Политики автоматического обслуживания объектов Iceberg#

Основные понятия#

Политики автоматического обслуживания#

Таблицы Iceberg требуют регулярного обслуживания: удаления устаревших snapshot и файлов, не связанных ни с одним актуальным snapshot, перезаписи мелких файлов данных и объединения манифестов. При большом количестве таблиц ручной запуск этих операций становится непрактичным: отдельные cron-задачи сами становятся инфраструктурой, требующей поддержки, а глобальное расписание не учитывает разную интенсивность изменений в разных таблицах. Политики автоматического обслуживания решают эту задачу декларативно: вы описываете, при каком состоянии таблицы обслуживание имеет смысл, а каталог сам отслеживает это состояние и запускает нужные операции в подходящий момент.

Политика автоматического обслуживания — это именованная конфигурация, которая объединяет:

  • Целевую группу объектов — набор таблиц и материализованных представлений, для которых выполняется обслуживание (см. Группы объектов обслуживания).

  • Правила — одно или несколько правил, каждое из которых описывает, какую операцию выполнять и при каких условиях.

  • Временное окно — опциональное расписание, ограничивающее время, в которое политика может выполняться.

Каталог периодически проверяет все активные политики: для каждого объекта из целевой группы оценивает условия правил и, если условия выполнены, автоматически запускает соответствующие операции обслуживания.

Структура политики#

Политика задаётся следующими параметрами:

  • Имя политики (policy-name) — уникальное имя, по которому политика идентифицируется в командах CLI и API.

  • Целевая группа объектов (target-object-group-name) — группа объектов, к которой применяется политика.

  • Владелец (owner) — принципал, от имени которого выполняется политика. Устанавливается автоматически при создании.

  • Описание (description) — произвольный текстовый комментарий.

  • Статус активности (active) — флаг включения/выключения политики. Неактивная политика не проверяется и не запускает операции.

  • Временное окно (window) — расписание, в рамках которого политика может работать (см. Временное окно).

Политика содержит одно или несколько правил. Каждое правило определяет:

  • Конфигурацию операции — какую операцию выполнять, на каком вычислительном движке, с какими параметрами.

  • Триггер — условия, при которых правило срабатывает (см. Двухуровневая система триггеров).

  • Режим точки отсчёта — как инициализировать счётчики активности при добавлении правила (см. Режим точки отсчёта).

Типы операций#

Каждое правило указывает вычислительный движок, который будет выполнять операцию обслуживания. Доступны три движка:

  • local — встроенный движок каталога. Не требует внешней инфраструктуры, подходит для лёгких операций, таких как удаление устаревших snapshot и очистка orphan файлов.

  • Apache Spark — распределённый движок для ресурсоёмких операций: отптимизация файлов данных (compaction), перестроения манифестов и т.д. Требует доступ к кластеру Spark.

  • CedrusData — выполняет операции обслуживания через запросы к кластеру CedrusData (Trino).

Подробнее о доступных движках и поддерживаемых операциях: Движки для выполнения операций обслуживания, Движок обслуживания Apache Spark.

Двухуровневая система триггеров#

Каждое правило содержит триггер, который определяет, когда правило должно сработать для конкретного объекта. Триггер состоит из двух уровней проверки.

Счётчики активности#

Первый уровень — быстрая предварительная проверка на основе счётчиков активности. Каталог отслеживает для каждой пары (объект, правило) три счётчика, измеряющих изменения с момента добавления правила или последнего выполнения операции:

  • Количество коммитов (commits-added) — сколько коммитов Iceberg произошло.

  • Объём добавленных данных (bytes-added) — сколько байт данных было добавлено.

  • Прошедшее время (time-elapsed) — сколько миллисекунд прошло.

Для каждого счётчика задаётся пороговое значение. Необходимо указать хотя бы один счётчик, значение порога должно быть больше нуля. Правило становится кандидатом на проверку, когда хотя бы один из порогов превышен.

