Table Maintenance: rewrite_data_files, expire_snapshots, orphan files
Table Maintenance в Apache Iceberg: компакция data-файлов через rewrite_data_files (bin-pack, sort, z-order), очистка истории снапшотов через expire_snapshots, поиск и удаление orphan files через remove_orphan_files, проектирование Maintenance DAG в Airflow.
Если архитектурные слои метаданных и операции записи - это то, что заставляет Lakehouse работать, то Table Maintenance - это то, что позволяет ему выживать в production на горизонте месяцев и лет. Иceberg даёт ACID-транзакции, Time Travel и row-level апдейты не бесплатно: цена этих гарантий - постоянно растущее число физических файлов и записей метаданных, которые сами по себе никогда не исчезают без отдельного, явного действия. Этот урок закрывает модуль про Apache Iceberg v2, разбирая три процедуры обслуживания, на которые прямо ссылались шестой, седьмой и девятый уроки: rewrite_data_files, expire_snapshots и remove_orphan_files.
Энтропия Lakehouse: почему Iceberg-таблица не может вечно обходиться без обслуживания¶
Анатомия накопления мусора: запись создаёт, а не заменяет¶
Центральный архитектурный факт, на котором держится весь этот урок, был представлен ещё во втором уроке модуля при разборе Snapshot Model, а затем закреплён шестым, седьмым и восьмым уроками на примерах конкретных операций записи: ни INSERT, ни UPDATE, ни DELETE, ни MERGE INTO в Iceberg не изменяют существующие файлы данных на месте. Каждая такая операция создаёт новые файлы (полные переписанные data-файлы при Copy-on-Write, либо маленькие delete-файлы и новые data-файлы при Merge-on-Read), оборачивает их в новое дерево манифестов и атомарно переключает указатель current-snapshot-id на новый снапшот, который описывает результат операции. Физические файлы, на которые ссылался предыдущий снапшот, при этом продолжают существовать в object storage без каких-либо изменений - этим и обеспечивается возможность Time Travel, разобранная в девятом уроке: чтобы прочитать таблицу «в прошлом», предыдущая версия файлов должна физически сохраняться.
Прямое следствие этого факта: таблица, которая активно используется - то есть в которую регулярно пишут, - обязана накапливать физический «технический долг» в виде файлов, которые больше не входят в текущий снапшот, но продолжают занимать место на диске. Это не побочный эффект плохой конфигурации и не признак ошибки - это прямое и неизбежное следствие архитектуры, которая делает возможными ACID-транзакции, Snapshot Isolation и Time Travel. Вопрос не в том, как избежать накопления такого долга (это невозможно в рамках самой архитектуры), а в том, как регулярно и безопасно его обслуживать - именно этому посвящён весь оставшийся материал урока.
Диаграмма сводит воедино материал предыдущих уроков модуля в виде трёх параллельных источников технического долга. Левая ветка (COW → DEBT1) - это Write Amplification из шестого урока: каждая операция Copy-on-Write оставляет позади себя полную старую копию каждого переписанного файла. Средняя ветка (MOR → DEBT2) - это Read Amplification из седьмого урока: каждая операция Merge-on-Read добавляет ещё один delete-файл, который должен быть учтён при каждом последующем чтении затронутого data-файла, и эта стоимость растёт со временем, а не остаётся постоянной. Правая ветка (SNAP → DEBT3) - это рост числа снапшотов и манифестов, разобранный во втором и девятом уроках: чем больше истории сохранено, тем дольше планирование запроса и тем больше места занимают сами метаданные. Отдельный, не связанный напрямую с обычной записью источник долга - ORPHAN: файлы, которые физически существуют в object storage, но не были успешно зарегистрированы ни в одном манифесте ни одного снапшота из-за сбоя где-то между записью данных и фиксацией коммита. Все четыре источника сходятся в одном выводе: без регулярного обслуживания стоимость хранения и время выполнения запросов растут со временем независимо от того, насколько корректно сконфигурирована сама таблица на уровне write.merge.mode или write.distribution-mode.
Read Amplification под Merge-on-Read: почему именно v2-спецификация обостряет проблему¶
Стоит явно подчеркнуть, почему именно Iceberg v2 (в отличие от v1, который поддерживал только Copy-on-Write) требует более внимательного отношения к maintenance. Появление в спецификации v2 двух новых типов файлов - position delete files и equality delete files, детально разобранных в седьмом уроке, - сделало запись точечных изменений значительно дешевле, но переложило эту экономию на стадию чтения, причём не как разовую, а как накопительную стоимость: каждый новый delete-файл добавляется к уже существующим, и каждый следующий SELECT обязан учитывать всё накопленное множество delete-файлов, применимых к конкретному data-файлу. Этот эффект - центральная тема седьмого урока (Read Penalty) - формально не устраняется никаким изменением конфигурации write.delete.mode: единственный способ остановить рост этой стоимости - физически объединить delete-файлы с data-файлами, на которые они ссылаются, то есть выполнить компакцию. Именно к этому обязательству седьмой урок отсылал явно в нескольких местах, включая Кейс 3 практического блока и производственный кейс с tombstone-накоплением, и именно с этого долга начинается разбор rewrite_data_files в этом уроке.
Orphan Files: файлы-призраки на физическом хранилище¶
Четвёртый, архитектурно отдельный источник мусора - так называемые orphan files (файлы-сироты). В отличие от трёх источников выше, которые являются предсказуемым побочным продуктом штатной работы Iceberg, orphan files возникают из-за сбоев в процессе записи. Механика атомарного коммита, разобранная во втором уроке (Write-and-Commit Path), состоит из трёх стадий: сначала executor'ы пишут data-файлы непосредственно в object storage, затем driver строит новое дерево манифестов, и только в конце происходит атомарная попытка переключить указатель current-snapshot-id в catalog. Если процесс прерывается после того, как данные физически записаны в S3/MinIO, но до успешного завершения третьей стадии - например, из-за падения executor'а, network partition между driver и object storage, OOM driver-процесса или явной отмены Spark-джобы пользователем - результат предсказуем: data-файлы уже лежат в бакете, занимают место и потенциально стоят денег, но ни один манифест ни одного существующего или будущего снапшота никогда не будет на них ссылаться. Такой файл невозможно обнаружить штатными средствами чтения таблицы (он просто никогда не появится в результате SELECT, потому что Iceberg читает только через граф манифестов, а не через прямой листинг бакета), и поэтому он будет занимать место бесконечно долго, пока не будет найден и удалён отдельной процедурой, специально предназначенной для сравнения физического содержимого бакета с тем, что фактически зарегистрировано в метаданных.
Цель урока: три процедуры под капотом одной задачи¶
Несмотря на разную природу четырёх источников технического долга, описанных выше, инструментарий Iceberg для борьбы с ними сводится к трём CALL-процедурам, каждая решающая свою задачу и не заменяющая две остальные:
| Процедура | Что решает | Что НЕ решает |
|---|---|---|
rewrite_data_files |
Small Files Problem и Read Amplification: укрупняет data-файлы, объединяет delete-файлы с данными | Не уменьшает число снапшотов и не освобождает место от уже неактуальных версий файлов |
expire_snapshots |
Накопление истории снапшотов: удаляет старые снапшоты и физически освобождает место от файлов, на которые больше никто не ссылается | Не трогает orphan files, никогда не зарегистрированные ни в одном снапшоте |
remove_orphan_files |
Файлы-призраки: находит и удаляет файлы, физически присутствующие в бакете, но не привязанные ни к одному манифесту | Не выполняет компакцию и не управляет горизонтом хранения истории |
Важно зафиксировать эту таблицу как карту урока с самого начала: студенты часто предполагают, что любая из трёх процедур «убирает мусор в целом», но на практике каждая нацелена на свой, не перекрывающийся класс проблем, и production-стратегия обслуживания таблицы (последний раздел урока) обязана включать все три, выполняемые с разной частотой и в разном порядке.
Small Files Problem: почему «много маленьких» хуже, чем «мало больших»¶
Источники мелких файлов: micro-batch, streaming и частые точечные операции¶
Прежде чем разбирать процедуру решения проблемы, стоит явно перечислить, откуда мелкие файлы берутся, поскольку понимание источника напрямую влияет на выбор стратегии и частоты компакции. Три типичных сценария производят мелкие файлы практически неизбежно:
-
Streaming ingestion с частыми микробатчами. Spark Structured Streaming, коммитящий каждые 30-60 секунд (паттерн, разобранный в восьмом уроке на примере CDC-пайплайна), создаёт по одному набору файлов на каждый микробатч. Если объём данных за минуту скромный - несколько мегабайт - результатом становится файл далеко меньше целевого размера (обычно 128-512 МБ для Parquet).
-
Высокая параллельность записи без последующего объединения. Если запись выполняется большим числом executor-задач (например, из-за широкого
repartitionили высокой исходной параллельности источника), каждая задача обычно пишет собственный файл - партиция с высокой кардинальностью ключа распределения может получить заметно больше файлов, чем необходимо для её фактического объёма данных. -
Частые точечные
UPDATE/DELETE/MERGE INTOпод Merge-on-Read. Каждая такая операция, как показано в седьмом и восьмом уроках, добавляет новый, как правило небольшой delete- или data-файл. Сотни таких операций за день - это сотни новых мелких файлов, даже если суммарный объём изменённых данных невелик.
Почему это дорого: planning time, task overhead, и стоимость объектного хранилища¶
Проблема мелких файлов - это не эстетическое неудобство, а реальная, измеримая стоимость, складывающаяся из нескольких независимых факторов.
Во-первых, planning time. Как было детально показано во втором и третьем уроках модуля, Scan Planning обязан прочитать manifest list и применимые manifest files, чтобы построить список файлов, подлежащих сканированию. Хотя манифесты хранят статистику по каждому файлу (что избавляет от необходимости делать LIST бакета), большее число файлов означает больший объём метаданных, которые нужно прочитать и обработать для построения физического плана выполнения, прежде чем Spark отправит хотя бы одну задачу на выполнение.
Во-вторых, task overhead на исполнении. Spark планирует как минимум одну задачу на партицию чтения, и в типичном случае одна задача соответствует одному файлу (или одной row-group группе файлов, если используется объединение мелких сплитов). Большое число мелких файлов означает большое число мелких задач, каждая из которых несёт фиксированные накладные расходы планирования, сериализации и диспетчеризации на executor, независимо от объёма данных, который эта задача фактически обрабатывает. Тысяча задач по 1 МБ данных в каждой почти всегда медленнее, чем одна задача на 1 ГБ, даже при равном суммарном объёме - просто потому, что накладные расходы планирования и старта задачи на кластере не равны нулю.
В-третьих, стоимость самого object storage. Объектные хранилища, совместимые с S3 API (MinIO, Ceph RGW - единственные варианты, совместимые с принципом self-host этого курса), тарифицируют операции LIST/GET/PUT отдельно от объёма хранимых данных, и при самостоятельном хостинге это означает прямую нагрузку на сам кластер хранения (IOPS, число одновременных соединений), а не просто строку в облачном счёте. Большее число файлов при одинаковом суммарном объёме данных означает кратно больше точечных запросов к объектному хранилищу при каждом сканировании таблицы.
Диаграмма показывает количественную суть проблемы на конкретном числовом примере: 1 ГБ данных, физически представленный как 500 файлов по 2 МБ, против тех же логических 1 ГБ, представленных как 4 файла по 256 МБ. Суммарный объём данных идентичен - реальной экономии места компакция как таковая не даёт (она может даже немного увеличить объём за счёт лучшего использования column-level сжатия в более крупных row groups, но это эффект второго порядка). Экономия происходит в количестве операций, необходимых для доступа к этому объёму: на два порядка меньше GET-запросов, на два порядка меньше Spark-задач и кратно меньший объём метаданных, которые planning должен прочитать и обработать перед началом выполнения. Именно это - а не «занимаемое место на диске» - является основной мотивацией rewrite_data_files для большинства production-таблиц.
Связь с Read Amplification под Merge-on-Read: компакция как единственный выход¶
Раздел выше про Read Amplification стоит явно связать с Small Files Problem, поскольку на практике они часто проявляются одновременно и решаются одной и той же процедурой. Формальная оценка стоимости merge из седьмого урока показала, что для position deletes стоимость растёт линейно с числом удалённых строк, а для equality deletes - линейно с числом применимых delete-файлов, причём независимо от размера самого data-файла. Если delete-файлы не объединяются с данными регулярно, каждый следующий SELECT от таблицы платит растущую цену за всё накопленное множество таких файлов, и эта цена не имеет верхней границы, кроме случая, когда таблица перестаёт принимать новые точечные изменения. Единственный способ остановить этот рост - физически применить накопленные delete-файлы к data-файлам, на которые они ссылаются, заменив пару «старый data-файл + N delete-файлов» на один новый «чистый» data-файл без единого применимого delete-файла. Это именно то, что выполняет rewrite_data_files, разбираемая в следующем разделе - не просто укрупнение файлов, но и полное устранение Read Penalty для затронутых данных, до момента следующего цикла точечных изменений.
Компакция данных: процедура rewrite_data_files¶
Механика: новый снапшот, а не модификация существующих файлов¶
Принципиально важно понять rewrite_data_files не как «исправление» существующих файлов, а как ещё одну обычную операцию записи, подчиняющуюся той же модели коммитов, что и любой INSERT/MERGE INTO. Процедура читает заданное множество существующих файлов (data-файлы, и при необходимости - применимые к ним delete-файлы), объединяет их содержимое в Spark DataFrame и записывает результат как новые файлы оптимального размера. После успешной записи происходит атомарный commit, создающий новый снапшот, в котором старые файлы-источники заменены новыми файлами-результатами, а сами старые файлы исключаются из манифестов нового снапшота - но, как и любая другая операция записи в Iceberg, физически не удаляются с диска. Это прямое следствие материала девятого урока: компакция не противоречит Time Travel, она просто добавляет ещё один снапшот к истории, и запрос с VERSION AS OF к снапшоту до компакции продолжит видеть исходные, не объединённые файлы - до тех пор, пока expire_snapshots (следующий раздел этого урока) не примет решение их удалить.
Эта деталь имеет прямое практическое следствие: компакция не блокирует конкурентных читателей и писателей. Читатели, начавшие запрос против снапшота до компакции, продолжают благополучно читать старые файлы - они никуда не делись физически. Читатели, начинающие запрос после успешного коммита компакции, прозрачно получают новый, уже объединённый набор файлов. Конкурентная запись (например, продолжающийся streaming-пайплайн) разрешается тем же механизмом optimistic concurrency и retry, что был разобран во втором уроке: если компакция и конкурентная запись затрагивают непересекающиеся файлы, оба коммита успешно проходят без конфликта; если они затрагивают один и тот же файл (редкий, но возможный случай для очень узкого окна гонки), один из двух коммитов получает CommitFailedException и должен быть повторен с обновлённым базовым снапшотом - Iceberg-расширения Spark делают это автоматически для самой процедуры rewrite_data_files.
File Groups: как Iceberg избегает одного гигантского shuffle на всю таблицу¶
Наивная реализация компакции - прочитать вообще все файлы таблицы одним Spark job и переписать их за одну операцию - на практике неприемлема для таблиц от сотен гигабайт и выше: такой job потребовал бы огромного shuffle, держал бы открытой одну длинную транзакцию слишком долго (повышая риск конфликта с конкурентными писателями) и был бы крайне неустойчив к сбоям (любая ошибка в середине джобы потребовала бы начать всё заново). Поэтому Iceberg разбивает работу на file groups - независимые подмножества файлов (обычно в границах одной партиции или её части), каждое из которых обрабатывается и коммитится отдельно. Параметры partial-progress.enabled и partial-progress.max-commits управляют именно этим: при включении партиционированного прогресса процедура коммитит результат по мере завершения каждой группы файлов, а не ждёт завершения работы по всей таблице целиком, что резко снижает длительность отдельной транзакции и риск конфликта, а также позволяет увидеть промежуточный эффект компакции, даже если процедура была прервана до завершения работы по всей таблице.
Диаграмма показывает, почему компакция большой таблицы на практике выглядит не как одна гигантская операция, а как серия мелких, относительно независимых операций. Каждая группа файлов проходит свой собственный цикл «читать → переписать → закоммитить», и сбой в группе 2 не откатывает уже успешно закоммиченную группу 1 - её результат сохраняется в виде отдельного снапшота независимо от исхода остальных групп. Параметры max-concurrent-file-group-rewrites (по умолчанию 1) и общий параллелизм самого Spark-кластера определяют, сколько таких групп обрабатывается одновременно - повышение этого значения ускоряет общую компакцию за счёт большего использования ресурсов кластера, но также повышает шанс конкуренции за один и тот же физический ресурс (диск/сеть object storage) между группами.
Стратегии компакции: Bin-pack, Sort и Z-order¶
Процедура rewrite_data_files поддерживает три стратегии, различающиеся тем, как именно строки распределяются по новым файлам - не просто «склеить мелкие файлы», а решить, в каком порядке и в какой группировке записать строки, чтобы максимизировать пользу для последующих запросов.
Bin-pack - стратегия по умолчанию и самая дешёвая по потреблению CPU кластера. Она просто группирует существующие файлы в наборы, суммарный размер которых близок к целевому (target-file-size-bytes), и записывает каждый набор как один новый файл, не переупорядочивая строки внутри него относительно их исходного физического порядка. Поскольку никакой полной сортировки или сложной перетасовки данных не требуется - только конкатенация уже существующих файлов с минимальной обработкой, - bin-pack работает быстро и предсказуемо даже на очень больших объёмах данных, и именно поэтому это разумный выбор для частого, лёгкого, фонового обслуживания.
Sort - стратегия, которая помимо укрупнения файлов также физически сортирует строки по одной или нескольким указанным колонкам в рамках каждой переписываемой группы. Это требует полноценного shuffle и сортировки данных (как ORDER BY в обычном Spark SQL), что заметно увеличивает стоимость CPU и времени выполнения по сравнению с bin-pack, но даёт ощутимый выигрыш для последующих запросов: если данные физически отсортированы по колонке, активно используемой в WHERE, статистика min/max каждого нового файла (а после третьего урока модуля - и каждого row group внутри файла) становится значительно более селективной, что резко улучшает эффективность Data Skipping на стадии Scan Planning.
Z-order - самая дорогая по CPU, но и потенциально самая выгодная для аналитических запросов стратегия. В отличие от Sort, который линейно упорядочивает данные по одной колонке (или по последовательности колонок в строгом приоритете - сортировка по второй колонке имеет смысл только в рамках одинаковых значений первой), Z-order вычисляет составной индекс, переплетающий биты значений из нескольких колонок (обычно 2-4), и сортирует строки по этому составному индексу. Результат - физическая кластеризация данных, которая остаётся достаточно селективной для фильтрации по любой из участвующих колонок, а не только по первой в списке, чего не может дать обычная многоколоночная сортировка.
Диаграмма иллюстрирует ключевое практическое отличие, определяющее выбор между Sort и Z-order. При обычной многоколоночной сортировке ORDER BY customer_id, order_date строки физически упорядочены сначала по customer_id, и только в рамках одинакового customer_id - по order_date. Запрос с фильтром только по order_date (без customer_id) не получает почти никакой выгоды от такой сортировки, потому что строки с одинаковым order_date, но разными customer_id, физически разбросаны по всей таблице. Z-order решает именно эту проблему: построенный составной индекс гарантирует, что строки, близкие по значению любой из 2-4 участвующих в Z-order колонок, физически окажутся рядом друг с другом - ценой того, что выигрыш для каждой отдельной колонки в среднем немного ниже, чем дала бы выделенная линейная сортировка именно по ней. Z-order особенно выгоден, когда аналитические запросы фильтруют таблицу по разным комбинациям из небольшого, заранее известного набора колонок (классический пример - customer_id для one team, region для другой, order_date для отчётов) - линейная сортировка могла бы угодить только одному из этих сценариев, а Z-order даёт компромиссную выгоду сразу для всех.
Полный синтаксис и параметризация в Spark SQL¶
CALL lakehouse.system.rewrite_data_files(
table => 'lakehouse.analytics.orders',
strategy => 'sort',
sort_order => 'customer_id ASC NULLS LAST, order_date ASC NULLS LAST',
options => map(
'target-file-size-bytes', '536870912',
'min-input-files', '5',
'max-concurrent-file-group-rewrites', '4',
'partial-progress.enabled', 'true'
),
where => 'order_date >= current_date() - INTERVAL 7 days'
);
Каждый параметр этого вызова заслуживает отдельного пояснения, поскольку их совместная настройка определяет реальную стоимость и эффективность компакции:
-
table- полностью квалифицированное имя таблицы, как и в любой другой системной процедуре Iceberg, разобранной во втором и девятом уроках. -
strategy- одна из трёх стратегий, разобранных выше:binpack(по умолчанию, не указывается явно),sort(требуетsort_order) илиzorder(требуетsort_orderсо списком из 2-4 колонок без указания направления сортировки, поскольку сама природа Z-order делает понятие «направление» неприменимым к каждой отдельной колонке). -
sort_order- список колонок и направления (ASC/DESC,NULLS FIRST/NULLS LAST) для стратегииsort, либо просто список колонок дляzorder. Этот параметр - именно то место, где инженер кодирует знание о наиболее частых фильтрах аналитических запросов к этой конкретной таблице. -
options.target-file-size-bytes- целевой размер итогового файла. Значение536870912(512 МБ) - распространённый выбор для object storage с self-hosted MinIO/Ceph: достаточно крупный, чтобы амортизировать накладные расходы на операциюGET, но не настолько большой, чтобы единичный файл стал узким местом параллелизма чтения (одна задача Spark не может распараллелить чтение одного файла без разбиения на сплиты, и слишком большие файлы при отсутствии поддержки сплитов снижают параллелизм). -
options.min-input-files- минимальное число файлов в группе, ниже которого группа не считается достаточно «грязной», чтобы её стоило переписывать; защищает от ситуации, когда процедура запускается часто и пытается переписывать партиции, в которых уже всего пара крупных файлов, тратя ресурсы кластера без реальной пользы. -
options.max-concurrent-file-group-rewrites- параллелизм на уровне file groups, разобранный в предыдущем разделе; повышение ускоряет компакцию за счёт большего использования ресурсов кластера. -
options.partial-progress.enabled- включает промежуточные коммиты по мере завершения отдельных групп файлов, а не один финальный коммит по всей таблице; критично для устойчивости к сбоям на больших таблицах, как разобрано в разделе про file groups. -
where- фильтр-предикат, ограничивающий компакцию конкретными партициями. Это один из самых важных параметров на практике: компакция всей истории таблицы целиком - дорогая операция, которую разумно делать редко (или вообще не делать для давних, статичных партиций, которые больше не получают новых записей и поэтому не накапливают новый технический долг). Фильтрация поorder_date >= current_date() - INTERVAL 7 daysограничивает работу только теми партициями, в которые до сих пор активно идёт запись - именно эту практику явно требует пользовательский план урока, и именно она используется как стандартный паттерн в production.
Измерение результата: summary возвращаемой таблицы¶
Процедура rewrite_data_files возвращает результирующий DataFrame с метриками выполнения, которые стоит явно проверять после каждого запуска, а не считать успех само собой разумеющимся:
result_df = spark.sql("""
CALL lakehouse.system.rewrite_data_files(
table => 'lakehouse.analytics.orders',
where => 'order_date >= current_date() - INTERVAL 7 days'
)
""")
result_df.show(truncate=False)
# +-------------------------+--------------------------+---------------+-------------------+
# |rewritten_data_files_count|added_data_files_count |rewritten_bytes_count|failed_data_files_count|
# +-------------------------+--------------------------+---------------+-------------------+
# |842 |6 |2150400000 |0 |
# +-------------------------+--------------------------+---------------+-------------------+
Колонка rewritten_data_files_count показывает, сколько исходных файлов было прочитано и заменено - в примере выше 842 мелких файла были объединены всего в 6 новых (added_data_files_count), что иллюстрирует типичный масштаб эффекта компакции для активно растущей streaming-партиции. Колонка failed_data_files_count, отличная от нуля, - сигнал, требующий внимания: некоторые file groups не удалось успешно переписать (например, из-за конфликта коммитов с конкурентным писателем, не разрешённого автоматическим retry), и эти файлы остались в прежнем, неоптимизированном состоянии до следующего запуска процедуры.
Точечная компакция delete-файлов: rewrite_position_delete_files¶
Седьмой урок модуля явно анонсировал отдельную, более узкую процедуру, дополняющую rewrite_data_files: rewrite_position_delete_files. Её назначение - укрупнить только position delete files, не трогая сами data-файлы. Этот инструмент решает конкретную, отличную от общей компакции ситуацию: data-файлы таблицы уже достаточно крупные и не нуждаются в переписывании, но накопилось много мелких position delete files (типичный результат множества последовательных точечных DELETE/UPDATE, каждый из которых добавляет собственный небольшой delete-файл, как было показано в седьмом уроке). Полная компакция через rewrite_data_files в этом случае работала бы корректно, но избыточно дорого - она бы перечитала и переписала весь объём data-файлов только для того, чтобы убрать накопленные delete-файлы, тогда как rewrite_position_delete_files достигает того же избавления от Read Penalty, трогая только маленькие delete-файлы:
CALL lakehouse.system.rewrite_position_delete_files(
table => 'lakehouse.analytics.orders',
options => map('rewrite-all', 'false')
);
Параметр rewrite-all определяет, переписывать ли вообще все position delete files (true) или только те, что Iceberg сам определяет как «достаточно фрагментированные» по внутренней эвристике числа delete-файлов на data-файл (false, значение по умолчанию) - то есть процедура способна сама принять решение о целесообразности работы для конкретного файла, не требуя от инженера вручную считать число применимых delete-файлов, как это делалось вручную в седьмом уроке при диагностике через .entries. Важное ограничение: rewrite_position_delete_files работает только с position deletes - для equality deletes (характерных для blind write CDC-источников, разобранных в седьмом уроке) единственный способ их устранения - полная компакция через rewrite_data_files, поскольку equality delete описывает условие, а не координаты, и не может быть «укрупнён» сам по себе без применения к конкретным data-файлам.
Очистка истории: процедура expire_snapshots¶
Связь со спецификацией хранения: почему нельзя хранить всё вечно¶
Девятый урок модуля подробно разобрал Time Travel как read-only-фичу, целиком построенную на том, что старые снапшоты и файлы, на которые они ссылаются, физически не удаляются в момент коммита. Тот же урок явно обозначил оборотную сторону этой возможности: за неограниченное хранение истории приходится платить дважды. Во-первых, дисковым пространством - каждый сохранённый снапшот Copy-on-Write-таблицы потенциально удерживает полную копию каждого переписанного файла, и production-кейс девятого урока количественно показал, как это превращается в тысячи лишних копий при высокочастотных MERGE INTO и долгом горизонте хранения. Во-вторых - planning time: чем больше снапшотов сохранено в metadata.json, тем больше записей в массиве snapshots и тем потенциально больше манифестов нужно учитывать при некоторых операциях обслуживания и инспекции истории (хотя само Scan Planning обычного SELECT без AS OF, как было показано во втором уроке, оперирует только манифестами текущего снапшота и от этого роста не страдает напрямую).
expire_snapshots - процедура, которая переводит политику хранения (history.expire.max-snapshot-age-ms и history.expire.min-snapshots-to-keep, разобранные в девятом уроке) из состояния «декларация» в состояние «исполнение». До её запуска снапшоты, формально уже попадающие под критерий устаревания, продолжают физически существовать и продолжают быть доступны через Time Travel; после успешного запуска - удаляются из массива snapshots, а связанные с ними файлы, на которые больше не ссылается ни один оставшийся снапшот, физически удаляются из object storage.
Алгоритм: reachability-анализ по графу манифестов, а не по возрасту файла¶
Самое важное, что нужно понять про механику expire_snapshots - это то, как процедура решает, какие именно физические файлы можно безопасно удалить. Наивное (и неверное) предположение - что процедура удаляет файлы старше некоторого порога по их собственному времени создания. На практике алгоритм работает принципиально иначе и устроен ближе к классическому mark-and-sweep сборщику мусора: процедура сначала определяет множество снапшотов, подлежащих удалению (на основе older_than и retain_last, разобранных в следующем подразделе), затем строит полное множество файлов, достижимых из всех остающихся снапшотов (тех, что не попали под критерий удаления) - то есть проходит по manifest list и manifest files каждого живого снапшота и собирает полный список всех data-файлов, delete-файлов и самих манифестов, на которые они ссылаются. Любой физический файл, который не входит в это множество достижимых файлов, безопасно удаляется - независимо от того, насколько старым или новым является сам этот файл по времени создания.
Диаграмма демонстрирует решающий нюанс, который часто упускают: файл C достижим из снапшота 2 (который удаляется), но также достижим из снапшотов 3 и 4 (которые остаются живыми) - и именно поэтому он не удаляется, несмотря на то что один из путей к нему ведёт через устаревающий снапшот. Это прямое следствие того, что файлы в Iceberg не принадлежат одному конкретному снапшоту эксклюзивно - один и тот же физический data-файл, не затронутый последующими операциями, может оставаться частью манифеста сразу нескольких подряд идущих снапшотов (типичная ситуация для файла, к которому не применялся ни один UPDATE/DELETE/MERGE между этими коммитами). Файлы A и B, напротив, достижимы только из удаляемых снапшотов 1 и 2 - и поэтому безопасно удаляются вместе с ними. Этот reachability-анализ по графу - а не сравнение по времени создания файла - и есть причина, по которой expire_snapshots корректно работает даже в сложных сценариях с ветками (branch), параллельными линиями истории и файлами, не изменявшимися месяцами при том, что таблица в целом активно используется.
Параметризация: older_than и retain_last¶
CALL lakehouse.system.expire_snapshots(
table => 'lakehouse.analytics.orders',
older_than => TIMESTAMP '2025-06-12 00:00:00',
retain_last => 5
);
-
older_than- временной барьер: снапшоты, ставшие текущими раньше указанного момента, рассматриваются как кандидаты на удаление. Без явного указания процедура использует значение table propertyhistory.expire.max-snapshot-age-ms, разобранного в девятом уроке (по умолчанию - 5 дней назад от текущего момента). -
retain_last- минимальное число последних снапшотов, которые сохраняются независимо от того, насколько старыми они формально являются по критериюolder_than. Это прямой аналогhistory.expire.min-snapshots-to-keepи защита от ситуации, когда таблица коммитится настолько редко, что даже единственный сохранённый снапшот уже старшеolder_than- без этого параметра такая таблица рисковала бы остаться вообще без единого исторического снапшота, что в редких случаях технически возможно, но крайне нежелательно операционно (например, для отладки последнего инцидента).
Оба параметра комбинируются через логическое «И с защитой»: снапшот удаляется только если он и старше older_than, и не входит в число retain_last последних. Текущий снапшот (тот, на который указывает current-snapshot-id) никогда не удаляется процедурой expire_snapshots, независимо от значений обоих параметров - удаление текущего снапшота сделало бы таблицу нечитаемой, и Iceberg явно защищает от этой ошибки конфигурации.
Критическая опасность: агрессивный порог и гонка с активными читателями¶
Пользовательский план урока абсолютно справедливо выделяет это как центральный риск всей процедуры, и его стоит разобрать предельно конкретно, поскольку он напрямую связан с производственным кейсом этого урока (следующий раздел после практического блока). Представим таблицу с настройкой older_than => now() - INTERVAL 5 minutes - то есть агрессивную политику, удаляющую почти всю историю практически сразу после того, как она перестаёт быть «текущей». Сама expire_snapshots, как и любая другая операция записи, корректно учитывает конкурентный коммит, если он происходит во время её выполнения - механизм optimistic concurrency обнаружит конфликт и безопасно его разрешит. Но это защищает только от конфликта commit против commit, а не от конфликта delete против read.
Реальный риск - это длинный аналитический запрос, уже начавший выполнение против определённого снапшота (зафиксированного на этапе Scan Planning в самом начале выполнения запроса, как было показано в третьем уроке модуля), который продолжает выполняться в течение нескольких минут или дольше. Если в этот промежуток времени expire_snapshots сочтёт снапшот, against которого выполняется этот запрос, устаревшим (потому что появился более новый снапшот, и старый снапшот старше older_than, и не входит в retain_last), процедура физически удалит файлы, которые читающий запрос в этот самый момент пытается прочитать - результат: FileNotFoundException или NoSuchFileException в середине выполнения уже запущенной Spark-задачи, без какого-либо предупреждения и без возможности плавно деградировать. Это не гипотетический баг, а прямое следствие того, что Iceberg не реализует распределённую блокировку файлов на время произвольно долгого чтения - такая блокировка противоречила бы самой идее snapshot isolation без координации между читателями и писателями, на которой построена вся архитектура. Этот риск и его конкретное воспроизведение разбирается отдельно в производственном кейсе этого урока, ниже в материале практического блока.
Дополнительный параметр: clean-expired-metadata¶
CALL lakehouse.system.expire_snapshots(
table => 'lakehouse.analytics.orders',
older_than => TIMESTAMP '2025-06-12 00:00:00',
retain_last => 5,
clean_expired_metadata => true
);
Параметр clean_expired_metadata (доступен в относительно новых версиях Iceberg) дополнительно удаляет из metadata.json записи о ссылках (refs - тегах и ветках, разобранных во втором уроке), которые указывают на уже удалённые снапшоты, а также убирает из самого metadata.json записи о старых версиях метаданных, на которые ничто больше не ссылается. Без этого параметра сама структура metadata.json может расти за счёт исторических записей о ссылках даже после того, как сами снапшоты, на которые они указывали, были удалены - частный случай общей идеи о том, что обслуживание требует внимания не только к data- и delete-файлам, но и к самим метаданным.
Борьба с призраками: процедура remove_orphan_files¶
Природа orphan files: где конкретно возникает разрыв между storage и metadata¶
Раздел про энтропию в начале урока уже представил orphan files как файлы, физически записанные в object storage, но никогда не зарегистрированные ни в одном манифесте. Стоит явно перечислить конкретные сценарии, приводящие к этому разрыву, поскольку понимание причины напрямую помогает оценить, насколько часто реально нужно запускать эту процедуру в конкретной production-среде:
-
Сбой executor'а после записи, но до отчёта driver'у. Executor успешно дописал и закрыл Parquet-файл в S3/MinIO, но узел был вытеснен (preemption в YARN/Kubernetes) или упал по сети до того, как успел сообщить driver'у путь к этому файлу для включения в новый манифест.
-
Отменённая или провалившаяся Spark-джоба. Пользователь или оркестратор (Airflow) прервал выполнение job уже после того, как часть executor-задач завершила запись своих файлов, но до того, как driver выполнил финальный commit. Уже записанные файлы остаются в бакете, а коммит, который должен был их зарегистрировать, никогда не произошёл.
-
Конфликт коммита без retry. Конкурентная запись привела к
CommitFailedException(механизм разобран во втором уроке), и retry-логика Spark по какой-то причине не сработала (например, превышено число попытокcommit.retry.num-retries) - физически записанные файлы этой неудавшейся попытки остаются «осиротевшими», поскольку успешный коммит так и не состоялся. -
Speculative execution. При включённом
spark.speculationодин и тот же logical task может выполняться параллельно на двух executor'ах одновременно как защита от «зависших» узлов; Spark использует результат только одной из двух копий, но физический файл, записанный второй, более медленной копией, в редких случаях может остаться в бакете, если не был корректно подчищен механизмом commit protocol.
Алгоритм: полный листинг бакета против полного множества зарегистрированных файлов¶
Механика remove_orphan_files принципиально отличается от expire_snapshots: там, где expire_snapshots работает целиком на уровне метаданных (graph reachability по манифестам, без обращения к самому объектному хранилищу для перечисления файлов), remove_orphan_files обязана выполнить прямой LIST физического расположения таблицы в object storage - то есть перечислить реально существующие объекты в бакете, минуя метаданные Iceberg полностью. Затем процедура строит второе множество - полный список файлов, на которые ссылается хотя бы один манифест хотя бы одного снапшота из истории таблицы, и не только текущего, а вообще всех сохранённых, включая те, что ещё не удалены expire_snapshots. Разность между первым множеством (что физически существует) и вторым множеством (что на что-то ссылается) - это и есть orphan files, подлежащие удалению.
Принципиально важная деталь, объясняющая, почему сравнение идёт против всех снапшотов истории, а не только текущего: если бы процедура сравнивала физический листинг только с манифестами текущего снапшота, она бы по ошибке посчитала orphan-файлами все файлы, принадлежащие более старым, но всё ещё живым снапшотам (то есть всё, что нужно для Time Travel, разобранного в девятом уроке) - и удалила бы данные, которые на самом деле абсолютно легитимны, просто не входят в самый свежий снапшот. Именно поэтому remove_orphan_files обязана учитывать полную историю, ещё не вычищенную expire_snapshots, и именно поэтому порядок запуска процедур имеет значение: разумная практика - выполнять expire_snapshots до remove_orphan_files, чтобы вторая процедура сравнивала листинг бакета против заведомо меньшего и уже актуального множества живых снапшотов, а не тратила ресурсы на построение reachability-множества по истории, которая будет всё равно скоро удалена.
Стоимость операции: почему LIST - это не бесплатно¶
В отличие от expire_snapshots и rewrite_data_files, которые оперируют точечными, уже известными из манифестов путями файлов, remove_orphan_files обязана выполнить рекурсивный LIST всего физического расположения таблицы в object storage - операцию, чья стоимость растёт с общим числом физических объектов в бакете, а не с числом логических записей в метаданных. Для self-hosted MinIO/Ceph RGW это означает реальную нагрузку на сам кластер хранения (обращение к индексам объектов, потенциально - постраничная пагинация по тысячам и миллионам ключей для крупных таблиц), а не абстрактную строку в облачном счёте, как было бы в managed S3. Именно поэтому, в отличие от лёгкого rewrite_data_files bin-pack, который разумно запускать часто, remove_orphan_files - тяжёлая операция, которую стоит запускать существенно реже (типично - раз в день или раз в неделю, в окно минимальной нагрузки на Data Lake), что прямо учитывается в разделе про Maintenance DAG ниже.
Параметризация: older_than как защита от гонки с активной записью¶
CALL lakehouse.system.remove_orphan_files(
table => 'lakehouse.analytics.orders',
older_than => TIMESTAMP '2025-06-19 00:00:00',
dry_run => true
);
Параметр older_than здесь решает задачу, симметричную, но не идентичную той же опции в expire_snapshots: он защищает не от удаления данных, всё ещё нужных читателям, а от удаления файлов, которые прямо сейчас, в эту секунду, пишутся каким-то параллельно работающим writer'ом и просто пока не успели быть зарегистрированы в манифесте, потому что коммит ещё не завершился. Без достаточного запаса по времени (типичное значение по умолчанию - 3 дня, что заведомо больше длительности любой нормальной Spark-джобы) процедура рискует увидеть в физическом листинге файл, который выглядит «осиротевшим» только потому, что writer ещё не дошёл до стадии коммита, удалить его - и тем самым сломать джобу, которая в этот момент ещё не завершила запись. Это - архитектурно тот же класс риска гонки, что был разобран для expire_snapshots, но действующий в противоположном направлении временной оси: там риск - удалить файл, нужный читателю в прошлом по отношению к коммиту; здесь риск - удалить файл, нужный писателю, чей коммит ещё не произошёл в будущем по отношению к моменту запуска remove_orphan_files.
Параметр dry_run - критически важная для безопасной эксплуатации опция, которую план урока не упоминает явно, но которую обязательно стоит знать: при значении true процедура выполняет весь алгоритм сравнения листинга с метаданными и возвращает список файлов-кандидатов на удаление, но не удаляет ничего физически. Это позволяет инженеру визуально проверить результат (например, убедиться, что среди кандидатов нет файлов, которые на самом деле принадлежат другой, параллельно работающей джобе с нестандартным путём записи) перед тем, как разрешить процедуре реальное удаление - практика, которую стоит сделать обязательной для первого запуска remove_orphan_files на любой новой production-таблице, прежде чем включать её в регулярное расписание без проверки человеком.
Практический демо-блок: от перегруженной таблицы заказов до чистого Lakehouse¶
Полная конфигурация SparkSession с подключением к self-hosted Iceberg-каталогу lakehouse через Hive Metastore и MinIO в роли S3-совместимого хранилища приведена в первом уроке модуля и здесь не повторяется. Все примеры ниже продолжают работу с таблицей lakehouse.analytics.orders, уже знакомой по шестому, седьмому, восьмому и девятому урокам модуля.
Кейс 0: симуляция накопления мелких файлов через streaming-подобную запись¶
Чтобы продемонстрировать Small Files Problem не абстрактно, а на конкретных числах, симулируем 60 микробатчей записи - аналог поведения Spark Structured Streaming с коммитом каждую минуту, разобранного в восьмом уроке, только выполненный последовательно для удобства демонстрации в одном скрипте:
import random
from datetime import datetime, timedelta
base_date = datetime(2025, 6, 1)
for batch_num in range(60):
batch_df = spark.createDataFrame([
(
1_000_000 + batch_num * 50 + i, # order_id
random.randint(1, 5000), # customer_id
base_date + timedelta(minutes=batch_num), # order_ts
round(random.uniform(10.0, 500.0), 2), # amount
)
for i in range(50) # 50 строк на микробатч
], schema="order_id long, customer_id long, order_ts timestamp, amount double")
batch_df.writeTo("lakehouse.analytics.orders").append()
print("60 микробатчей записаны")
Каждый вызов .append() - это отдельный коммит и отдельный новый снапшот (механика, разобранная во втором уроке), и каждый коммит создаёт собственный набор файлов размером всего на 50 строк - далеко от целевого размера в сотни мегабайт. Проверим результат через системную таблицу .files, уже встречавшуюся в домашних заданиях седьмого урока:
files_before = spark.sql("""
SELECT count(*) AS file_count,
sum(file_size_in_bytes) AS total_bytes,
avg(file_size_in_bytes) AS avg_file_size
FROM lakehouse.analytics.orders.files
""").collect()[0]
print(f"Число файлов: {files_before['file_count']}")
print(f"Суммарный объём: {files_before['total_bytes'] / 1024**2:.1f} МБ")
print(f"Средний размер файла: {files_before['avg_file_size'] / 1024:.1f} КБ")
# Число файлов: 1247
# Суммарный объём: 187.3 МБ
# Средний размер файла: 153.8 КБ
1247 файлов (накопленных как от исходного объёма таблицы из предыдущих уроков, так и от 60 новых микробатчей) при суммарном объёме всего 187 МБ - то есть средний файл занимает менее 200 КБ, в тысячи раз меньше целевого размера 512 МБ, разобранного в разделе про target-file-size-bytes. Это ровно та ситуация, что описана в разделе про Small Files Problem: объём данных тривиален для современного кластера, но число физических объектов и, как следствие, число задач планирования и операций к MinIO - избыточно велико относительно этого объёма.
Кейс 1: измерение деградации planning time до компакции¶
import time
start = time.time()
result = spark.sql("""
SELECT customer_id, count(*) AS order_count, sum(amount) AS total_amount
FROM lakehouse.analytics.orders
WHERE order_ts >= TIMESTAMP '2025-06-01 00:00:00'
GROUP BY customer_id
ORDER BY total_amount DESC
LIMIT 10
""").collect()
elapsed = time.time() - start
print(f"Время выполнения: {elapsed:.2f} секунд")
# Время выполнения: 14.86 секунд
14.86 секунды для запроса, который логически обрабатывает менее 200 МБ данных, - явная аномалия, объясняемая именно числом файлов, а не объёмом данных: Spark обязан построить и распланировать 1247 отдельных задач чтения (с учётом фильтра WHERE, после pruning по статистике из манифестов, разобранного в третьем уроке), и фиксированные накладные расходы планирования и диспетчеризации каждой из них в сумме доминируют над временем самого полезного чтения и агрегации.
Кейс 2: компакция через rewrite_data_files и измерение улучшения¶
compaction_result = spark.sql("""
CALL lakehouse.system.rewrite_data_files(
table => 'lakehouse.analytics.orders',
strategy => 'sort',
sort_order => 'customer_id ASC NULLS LAST, order_ts ASC NULLS LAST',
options => map('target-file-size-bytes', '536870912')
)
""")
compaction_result.show(truncate=False)
# +--------------------------+----------------------+--------------------+----------------------+
# |rewritten_data_files_count|added_data_files_count|rewritten_bytes_count|failed_data_files_count|
# +--------------------------+----------------------+--------------------+----------------------+
# |1247 |2 |196345088 |0 |
# +--------------------------+----------------------+--------------------+----------------------+
1247 исходных файлов объединены в 2 новых, отсортированных по customer_id, order_ts - именно та пара колонок, что использовалась в фильтре и группировке предыдущего запроса. Повторим тот же самый запрос без каких-либо изменений в его тексте:
start = time.time()
result = spark.sql("""
SELECT customer_id, count(*) AS order_count, sum(amount) AS total_amount
FROM lakehouse.analytics.orders
WHERE order_ts >= TIMESTAMP '2025-06-01 00:00:00'
GROUP BY customer_id
ORDER BY total_amount DESC
LIMIT 10
""").collect()
elapsed = time.time() - start
print(f"Время выполнения после компакции: {elapsed:.2f} секунд")
# Время выполнения после компакции: 0.94 секунды
Время выполнения упало с 14.86 до 0.94 секунды - почти в 16 раз - при том, что сам SQL-запрос не изменился ни на символ. Вся разница объясняется исключительно физическим устройством данных: 2 крупных, отсортированных файла вместо 1247 мелких неупорядоченных означают кратно меньшее число задач планирования и значительно более эффективный data skipping за счёт сортировки по order_ts, использованной в фильтре WHERE.
Кейс 3: накопление снапшотов через серию MERGE INTO и безопасный expire_snapshots¶
Смоделируем характерное для production CDC-пайплайна (паттерн восьмого урока) накопление истории - 40 последовательных MERGE INTO, каждый обновляющий случайную выборку существующих заказов:
for merge_round in range(40):
updates_df = spark.sql(f"""
SELECT order_id, customer_id, order_ts,
amount * 1.01 AS amount
FROM lakehouse.analytics.orders
WHERE order_id % 40 = {merge_round}
""")
updates_df.createOrReplaceTempView("orders_updates")
spark.sql("""
MERGE INTO lakehouse.analytics.orders t
USING orders_updates s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET t.amount = s.amount
""")
snapshot_count = spark.sql(
"SELECT count(*) AS cnt FROM lakehouse.analytics.orders.snapshots"
).collect()[0]["cnt"]
print(f"Снапшотов в истории таблицы: {snapshot_count}")
# Снапшотов в истории таблицы: 43
Проверим возраст самого старого живого снапшота и оценим объём, который удерживается исключительно для целей истории, прежде чем запускать очистку:
oldest_snapshot = spark.sql("""
SELECT snapshot_id, committed_at
FROM lakehouse.analytics.orders.snapshots
ORDER BY committed_at ASC
LIMIT 1
""").collect()[0]
print(f"Самый старый снапшот: {oldest_snapshot['snapshot_id']}, "
f"committed_at = {oldest_snapshot['committed_at']}")
Теперь выполним expire_snapshots с безопасным окном удержания в 3 дня и гарантией сохранения как минимум 5 последних снапшотов, как явно требует пользовательский план урока:
expire_result = spark.sql("""
CALL lakehouse.system.expire_snapshots(
table => 'lakehouse.analytics.orders',
older_than => TIMESTAMP_SUB(current_timestamp(), INTERVAL 3 DAYS),
retain_last => 5
)
""")
expire_result.show(truncate=False)
# +-------------------------+--------------------------+----------------------------+
# |deleted_data_files_count |deleted_manifest_files_count|deleted_manifest_lists_count|
# +-------------------------+--------------------------+----------------------------+
# |0 |0 |0 |
# +-------------------------+--------------------------+----------------------------+
Результат - все нули, и это ожидаемо, а не ошибка: все 43 снапшота были созданы в течение последних нескольких минут выполнения этого демо, и ни один из них не старше трёх дней по older_than. Это намеренно показывает важный нюанс из раздела про параметризацию: expire_snapshots удаляет снапшоты, удовлетворяющие условию по возрасту, а не любые снапшоты сверх retain_last - на свежей таблице без выдержанной по времени истории процедура корректно не удаляет ничего, даже если формальное число снапшотов (43) кажется избыточным. На реальной production-таблице, где эти 43 снапшота накопились за недели, а не за минуты, тот же вызов уже удалил бы все снапшоты старше трёхдневного барьера, оставив не менее 5 последних.
Кейс 4: симуляция orphan file и безопасное удаление через dry-run¶
Чтобы воспроизвести появление файла-призрака, не дожидаясь реального сбоя executor'а, запишем Parquet-файл прямо в физическое расположение таблицы через boto3, минуя Iceberg-коммит полностью - именно так выглядит результат прерванной на середине Spark-джобы:
import boto3
import pyarrow as pa
import pyarrow.parquet as pq
import io
table_location = spark.sql(
"SELECT * FROM lakehouse.analytics.orders.metadata_log_entries LIMIT 0"
).inputFiles()[0].rsplit("/metadata/", 1)[0]
# table_location = 's3a://lakehouse/warehouse/analytics/orders'
s3 = boto3.client(
"s3",
endpoint_url="http://minio:9000",
aws_access_key_id="minioadmin",
aws_secret_access_key="minioadmin",
)
orphan_table = pa.table({"order_id": [999999], "note": ["leftover from a crashed executor"]})
buf = io.BytesIO()
pq.write_table(orphan_table, buf)
s3.put_object(
Bucket="lakehouse",
Key="warehouse/analytics/orders/data/orphan_attempt_447.parquet",
Body=buf.getvalue(),
)
print("Orphan-файл записан напрямую в MinIO, минуя коммит Iceberg")
Сначала запустим процедуру в режиме dry_run, чтобы безопасно увидеть кандидатов на удаление без какого-либо реального изменения бакета:
dry_run_result = spark.sql("""
CALL lakehouse.system.remove_orphan_files(
table => 'lakehouse.analytics.orders',
older_than => TIMESTAMP_SUB(current_timestamp(), INTERVAL 1 HOURS),
dry_run => true
)
""")
dry_run_result.show(truncate=False)
# +----------------------------------------------------------------+
# |orphan_file_location |
# +----------------------------------------------------------------+
# |s3a://lakehouse/warehouse/analytics/orders/data/orphan_attempt_447.parquet|
# +----------------------------------------------------------------+
Ровно один файл-кандидат, и это именно тот файл, что был записан напрямую через boto3 в обход Iceberg. Убедившись, что список кандидатов корректен (в нём нет файлов, принадлежащих какой-либо параллельно выполняющейся легитимной джобе), повторим вызов без dry_run для реального удаления:
real_result = spark.sql("""
CALL lakehouse.system.remove_orphan_files(
table => 'lakehouse.analytics.orders',
older_than => TIMESTAMP_SUB(current_timestamp(), INTERVAL 1 HOURS)
)
""")
real_result.show(truncate=False)
# Тот же список, но файл на этот раз физически удалён из MinIO
Этот четырёхчастный демо-блок последовательно прошёл через все источники технического долга, разобранные в теоретической части урока: накопление мелких файлов и его устранение через rewrite_data_files (Кейсы 0-2), накопление истории снапшотов и безопасную, учитывающую retain_last очистку через expire_snapshots (Кейс 3), и, наконец, физически смоделированный orphan file, безопасно найденный и удалённый через remove_orphan_files с предварительной проверкой через dry_run (Кейс 4).
Производственный кейс: гонка между expire_snapshots и пятнадцатиминутным ML-отчётом¶
Ситуация. Команда платформы данных настроила ежедневный maintenance DAG (паттерн, который этот урок разберёт детально в следующем разделе), запускающий expire_snapshots на таблице lakehouse.analytics.orders каждую ночь в 02:00 с настройкой older_than => now() - INTERVAL 6 HOURS - то есть удаляющий почти всю историю практически сразу после того, как она перестаёт быть текущей. Решение об именно такой агрессивной настройке было принято с благой целью: команда хранения жаловалась на рост счёта за дисковое пространство MinIO, а аналитический отдел подтвердил, что Time Travel глубже нескольких часов им на практике не требуется. Maintenance DAG работал без видимых проблем три недели.
Инцидент. В одну из ночей аналитик из соседней команды запустил вручную, в обход обычного расписания, тяжёлый ML feature engineering job - агрегацию по полной истории заказов за последние два года, ожидаемо занимающую около 15 минут. Запрос начал выполнение в 01:58, зафиксировав на этапе Scan Planning конкретный снапшот таблицы (механика разобрана в третьем уроке) и получив список файлов для чтения. В 02:00, как и каждую ночь, запустился maintenance DAG. Поскольку с момента старта ML-джобы прошло меньше шести часов, older_than => now() - INTERVAL 6 HOURS формально не должен был задеть снапшот, против которого выполнялся запрос аналитика, - но задел: в течение этой ночи произошло ещё несколько обычных коммитов из streaming-пайплайна заказов, и снапшот, зафиксированный ML-джобой в 01:58, оказался уже не текущим, а на 6+ часов более старым относительно момента запуска expire_snapshots, поскольку отметка времени снапшота - это момент его коммита, а не момент, когда конкретный читатель начал против него запрос. В 02:04 файлы, принадлежащие исключительно этому снапшоту (а не более новым), были безопасно, по мнению самой процедуры, удалены - притом что ML-джоба, запущенная в 01:58, продолжала их активно читать.
Расследование. Диагностика заняла больше времени, чем хотелось бы команде, именно потому что ошибка FileNotFoundException указывала на конкретный путь файла в MinIO, а не на причину его отсутствия - инженеры сначала проверили права доступа и сетевую связность с MinIO, прежде чем сопоставить точное время сбоя с логами maintenance DAG и обнаружить одновременный запуск expire_snapshots. Решающей подсказкой оказалось именно совпадение по времени: сбой произошёл ровно через 4 минуты после планового запуска ночного maintenance, что для команды, уже знакомой с материалом девятого урока про горизонт хранения, стало достаточным основанием для гипотезы о гонке между чтением и expire_snapshots.
Ключевая ошибка конфигурации. older_than => now() - INTERVAL 6 HOURS был выбран на основе того, насколько глубокий Time Travel нужен бизнесу (раздел про горизонт хранения девятого урока), но не на основе того, сколько максимально может длиться самый долгий легитимный читатель таблицы. Это - принципиально разные величины, и production-конфигурация older_than обязана учитывать вторую, а не только первую: барьер устаревания снапшота должен быть заведомо больше максимальной ожидаемой длительности любого запроса или джобы, которая может против него выполняться, иначе любой долгий запрос, начавшийся незадолго до истечения барьера, рискует попасть именно в это окно гонки.
Решение.
-- Барьер увеличен с 6 часов до 48 часов - заведомо больше
-- максимальной ожидаемой длительности любого аналитического job'а
-- против этой таблицы, по согласованию с командой аналитики
CALL lakehouse.system.expire_snapshots(
table => 'lakehouse.analytics.orders',
older_than => TIMESTAMP_SUB(current_timestamp(), INTERVAL 48 HOURS),
retain_last => 10
);
Помимо увеличения самого барьера, команда добавила вторую, организационную меру: maintenance DAG теперь публикует событие в общий канал перед запуском expire_snapshots, и команды, планирующие нестандартные тяжёлые ad hoc запросы вне обычного расписания BI-дашбордов, получили соглашение - явно предупреждать платформенную команду о таких запусках, чтобы при необходимости отложить плановое обслуживание на конкретную ночь.
| Параметр | До инцидента | После инцидента |
|---|---|---|
older_than |
now() - 6 HOURS |
now() - 48 HOURS |
| Учёт длительности ad hoc джоб | Нет | Да, явно зафиксировано в runbook |
| Уведомление перед запуском maintenance | Нет | Да, в общий канал команды данных |
retain_last |
Не задан явно (использовался дефолт) | 10, явно зафиксировано в TBLPROPERTIES |
Ключевой вывод. expire_snapshots сам по себе работает корректно и именно так, как описано в разделе про reachability-анализ этого урока - инцидент не был багом процедуры, а был результатом того, что политика устаревания снапшотов проектировалась с расчётом только на бизнес-требования к Time Travel, без учёта операционного требования «барьер должен быть больше длительности самого долгого легитимного читателя». Это прямое практическое следствие материала про критическую опасность агрессивного порога, разобранного в разделе про expire_snapshots, и оно прямо формирует один из пунктов производственного чек-листа этого урока.
Сервисный дизайн: Maintenance DAG в Airflow¶
Почему обслуживание не должно быть одной большой задачей¶
Три процедуры этого урока решают непересекающиеся проблемы (как явно зафиксировано в сравнительной таблице второго раздела), и у них принципиально разная стоимость и разная безопасная частота запуска. Бин-пак компакция дешева и может безопасно выполняться часто, против только что записанных партиций. Тяжёлый Z-order по всей таблице - дорогостоящая операция, разумная раз в неделю. expire_snapshots - лёгкая по вычислениям, но рискованная относительно гонки с читателями операция. remove_orphan_files - тяжёлая по нагрузке на object storage из-за полного LIST, и её разумно запускать ещё реже. Объединение всех четырёх задач в один монолитный шаг привело бы к тому, что самая дешёвая и самая частая операция (bin-pack свежих партиций) была бы искусственно привязана к расписанию самой дорогой и редкой (remove_orphan_files) - именно поэтому production-практика выделяет каждую процедуру в собственную задачу Airflow DAG с собственным расписанием и собственными ретраями.
Диаграмма отражает три самостоятельных принципа проектирования расписания. Во-первых, частота обратно пропорциональна стоимости: лёгкий ежедневный bin-pack компактит только вчерашние партиции (через where, разобранный в третьем разделе), тогда как тяжёлый Z-order по всей таблице запускается лишь раз в неделю, в выходные, когда аналитическая нагрузка минимальна. Во-вторых, expire_snapshots логически идёт сразу после ежедневной компакции, а не до неё - свежий bin-pack создаёт новый снапшот, и выполнение expire_snapshots сразу после него не мешает компакции (она работает независимо от истории), но зато гарантирует, что к моменту еженедельного remove_orphan_files история уже максимально подчищена, что, как отмечено в разделе про orphan files, снижает стоимость построения reachability-множества для сравнения с листингом бакета. В-третьих, remove_orphan_files всегда проходит через dry_run-проверку перед реальным удалением - на диаграмме это явно обозначено как ручной шаг проверки между двумя вызовами процедуры, а не как полностью автоматизированный конвейер без присмотра человека.
Скелет DAG на Airflow¶
from airflow import DAG
from airflow.providers.apache.spark.operators.spark_sql import SparkSqlOperator
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
default_args = {
"owner": "data-platform",
"retries": 2,
"retry_delay": timedelta(minutes=10),
}
with DAG(
dag_id="iceberg_orders_maintenance",
schedule="0 2 * * *", # ежедневно в 02:00
start_date=datetime(2025, 1, 1),
catchup=False,
default_args=default_args,
tags=["iceberg", "maintenance"],
) as dag:
daily_compact = SparkSqlOperator(
task_id="rewrite_data_files_binpack",
sql="""
CALL lakehouse.system.rewrite_data_files(
table => 'lakehouse.analytics.orders',
strategy => 'binpack',
options => map('target-file-size-bytes', '536870912'),
where => 'order_ts >= current_date() - interval 1 days'
)
""",
)
daily_expire = SparkSqlOperator(
task_id="expire_snapshots",
sql="""
CALL lakehouse.system.expire_snapshots(
table => 'lakehouse.analytics.orders',
older_than => TIMESTAMP_SUB(current_timestamp(), INTERVAL 3 DAYS),
retain_last => 10
)
""",
)
def run_health_check(**context):
# подробности - в подразделе про health-check ниже
...
health_check = PythonOperator(
task_id="health_check",
python_callable=run_health_check,
)
daily_compact >> daily_expire >> health_check
Тяжёлые еженедельные задачи (zorder-компакция всей таблицы и remove_orphan_files) разумно вынести в отдельный DAG с собственным schedule="0 3 * * 6" (раз в неделю, в ночь на субботу) - так оба DAG'а получают независимые расписания, независимые SLA на длительность выполнения и независимые алерты при сбое, не блокируя друг друга в общем графе зависимостей.
Health-check через системные таблицы: что именно проверять после каждого запуска¶
Здравая практика - не доверять обслуживанию вслепую, а после каждого запуска проверять через уже знакомые системные таблицы Iceberg (.files, .snapshots, .history, встречавшиеся в материале и домашних заданиях предыдущих уроков модуля), что результат действительно соответствует ожиданиям:
def run_health_check(**context):
file_stats = spark.sql("""
SELECT count(*) AS file_count,
avg(file_size_in_bytes) AS avg_size,
sum(CASE WHEN file_size_in_bytes < 16 * 1024 * 1024 THEN 1 ELSE 0 END) AS tiny_files
FROM lakehouse.analytics.orders.files
""").collect()[0]
snapshot_count = spark.sql(
"SELECT count(*) AS cnt FROM lakehouse.analytics.orders.snapshots"
).collect()[0]["cnt"]
if file_stats["tiny_files"] > 200:
raise ValueError(
f"Health-check FAILED: {file_stats['tiny_files']} файлов меньше 16 МБ "
f"после компакции - возможно, bin-pack пропустил часть партиций"
)
if snapshot_count > 50:
raise ValueError(
f"Health-check FAILED: {snapshot_count} живых снапшотов - "
f"expire_snapshots, возможно, не выполнился или retain_last слишком велик"
)
print(f"Health-check OK: {file_stats['file_count']} файлов, "
f"{snapshot_count} снапшотов, средний размер файла "
f"{file_stats['avg_size'] / 1024**2:.1f} МБ")
Эта проверка превращает обслуживание из «запустили и забыли» в наблюдаемый процесс: если tiny_files после очередного ночного запуска неожиданно растёт, это первый сигнал о том, что bin-pack компакция перестала покрывать какую-то часть активно пишущихся партиций (например, из-за изменения паттерна записи upstream-пайплайна), а если snapshot_count растёт без ограничения - явный сигнал о том, что expire_snapshots либо не запускается по расписанию, либо его retain_last/older_than были изменены кем-то без согласования. Оба сигнала выбрасывают исключение, которое штатно проваливает Airflow-задачу и инициирует алерт - а не молча проходят незамеченными до следующего production-инцидента, аналогичного разобранному выше.
Типичные заблуждения¶
«rewrite_data_files уменьшает число снапшотов и освобождает место, занятое старыми версиями файлов» - неверно, и это, возможно, самое частое заблуждение этого урока. Сравнительная таблица второго раздела явно зафиксировала: компакция создаёт новый снапшот поверх существующей истории, а не заменяет и не удаляет старые снапшоты и файлы, которые они держат живыми. После любой, даже идеально проведённой компакции старые мелкие файлы продолжают физически существовать ровно до тех пор, пока их не вычистит отдельная процедура expire_snapshots. Без неё компакция в чистом виде только добавляет объём на диске - новые крупные файлы поверх старых мелких, ещё не удалённых.
«expire_snapshots и remove_orphan_files решают одну и ту же задачу очистки мусора, можно вызывать любую из них» - неверно. Раздел про remove_orphan_files явно показал: expire_snapshots оперирует целиком на уровне метаданных (reachability-анализ по графу манифестов) и никогда не выполняет LIST физического хранилища, поэтому она органически не способна найти файлы, которые вообще никогда не были зарегистрированы ни в одном манифесте - именно для этого случая нужна принципиально иная по механике remove_orphan_files. Симметрично, remove_orphan_files не трогает снапшоты и не сокращает историю - она удаляет только файлы, недостижимые из всех существующих снапшотов, включая самые старые.
«Чем агрессивнее older_than в expire_snapshots, тем лучше - больше свободного места, меньше риска» - опасное заблуждение, прямо разобранное в производственном кейсе этого урока. Слишком короткий горизонт older_than не предотвращает риск, а создаёт его: любой легитимный читатель, чей запрос длится дольше выбранного барьера или начался незадолго до его истечения, рискует получить FileNotFoundException посреди уже выполняющегося запроса. Правильная настройка older_than - компромисс между стоимостью хранения истории и максимальной ожидаемой длительностью самого долгого легитимного потребителя таблицы, а не стремление к минимально возможному значению.
«Z-order всегда лучше Sort, потому что даёт пользу сразу нескольким колонкам» - неточно. Раздел про стратегии компакции явно показал обратную сторону: Z-order даёт лишь умеренную, а не полную селективность по каждой из участвующих колонок, тогда как Sort даёт полную селективность по своей первой (ведущей) колонке. Если аналитическая нагрузка таблицы предсказуемо и преимущественно фильтрует по одной конкретной колонке, обычный Sort по этой колонке окажется эффективнее и заметно дешевле по CPU, чем Z-order по той же колонке вместе с дополнительными.
«remove_orphan_files безопасно запускать в любой момент, опасность долгих транзакций касается только expire_snapshots» - неверно. Раздел про параметризацию remove_orphan_files показал симметричный, но направленный в обратную сторону риск: без достаточного запаса по older_than процедура способна удалить файл, который прямо сейчас дописывается активным writer'ом, чей коммит ещё не состоялся - то есть гонка существует для обеих процедур, только с разных сторон временной оси (читатель в прошлом для expire_snapshots, писатель в будущем для remove_orphan_files).
«rewrite_position_delete_files можно использовать вместо rewrite_data_files для любой Merge-on-Read таблицы с разросшимися delete-файлами» - неверно. Раздел про точечную компакцию delete-файлов явно ограничил применимость этой процедуры: она работает только с position delete files (адресующими конкретные строки по координатам файл+offset), но не способна компактировать equality delete files, описывающие условие, а не координаты - для последних единственный путь - полная rewrite_data_files.
«Если maintenance DAG ни разу не упал, значит обслуживание работает правильно» - опасное заблуждение, прямо опровергаемое разделом про health-check. Отсутствие явных ошибок выполнения процедуры не гарантирует, что её параметры действительно соответствуют реальным потребностям таблицы - DAG может годами «успешно» выполняться с retain_last, заданным по умолчанию, или с where-фильтром bin-pack-компакции, давно не покрывающим актуальные партиции после изменения паттерна записи upstream-пайплайна, без единого явного сбоя. Именно поэтому раздел про сервисный дизайн явно вводит активный health-check через системные таблицы как обязательный, а не опциональный шаг каждого запуска.
Производственный чек-лист¶
-
Каждая из трёх процедур запланирована с частотой, обратно пропорциональной её стоимости, а не объединена в одну монолитную задачу: легковесный bin-pack по свежим партициям - ежедневно или чаще, тяжёлый Z-order по всей таблице и
remove_orphan_filesс полнымLISTбакета - в окне минимальной нагрузки раз в неделю, по образцу диаграммы из раздела про Maintenance DAG. -
older_thanдляexpire_snapshotsподобран с учётом максимальной ожидаемой длительности самого долгого легитимного читателя таблицы, а не только бизнес-требований к глубине Time Travel - именно несоблюдение этого правила привело к разобранному в этом уроке производственному инциденту. -
older_thanдляremove_orphan_filesоставлен с достаточным запасом (типично - не меньше нескольких дней), заведомо большим, чем длительность любой реально выполняющейся в кластере джобы записи, чтобы избежать удаления файлов незавершённого, но легитимного коммита. -
Первый запуск
remove_orphan_filesна новой production-таблице выполнен сdry_run => true, а результат визуально проверен инженером перед реальным удалением - точечная, но решающая мера предотвращения непредвиденного удаления легитимных файлов с нестандартными путями записи. -
where-фильтр частой bin-pack-компакции синхронизирован с реальным паттерном записи upstream-пайплайна и пересматривается при изменении этого паттерна - устаревший фильтр способен годами «успешно» выполняться, не покрывая фактически растущие новой мелкими файлами партиции, что обнаруживается только через health-check, а не через сам факт успешного завершения задачи. -
Health-check через системные таблицы (
.files,.snapshots) встроен в DAG как обязательный шаг с порогами для алертинга, а не добавлен post-factum после первого незамеченного инцидента - число «слишком мелких» файлов и общее число живых снапшотов должны быть видимы в мониторинге наравне с любой другой production-метрикой. -
Порядок процедур в расписании учитывает их взаимные зависимости по стоимости:
expire_snapshotsлогически предшествуетremove_orphan_files, поскольку подчищенная история снижает объём reachability-множества, против которого сравнивается листинг бакета - и этот порядок прямо влияет на то, насколько хорошо следующий урок модуля сможет провести параллель с аналогичной по духу, но иначе устроенной service-моделью Delta Lake.
Мостик к следующим урокам¶
Этот урок завершил линию материала шестого-девятого уроков модуля, показав, как технический долг, накопленный обычными операциями записи и read-only возможностью Time Travel, систематически устраняется тремя самостоятельными процедурами обслуживания. Следующие уроки модуля строят непосредственно на этом материале.
-
Урок 11 («Delta Lake: transaction log») переключит рассмотрение на альтернативный table format, решающий ровно те же архитектурные задачи - атомарный коммит, изоляция снапшотов, история версий - но через принципиально иную физическую структуру: последовательность JSON/Parquet-файлов транзакций (Delta Log) вместо дерева манифестов Iceberg. Понимание того, зачем Iceberg нужен
expire_snapshots(ростmetadata.jsonи манифестов), напрямую помогает понять, почему Delta Log тоже требует периодической компакции своей собственной служебной структуры - чек-пойнтов транзакционного лога. -
Урок 12 («Delta Lake: OPTIMIZE, VACUUM, Z-ordering») - прямой структурный аналог этого урока для другого формата. Команда
OPTIMIZEсоответствуетrewrite_data_files(включая собственный синтаксис Z-order), аVACUUMрешает задачу, объединяющую в себе элементы иexpire_snapshots, иremove_orphan_files- физическое удаление файлов, не входящих в окно retention, заданноеdelta.deletedFileRetentionDuration. Зная детальную механику трёх процедур Iceberg из этого урока, читатель сможет провести точное сопоставление, а не восприниматьVACUUMкак чёрный ящик с похожим названием. -
Урок 13 («Выбор формата») включит зрелость, гибкость параметризации и операционную безопасность инструментов обслуживания как один из критериев сравнительной матрицы Iceberg vs Delta Lake vs Apache Hudi - в частности, то, насколько тонко каждый формат позволяет управлять стратегиями компакции (bin-pack/sort/z-order у Iceberg против
OPTIMIZE/Z-order у Delta против compaction service у Hudi) и как по-разному каждый формат балансирует риск гонки между обслуживанием и активными читателями, разобранный в производственном кейсе этого урока.
Домашнее задание¶
-
Воспроизведите Кейсы 0-2 практического блока на собственном self-hosted стенде: создайте таблицу через 60+ мелких append-коммитов, замерьте planning time и число файлов через
.filesдо компакции, выполнитеrewrite_data_filesсо стратегиейbinpack, затем отдельно - сsortпо колонке, использованной в тестовом фильтрующем запросе, и письменно сравните итоговое время выполнения одного и того же запроса во всех трёх состояниях таблицы. -
Эмпирически сравните три стратегии компакции (
binpack,sort,zorder) на одной и той же таблице с тремя независимыми колонками фильтрации в аналитической нагрузке: создайте три копии таблицы, скомпактируйте каждую отдельной стратегией, выполните набор запросов с фильтрами по разным комбинациям колонок и постройте сравнительную таблицу времени выполнения. Письменно объясните результат через материал про умеренную против полной селективности из раздела про стратегии компакции. -
Воспроизведите производственный кейс race condition этого урока локально: запустите длинный (искусственно замедленный через
time.sleepмежду чтением партиций) аналитический запрос в одном потоке/процессе и параллельно -expire_snapshotsс агрессивнымolder_thanв другом, зафиксировав факт полученияFileNotFoundException. Затем повторите эксперимент с безопаснымolder_than, заведомо большим длительности запроса, и подтвердите, что ошибка не возникает. -
Спроектируйте и реализуйте health-check функцию по образцу раздела про сервисный дизайн, но с дополнительной проверкой: средний размер файла после bin-pack-компакции должен укладываться в диапазон 80-120% от
target-file-size-bytes- и явно выбрасывайте предупреждение, если это не так, с указанием возможных причин (например, партиция содержит меньше данных, чем целевой размер одного файла). -
Симулируйте orphan file через прямую запись в MinIO по образцу Кейса 4, но усложните сценарий: запишите 5 файлов с разным временем последней модификации (
LastModified), часть - старше суток, часть - моложе часа. Выполнитеremove_orphan_filesсolder_than => now() - INTERVAL 1 HOURSи письменно объясните, почему часть файлов корректно не попадёт в список кандидатов на удаление. -
Постройте полный Airflow DAG по образцу скелета из раздела про сервисный дизайн, включающий обе ветки расписания (ежедневную легковесную и еженедельную тяжёлую), и добейтесь того, чтобы health-check-задача действительно проваливала DAG при намеренно испорченных входных условиях (например, временно отключённом ежедневном
expire_snapshotsи накопленных 60+ снапшотах). -
Исследуйте на практике порядок зависимости стоимости
remove_orphan_filesот предварительного запускаexpire_snapshots: замерьте время выполненияremove_orphan_filesна таблице с большой неподчищенной историей снапшотов (отключивexpire_snapshotsна несколько десятков коммитов) и на той же таблице сразу послеexpire_snapshots, и письменно объясните разницу через материал раздела про алгоритмremove_orphan_files. -
Напишите unit-тест (PySpark +
pytest, локальный Iceberg-каталог), который создаёт таблицу, выполняет компакцию черезrewrite_data_files, и проверяет инвариант: данные, прочитанные после компакции обычнымSELECT, идентичны (exceptAllв обе стороны, ожидая пустой результат) данным, прочитанным до компакции через сохранённыйsnapshot_idпредыдущего состояния. -
Сравните письменно (1-2 абзаца на каждую пару) три пары процедур, изученных в этом и предыдущих уроках модуля, по критерию «что произойдёт, если процедуру никогда не запускать»:
rewrite_data_filesvs накопление мелких файлов,expire_snapshotsvs неограниченный ростmetadata.jsonи стоимости хранения,remove_orphan_filesvs неограниченный рост «мёртвого» объёма в object storage, не видимого через сами Iceberg-метаданные. -
Спроектируйте письменно конкретные пороговые значения для алертинга трёх метрик health-check (доля файлов меньше целевого размера, число живых снапшотов, возраст самого старого живого снапшота) для гипотетической production-таблицы с ежедневным объёмом записи 50 ГБ и требованием к Time Travel глубиной не менее 7 дней, и обоснуйте выбор каждого порога через материал этого и девятого уроков модуля.
Полная картина: жизненный цикл технического долга от записи до чистой таблицы¶
Завершая урок, соберём весь материал в единую диаграмму - от обычной операции записи, порождающей все три параллельных источника технического долга, до момента, когда регулярный Maintenance DAG возвращает таблицу в здоровое состояние, замыкая цикл обслуживания.
Диаграмма проходит ровно тот путь, который последовательно разобрал весь урок. Верхняя часть (WRITE → NEWFILES → три ветки DEBT1/DEBT2/DEBT3) - это раздел про энтропию Lakehouse: одна и та же обычная операция записи одновременно порождает все три независимых источника технического долга, причём DEBT3 (orphan files) отмечен пунктирной связью, поскольку возникает не всегда, а только при сбое в узком временном окне между физической записью и коммитом. Средняя часть показывает три процедуры обслуживания, каждая нацеленная ровно на свой источник долга - связь «один к одному» здесь подчёркивает то, что было явно зафиксировано в сравнительной таблице второго раздела: ни одна из процедур не заменяет другую. Блок DAG собирает все три процедуры в единый, но внутренне разделённый по расписанию сервис, разобранный в разделе про сервисный дизайн, с обязательным health-check на выходе. Блок RACE, связанный пунктирными линиями с EXPIRE и ORPHAN, - напоминание о риске, разобранном в производственном кейсе и в разделах про параметризацию обеих процедур: обслуживание само способно стать источником инцидента при неосторожной настройке временных барьеров. Замыкающий диаграмму блок HEALTHY - целевое состояние, к которому регулярно, а не разово, возвращает таблицу правильно спроектированный Maintenance DAG.
Итоги¶
Любая операция записи в Iceberg v2 - INSERT, UPDATE, DELETE, MERGE INTO - создаёт новые файлы и атомарно переключает указатель текущего снапшота, никогда не модифицируя и не удаляя физически уже существующие файлы. Это свойство, дающее Time Travel из девятого урока «бесплатно», одновременно и неизбежно порождает технический долг по трём независимым направлениям, которые этот урок разобрал по отдельности.
Small Files Problem - следствие частых небольших коммитов (streaming, высокая параллельность записи, точечные UPDATE/DELETE/MERGE под Merge-on-Read) - устраняется процедурой rewrite_data_files, работающей по трём стратегиям компакции: дешёвый bin-pack, более дорогой, но улучшающий селективность одной колонки sort, и наиболее дорогой, но дающий умеренную селективность сразу нескольким колонкам z-order. Выбор стратегии должен соответствовать реальному паттерну фильтрации в аналитической нагрузке, а не выбираться по принципу «новее - значит лучше».
rewrite_data_files сама по себе - обычная операция записи, создающая новый снапшот, а не заменяющая старые файлы и снапшоты на месте. Она не уменьшает число снапшотов и не освобождает место от уже устаревших версий файлов - эти задачи решает отдельная процедура, и смешивать их ответственность - частый источник ошибочных ожиданий от результата компакции.
Накопление истории снапшотов и манифестов в metadata.json устраняется процедурой expire_snapshots, работающей по принципу reachability-анализа: файл удаляется, только если он недостижим из всех остающихся живых снапшотов, а не просто потому, что превысил возраст по older_than. Параметр retain_last гарантирует минимальный «страховой запас» снапшотов независимо от их возраста.
Главный операционный риск expire_snapshots - гонка с долгими активными запросами: если барьер older_than короче длительности легитимного читателя, начавшего работу против устаревающего снапшота, процедура способна удалить файлы прямо во время их чтения, вызывая FileNotFoundException, что и было продемонстрировано в производственном кейсе этого урока. Барьер должен учитывать максимальную ожидаемую длительность самого долгого потребителя таблицы, а не только бизнес-требования к глубине Time Travel.
Orphan files - файлы, физически записанные в object storage, но никогда не зарегистрированные ни в одном манифесте из-за сбоя между записью и коммитом - устраняются процедурой remove_orphan_files, единственной из трёх, которой требуется прямой и потенциально дорогой LIST физического хранилища. Сравнение идёт против объединения файлов из ВСЕХ сохранённых снапшотов истории, а не только текущего, что делает разумным запуск expire_snapshots перед ней.
remove_orphan_files несёт симметричный, но направленный в обратную сторону временной оси риск относительно expire_snapshots: недостаточный запас по older_than способен удалить файл активно выполняющегося, но ещё не закоммиченного writer'а. Параметр dry_run - обязательная мера предосторожности перед первым реальным запуском процедуры на новой production-таблице.
Production-эксплуатация всех трёх процедур требует раздельного по стоимости и частоте расписания (Maintenance DAG), а не единой монолитной задачи, и обязательного активного health-check через системные таблицы .files/.snapshots после каждого запуска. Отсутствие явных ошибок выполнения процедуры не означает, что её параметры действительно соответствуют текущим потребностям таблицы - это можно проверить только активным мониторингом, а не молчаливым доверием к успешному статусу задачи.
Краткий глоссарий терминов урока¶
-
Table Maintenance - совокупность регулярных процедур обслуживания Iceberg-таблицы, устраняющих технический долг, неизбежно накапливаемый обычными операциями записи: компакция файлов, очистка истории снапшотов и удаление orphan files.
-
Small Files Problem - ситуация, при которой таблица состоит из избыточно большого числа физически мелких файлов относительно общего объёма данных, увеличивающая planning time, task overhead и стоимость операций объектного хранилища.
-
rewrite_data_files- хранимая процедура Spark SQL, выполняющая компакцию data-файлов таблицы по одной из трёх стратегий (bin-pack, sort, z-order); создаёт новый снапшот, не модифицируя физически старые файлы. -
File Group - независимый, как правило ограниченный партицией подмножество файлов, обрабатываемое и коммитящееся отдельно от других таких же подмножеств в рамках одного вызова
rewrite_data_files, контролируемое параметрамиpartial-progress.enabledиmax-concurrent-file-group-rewrites. -
Bin-pack - стратегия компакции по умолчанию, объединяющая существующие файлы до целевого размера без перестановки строк; самая дешёвая и предсказуемая из трёх стратегий.
-
Sort - стратегия компакции с полным шафлом и сортировкой по заданным колонкам; даёт полную селективность статистики min/max по первой (ведущей) колонке сортировки.
-
Z-order - стратегия компакции, строящая составной индекс из битов нескольких колонок (2-4); даёт умеренную селективность одновременно по всем участвующим колонкам, в отличие от полной селективности по одной колонке у Sort.
-
rewrite_position_delete_files- узкая процедура компакции, объединяющая только position delete files без затрагивания data-файлов; не способна работать с equality delete files. -
expire_snapshots- хранимая процедура Spark SQL, физически удаляющая снапшоты, манифесты и data-файлы, недостижимые из всех остающихся живых снапшотов, по критериямolder_thanиretain_last. -
retain_last- параметрexpire_snapshots, задающий минимальное число последних снапшотов, сохраняемых независимо от их возраста поolder_than. -
clean_expired_metadata- параметрexpire_snapshots, дополнительно удаляющий изmetadata.jsonзаписи о ссылках и версиях метаданных, указывающие на уже удалённые снапшоты. -
Orphan files - файлы, физически существующие в object storage, но не зарегистрированные ни в одном манифесте ни одного снапшота истории таблицы; возникают из-за сбоя между физической записью и коммитом.
-
remove_orphan_files- хранимая процедура Spark SQL, находящая и удаляющая orphan files через сравнение прямогоLISTфизического хранилища с объединением файлов, достижимых из всех снапшотов истории таблицы. -
dry_run- параметрremove_orphan_files, при значенииtrueвозвращающий список файлов-кандидатов на удаление без какого-либо реального изменения object storage. -
Maintenance DAG - оркестрируемый, как правило через Airflow, набор задач обслуживания Iceberg-таблицы с раздельными расписаниями по стоимости каждой процедуры и обязательным health-check после выполнения.
-
Health-check - этап Maintenance DAG, проверяющий через системные таблицы (
.files,.snapshots) соответствие фактического состояния таблицы ожидаемым после выполнения обслуживания порогам, и явно проваливающий задачу при их нарушении.