Delta Lake: OPTIMIZE, VACUUM, Z-ordering и auto-optimize - production-операции для managed tables
Эксплуатация Delta Lake в production: компакция мелких файлов через OPTIMIZE, многомерная кластеризация Z-ORDER BY, физическое удаление устаревших файлов через VACUUM, Optimized Write и Auto Compact, Deletion Vectors как Merge-on-Read для Delta, и проектирование Maintenance DAG в Airflow.
Десятый урок модуля разобрал три процедуры обслуживания Iceberg-таблиц - rewrite_data_files, expire_snapshots, remove_orphan_files - и зафиксировал главный тезис: транзакционность и Time Travel не достаются бесплатно, их цена - постоянно растущее число физических файлов, которое не исчезает само по себе без отдельного, явного действия. Одиннадцатый урок детально разобрал, как Delta Lake устроен внутри - _delta_log, JSON-коммиты, checkpoint’ы - и из этого устройства прямо следует, что у Delta Lake тот же класс проблем требует решения, просто под другими именами команд и с другой внутренней механикой.
Этот урок закрывает тему эксплуатации Delta Lake тем же способом, каким десятый урок закрыл тему эксплуатации Iceberg: разбирает OPTIMIZE, VACUUM, Z-ORDER BY, Optimized Write, Auto Compact и Deletion Vectors не как изолированные команды для запоминания, а как взаимосвязанные процедуры, каждая из которых решает конкретный, конкретно названный класс операционного риска. В конце урока эти процедуры собираются в единый production-регламент - Maintenance DAG в Airflow - применимый к self-hosted Lakehouse без какой-либо зависимости от управляемой платформы.
Проблема мелких файлов и уплотнение через OPTIMIZE¶
Откуда берутся мелкие файлы и почему объектное хранилище их не любит¶
Десятый урок модуля вводил термин Small Files Problem применительно к Iceberg, и природа этой проблемы у Delta Lake идентична, потому что коренится не в формате метаданных, а в самом способе, которым Spark пишет данные. Каждый микро-батч Structured Streaming, каждая параллельная задача инкрементального ETL, каждый отдельный INSERT/MERGE с малым числом строк создаёт собственный набор Parquet-файлов - и если объём данных в одной операции записи невелик, итоговые файлы получаются крошечными: счёт идёт на килобайты или единицы мегабайт вместо целевых сотен мегабайт - гигабайта.
Деградация от этого - не абстрактная угроза, а конкретный, измеримый эффект на двух разных уровнях. Во-первых, объектные хранилища (S3, MinIO, Ceph RGW) физически устроены как key-value сервисы поверх HTTP, и у каждого обращения к объекту - будь то LIST для перечисления файлов партиции или GET для скачивания конкретного файла - есть фиксированная сетевая накладная стоимость (latency на установление соединения, аутентификацию, обработку запроса), которая не зависит от размера самого объекта. Десять тысяч файлов по 50 КБ обходятся кластеру в десять тысяч отдельных HTTP-запросов, тогда как тот же объём данных в виде десяти файлов по 50 МБ - это в тысячу раз меньше обращений к хранилищу при идентичном суммарном объёме полезных данных.
Во-вторых, у каждого Parquet-файла, который Spark включает в физический план выполнения, есть собственная задача (task) в DAG приложения - и при тысячах мелких файлов планировщик Spark создаёт тысячи мелких задач, каждая из которых несёт фиксированный overhead на запуск executor-потока, сериализацию, сбор метрик и коммуникацию с драйвером. Время выполнения такой задачи может оказаться меньше, чем накладные расходы на её диспетчеризацию - то есть кластер тратит больше времени на управление работой, чем на саму работу, что прямо симметрично проблеме «накладные расходы на planning растут быстрее полезного объёма» из одиннадцатого урока этого модуля.
OPTIMIZE как процедура компакции: что происходит физически¶
Команда OPTIMIZE - основной, исторически первый инструмент Delta Lake для борьбы с этой энтропией. Механика предельно прямолинейна по сравнению с трёхстадийным rewrite_data_files Iceberg, разобранным в десятом уроке: Spark считывает множество мелких Parquet-файлов в рамках указанной таблицы (или конкретной партиции, если задан предикат), перегруппировывает строки в новые файлы целевого размера и записывает их на диск, после чего фиксирует один новый коммит в _delta_log - JSON-файл, который одной транзакцией помечает все старые мелкие файлы как remove, а все новые крупные - как add.
-- Простейший вызов: уплотнить всю таблицу целиком
OPTIMIZE delta.`s3a://lakehouse/warehouse/orders_delta`;
-- С фильтрацией по партиции - оптимизируем только свежие данные
OPTIMIZE delta.`s3a://lakehouse/warehouse/orders_delta`
WHERE order_date >= '2025-06-01';
Критически важная деталь, прямо следующая из модели _delta_log, разобранной в одиннадцатом уроке: OPTIMIZE не модифицирует существующие файлы на месте (что и невозможно для неизменяемого формата Parquet, как объяснялось в шестом уроке модуля про Copy-on-Write) и не трогает физические байты старых файлов - она лишь добавляет новый коммит, логически переключающий состояние таблицы на новый, более компактный набор файлов. Старые мелкие файлы продолжают физически существовать на диске после OPTIMIZE ровно по той же причине, по которой Iceberg не удаляет файлы физически при rewrite_data_files: чтобы не сломать Time Travel к версиям, существовавшим до компакции. Их физическое исчезновение - отдельная задача, решаемая VACUUM (раздел 3 этого урока).
Изоляция от основных пайплайнов записи и стратегия частичной оптимизации¶
OPTIMIZE - тяжёлая операция: она требует прочитать с диска и переписать заново весь объём затронутых данных, что означает полноценную нагрузку на сеть, CPU (распаковка и пересжатие Parquet) и executor-память кластера, сопоставимую по характеру с полным backfill. По этой причине OPTIMIZE принципиально нельзя запускать как часть того же Spark-job, который выполняет основную потоковую или пакетную запись в таблицу - совместное выполнение создаёт конкуренцию за ресурсы executor’ов между задачей, критичной по latency (доставка свежих данных), и задачей, критичной по throughput (переупаковка истории), и обычно проигрывает именно latency-критичная задача.
Практическая рекомендация, прямо аналогичная File Groups из десятого урока (где Iceberg разбивает компакцию по партициям, чтобы не создавать единый гигантский shuffle на всю таблицу): OPTIMIZE следует запускать как отдельный, изолированный job, с собственным расписанием, собственным пулом executor’ов, и - что особенно важно для production-таблиц с историей в сотни терабайт - с обязательной фильтрацией по партициям через WHERE. Оптимизировать партиции вчерашнего и сегодняшнего дня, в которые ещё идёт активная запись и которые чаще всего читаются аналитическими запросами, экономически оправдано почти всегда; повторно гонять OPTIMIZE по партициям трёхлетней давности, которые уже были уплотнены однажды и больше не получают новых записей, - чистая трата вычислительных ресурсов кластера без какого-либо выигрыша для текущих потребителей данных.
-- Production-паттерн: оптимизируем только последние N дней,
-- остальная история уже уплотнена предыдущими запусками
OPTIMIZE delta.`s3a://lakehouse/warehouse/orders_delta`
WHERE order_date >= date_sub(current_date(), 3);
Многомерная кластеризация: Z-ORDER BY¶
Почему обычное партиционирование не спасает при высокой кардинальности¶
Партиционирование по дате или по региону работает хорошо именно потому, что число уникальных значений таких колонок невелико - тысяча партиций по дням за три года или десяток партиций по регионам не создают проблем при листинге директорий. Но если аналитики систематически фильтруют запросы по колонкам с высокой кардинальностью - user_id с миллионами уникальных значений, device_id, session_id - наивная попытка партиционировать таблицу по такой колонке создаёт миллионы крошечных директорий-партиций, что напрямую возвращает Small Files Problem из раздела 1, только теперь как структурное, встроенное в саму схему хранения свойство таблицы, а не как побочный эффект частой записи.
Единственная альтернатива - оставить такие высококардинальные колонки вне партиционной схемы и положиться на data skipping по статистике файлов (min/max значения колонок, сохранённые в add-записях _delta_log, как было показано в одиннадцатом уроке) для отсечения нерелевантных файлов при чтении. Проблема в том, что data skipping эффективен только тогда, когда строки со схожими значениями физически сгруппированы в одних и тех же файлах - если значения user_id распределены по файлам полностью случайно (что естественно происходит, когда Spark пишет данные в порядке их появления в потоке, без какой-либо сортировки), то диапазон min..max практически любого файла будет покрывать весь спектр возможных значений, и ни один файл не получится безопасно отбросить.
Механика Z-Ordering: многомерная кластеризация без явного партиционирования¶
Z-Ordering решает именно эту задачу - физически переупорядочивает строки внутри новых файлов так, чтобы строки со схожими значениями указанных колонок оказывались в одних и тех же файлах, причём не по одной колонке, а сразу по нескольким одновременно. Название отсылает к Z-order curve (кривой Мортона) - математической конструкции, отображающей многомерное пространство значений в одномерный порядок так, что точки, близкие в многомерном пространстве, как правило остаются близкими и после линеаризации; именно этот порядок Spark использует для физической перестановки строк перед записью файлов.
OPTIMIZE delta.`s3a://lakehouse/warehouse/orders_delta`
ZORDER BY (user_id, device_id);
Выполнение этой команды совмещает обычную компакцию из раздела 1 (мелкие файлы объединяются в крупные) с дополнительным шагом перестановки: перед физической записью новых файлов Spark вычисляет Z-значение для каждой строки на основе указанных колонок и сортирует данные по этому значению, прежде чем разбивать их на финальные файлы целевого размера. В результате каждый итоговый файл получает узкий, плотный диапазон min..max сразу по обеим указанным колонкам - а не широкий, покрывающий весь спектр значений диапазон, который получился бы при случайном распределении строк.
Почему это ускоряет запросы: Data Skipping на уровне статистики, без скачивания файлов¶
Выгода от узких диапазонов min..max проявляется в момент Query Planning, разобранного в одиннадцатом уроке: когда Spark получает запрос WHERE user_id = 4200, движок сначала читает только метаданные файлов (статистику из add-записей в checkpoint’е или хвосте лога), и для каждого файла проверяет, попадает ли искомое значение в его диапазон min..max - без необходимости скачивать сам файл с S3 хотя бы для одной байты полезных данных. На диаграмме выше запрос user_id = 4200 после Z-Ordering мгновенно отбрасывает file_1 (диапазон 1..1500) и file_2 (диапазон 1501..3200), оставляя для реального чтения только file_3 - притом что без Z-Ordering все три файла содержали бы потенциально подходящие строки и не могли бы быть отброшены на основании одной только статистики.
На production-таблице с тысячами файлов этот эффект масштабируется ощутимо: типичный практический результат Z-Ordering по правильно выбранным, часто фильтруемым колонкам - отсечение 90-95% файлов ещё на этапе планирования, до единого обращения к данным. Эффект особенно заметен для ad-hoc аналитических запросов с точечными или узко-диапазонными фильтрами по высококардинальным колонкам - именно тот класс нагрузки, для которого обычное партиционирование органически не подходит.
Ограничения: почему не стоит указывать больше 3-4 колонок¶
Z-Ordering не масштабируется безгранично на произвольное число колонок, и это ограничение - прямое следствие самой математики кривой Мортона. С ростом числа измерений «близость» в многомерном пространстве всё хуже сохраняется при проецировании на одномерный порядок - эффект, известный в литературе по многомерному индексированию как «проклятие размерности» (curse of dimensionality): чем больше колонок участвует в Z-Order, тем менее плотными становятся итоговые диапазоны min..max по каждой из них, и тем слабее эффект data skipping для каждой отдельной колонки.
Практическая рекомендация, подтверждённая опытом эксплуатации Delta Lake в production, - ограничивать ZORDER BY тремя, в редких случаях четырьмя колонками, выбранными по критерию «по этой колонке аналитики действительно фильтруют запросы чаще всего» (это можно эмпирически подтвердить через анализ истории выполненных запросов в Spark History Server или внешнем query-логе). Добавление пятой, шестой колонки в ZORDER BY в попытке угодить сразу всем возможным паттернам фильтрации на практике обычно ухудшает, а не улучшает результат: эффективность кластеризации размывается настолько, что для каждой отдельной колонки она становится сопоставима с полным отсутствием Z-Ordering.
Сравнение с подходом Iceberg¶
Десятый урок модуля уже разбирал Z-order стратегию rewrite_data_files Iceberg (strategy => 'zorder') как один из трёх вариантов компакции наравне с bin-pack и sort. Идея в обоих форматах идентична - многомерная физическая кластеризация для ускорения data skipping - но оформлена по-разному: у Iceberg Z-order - один из параметров универсальной процедуры компакции, явно конфигурируемый через sort_order; у Delta Lake - отдельное, синтаксически выделенное расширение базовой команды OPTIMIZE. На уровне эксплуатации различие не носит принципиального характера: оба формата требуют той же дисциплины выбора колонок (не больше 3-4, по реальному паттерну фильтрации) и того же понимания, что Z-order - это компакция с дополнительным шагом сортировки, а не самостоятельная, отдельная от уплотнения файлов операция.
| Аспект | Delta Lake | Apache Iceberg |
|---|---|---|
| Синтаксис | OPTIMIZE table ZORDER BY (col1, col2) |
CALL ...rewrite_data_files(strategy => 'zorder', sort_order => 'col1,col2') |
| Совмещение с компакцией | Встроено - один и тот же вызов | Один и тот же вызов с параметром strategy |
| Рекомендуемое число колонок | 3-4 | 3-4 (та же математика curve, то же ограничение) |
| Хранение статистики результата | min/max в add-записи _delta_log |
min/max в Manifest File |
Куда движется эволюция: краткий взгляд на Liquid Clustering¶
OPTIMIZE ... ZORDER BY - не последняя точка эволюции многомерной кластеризации у Delta Lake. Более новая возможность, Liquid Clustering (объявляемая на уровне таблицы через CLUSTER BY (col1, col2) при создании таблицы, а не как параметр отдельного вызова OPTIMIZE), решает практическую проблему, оставшуюся за рамками классического ZORDER BY: набор колонок для Z-Order жёстко фиксируется в момент каждого вызова OPTIMIZE и никак не запоминается таблицей между вызовами, поэтому смена набора «горячих» колонок фильтрации требует ручного отслеживания и согласованного повторения одних и тех же аргументов командой эксплуатации при каждом следующем запуске. Liquid Clustering делает выбор кластеризующих колонок свойством самой таблицы, инкрементально поддерживаемым при каждой последующей записи и каждой компакции, без необходимости полного пересчёта кластеризации с нуля при добавлении новых данных.
Тринадцатый, завершающий урок модуля приводит конкретный производственный кейс, в котором именно эта более новая возможность создала проблему: Trino-коннектор определённой версии не умел интерпретировать метаданные Liquid Clustering и из-за этого деградировал до полного сканирования партиции вместо кластеризованного data skipping. Этот кейс - не случайность, а системная иллюстрация общего принципа, уже зафиксированного в производственном чек-листе этого урока: любая возможность формата, появившаяся позже базового протокола (Deletion Vectors, Liquid Clustering, Column Mapping из одиннадцатого урока), требует явной проверки поддержки во всех движках, фактически читающих таблицу, а не только в том движке, который её записывает.
Физическое уничтожение данных и очистка диска с помощью VACUUM¶
Логическое удаление - это не физическое удаление¶
Один из самых частых источников путаницы у инженеров, впервые эксплуатирующих Delta Lake в production, - предположение, что OPTIMIZE, UPDATE, DELETE или MERGE сами по себе освобождают место на диске. Одиннадцатый урок прямо показал механику: remove-запись в _delta_log помечает файл как логически удалённый из текущего состояния таблицы, но не удаляет ни единого байта с физического хранилища. Это - не недосмотр архитектуры, а намеренное, обязательное условие работы Time Travel: запрос VERSION AS OF 40 должен суметь прочитать файлы, актуальные на версии 40, даже если к версии 45 эти файлы давно вышли из текущего состояния таблицы через remove.
Без какого-либо регулярного процесса очистки итог предсказуем: объём физически занятого пространства на S3/MinIO растёт монотонно и без верхней границы, потому что каждая операция UPDATE/DELETE/MERGE/OPTIMIZE добавляет новые файлы, но никогда не убирает старые. Для таблицы с интенсивной потоковой записью и частыми точечными обновлениями это означает, что физический объём хранилища может в несколько раз превышать объём данных, реально нужных для чтения текущей версии таблицы - а на платформах с тарификацией по фактически занятому объёму object storage это прямая, измеримая статья расходов, которая не видна ни в одном бизнес-дашборде до тех пор, пока кто-то не сверит счёт от облачного провайдера с ожидаемым объёмом данных.
VACUUM: санитар, работающий по порогу возраста¶
Команда VACUUM решает эту задачу - физически удаляет с диска файлы данных, которые одновременно удовлетворяют двум условиям: (1) файл не входит в множество живых файлов ни одной версии таблицы новее порога удержания, и (2) возраст файла (по времени последней модификации) превышает этот порог. По умолчанию порог удержания (Retention Threshold) равен 168 часам - ровно семь дней, и это значение выбрано не произвольно: оно должно с запасом перекрывать типичную продолжительность долгих аналитических job’ов и Time Travel запросов, которые могут быть запущены против таблицы и оставаться активными в момент срабатывания VACUUM.
-- Стандартный безопасный вызов - используется retention по умолчанию (168 часов)
VACUUM delta.`s3a://lakehouse/warehouse/orders_delta`;
-- Явное указание порога удержания
VACUUM delta.`s3a://lakehouse/warehouse/orders_delta` RETAIN 168 HOURS;
-- Dry run: посмотреть, какие файлы будут удалены, без реального удаления
VACUUM delta.`s3a://lakehouse/warehouse/orders_delta` RETAIN 168 HOURS DRY RUN;
Внутренний алгоритм VACUUM концептуально похож на reachability-анализ expire_snapshots у Iceberg, разобранный в десятом уроке, хотя оперирует более простой структурой данных: Delta Lake перечисляет все файлы, физически присутствующие в директории таблицы (полный LIST по объектному хранилищу), вычитает из этого множества файлы, которые остаются «живыми» (присутствуют как add без последующего remove) в истории _delta_log не старше порога удержания, и удаляет всё, что осталось - то есть файлы, чей remove произошёл раньше границы retention, и которые поэтому больше не нужны ни одной версии таблицы, доступной для Time Travel.
Опасность агрессивного снижения порога удержания¶
Снижение порога удержания ниже значения, реально нужного для обслуживания самых долгих читателей таблицы, - один из самых разрушительных типов операционной ошибки во всей экосистеме Delta Lake, потому что его последствия проявляются не сразу и не локально: они бьют по job’ам, которые в момент запуска VACUUM ни сном ни духом не подозревают о его существовании. По умолчанию Delta Lake вообще не позволяет снизить порог ниже разумного минимума без явного, осознанного шага - попытка вызвать VACUUM ... RETAIN 0 HOURS напрямую завершится ошибкой защиты, если предварительно не отключить проверку через spark.databricks.delta.retentionDurationCheck.enabled = false.
-- Явное, осознанное отключение защиты - использовать только понимая последствия
SET spark.databricks.delta.retentionDurationCheck.enabled = false;
VACUUM delta.`s3a://lakehouse/warehouse/orders_delta` RETAIN 0 HOURS;
Если эта защита отключена и VACUUM запущен с нулевым или крайне малым retention, конкретный сценарий отказа выглядит так: аналитический Spark-job уже прочитал список файлов, относящихся к версии 50 таблицы, и приступил к их фактическому скачиванию и обработке - в этот момент параллельно срабатывает VACUUM с агрессивным порогом, обнаруживает, что версия 50 старше нового, заниженного retention, и физически удаляет файлы, которые, как ему кажется, больше никому не нужны. Job, уже держащий список путей к этим файлам, получает FileNotFoundException при попытке прочитать конкретный файл, которого уже физически не существует - падение job’а, которое выглядит как случайная, необъяснимая инфраструктурная нестабильность, если не знать о только что отработавшем VACUUM в логах соседнего DAG.
Отдельный, не связанный с безопасностью, параметр spark.databricks.delta.vacuum.parallelDelete.enabled управляет не порогом удержания, а тем, выполняется ли само физическое удаление файлов параллельно силами executor’ов (true) или последовательно одним потоком драйвера (значение по умолчанию false); включение параллельного удаления ускоряет VACUUM на таблицах с большим числом файлов-кандидатов на удаление, но не имеет никакого отношения к риску гонки с активными читателями - этот риск контролируется исключительно величиной retention threshold и retentionDurationCheck.enabled, и их не стоит путать между собой при настройке.
Проектирование безопасного окна для VACUUM¶
Практический вывод для production-эксплуатации: безопасный порог удержания должен определяться не техническим минимумом, который соглашается принять Delta Lake, а максимальной реальной продолжительностью долгих сессий, которые организация допускает против конкретной таблицы. Если в компании принято, что ML-инженеры могут держать открытую интерактивную Spark-сессию с Time Travel запросом до трёх дней, retention threshold обязан с запасом перекрывать этот сценарий - стандартные 168 часов (семь дней) на практике покрывают подавляющее большинство организационных паттернов использования и редко требуют ручной корректировки в сторону уменьшения.
Снижение retention ниже значения по умолчанию оправдано в узком наборе случаев - например, для staging-таблиц с заведомо коротким жизненным циклом потребителей, где быстрая очистка диска важнее долгого Time Travel - и в любом случае требует явного согласования с командами, фактически потребляющими таблицу, а не одностороннего решения команды эксплуатации платформы. Эта же мысль будет прямо повторена в разделе про Maintenance DAG: расписание VACUUM должно учитывать не только внутренние процедуры обслуживания, но и внешних потребителей, о существовании которых платформенная команда может даже не знать.
Автоматическое управление файлами: Auto Optimize и Optimized Write¶
От реактивного обслуживания к проактивному предотвращению проблемы¶
Всё, что было разобрано в разделах 1-3, - это реактивный подход: данные сначала пишутся как есть, в том виде и с тем размером файлов, который получается естественным образом из логики job’а записи, а проблема мелких файлов устраняется отдельно, отложенно, отдельной процедурой OPTIMIZE, запускаемой по независимому расписанию. У такого подхода есть встроенный недостаток - между моментом записи мелких файлов и моментом их компакции существует временное окно, в течение которого любой запрос к таблице вынужден работать с неоптимальной физической раскладкой, а на таблицах с высокой частотой записи (потоковые источники, события) это окно может практически не закрываться, если OPTIMIZE запускается реже, чем накапливается следующая порция мелких файлов.
Delta Lake предлагает альтернативу - проактивное управление размером файлов прямо в момент записи, без необходимости ждать отдельного запуска OPTIMIZE. Эта возможность называется Auto Optimize и фактически складывается из двух независимых, но взаимодополняющих механизмов: Optimized Write, который меняет поведение самой операции записи, и Auto Compact, который запускает мини-компакцию сразу после завершения транзакции записи, если она оставила за собой мелкие файлы.
Optimized Write: адаптивное распределение перед финальной записью¶
При обычной записи Spark создаёт ровно столько выходных файлов, сколько было партиций (partitions) в Spark RDD/DataFrame непосредственно перед операцией записи - а число этих партиций определяется предшествующими широкими трансформациями (shuffle), параметром spark.sql.shuffle.partitions или особенностями источника данных, и почти никогда не подбирается осознанно под целевой размер итогового Parquet-файла. Результат - типичная ситуация, когда из-за избыточного числа входных партиций при записи небольшого объёма данных получается избыточное число мелких выходных файлов, даже если сам объём данных в одной транзакции невелик.
Optimized Write вставляет дополнительный, адаптивный шаг перед финальной записью: Spark оценивает объём данных, подлежащих записи, и динамически выполняет локальный shuffle (репартиционирование), приводя число выходных партиций в соответствие с целевым размером файла - вместо записи, скажем, двухсот партиций по 5 МБ каждая получается, например, десять партиций по 100 МБ. Эта перегруппировка происходит за счёт дополнительного, но обычно недорогого по объёму данных шага shuffle прямо в рамках того же job’а записи, то есть стоимость переносится с отложенной, тяжёлой компакции (OPTIMIZE целой партиции постфактум) на small, постоянно работающий механизм адаптации, встроенный в каждую write-операцию.
Auto Compact: мини-компакция в рамках текущей сессии¶
Optimized Write улучшает раскладку данных в рамках одной write-транзакции, но не решает проблему, накапливающуюся через множество отдельных, небольших транзакций - например, серию из тысячи микро-батчей Structured Streaming, каждый из которых после применения Optimized Write всё равно создаёт несколько в меру крупных, но всё же не идеально упакованных файлов, потому что объём данных в отдельном микро-батче слишком мал для формирования файла целевого размера в принципе. Эту проблему закрывает второй механизм - Auto Compact: после завершения каждой транзакции записи Spark проверяет, не привела ли она к появлению избыточного числа файлов меньше определённого порога размера в затронутой партиции, и если да - запускает компактную мини-версию OPTIMIZE прямо в рамках текущей driver-сессии, без отдельного job’а и отдельного запуска по расписанию.
-- Включение Auto Optimize для конкретной таблицы:
-- Optimized Write + Auto Compact одновременно
ALTER TABLE orders_delta SET TBLPROPERTIES (
delta.autoOptimize.optimizeWrite = true,
delta.autoOptimize.autoCompact = true
);
-- Альтернатива - включить как настройку SparkSession,
-- которая будет применяться по умолчанию ко всем новым записям
SET spark.databricks.delta.autoCompact.enabled = true;
SET spark.databricks.delta.autoOptimize.optimizeWrite = true;
Тradeoff: где Auto Optimize выгоден, а где вреден¶
Auto Optimize - не безусловно полезная настройка, которую стоит включать на каждой таблице по умолчанию: оба её механизма платят за улучшение физической раскладки данных дополнительными ресурсами, потраченными именно в момент записи, а не отложенно в окне обслуживания. Optimized Write добавляет дополнительный shuffle к каждой write-операции - то есть дополнительную сетевую и CPU-нагрузку именно в тот момент, когда задержка записи (write latency) часто наиболее критична, особенно для потоковых пайплайнов с жёсткими SLA на end-to-end latency. Auto Compact добавляет к каждой транзакции дополнительную проверку условий компакции и, при срабатывании, дополнительный раунд чтения-перезаписи файлов прямо внутри той же сессии, что увеличивает время коммита транзакции (commit latency) - для систем, ожидающих от Delta Lake предсказуемо короткое время отклика на запись, эта дополнительная задержка может оказаться неприемлемой.
Практическое правило, вытекающее из этого tradeoff: Auto Optimize оправдан на таблицах с высокочастотной потоковой записью небольшими порциями, где альтернатива - быстрое накопление огромного числа мелких файлов между плановыми запусками OPTIMIZE - наносит больший вред, чем небольшое увеличение latency самой записи. На таблицах с редкими, но крупными батч-загрузками (например, ежедневный полный backfill витрины), где входные данные уже естественным образом формируют крупные файлы, Auto Optimize добавляет накладные расходы без соразмерной выгоды, и предпочтительнее оставить файлы как есть, опираясь на плановый OPTIMIZE из раздела 1 для тонкой доводки раскладки уже постфактум.
Современная оптимизация MoR-таблиц: Deletion Vectors¶
Историческая отправная точка: Delta как чистый Copy-on-Write¶
Тринадцатый урок этого модуля уже вкратце упоминал Deletion Vectors при сравнении трёх форматов, но здесь стоит восстановить полную картину с самого начала, потому что именно она объясняет, почему этот механизм появился и какую именно боль он снимает. Исторически, до версии 2.3, Delta Lake реализовывал исключительно модель Copy-on-Write, разобранную в шестом уроке модуля: любая операция DELETE или UPDATE, затрагивающая хотя бы одну строку внутри Parquet-файла, требовала полностью переписать весь этот файл - прочитать его целиком, исключить или изменить нужные строки, и записать результат как совершенно новый физический файл, после чего старый файл помечался как remove в _delta_log.
Формула Write Amplification из шестого урока модуля - отношение объёма физически перезаписанных данных к объёму данных, реально затронутых логической операцией, - в применении к Delta-таблицам докрутилась до того же самого экстремального случая, который уже был разобран на примере Iceberg: UPDATE, меняющий одну строку в файле объёмом 500 МБ, вынуждал переписать все 500 МБ ради изменения нескольких десятков байт. На таблицах с частыми точечными обновлениями (GDPR-удаления по user_id, исправление отдельных дефектных записей, late-arriving corrections) этот паттерн превращался в постоянный, дорогой источник нагрузки на кластер - именно ту проблему, для которой Merge-on-Read был изобретён в принципе (седьмой урок модуля).
Deletion Vectors: переход Delta к частичному Merge-on-Read¶
Начиная с Delta Lake 2.3 (и достигнув полноценной production-зрелости к версии 3.0), формат получил собственный механизм Merge-on-Read - Deletion Vectors, включаемый табличным свойством delta.enableDeletionVectors = true. Идея зеркально повторяет Position Delete из седьмого урока модуля (про Iceberg Merge-on-Read): вместо переписывания целого файла ради удаления нескольких строк, Delta Lake записывает отдельный, маленький побочный файл, который отмечает позиции удалённых строк внутри существующего, неизменного базового Parquet-файла - и при чтении эти позиции просто пропускаются на лету.
ALTER TABLE orders_delta SET TBLPROPERTIES (delta.enableDeletionVectors = true);
DELETE FROM orders_delta WHERE status = 'cancelled' AND order_date < '2025-01-01';
Физическая форма этого побочного файла - битовая карта (bitmap), закодированная компактным алгоритмом RoaringBitmap: каждый бит карты соответствует одной позиции (порядковому номеру строки) в базовом файле, и бит установлен в 1 ровно для тех строк, которые логически удалены. RoaringBitmap эффективно сжимает как плотные, так и разрежённые диапазоны установленных битов, что делает Deletion Vector файлы крайне компактными даже для базовых файлов с миллионами строк - типичный размер такого файла измеряется килобайтами, а не сопоставим по объёму с самими данными.
Read Amplification: цена, которую платит каждый последующий запрос¶
Как и любая форма Merge-on-Read, Deletion Vectors не делают удаление бесплатным - они переносят стоимость операции с момента записи на момент каждого последующего чтения, в точности по той же логике Read Penalty, что была формализована в седьмом уроке модуля. Каждый запрос, читающий файл с присоединённым Deletion Vector, обязан дополнительно прочитать сам Deletion Vector файл и применить его как маску, исключающую помеченные позиции из результата, прежде чем строки этого файла попадут в дальнейшую обработку - небольшая, но не нулевая постоянная накладная стоимость, которая, в отличие от единоразовой стоимости Copy-on-Write перезаписи, повторяется при каждом обращении к файлу, пока Deletion Vector не будет смыт обратно в данные.
Критически важное структурное ограничение, прямо унаследованное от семантики бит-карты: Deletion Vector способен только пометить существующую строку как удалённую - он не умеет менять значения внутри строки. Поэтому UPDATE для Delta-таблицы с включёнными Deletion Vectors реализован как гибридная операция: старая версия изменяемой строки помечается удалённой через Deletion Vector (дешёвая часть, как при DELETE), а новое значение строки записывается как новая строка в новый, отдельный файл (дорогая часть, как при классическом Copy-on-Write). Это - прямое структурное отличие от Iceberg, у которого помимо position delete существует ещё equality delete (восьмой урок модуля, через MERGE INTO), позволяющий пометить строки на удаление по значению предиката без необходимости сначала их физически прочитать; у Delta Lake такого механизма нет, и это явно зафиксированный, не временный архитектурный выбор формата, а не нереализованная пока возможность.
Возврат долга: компакция Deletion Vectors через OPTIMIZE¶
Read Amplification, накопленный Deletion Vectors, не растёт бесконечно сам по себе именно потому, что обычный OPTIMIZE из раздела 1 этого урока, переписывая файлы, одновременно вшивает результат применения всех накопленных Deletion Vector в новый базовый файл: строки, помеченные как удалённые, физически не попадают в новый файл, а сам Deletion Vector становится не нужен и удаляется вместе с заменённым базовым файлом через стандартный механизм remove/add. Это - в точности тот же принцип, который седьмой урок модуля формулировал для Iceberg как rewrite_position_delete_files/rewrite_data_files: Merge-on-Read даёт дешёвую запись здесь и сейчас, а периодическая компакция возвращает таблицу к состоянию с нулевым Read Amplification, конвертируя накопленный технический долг в разовую стоимость одной операции OPTIMIZE.
Практическое следствие для эксплуатации: таблицы с активно включёнными Deletion Vectors и интенсивным потоком DELETE/UPDATE обязаны попадать в регулярное расписание OPTIMIZE чаще, чем таблицы без MoR-нагрузки - иначе число накопленных Deletion Vector файлов и связанный с ними Read Penalty будут расти с каждым новым удалением, постепенно съедая весь выигрыш, который Deletion Vectors должны были дать по сравнению с прямым Copy-on-Write. Раздел 7 этого урока, посвящённый Maintenance DAG, учитывает это явно: таблицы с высокой частотой точечных удалений получают более частый цикл OPTIMIZE, чем таблицы, в которые данные только дописываются (append-only).
Особенности обслуживания Managed Tables в различных средах¶
Managed Tables: что меняется по сравнению с path-based таблицами¶
Все примеры предыдущих разделов адресовали таблицы по физическому пути (delta.\s3a://...`) - это **unmanaged** (External) таблицы, у которых жизненный цикл метаданных в каталоге (если он вообще существует) не связан с жизненным циклом самих файлов. Заголовок этого урока, однако, говорит конкретно про **managed tables** - таблицы, зарегистрированные в метасторе (Hive Metastore, Unity Catalog, или встроенныйspark_catalog, который Delta Lake переопределяет черезDeltaCatalog, как было показано в одиннадцатом уроке) без явного указанияLOCATION`, для которых Spark сам выбирает физический путь хранения и сам берёт на себя ответственность за его удаление, если таблица будет дропнута.
-- Managed table: путь к данным выбирает и контролирует сам Spark/метастор
CREATE TABLE orders_delta (
order_id BIGINT,
user_id BIGINT,
status STRING,
order_date DATE
)
USING DELTA;
-- Unmanaged (external) table: путь явно указан и не управляется метастором
CREATE TABLE orders_delta_external (
order_id BIGINT,
user_id BIGINT,
status STRING,
order_date DATE
)
USING DELTA
LOCATION 's3a://lakehouse/warehouse/orders_delta_external';
Сами процедуры OPTIMIZE, VACUUM, ZORDER BY и Auto Optimize работают идентично для managed и unmanaged таблиц - физическая механика компакции, кластеризации и удаления файлов не зависит от того, кто формально владеет записью о таблице в каталоге. Различие, которое реально имеет значение для эксплуатации, проявляется не в синтаксисе самих команд, а в том, кто и как берёт на себя ответственность за их регулярный, надёжный запуск - и именно здесь среды Databricks и самостоятельно собранного (self-hosted) кластера OSS Spark расходятся существенно.
Databricks: автоматизация обслуживания как часть платформы¶
Managed-платформа Databricks исторически продвигала идею, что само обслуживание Delta-таблиц должно быть прозрачным для пользователя - точно так же, как СУБД не требует от прикладного разработчика вручную запускать VACUUM в PostgreSQL. На уровне Databricks Runtime это выражается в нескольких практических механизмах: автоматический сбор и обновление статистики таблиц (без необходимости вручную вызывать ANALYZE TABLE), движок выполнения Photon, ускоряющий именно тяжёлые операции компакции и Z-Ordering за счёт векторизованного, написанного на C++ исполнения вместо JVM-байткода, и режим Serverless, в котором сама платформа управляет выделением вычислительных ресурсов под конкретную задачу обслуживания без необходимости вручную поднимать и настраивать кластер.
Отдельная, более автоматизированная форма этого подхода - функция OPTIMIZE и Auto Optimize, интегрированные непосредственно в Databricks Jobs с предсказуемой биллинговой моделью: платформа способна сама определять, какие таблицы нуждаются в обслуживании, основываясь на собственной телеметрии записи, и планировать запуск OPTIMIZE без явного DAG, написанного инженером (этот более продвинутый, предиктивный режим выходит за рамки данного урока, ориентированного на self-hosted Lakehouse, но важно знать о его существовании, чтобы не пытаться вручную воссоздавать то, что управляемая платформа уже умеет делать сама).
OSS Spark + self-hosted объектное хранилище: ручная настройка обязательна¶
Self-hosted стек на голом PySpark и MinIO/Ceph RGW, который этот курс рассматривает как основной целевой сценарий (в силу принципа self-host only из правил этого курса), не получает ничего из перечисленного бесплатно - каждая процедура обслуживания требует явно написанного и явно запланированного кода, и каждый параметр производительности требует осознанной настройки SparkSession под конкретный объём данных и доступные ресурсы кластера.
Наиболее чувствительная точка - OPTIMIZE ... ZORDER BY на крупных таблицах: фаза сортировки по Z-значению требует полного shuffle затрагиваемых данных, и если число партиций shuffle (spark.sql.shuffle.partitions) подобрано неудачно относительно объёма обрабатываемых данных и доступной executor-памяти, операция рискует упасть с OutOfMemoryError прямо в процессе компакции крупной исторической партиции. Поскольку OPTIMIZE и так является ресурсоёмкой, изолированной от основного пайплайна операцией (раздел 1), разумно выделять под неё отдельный SparkSession/job с настройками, специально подобранными под характер этой нагрузки, а не унаследованными от конфигурации обычного ETL-job’а.
# Пример SparkSession, специально сконфигурированной под тяжёлую
# maintenance-нагрузку (OPTIMIZE ZORDER BY) на self-hosted кластере
spark = (
SparkSession.builder
.appName("delta-maintenance-zorder")
.config("spark.sql.shuffle.partitions", "400")
.config("spark.executor.memory", "8g")
.config("spark.executor.memoryOverhead", "2g")
.config("spark.sql.adaptive.enabled", "true")
.config("spark.sql.adaptive.coalescePartitions.enabled", "true")
.config("spark.databricks.delta.vacuum.parallelDelete.enabled", "true")
.getOrCreate()
)
Параметр spark.sql.shuffle.partitions здесь подобран существенно выше значения, разумного для обычного ETL-job’а на той же таблице: цель maintenance-job’а - максимально равномерно распределить тяжёлую сортировку Z-Order между executor’ами, чтобы ни одна задача не получила непропорционально большой объём данных и не вызвала переполнение памяти конкретного executor’а. Включение Adaptive Query Execution (spark.sql.adaptive.enabled) позволяет Spark динамически перераспределять данные между partition’ами уже во время выполнения, частично компенсируя неточный выбор статического значения shuffle.partitions, что особенно полезно, поскольку реальный объём данных в конкретной партиции таблицы редко известен заранее с высокой точностью.
Практическое правило для самостоятельной эксплуатации¶
Главный операционный вывод этого раздела: на self-hosted Lakehouse ответственность за надёжность maintenance-процедур целиком лежит на инженерной команде, и эта ответственность не сводится к запоминанию правильного синтаксиса команд - она включает в себя осознанную настройку ресурсов кластера под конкретный профиль нагрузки maintenance-job’а, мониторинг успешности и продолжительности каждого запуска, и явное проектирование расписания, которое не конкурирует за ресурсы с основными ETL-пайплайнами. Раздел 7 формализует это расписание в виде конкретного Airflow DAG.
Проектирование расписания обслуживания: Maintenance Window в Airflow¶
Почему процедуры нельзя смешивать в одну задачу¶
Десятый урок модуля уже формулировал общий принцип проектирования maintenance-расписания для Iceberg: процедуры с разной стоимостью и разным профилем риска должны жить в отдельных задачах DAG с собственным расписанием, а не быть искусственно склеены в один монолитный шаг. Для Delta Lake этот принцип сохраняется целиком, хотя набор процедур короче - вместо трёх независимых вызовов (rewrite_data_files, expire_snapshots, remove_orphan_files) здесь всего два явных командных глагола (OPTIMIZE, VACUUM), но внутри одного OPTIMIZE скрываются два разных по стоимости режима: лёгкая компакция без сортировки и тяжёлая компакция с ZORDER BY.
Лёгкий OPTIMIZE без Z-Order, ограниченный недавними партициями через WHERE, дешев в вычислительном смысле (он не требует полного shuffle всей таблицы) и может безопасно запускаться часто - несколько раз в сутки для критичных, активно читаемых таблиц, чтобы не давать накапливаться мелким файлам между более редкими тяжёлыми циклами. Тяжёлый OPTIMIZE ZORDER BY по историческим партициям целиком - дорогая операция, требующая полного shuffle затронутых данных (раздел 6), и её разумное место - еженедельное технического окно по выходным, когда аналитическая нагрузка на кластер минимальна. VACUUM логически должен идти строго после завершения всех циклов оптимизации текущей недели - досрочный VACUUM, запущенный пока ещё идёт OPTIMIZE, рискует удалить файлы, на которые компакция в этот самый момент claims как на входные данные операции чтения.
Скелет DAG на Airflow¶
Структурно DAG для Delta Lake почти буквально повторяет паттерн, уже показанный для Iceberg в десятом уроке - отдельные задачи, явные зависимости через >>, ретраи на уровне default_args - с заменой конкретных вызовов CALL system.* на OPTIMIZE/VACUUM:
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="delta_orders_maintenance_light",
schedule="0 */6 * * *", # каждые 6 часов
start_date=datetime(2025, 1, 1),
catchup=False,
default_args=default_args,
tags=["delta", "maintenance"],
) as dag:
light_optimize = SparkSqlOperator(
task_id="optimize_light",
sql="""
OPTIMIZE delta.`s3a://lakehouse/warehouse/orders_delta`
WHERE order_date >= date_sub(current_date(), 1)
""",
)
def run_health_check(**context):
# подробности - в подразделе про health-check ниже
...
health_check = PythonOperator(
task_id="health_check_light",
python_callable=run_health_check,
)
light_optimize >> health_check
with DAG(
dag_id="delta_orders_maintenance_deep",
schedule="0 3 * * 6", # раз в неделю, ночь на субботу
start_date=datetime(2025, 1, 1),
catchup=False,
default_args=default_args,
tags=["delta", "maintenance"],
) as deep_dag:
deep_optimize = SparkSqlOperator(
task_id="optimize_zorder",
sql="""
OPTIMIZE delta.`s3a://lakehouse/warehouse/orders_delta`
ZORDER BY (user_id, device_id)
""",
)
vacuum = SparkSqlOperator(
task_id="vacuum",
sql="""
VACUUM delta.`s3a://lakehouse/warehouse/orders_delta`
RETAIN 168 HOURS
""",
)
health_check_deep = PythonOperator(
task_id="health_check_deep",
python_callable=run_health_check,
)
deep_optimize >> vacuum >> health_check_deep
Разделение на два отдельных DAG’а (..._light и ..._deep) - не косметическое решение, а прямое следствие принципа из начала этого раздела: лёгкая и тяжёлая процедуры получают независимые расписания, независимые SLA на длительность выполнения и независимые алерты при сбое, и сбой тяжёлого еженедельного цикла (например, временная нехватка ресурсов кластера в выходные) не должен блокировать или задерживать лёгкий ежедневный цикл, критичный для свежих, активно читаемых партиций.
Health-check через DESCRIBE HISTORY: что проверять после каждого запуска¶
Одиннадцатый урок модуля уже представил команду DESCRIBE HISTORY как способ просматривать историю коммитов Delta-таблицы. Для целей health-check после maintenance-запуска особенно важно поле operationMetrics, которое Delta Lake заполняет для каждой операции OPTIMIZE и VACUUM конкретными числовыми показателями результата - именно эти числа дают объективный, не основанный на интуиции ответ на вопрос «действительно ли последний запуск обслуживания принёс пользу»:
history_df = spark.sql(
"DESCRIBE HISTORY delta.`s3a://lakehouse/warehouse/orders_delta`"
)
last_optimize = (
history_df
.filter("operation = 'OPTIMIZE'")
.orderBy("version", ascending=False)
.limit(1)
)
last_optimize.select("version", "timestamp", "operationMetrics").show(truncate=False)
# operationMetrics содержит, среди прочего:
# numAddedFiles, numRemovedFiles,
# numAddedBytes, numRemovedBytes,
# minFileSize, p25FileSize, p50FileSize, p75FileSize, maxFileSize
Практический критерий успеха для OPTIMIZE: numRemovedFiles должно заметно превышать numAddedFiles (множество мелких файлов схлопнулось в существенно меньшее число крупных), а p50FileSize (медианный размер файла после операции) должен приближаться к целевому размеру файла (по умолчанию около 1 ГБ, настраивается через spark.databricks.delta.optimize.maxFileSize) - если после нескольких запусков подряд p50FileSize остаётся существенно ниже цели, это сигнал, что входящие данные продолжают поступать быстрее, чем успевает работать текущая частота OPTIMIZE, и расписание нужно учащать. Для VACUUM аналогичным критерием служит прямое сравнение объёма, фактически занятого таблицей на S3/MinIO (через aws s3 ls --summarize --recursive или встроенные метрики объектного хранилища), до и после запуска - устойчивое отсутствие сокращения объёма при растущем числе версий таблицы означает, что VACUUM либо не запускается достаточно часто, либо retention threshold выставлен заведомо избыточно консервативно для реального профиля читателей этой конкретной таблицы.
Сводная таблица: частота процедур по профилю таблицы¶
Решение о конкретной частоте каждой процедуры не должно приниматься интуитивно для каждой новой таблицы с нуля - ниже сведена практическая отправная точка, основанная на профиле нагрузки, которую можно скорректировать по результатам собственного health-check, а не на пустом месте:
| Профиль таблицы | Лёгкий OPTIMIZE |
OPTIMIZE ZORDER BY |
VACUUM |
|---|---|---|---|
| Потоковый append, без точечных изменений | каждые 1-6 часов | раз в неделю | раз в неделю, после оптимизации |
Частые точечные DELETE/UPDATE (Deletion Vectors) |
каждые 1-6 часов | раз в неделю | раз в неделю, после оптимизации |
| Редкий, крупнообъёмный батч (ежедневный backfill) | не требуется отдельно от записи | раз в месяц или по необходимости | раз в месяц |
| Staging / временные таблицы короткого жизненного цикла | по необходимости | обычно не требуется | агрессивный retention возможен по согласованию |
Практический демо-блок: от захламлённой таблицы заказов до чистого managed table¶
Семь теоретических разделов выше были построены на постоянном чередовании «механика - риск - production-рекомендация». Эта часть урока переводит то же чередование в воспроизводимые команды против одной и той же managed Delta-таблицы заказов: имитация захламления мелкими файлами, измерение деградации, последовательное применение OPTIMIZE, ZORDER BY, VACUUM, Deletion Vectors и Auto Optimize, и измерение результата на каждом шаге.
Настройка окружения¶
SparkSession для этого демо-блока сфокусирован на одном Delta Lake, без параллельного подключения Iceberg-каталога, который использовался в демо-блоке одиннадцатого урока - тема этого урока не требует сравнения архитектур бок-о-бок, только глубокого, изолированного погружения в эксплуатацию Delta Lake.
from pyspark.sql import SparkSession
spark = (
SparkSession.builder
.appName("delta-maintenance-deep-dive")
.config("spark.jars.packages", "io.delta:delta-spark_2.12:3.1.0")
.config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
.config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")
.config("spark.hadoop.fs.s3a.endpoint", "http://minio:9000")
.config("spark.hadoop.fs.s3a.access.key", "minioadmin")
.config("spark.hadoop.fs.s3a.secret.key", "minioadmin")
.config("spark.hadoop.fs.s3a.path.style.access", "true")
.getOrCreate()
)
# Managed table - путь к данным выбирает спарк/метастор, как описано в разделе 6
spark.sql("""
CREATE TABLE IF NOT EXISTS orders_delta (
order_id BIGINT,
user_id BIGINT,
device_id INT,
amount DOUBLE,
status STRING,
order_date DATE
)
USING DELTA
""")
Кейс 0: имитация Small Files Problem через серию микро-батчей¶
Чтобы наблюдать реальную деградацию, а не верить в неё абстрактно, нужно сначала её воспроизвести. Цикл из ста небольших INSERT, каждый по 50-200 строк, имитирует поведение потокового источника, дописывающего данные мелкими порциями - именно тот паттерн, который раздел 1 назвал основным источником мелких файлов в production.
import random
from datetime import date, timedelta
for batch in range(100):
rows = [
(
batch * 1000 + i,
random.randint(1, 50_000), # user_id - высокая кардинальность
random.randint(1, 500), # device_id
round(random.uniform(5, 500), 2),
random.choice(["new", "shipped", "cancelled"]),
date(2025, 1, 1) + timedelta(days=random.randint(0, 200)),
)
for i in range(random.randint(50, 200))
]
df = spark.createDataFrame(
rows, ["order_id", "user_id", "device_id", "amount", "status", "order_date"]
)
df.write.format("delta").mode("append").saveAsTable("orders_delta")
print(spark.sql("SELECT COUNT(*) AS total_rows FROM orders_delta").collect())
Кейс 1: измерение деградации до компакции¶
После ста микро-батчей таблица физически состоит из ста (или более - в зависимости от внутреннего числа партиций Spark на запись) отдельных мелких Parquet-файлов. Прямой способ убедиться в этом без догадок - запрос к системной таблице DESCRIBE DETAIL, которая, в отличие от DESCRIBE HISTORY (история коммитов), показывает текущий физический срез таблицы, включая число файлов и суммарный размер:
detail_before = spark.sql("DESCRIBE DETAIL orders_delta").select(
"numFiles", "sizeInBytes"
).collect()[0]
print(f"Файлов до OPTIMIZE: {detail_before['numFiles']}")
print(f"Средний размер файла: {detail_before['sizeInBytes'] / detail_before['numFiles'] / 1024:.1f} KB")
# Замеряем время точечного запроса по высококардинальной колонке -
# без какой-либо физической кластеризации каждый файл придётся открыть
import time
start = time.time()
spark.sql("SELECT * FROM orders_delta WHERE user_id = 4200").collect()
print(f"Время запроса без оптимизации: {time.time() - start:.2f} сек")
Типичный результат на таком эксперименте - несколько сотен файлов размером по 30-80 КБ каждый, и точечный запрос по user_id, вынужденный открыть подавляющее большинство из них, поскольку диапазоны min..max каждого файла из-за случайного распределения значений перекрывают друг друга почти полностью.
Кейс 2: компакция через OPTIMIZE и измерение результата¶
result = spark.sql("OPTIMIZE orders_delta")
result.show(truncate=False)
detail_after = spark.sql("DESCRIBE DETAIL orders_delta").select(
"numFiles", "sizeInBytes"
).collect()[0]
print(f"Файлов после OPTIMIZE: {detail_after['numFiles']}")
print(f"Средний размер файла: {detail_after['sizeInBytes'] / detail_after['numFiles'] / (1024**2):.1f} MB")
Возвращаемая OPTIMIZE таблица содержит вложенную структуру metrics с теми же полями, что были разобраны в разделе про health-check (numFilesAdded, numFilesRemoved), - это тот же набор показателей, который Airflow-задача из раздела 7 проверяет автоматически после каждого запуска, только здесь он наблюдается интерактивно, сразу после выполнения команды.
Кейс 3: ZORDER BY и наблюдение Data Skipping¶
spark.sql("OPTIMIZE orders_delta ZORDER BY (user_id, device_id)")
# Тот же точечный запрос - теперь по плотно кластеризованным файлам
start = time.time()
plan = spark.sql("SELECT * FROM orders_delta WHERE user_id = 4200")
plan.collect()
print(f"Время запроса после Z-Ordering: {time.time() - start:.2f} сек")
# Просмотр физического плана: PartitionFilters / PushedFilters
# показывают, сколько файлов реально было прочитано планировщиком
plan.explain(mode="formatted")
В выводе explain(mode="formatted") после Z-Ordering ожидаемо заметна существенно меньшая цифра в строке, отвечающей за число просканированных файлов (number of files read), по сравнению с аналогичным планом для запроса из Кейса 1 - прямое, наблюдаемое в выводе Spark подтверждение эффекта data skipping, разобранного в разделе 2 теоретической части.
Кейс 4: VACUUM - от логического удаления до физического¶
# Удаляем отменённые заказы старше полугода - логическое удаление,
# создающее remove-записи в _delta_log, но не освобождающее диск
spark.sql("""
DELETE FROM orders_delta
WHERE status = 'cancelled' AND order_date < '2025-03-01'
""")
# Dry run - смотрим, что VACUUM удалил бы, без реального удаления
dry_run_result = spark.sql("VACUUM orders_delta RETAIN 168 HOURS DRY RUN")
dry_run_result.show(truncate=False)
# Реальный запуск с безопасным retention по умолчанию
spark.sql("VACUUM orders_delta RETAIN 168 HOURS")
Поскольку демонстрационная таблица только что создана и все её файлы моложе 168 часов, реальный VACUUM в рамках этого конкретного эксперимента не найдёт файлов-кандидатов на удаление - это ожидаемое, корректное поведение защиты, разобранной в разделе 3, а не ошибка демо: чтобы увидеть реальное физическое удаление, нужно либо подождать семь дней, либо (исключительно в учебных целях, на тестовой, ни для кого не критичной таблице) сознательно снизить порог через retentionDurationCheck.enabled = false и RETAIN 0 HOURS, воспроизводя ровно тот риск, который раздел 3 описывал как причину FileNotFoundException для параллельных читателей.
Кейс 5: жизненный цикл Deletion Vector - от DELETE до компакции¶
spark.sql("ALTER TABLE orders_delta SET TBLPROPERTIES (delta.enableDeletionVectors = true)")
# Точечное удаление - теперь обслуживается Deletion Vector, а не полной перезаписью
spark.sql("DELETE FROM orders_delta WHERE status = 'cancelled'")
# Подтверждение через DESCRIBE HISTORY: операция отмечена, но базовые файлы не тронуты
spark.sql("DESCRIBE HISTORY orders_delta") \
.select("version", "operation", "operationMetrics") \
.filter("operation = 'DELETE'") \
.show(truncate=False)
# operationMetrics для DELETE с Deletion Vectors содержит numDeletionVectorsAdded
# вместо привычных numRemovedFiles/numAddedFiles - прямое числовое подтверждение
# того, что физические файлы не переписывались
# Компакция - Deletion Vector "вшивается" в данные, накопленный Read Penalty снимается
spark.sql("OPTIMIZE orders_delta")
Поле operationMetrics для версии с DELETE - самое прямое доступное подтверждение теоретического материала раздела 5: значение numDeletionVectorsAdded больше нуля и около-нулевые numRemovedFiles/numAddedFiles однозначно показывают, что операция отработала через Merge-on-Read, а не через классическую Copy-on-Write перезапись, которая была бы видна как большое число numRemovedFiles, равное числу затронутых файлов.
Кейс 6: Auto Optimize в действии - сравнение поведения записи¶
spark.sql("""
ALTER TABLE orders_delta SET TBLPROPERTIES (
delta.autoOptimize.optimizeWrite = true,
delta.autoOptimize.autoCompact = true
)
""")
# Повторяем имитацию микро-батчей из Кейса 0, на этот раз с включённым Auto Optimize
for batch in range(100, 130):
rows = [
(
batch * 1000 + i,
random.randint(1, 50_000),
random.randint(1, 500),
round(random.uniform(5, 500), 2),
random.choice(["new", "shipped", "cancelled"]),
date(2025, 1, 1) + timedelta(days=random.randint(0, 200)),
)
for i in range(random.randint(50, 200))
]
df = spark.createDataFrame(
rows, ["order_id", "user_id", "device_id", "amount", "status", "order_date"]
)
df.write.format("delta").mode("append").saveAsTable("orders_delta")
detail_auto = spark.sql("DESCRIBE DETAIL orders_delta").select(
"numFiles", "sizeInBytes"
).collect()[0]
print(f"Средний размер файла с Auto Optimize: {detail_auto['sizeInBytes'] / detail_auto['numFiles'] / (1024**2):.1f} MB")
Сравнение среднего размера файла, полученного в этом кейсе, со средним размером файла из Кейса 0 (записанного без Auto Optimize) - наглядная демонстрация раздела 4: даже без единого явного вызова OPTIMIZE после записи новых тридцати микро-батчей, средний размер файла остаётся заметно выше, чем при «голой» записи, потому что Optimized Write и Auto Compact уже сделали часть работы компакции проактивно, в момент самой записи.
Производственный кейс¶
Платформенная команда retail-аналитики эксплуатировала Delta-таблицу events.clickstream объёмом около 40 ТБ - основной источник для ML-команды, которая раз в сутки строила фичи для модели рекомендаций через интерактивную PySpark-сессию с явным VERSION AS OF, фиксируя версию таблицы в начале своего шестичасового batch job’а, чтобы все вычисления в рамках одного запуска опирались на консистентный, неизменный срез данных, не подверженный новым записям, приходящим параллельно от потокового источника.
Объём, физически занятый таблицей на S3-совместимом хранилище, рос быстрее, чем ожидала команда платформы, и за квартал превысил бюджет, согласованный с финансовым отделом, в полтора раза - расследование показало, что VACUUM на этой таблице запускался нерегулярно, вручную, по случаю, и накопил больше года невычищенных, логически удалённых файлов. Под давлением срочного запроса «снизить расходы на хранение в этом же квартале» инженер платформенной команды решил выполнить агрессивную одноразовую очистку: отключил защитную проверку (spark.databricks.delta.retentionDurationCheck.enabled = false) и запустил VACUUM events.clickstream RETAIN 24 HOURS - выбрав суточный порог как компромисс между «быстро» и «не совсем ноль», без согласования этого порога с ML-командой и без проверки, какие job’ы в этот момент реально читают таблицу.
VACUUM отработал успешно и физически освободил заметный объём хранилища - но ровно в то же окно времени уже шёл шестичасовой job ML-команды, зафиксировавший версию таблицы за два дня до запуска VACUUM. На третьем часу выполнения job упал с FileNotFoundException: один из файлов, на который ссылалась зафиксированная версия, физически перестал существовать, потому что его возраст (по времени последней модификации) превысил выбранный суточный порог удержания, и VACUUM счёл его безопасным к удалению, не зная (и не имея никакого механизма узнать) о существовании активного, уже выполняющегося чтения именно этой версии.
Постфактум-анализ показал две накладывающиеся ошибки: во-первых, платформенная команда не вела реестра реальных потребителей таблицы и их паттернов доступа (как долго длятся их job’ы, фиксируют ли они конкретную версию), что сделало невозможным осознанный выбор безопасного retention threshold; во-вторых, решение о снижении retention принималось изолированно, под давлением бюджетного срока, без процесса согласования изменений конфигурации, влияющих на гарантии Time Travel, с командами-потребителями. Исправление включало восстановление стандартного RETAIN 168 HOURS, формальное соглашение об уровне обслуживания (SLA) с ML-командой, ограничивающее максимальную длительность их job’ов шестью часами и формализующее это число как нижнюю границу для любого будущего изменения retention threshold, и перенос всех запусков VACUUM в формальный Maintenance DAG (раздел 7) с фиксированным окном по выходным, когда ML-команда заведомо не запускает интерактивные сессии.
Отдельный практический урок, зафиксированный командой во внутреннем post-mortem: сам факт, что VACUUM физически успешно выполнился и не выдал ни одной ошибки в момент собственного запуска, создал ложное ощущение безопасности - симптом проявился не у VACUUM, а у совершенно другого, формально не связанного job’а, спустя несколько часов. Это - тот же класс «отложенного, не локального отказа», что уже разбирался для expire_snapshots в производственном кейсе десятого урока модуля: операция очистки в обоих форматах может рапортовать полный успех, при этом создавая отказ где-то ещё, у читателя, о существовании которого инициатор очистки может даже не подозревать.
Долгосрочное организационное изменение, которое команда внедрила по итогам этого инцидента, оказалось проще, чем любое техническое исправление: реестр потребителей (consumer registry) - простая таблица в внутренней wiki-системе, куда каждая команда, читающая production-таблицу через Time Travel или через долгие интерактивные сессии, обязана была вписать максимальную ожидаемую продолжительность своих job’ов и характер используемого доступа. Любое изменение retention threshold для конкретной таблицы стало требовать явной сверки с этим реестром, а не быть единоличным техническим решением инженера, ответственного за бюджет хранения - простой процессный барьер, не требующий новой инфраструктуры, но устраняющий именно тот класс ошибки, который привёл к инциденту: решение, влияющее на гарантии для других команд, принималось изолированно, без видимости их фактических паттернов использования.
Типичные заблуждения¶
«OPTIMIZE освобождает место на диске». Неверно: OPTIMIZE лишь логически переключает таблицу на новый, более компактный набор файлов через коммит в _delta_log; старые файлы продолжают физически занимать место до тех пор, пока их не удалит VACUUM после истечения порога удержания - это прямо разобрано в разделах 1 и 3 и является самым частым источником путаницы у тех, кто впервые сталкивается с двухкомандной моделью обслуживания Delta Lake.
«Z-Ordering - это то же самое, что партиционирование, просто под другим названием». Неверно: партиционирование - структурное свойство схемы хранения, определяющее физическую директорию файла и видимое в предикате запроса; Z-Ordering - результат компакции, влияющий только на порядок строк внутри файлов в рамках уже существующей партиционной структуры (или её полного отсутствия) и работающий исключительно через статистику min/max для data skipping, а не через директории.
«Чем больше колонок указано в ZORDER BY, тем лучше». Неверно и прямо противоположно происходящему на практике: раздел 2 показал, что эффективность Z-Ordering деградирует уже при превышении 3-4 колонок из-за «проклятия размерности» кривой Мортона - добавление пятой-шестой колонки чаще ухудшает, чем улучшает результат для каждой отдельно взятой колонки.
«VACUUM RETAIN 0 HOURS - это просто более агрессивная, но в остальном безопасная версия обычного VACUUM». Категорически неверно, и производственный кейс этого урока - прямая иллюстрация последствий именно этого заблуждения: снижение порога удержания ниже реальной продолжительности активных читателей создаёт FileNotFoundException не для самого VACUUM, а для произвольного другого job’а, который может упасть спустя часы после успешного завершения очистки.
«Deletion Vectors полностью заменяют необходимость в OPTIMIZE». Неверно: Deletion Vectors снижают стоимость записи при DELETE (раздел 5), но создают Read Penalty, накапливающийся при каждом новом удалении, и не исчезающий сам по себе - регулярный OPTIMIZE остаётся обязательным элементом эксплуатации даже на таблице с включёнными Deletion Vectors, просто решает другую, отложенную проблему.
«Auto Optimize - это бесплатная, не имеющая недостатков настройка, которую стоит включить на всех таблицах сразу». Неверно: раздел 4 показал прямой tradeoff между улучшением физической раскладки данных и дополнительной задержкой самой write-операции (commit latency) - на таблицах с редкой, но крупнообъёмной батч-записью эта задержка не компенсируется соразмерной выгодой.
«Managed и Unmanaged (External) Delta-таблицы обслуживаются принципиально по-разному». Неверно: раздел 6 показал, что сами процедуры OPTIMIZE/VACUUM/ZORDER BY идентичны для обоих типов регистрации таблицы в каталоге; различие касается только того, кто управляет физическим путём хранения и берёт на себя удаление файлов при DROP TABLE, а не механики самого обслуживания.
«На self-hosted Spark достаточно просто запускать те же команды OPTIMIZE/VACUUM, что и в документации Databricks, без дополнительной настройки». Неточно: раздел 6 явно показал, что синтаксис команд идентичен, но ресурсная конфигурация SparkSession (shuffle.partitions, executor memory) на самостоятельно собранном кластере требует осознанной настройки под конкретный объём данных - то, что на managed-платформе автоматизировано платформой (Photon, Serverless), на self-hosted кластере остаётся прямой ответственностью инженерной команды.
Производственный чек-лист¶
-
Расписание
OPTIMIZEразделено на лёгкий частый цикл (без Z-Order, с фильтром по свежим партициям) и тяжёлый редкий цикл (ZORDER BY, по историческим партициям) - раздел 7 показал, что смешение этих двух режимов в одной задаче привязывает дешёвую, частую операцию к расписанию дорогой, редкой, что не оправдано ни для одной из них. -
VACUUMзапускается строго после завершения всех цикловOPTIMIZEтекущего окна обслуживания, а не параллельно или до них - раздел 3 и производственный кейс этого урока показывают, что нарушение этого порядка создаёт прямой риск гонки между удалением файлов и операциями, которые их ещё используют. -
Retention threshold для
VACUUMзафиксирован не меньше, чем максимальная задокументированная продолжительность сессий реальных потребителей таблицы (Time Travel запросы, долгие интерактивные сессии, batch job’ы с фиксацией версии) - значение по умолчанию (168 часов) разумно почти всегда, и любое отклонение от него требует явного согласования с командами-потребителями, а не одностороннего решения платформенной команды. -
spark.databricks.delta.retentionDurationCheck.enabled = falseне используется как стандартная практика - это осознанный, разовый инструмент для узкого набора обоснованных случаев (например, staging-таблицы с коротким жизненным циклом потребителей), а не настройка, которую можно безопасно «включить и забыть» на production-таблице. -
Колонки для
ZORDER BYвыбраны на основе реального анализа паттернов фильтрации запросов (через Spark History Server или внешний query-лог), а не на основе интуитивного предположения «это важная колонка» - и их число не превышает 3-4 по причинам, разобранным в разделе 2. -
Auto Optimize включён осознанно, по результатам анализа профиля нагрузки конкретной таблицы, а не по умолчанию на всех таблицах сразу - таблицы с частыми, малыми write-транзакциями (потоковые источники) выигрывают от него больше, чем таблицы с редкими, крупнообъёмными батч-загрузками.
-
Deletion Vectors включены с пониманием их структурного ограничения (только пометка удаления, не изменение значения) и с регулярным
OPTIMIZEв расписании для компакции накопленных bitmap-файлов - раздел 5 показал, что без этой компакции Read Penalty продолжает расти с каждым новымDELETE. -
Совместимость со сторонними инструментами (BI-движки, Trino/Presto-коннекторы, инструменты каталогизации) проверена до включения Deletion Vectors или Z-Ordering на production-таблице, особенно за пределами Databricks Runtime - тринадцатый урок модуля приводит конкретный production-кейс несовместимости подобного рода для Liquid Clustering, и тот же класс риска применим к любой новой возможности формата, появившейся позже базового протокола.
-
На self-hosted кластере ресурсная конфигурация SparkSession для
OPTIMIZE ZORDER BYподобрана под реальный объём затрагиваемых партиций, а не унаследована от конфигурации обычного ETL-job’а - раздел 6 показал, что неудачно подобранныйshuffle.partitionsспособен привести кOutOfMemoryErrorпрямо в процессе компакции крупной исторической партиции.
Мостик к следующему уроку¶
Этот урок закрыл практическую эксплуатацию Delta Lake тем же способом, каким десятый урок модуля закрыл практическую эксплуатацию Iceberg - не как набор изолированных команд для запоминания, а как взаимосвязанную систему процедур, каждая из которых решает конкретный, измеримый класс операционного риска: Small Files Problem решается компакцией, неэффективный data skipping на высококардинальных колонках - многомерной кластеризацией, бесконечный рост физического хранилища - удержанием и удалением по порогу, реактивное отставание от потока записи - проактивным управлением размером файлов на лету, а растущая стоимость Merge-on-Read - периодическим возвратом накопленного технического долга в базовые файлы.
Тринадцатый, завершающий урок модуля поднимается на уровень выше обеих архитектур, разобранных по отдельности (Iceberg в уроках 1-10, Delta Lake в уроках 11-12), и добавляет третий формат - Apache Hudi - чтобы дать инженеру не очередной список фактов о том, «что есть в каждом формате», а конкретный, применимый на практике инструмент принятия решения: какой формат выбрать под заданный профиль нагрузки, существующую экосистему движков и операционную зрелость команды, и какую цену придётся платить за этот выбор в горизонте нескольких лет эксплуатации. Всё, что было разобрано в этом уроке - двухкомандная модель обслуживания, Z-Ordering вместо Hidden Partitioning, Deletion Vectors как единственная форма Merge-on-Read у Delta Lake - станет одной из опорных точек сравнения в финальной матрице выбора.
Домашнее задание¶
-
Разверните self-hosted Delta Lake поверх локального MinIO и воспроизведите Кейс 0 из демо-блока этого урока (сто микро-батчей по 50-200 строк). Замерьте число файлов и средний размер файла до какой-либо оптимизации, затем зафиксируйте те же метрики после
OPTIMIZEбезZORDER BY- убедитесь, чтоnumRemovedFilesвoperationMetricsсоответствует реальному числу файлов, существовавших до операции. -
На той же таблице выполните
OPTIMIZE ... ZORDER BY (user_id), затем сравните выводexplain(mode="formatted")для точечного запросаWHERE user_id = <конкретное значение>до и после Z-Ordering. Зафиксируйте конкретную числовую разницу в количестве просканированных файлов (number of files read). -
Повторите эксперимент из задания 2, но используя для
ZORDER BYсразу шесть колонок вместо одной-двух. Сравните полученный эффект data skipping для конкретной, заранее выбранной «важной» колонки с результатом, полученным приZORDER BYтолько по этой одной колонке - убедитесь на собственных числах в эффекте «проклятия размерности», разобранном в разделе 2. -
Создайте тестовую (не критичную ни для одного реального потребителя) Delta-таблицу, выполните
DELETEчасти строк, затем сразу запуститеVACUUM ... RETAIN 0 HOURSс предварительно отключённой проверкойretentionDurationCheck.enabled. Параллельно, в отдельном процессе, запустите Time Travel запрос (VERSION AS OFк версии, существовавшей доDELETE) и зафиксируйте получаемое исключение - это управляемое, безопасное воспроизведение риска, разобранного в разделе 3 и в производственном кейсе урока. -
Включите
delta.enableDeletionVectors = trueна тестовой таблице, выполнитеDELETEпо предикату, затрагивающему заметную долю строк одного крупного файла, и сравните черезDESCRIBE DETAILфизический размер файлов до и после операции. Убедитесь, что размер базового файла не изменился, а появился отдельный, существенно меньший по объёму файл Deletion Vector. -
На той же таблице из задания 5 выполните
UPDATE(а неDELETE) по аналогичному предикату и сравните набор файлов, созданных операцией, с набором, созданнымDELETE. Объясните своими словами, почемуUPDATEсоздаёт новый файл данных, аDELETE- только Deletion Vector, опираясь на структурное ограничение, разобранное в разделе 5. -
Включите Auto Optimize (
delta.autoOptimize.optimizeWriteиdelta.autoOptimize.autoCompact) на новой тестовой таблице и повторите имитацию микро-батчей из Кейса 0. Сравните средний размер итогового файла с результатом, полученным без Auto Optimize, и оцените разницу в суммарном времени выполнения всех ста write-операций - зафиксируйте конкретные числа, иллюстрирующие tradeoff между write latency и итоговым качеством физической раскладки данных, разобранный в разделе 4. -
Спроектируйте на бумаге (или в виде реального Airflow DAG) полный Maintenance DAG для гипотетической таблицы с профилем нагрузки «потоковая запись каждые 30 секунд, точечные
DELETEпо GDPR-запросам раз в сутки, аналитические запросы с фильтрами поuser_idиevent_type». Обоснуйте выбранную частоту лёгкогоOPTIMIZE, состав колонок дляZORDER BY, день и час запуска тяжёлого цикла, и выбранный retention threshold дляVACUUM. -
Сравните операционную модель
OPTIMIZE/VACUUM/ZORDER BYDelta Lake с операционной модельюrewrite_data_files/expire_snapshots/remove_orphan_filesIceberg, разобранной в десятом уроке модуля, заполнив таблицу из 5 строк: для каждой строки укажите команду Delta Lake, команду Iceberg, и одно предложение о практическом различии в эксплуатации (частота запуска, типичный риск, типичная стоимость). -
Найдите в официальной документации Delta Lake текущее значение
spark.databricks.delta.optimize.maxFileSizeпо умолчанию для используемой вами версии (io.delta:delta-spark) и объясните, почему это значение выбрано именно таким - какой баланс между числом файлов, накладными расходами на их перечисление (раздел 1) и стоимостью полной перезаписи при будущемOPTIMIZEоно отражает. -
Постройте сводную таблицу из раздела 7 («Сводная таблица: частота процедур по профилю таблицы») для одной из реальных таблиц, с которыми вы уже работали (или для гипотетической таблицы из вашего текущего проекта), и явно укажите, какие конкретные числа в вашем расписании отличаются от приведённых в уроке значений по умолчанию, и почему - какой конкретный факт о профиле нагрузки этой таблицы оправдывает отклонение.
-
Создайте табличный набор данных с явно высококардинальной колонкой (например, синтетический
device_idс миллионом уникальных значений на десяти миллионах строк) и сравните на нём три стратегии физической организации: обычное партиционирование по этой колонке, отсутствие какой-либо кластеризации, иOPTIMIZE ZORDER BYпо этой колонке. Зафиксируйте число физических партиций/директорий в первом случае и сравните время выполнения идентичного точечного запроса во всех трёх случаях - это прямая практическая проверка тезиса раздела 2 о том, почему классическое партиционирование не подходит для высокой кардинальности.
Полная картина: от записи до устойчивого, обслуживаемого managed table¶
Диаграмма ниже сводит весь урок в единый жизненный цикл: от момента, когда поток событий создаёт первый мелкий файл, через все процедуры обслуживания, разобранные в семи разделах, до состояния таблицы, готового к следующему циклу записи без накопления необслуживаемого технического долга.
Три параллельные ветви диаграммы (Проактивный слой, Реактивный слой, Merge-on-Read нагрузка) сходятся в одной и той же финальной процедуре - OPTIMIZE - что отражает центральный тезис всего урока: независимо от того, откуда возникла фрагментация физических файлов - из потока записи, из плановой компакции или из накопленных Deletion Vector, - механизм её устранения один и тот же, и именно поэтому раздел 7 строит вокруг OPTIMIZE (в двух его режимах - лёгком и тяжёлом) и VACUUM единое расписание Maintenance DAG, а не отдельный конвейер на каждый источник проблемы.
Итоги¶
Small Files Problem - центральная мотивация всего урока: микро-батчи потоковой и пакетной записи создают множество мелких физических файлов, каждый из которых добавляет фиксированную накладную стоимость на уровне объектного хранилища (LIST/GET) и на уровне планировщика Spark; OPTIMIZE решает эту проблему компакцией - переписывает множество мелких файлов в меньшее число крупных, фиксируя результат одним коммитом remove/add в _delta_log.
Z-Ordering (OPTIMIZE ... ZORDER BY) решает отдельную, но смежную задачу - эффективный data skipping на высококардинальных колонках, для которых классическое партиционирование органически не подходит. Многомерная кластеризация по кривой Мортона физически группирует строки со схожими значениями в одних и тех же файлах, сужая диапазоны min..max статистики и позволяя планировщику отбрасывать файлы без их скачивания - но эффект деградирует за пределами 3-4 колонок.
OPTIMIZE и любые другие операции изменения файлов никогда не освобождают физическое хранилище - они лишь логически переключают состояние таблицы через коммит в _delta_log. Физическое удаление выполняет только VACUUM, работающий по порогу удержания (Retention Threshold, по умолчанию 168 часов), и снижение этого порога ниже реальной продолжительности активных читателей создаёт прямой риск FileNotFoundException для job’ов, не имеющих никакого отношения к самому VACUUM.
Auto Optimize (Optimized Write + Auto Compact) переносит часть работы по управлению размером файлов с отложенной плановой компакции на момент самой записи, проактивно распределяя данные перед финальной записью и запуская мини-компакцию сразу после транзакции при необходимости - ценой дополнительной задержки самой записи, что оправдано для высокочастотных потоковых нагрузок и не оправдано для редких крупнообъёмных батчей.
Deletion Vectors - собственная реализация Merge-on-Read у Delta Lake, появившаяся с версии 2.3 и достигшая зрелости к 3.0: вместо полной перезаписи файла при DELETE создаётся компактная битовая карта (RoaringBitmap) удалённых позиций, что резко снижает Write Amplification ценой растущего Read Penalty, который снимается последующим OPTIMIZE. Структурное ограничение - Deletion Vector умеет только помечать удаление, не изменение значения, поэтому UPDATE остаётся гибридной, частично Copy-on-Write операцией.
Managed-платформы (Databricks) автоматизируют значительную часть обслуживания (сбор статистики, ускоренное исполнение через Photon, предиктивный запуск OPTIMIZE), тогда как self-hosted PySpark поверх MinIO/S3 требует явной, осознанной настройки ресурсов SparkSession под профиль конкретной maintenance-задачи - в первую очередь для тяжёлого OPTIMIZE ZORDER BY, рискующего OutOfMemoryError при неудачно подобранном shuffle.partitions.
Production-расписание обслуживания строится как несколько независимых задач Airflow с разной частотой и разным профилем риска, а не как одна монолитная процедура: лёгкий частый OPTIMIZE для свежих партиций, тяжёлый редкий OPTIMIZE ZORDER BY для исторических данных в техническое окно, и VACUUM строго после завершения обоих циклов компакции - с health-check через DESCRIBE HISTORY/operationMetrics и прямое сравнение объёма хранилища как объективным критерием успеха каждого запуска.
Краткий глоссарий терминов урока¶
Small Files Problem - накопление большого числа физически мелких файлов данных вследствие частых, небольших по объёму write-транзакций; основной мотивирующий контекст для компакции.
OPTIMIZE - команда Delta Lake, переписывающая множество мелких Parquet-файлов в меньшее число файлов целевого размера через единый коммит remove/add в _delta_log.
Z-Ordering / ZORDER BY - расширение OPTIMIZE, физически переупорядочивающее строки по кривой Мортона перед записью, чтобы сгруппировать строки со схожими значениями указанных колонок в одних и тех же файлах для эффективного data skipping.
Data Skipping - отсечение нерелевантных для запроса файлов на этапе планирования, основанное только на статистике файла (min/max значения колонок), без скачивания самого файла.
Curse of dimensionality (в контексте Z-Order) - деградация эффективности многомерной кластеризации при росте числа измерений (колонок), из-за которой ZORDER BY не рекомендуется указывать для более чем 3-4 колонок.
VACUUM - команда Delta Lake, физически удаляющая с хранилища файлы, помеченные remove и старше порога удержания (Retention Threshold).
Retention Threshold (RETAIN n HOURS) - порог возраста файла, ниже которого VACUUM не удаляет файл, даже если он логически больше не входит в текущее состояние таблицы; по умолчанию 168 часов (семь дней).
spark.databricks.delta.retentionDurationCheck.enabled - защитный флаг, предотвращающий запуск VACUUM с порогом удержания ниже безопасного минимума; его отключение - явный, осознанный шаг повышения риска.
spark.databricks.delta.vacuum.parallelDelete.enabled - параметр производительности (не безопасности), управляющий тем, выполняется ли физическое удаление файлов VACUUM параллельно на executor’ах или последовательно на драйвере.
Auto Optimize - объединённое название для Optimized Write и Auto Compact - проактивного управления размером файлов в момент записи, в противовес реактивному плановому OPTIMIZE.
Optimized Write - механизм адаптивного локального shuffle перед финальной записью, приводящий число выходных файлов в соответствие с целевым размером.
Auto Compact - механизм автоматического запуска мини-компакции сразу после завершения транзакции записи, если она оставила за собой избыточное число мелких файлов.
Deletion Vector - компактный битовый файл (бинарная карта позиций, обычно в кодировке RoaringBitmap), помечающий логически удалённые строки внутри неизменного базового Parquet-файла; реализация Merge-on-Read у Delta Lake.
delta.enableDeletionVectors - табличное свойство, включающее использование Deletion Vectors для операций DELETE/UPDATE/MERGE.
Read Amplification / Read Penalty (в контексте Deletion Vectors) - дополнительная стоимость чтения, возникающая из-за необходимости применять Deletion Vector как маску к каждому обращению к затронутому файлу до момента его компакции.
Managed table - таблица, зарегистрированная в каталоге без явного LOCATION; путь физического хранения и ответственность за удаление файлов при DROP TABLE берёт на себя сама платформа/метастор.
Unmanaged (External) table - таблица, зарегистрированная с явным LOCATION; жизненный цикл физических файлов не связан с записью в каталоге.
DESCRIBE DETAIL - команда Delta Lake, показывающая текущий физический срез таблицы (число файлов, суммарный размер, формат), в противовес DESCRIBE HISTORY, показывающей историю коммитов.
operationMetrics - поле записи DESCRIBE HISTORY, содержащее конкретные числовые показатели результата операции (numAddedFiles, numRemovedFiles, numDeletionVectorsAdded и аналогичные) - основа объективного health-check после maintenance-запуска.
Maintenance DAG - набор независимых, по-разному запланированных задач Airflow, реализующих регулярное обслуживание таблицы: лёгкая компакция, тяжёлая Z-order компакция, физическая очистка через VACUUM.
Photon - векторизованный движок выполнения Databricks, ускоряющий, среди прочего, операции OPTIMIZE и ZORDER BY за счёт исполнения, написанного на C++ вместо JVM-байткода.
Serverless (Databricks) - режим, в котором платформа сама управляет выделением вычислительных ресурсов под конкретную задачу обслуживания, без необходимости вручную поднимать и настраивать постоянный кластер.
Liquid Clustering (CLUSTER BY) - более новая возможность Delta Lake, делающая выбор кластеризующих колонок постоянным свойством таблицы, инкрементально поддерживаемым при каждой записи и компакции, в противовес ZORDER BY, который требует повторного явного указания тех же колонок при каждом отдельном вызове OPTIMIZE.
RoaringBitmap - алгоритм компактного сжатия битовых карт, используемый для физического представления Deletion Vector; эффективно сжимает как плотные, так и разрежённые диапазоны установленных битов.
spark.databricks.delta.optimize.maxFileSize - параметр, определяющий целевой размер файла, к которому стремится OPTIMIZE при компакции (по умолчанию около 1 ГБ).
Write Amplification (применительно к Deletion Vectors) - объём данных, физически перезаписанных операцией изменения по отношению к объёму данных, реально затронутых логической операцией; Deletion Vectors радикально снижают эту величину для DELETE, но не для UPDATE, который остаётся гибридной операцией (раздел 5).