Проверка счётчиков выполняется на уровне базы данных каталога и не требует обращения к метаданным Iceberg, что делает её очень быстрой даже при большом количестве объектов.

Таким образом счётчики работают как фильтр: если таблица не изменялась с прошлого запуска, нет смысла считывать её метаданные. Это экономит ресурсы при большом количестве таблиц в группе.

Условие запуска#

Второй уровень — вычисление выражения на языке DSL. Оно выполняется только для объектов, прошедших проверку счётчиков активности, и работает как дополнительный фильтр: срабатывание счётчиков говорит лишь о том, что таблица изменилась, но ещё не означает, что её действительно нужно обслуживать.

Условие запуска — это выражение, которое имеет доступ к метаданным таблицы Iceberg: количеству и размеру файлов данных, количеству snapshot, времени с последнего обновления и другим характеристикам. Выражение должно возвращать логическое значение (true / false). Такая дополнительная проверка экономит ресурсы: сама операция обслуживания (compaction, удаление snapshot и т.п.) обычно заметно дороже вычисления условия, поэтому отсеять случаи, где она не принесёт пользы, выгоднее, чем запускать её вслепую.

Пример условия: запуск оптимизации файлов данных, когда в таблице более 100 файлов и средний размер файла меньше 128 МБ:

dataFiles > 100 AND avgDataFileSize < 128MB

Подробное описание языка DSL, доступных переменных и функций приведено в разделе Справочник DSL.

Режим точки отсчёта#

При добавлении правила в политику необходимо определить, как инициализировать счётчики активности для существующих объектов. Этот выбор задаётся режимом точки отсчёта (baseline-mode):

  • catch-up (по умолчанию) — счётчики начинают с нуля. Вся ранее накопленная активность учитывается при первой проверке. Используйте этот режим, когда нужно обработать объекты, для которых обслуживание ранее не выполнялось.

  • going-forward — счётчики инициализируются текущими значениями метрик объекта. Учитывается только новая активность, возникшая после добавления правила. Используйте этот режим, когда существующие объекты не требуют немедленного обслуживания и вы хотите отслеживать только будущие изменения.

Пример. В каталоге есть таблица orders, в которой за время работы накопилось 500 коммитов. Вы создаёте политику с правилом оптимизации файлов и порогом счётчика commits-added = 100.

  • В режиме catch-up счётчик коммитов начнёт с нуля. Разница между текущим количеством коммитов таблицы и нулевой точкой отсчёта (500) сразу превысит порог 100, и правило сработает при первой же проверке — файлы таблицы будут оптимизированы.

  • В режиме going-forward точка отсчёта установится на текущее значение (500). Правило сработает только после того, как таблица получит ещё 100 новых коммитов (то есть когда общее число коммитов достигнет 600).

Примечание

В режиме going-forward, если объект добавляется в группу после создания правила, то первая проверка этого объекта только устанавливает точку отсчёта и не приводит к запуску операции.

Временное окно#

Временное окно ограничивает период, в течение которого политика может выполнять проверки и запускать операции. Это позволяет, например, проводить обслуживание только в ночное время или в выходные дни, чтобы не нагружать кластер в рабочие часы.

Временное окно задаётся следующими параметрами:

  • Часовой пояс (time-zone-id) — идентификатор часового пояса (например, Europe/Moscow, UTC).

  • Периоды (periods) — один или несколько временных интервалов. Каждый период определяется:

    • Cron-выражением (cron) — стандартное cron-выражение, задающее момент начала окна (минута, час, день месяца, месяц, день недели).

    • Длительностью (duration) — продолжительность окна в миллисекундах от момента, определённого cron-выражением.

Политика может выполняться, когда соблюдены оба условия:

  1. Политика активна (active = true).

  2. Временное окно не задано, либо текущее время попадает хотя бы в один из сконфигурированных периодов.

Пример: обслуживание только по будням с 02:00 до 06:00 по московскому времени:

{
  "time-zone-id": "Europe/Moscow",
  "periods": [
    { "cron": "0 2 * * 1-5", "duration": 14400000 }
  ]
}

Если временное окно не задано и политика активна, то операции могут выполняться в любое время. Для промышленного использования рекомендуется всегда задавать временное окно.

Выполнение политик#

Каталог выполняет политики в фоновом цикле, разбитом на несколько фаз: поиск кандидатов, параллельная проверка условий запуска и отправка отобранных кандидатов в подсистему выполнения операций. Такая структура рассчитана на большие группы объектов: подавляющее большинство таблиц отсеивается дешёвой проверкой счётчиков ещё до того, как потребуется загружать их метаданные, а оставшиеся кандидаты проверяются параллельно.

Цикл работы политик#

Цикл оценки запускается периодически планировщиком каталога. На каждой итерации для каждой активной политики, для которой текущее время попадает в рамки временного окна (см. Временное окно), выполняются следующие шаги:

  1. Аутентификация. Аутентифицируется заданный для политики run-as пользователь, от имени которого выполняются все последующие действия.

  2. Поиск кандидатов. Каталог находит пары (объект, правило), для которых значение хотя бы одного из счётчиков активности (commits-added, bytes-added, time-elapsed) превысило заданный порог. Этот шаг не требует загрузки метаданных Iceberg, поэтому масштабируется на большие группы объектов (см. Двухуровневая система триггеров).

  3. Проверка условий запуска. Для каждого кандидата загружаются метаданные таблицы, вычисляется выражение DSL из конфигурации триггера, результат (сработало или нет) записывается в журнал активности.

  4. Запуск операций. Кандидаты, для которых условие DSL вернуло true, постепенно отправляются в подсистему выполнения операций обслуживания с учётом общего лимита одновременно выполняющихся операций.

Важно: как только DSL-условие вычислено в true для конкретной пары (объект, правило), повторное вычисление не производится до завершения запланированной операции. Это гарантирует, что одна и та же операция не будет поставлена в очередь дважды, и исключает лишнюю нагрузку на чтение метаданных.

Параллельная проверка условий#

Проверка условий запуска для кандидатов, прошедших фильтр счётчиков активности, выполняется параллельно в выделенном пуле потоков. Каждая пара (объект, правило) обрабатывается как независимая задача: загружаются метаданные таблицы, вычисляется заранее скомпилированное выражение DSL и результат сохраняется в журнал активности.

При большом количестве таблиц в группах увеличение размера пула позволяет ускорить фазу вычисления условий, но увеличивает нагрузку на чтение метаданных. Рекомендуется подбирать значение исходя из доступных ресурсов и количества таблиц.

Журнал активности#

Журнал активности (evaluation log) — ключевой инструмент для понимания поведения политик. В журнале фиксируется каждое вычисление условия: как успешные запуски, так и пропуски с причиной.

Запись журнала содержит:

  • идентификатор правила политики;

  • сведения об объекте Iceberg;

  • конфигурацию триггера на момент оценки (пороговые значения счётчиков и выражение DSL);

  • фактические значения счётчиков активности (commits-added, bytes-added, time-elapsed) на момент оценки;

  • результат проверки условия: сработало правило или нет, с подробностями вычисления;

  • ссылку на запущенную операцию обслуживания, если проверка привела к её запуску.

Журнал предназначен для диагностики поведения политик: он позволяет ответить на вопросы: «почему операция не выполнилась на этой таблице?», «какое значение метрики помешало?», «как часто условие срабатывает?».

Пример. Для таблицы orders настроено правило compaction со следующим триггером:

  • счётчик commits-added с порогом 100;

  • условие запуска dataFiles > 100 AND avgDataFileSize < 128MB.

Инженер замечает, что compaction не запускается, хотя в таблицу активно пишутся данные. Запись в журнале активности показывает:

  • commits-added = 150 — порог пройден;

  • фактические метрики таблицы: dataFiles = 80, avgDataFileSize = 210MB;

  • результат условия DSL: false, операция не запущена.

Причина видна сразу: счётчики сработали, таблица действительно менялась, но условие DSL отсеяло её — файлов всего 80, и они уже крупнее целевого размера. Без журнала разобрать эту ситуацию было бы значительно сложнее: к моменту расследования метаданные таблицы могли бы уже измениться, и восстановить, на каких именно данных сработал или не сработал триггер, было бы невозможно. Журнал фиксирует и значения метрик, и результат проверки именно в тот момент, когда она выполнялась.

Управление размером журнала. Журнал может расти быстро, особенно при большом количестве таблиц и частых циклах работы политик. Для предотвращения проблем с дисковым пространством каталог автоматически удаляет наиболее старые записи. Время жизни записей настраивается параметром:

maintenance.policy.evaluation-log-retention = 14d

При увеличении числа таблиц или количества правил рекомендуется следить за объёмом журнала и при необходимости уменьшать время хранения.

Параметры конфигурации#

Работа политик настраивается следующими параметрами каталога:

Параметр

Описание

По умолчанию

Минимум

maintenance.policy.evaluation-cycle-interval

Интервал между циклами работы политик.

60s

1s

maintenance.max-active-operations

Максимальное количество одновременно выполняющихся операций обслуживания.

8

1

maintenance.policy.evaluation-pool-size

Размер пула потоков, выполняющих проверку условий запуска.

количество ядер CPU × 2

1

maintenance.policy.evaluation-log-retention

Срок хранения записей журнала активности.

14d

1d

Справочник DSL#

Типы данных#

DSL поддерживает следующие типы данных:

Тип

Описание

Примеры значений

BOOLEAN

Логический тип

true, false

LONG

Целое число (64-битное)

0, 42, 1024

DOUBLE

Число с плавающей точкой

0.5, 3.14, 100.0

STRING

Строка

'hello', "world"

Тип каждого выражения определяется на этапе компиляции. Несовместимость типов (например, сравнение строки с числом) приводит к ошибке компиляции.

В арифметических операциях, если один из операндов имеет тип DOUBLE, результат будет DOUBLE. Операция деления (/) всегда возвращает DOUBLE.

Литералы#

Числовые литералы#

Целые числа записываются без суффиксов:

100
0
1048576

Числа с плавающей точкой содержат десятичную точку:

0.5
3.14
100.0

Литералы размера данных#

Для удобства работы с размерами файлов DSL поддерживает суффиксы, которые автоматически преобразуют значение в байты:

Суффикс

Множитель

Пример

Значение в байтах

KB / kb

1024

100KB

102 400

MB / mb

1024 × 1024

512MB

536 870 912

GB / gb

1024³

1GB

1 073 741 824

Примеры:

1KB == 1024
1MB == 1048576
1GB == 1073741824
512MB
100KB

Литералы длительности#

Для работы с временными интервалами DSL поддерживает суффиксы, которые преобразуют значение в миллисекунды:

Суффикс

Множитель

Пример

Значение в миллисекундах

m

60 000 (минуты)

30m

1 800 000

h

3 600 000 (часы)

2h

7 200 000

d

86 400 000 (дни)

7d

604 800 000

Примеры:

30m
24h
7d

Строковые литералы#

Строки заключаются в одинарные или двойные кавычки. Обратный слэш (\) используется для экранирования:

'hello'
"world"
'it\'s a test'

Логические литералы#

true
false

Переменные#

Переменные содержат характеристики текущей таблицы Iceberg.

Переменные, основанные на файлах данных#

Следующие переменные требуют загрузки метаданных таблицы и чтения всех манифестов текущего snapshot:

Переменная

Тип

Описание

dataFiles

LONG

Количество файлов данных

dataFilesSize

LONG

Суммарный размер файлов данных в байтах

dataRecords

LONG

Суммарное количество записей в файлах данных

equalityDeleteFiles

LONG

Количество файлов equality-delete

equalityDeleteFilesSize

LONG

Суммарный размер файлов equality-delete в байтах

equalityDeleteRecords

LONG

Суммарное количество записей в файлах equality-delete

positionDeleteFiles

LONG

Количество файлов position-delete

positionDeleteFilesSize

LONG

Суммарный размер файлов position-delete в байтах

positionDeleteRecords

LONG

Суммарное количество записей в файлах position-delete

dataFilesWithDeletes

LONG

Количество файлов данных, к которым применяются delete-файлы (position или equality)

dataFilesWithPositionDeletes

LONG

Количество файлов данных, к которым применяются position-delete файлы

dataFilesWithEqualityDeletes

LONG

Количество файлов данных, к которым применяются equality-delete файлы

avgDataFileSize

DOUBLE

Средний размер файла данных в байтах; для пустой таблицы значение равно 0

Переменные dataFilesWithDeletes, dataFilesWithPositionDeletes и dataFilesWithEqualityDeletes дополнительно требуют планирования сканирования таблицы, чтобы сопоставить delete-файлы с файлами данных, к которым они применяются. Их вычисление дороже остальных переменных этой группы.

Прочие переменные#

Переменная

Тип

Описание

snapshots

LONG

Количество snapshot в таблице

oldestSnapshotAge

LONG

Возраст самого старого snapshot в миллисекундах; для таблицы без snapshot значение равно 0

manifests

LONG

Количество файлов манифестов в текущем snapshot; для таблицы без snapshot значение равно 0

timeSinceLastUpdate

LONG

Время с момента последнего обновления таблицы в миллисекундах

Операторы#

Операторы сравнения#

Оператор

Описание

Пример

==

Равно

dataFiles == 100

!=

Не равно

snapshots != 0

>

Больше

dataFiles > 10

>=

Больше или равно

dataFilesSize >= 1GB

<

Меньше

avgDataFileSize < 1MB

<=

Меньше или равно

snapshots <= 5

Сравниваемые значения должны быть совместимых типов. Числовые типы LONG и DOUBLE совместимы между собой.

Арифметические операторы#

Оператор

Описание

Пример

+

Сложение

dataFiles + equalityDeleteFiles

-

Вычитание

dataFilesSize - 100MB

*

Умножение

dataFiles * 2

/

Деление

dataFilesSize / dataFiles

Деление всегда возвращает DOUBLE. При делении на ноль условие считается невыполненным (false).

Унарный минус поддерживается для числовых значений: -dataFiles.

Логические операторы#

Оператор

Описание

Пример

AND

Логическое И

dataFiles > 10 AND equalityDeleteFiles > 0

OR

Логическое ИЛИ

dataFiles > 1000 OR dataFilesSize > 10GB

!

Логическое НЕ

!property('maintenance.disabled', false)

Приоритет операторов#

От высшего к низшему:

  1. !, унарный -

  2. *, /

  3. +, -

  4. >, >=, <, <=

  5. ==, !=

  6. AND

  7. OR

Скобки () позволяют изменить порядок вычисления:

dataFiles > 100 AND (equalityDeleteFiles > 0 OR positionDeleteFiles > 0)

Функции#

Функции преобразования единиц#

Для фиксированных констант обычно удобнее и короче использовать литералы размера и длительности: 128MB, 1GB, 24h, 7d.

Функции KB, MB, GB, hours и days нужны, когда единицу измерения надо применить не к литералу, а к результату другого выражения типа LONG, например к значению из property(...) или к результату арифметики.

KB(n)#

Преобразует значение в килобайтах в байты.

KB(1) == 1024
countDataFiles(f -> f.size < KB(property('small-file-threshold-kb', 512))) > 10

MB(n)#

Преобразует значение в мегабайтах в байты.

MB(1) == 1048576
avgDataFileSize < MB(property('target-file-size-mb', 128))

GB(n)#

Преобразует значение в гигабайтах в байты.

GB(1) == 1073741824
dataFilesSize > GB(property('compaction-threshold-gb', 10))

hours(n)#

Преобразует значение в часах в миллисекунды.

hours(1) == 3600000
timeSinceLastUpdate > hours(property('stale-hours', 24))

days(n)#

Преобразует значение в днях в миллисекунды.

days(1) == 86400000
timeSinceLastUpdate > days(property('stale-days', 7))

countDataFiles(predicate)#

Подсчитывает количество файлов данных, удовлетворяющих заданному предикату.

Аргумент — лямбда-выражение с предикатом. Параметр лямбды имеет следующие поля:

Поле

Тип

Описание

size

LONG

Размер файла данных в байтах

records

LONG

Количество записей в файле данных

hasDeletes

BOOLEAN

К файлу данных применяются delete-файлы (position или equality)

hasPositionDeletes

BOOLEAN

К файлу данных применяются position-delete файлы

hasEqualityDeletes

BOOLEAN

К файлу данных применяются equality-delete файлы

Примеры:

countDataFiles(f -> f.size < 100MB)
countDataFiles(f -> f.size > 1GB)
countDataFiles(f -> f.records < 1000)
countDataFiles(f -> f.size < 512MB AND f.records < 10000)
countDataFiles(f -> f.size < MB(64) AND f.hasDeletes)
countDataFiles(f -> true)

Внутри предиката можно ссылаться на переменные и функции из внешнего контекста:

countDataFiles(f -> f.size < property('write.target-file-size-bytes', 512MB))

property(name, default)#

Возвращает значение свойства (property) Iceberg-таблицы.

Параметр

Описание

name

Имя свойства таблицы (строка)

default

Значение по умолчанию, если свойство отсутствует или невалидно

Тип возвращаемого значения определяется типом аргумента default. Если строковое значение свойства не может быть преобразовано в требуемый тип, то возвращается значение по умолчанию.

Примеры:

property('write.target-file-size-bytes', 512MB)
property('maintenance.disabled', false)
property('compaction.max-files', 100)
property('owner.team', 'default')

Примеры условий для политик обслуживания#

Compaction: слишком много мелких файлов#

Запуск compaction, когда в таблице более 100 файлов данных и средний размер файла меньше 128 МБ:

dataFiles > 100 AND avgDataFileSize < 128MB

Compaction: на основе целевого размера файлов из свойств таблицы#

Запуск compaction, когда более половины файлов данных меньше целевого размера, заданного в свойствах таблицы:

dataFiles > 10 AND countDataFiles(f -> f.size < property('write.target-file-size-bytes', 512MB)) / dataFiles > 0.5

Удаление delete-файлов#

Запуск compaction при накоплении equality-delete файлов:

equalityDeleteFiles > 10 OR equalityDeleteRecords > 100000

Запуск compaction, когда суммарный размер delete-файлов превышает 1 ГБ:

equalityDeleteFilesSize + positionDeleteFilesSize > 1GB

Запуск compaction, когда delete-файлы применяются к большому числу файлов данных:

dataFilesWithDeletes > 100000

Запуск compaction для мелких файлов данных, к которым применяются delete-файлы:

countDataFiles(f -> f.size < MB(64) AND f.hasDeletes) > 5

Комбинированное условие#

Запуск compaction при выполнении хотя бы одного из условий: много мелких файлов, большое количество delete-файлов или таблица давно не обновлялась — но только если maintenance не отключено через свойство таблицы:

!property('maintenance.disabled', false) AND (
    (dataFiles > 100 AND avgDataFileSize < 128MB)
    OR equalityDeleteFiles > 50
    OR (timeSinceLastUpdate > 3d AND snapshots > 50)
)

Expire snapshots: удаление устаревших snapshot#

Запуск expire-snapshots, когда в таблице более одного snapshot и самый старый snapshot старше 5 дней:

snapshots > 1 AND oldestSnapshotAge > days(5)

Условие с подсчётом файлов по размеру#

Запуск, если количество файлов данных размером более 1 ГБ превышает 10:

countDataFiles(f -> f.size > 1GB) > 10

Запуск, если количество файлов с малым количеством записей превышает 50:

countDataFiles(f -> f.records < 1000) > 50