Политики автоматического обслуживания объектов 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-выражением.
Политика может выполняться, когда соблюдены оба условия:
Политика активна (
active=true).Временное окно не задано, либо текущее время попадает хотя бы в один из сконфигурированных периодов.
Пример: обслуживание только по будням с 02:00 до 06:00 по московскому времени:
{
"time-zone-id": "Europe/Moscow",
"periods": [
{ "cron": "0 2 * * 1-5", "duration": 14400000 }
]
}
Если временное окно не задано и политика активна, то операции могут выполняться в любое время. Для промышленного использования рекомендуется всегда задавать временное окно.
Выполнение политик#
Каталог выполняет политики в фоновом цикле, разбитом на несколько фаз: поиск кандидатов, параллельная проверка условий запуска и отправка отобранных кандидатов в подсистему выполнения операций. Такая структура рассчитана на большие группы объектов: подавляющее большинство таблиц отсеивается дешёвой проверкой счётчиков ещё до того, как потребуется загружать их метаданные, а оставшиеся кандидаты проверяются параллельно.
Цикл работы политик#
Цикл оценки запускается периодически планировщиком каталога. На каждой итерации для каждой активной политики, для которой текущее время попадает в рамки временного окна (см. Временное окно), выполняются следующие шаги:
Аутентификация. Аутентифицируется заданный для политики
run-asпользователь, от имени которого выполняются все последующие действия.Поиск кандидатов. Каталог находит пары (объект, правило), для которых значение хотя бы одного из счётчиков активности (
commits-added,bytes-added,time-elapsed) превысило заданный порог. Этот шаг не требует загрузки метаданных Iceberg, поэтому масштабируется на большие группы объектов (см. Двухуровневая система триггеров).Проверка условий запуска. Для каждого кандидата загружаются метаданные таблицы, вычисляется выражение DSL из конфигурации триггера, результат (сработало или нет) записывается в журнал активности.
Запуск операций. Кандидаты, для которых условие 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При увеличении числа таблиц или количества правил рекомендуется следить за объёмом журнала и при необходимости уменьшать время хранения.
Параметры конфигурации#
Работа политик настраивается следующими параметрами каталога:
Параметр |
Описание |
По умолчанию |
Минимум |
|---|---|---|---|
|
Интервал между циклами работы политик. |
|
|
|
Максимальное количество одновременно выполняющихся операций обслуживания. |
|
|
|
Размер пула потоков, выполняющих проверку условий запуска. |
|
|
|
Срок хранения записей журнала активности. |
|
|
Справочник DSL#
Типы данных#
DSL поддерживает следующие типы данных:
Тип |
Описание |
Примеры значений |
|---|---|---|
|
Логический тип |
|
|
Целое число (64-битное) |
|
|
Число с плавающей точкой |
|
|
Строка |
|
Тип каждого выражения определяется на этапе компиляции. Несовместимость типов (например, сравнение строки с числом) приводит к ошибке компиляции.
В арифметических операциях, если один из операндов имеет тип DOUBLE, результат будет DOUBLE.
Операция деления (/) всегда возвращает DOUBLE.
Литералы#
Числовые литералы#
Целые числа записываются без суффиксов:
100
0
1048576
Числа с плавающей точкой содержат десятичную точку:
0.5
3.14
100.0
Литералы размера данных#
Для удобства работы с размерами файлов DSL поддерживает суффиксы, которые автоматически преобразуют значение в байты:
Суффикс |
Множитель |
Пример |
Значение в байтах |
|---|---|---|---|
|
1024 |
|
102 400 |
|
1024 × 1024 |
|
536 870 912 |
|
1024³ |
|
1 073 741 824 |
Примеры:
1KB == 1024
1MB == 1048576
1GB == 1073741824
512MB
100KB
Литералы длительности#
Для работы с временными интервалами DSL поддерживает суффиксы, которые преобразуют значение в миллисекунды:
Суффикс |
Множитель |
Пример |
Значение в миллисекундах |
|---|---|---|---|
|
60 000 (минуты) |
|
1 800 000 |
|
3 600 000 (часы) |
|
7 200 000 |
|
86 400 000 (дни) |
|
604 800 000 |
Примеры:
30m
24h
7d
Строковые литералы#
Строки заключаются в одинарные или двойные кавычки. Обратный слэш (\) используется для экранирования:
'hello'
"world"
'it\'s a test'
Логические литералы#
true
false
Переменные#
Переменные содержат характеристики текущей таблицы Iceberg.
Переменные, основанные на файлах данных#
Следующие переменные требуют загрузки метаданных таблицы и чтения всех манифестов текущего snapshot:
Переменная |
Тип |
Описание |
|---|---|---|
|
|
Количество файлов данных |
|
|
Суммарный размер файлов данных в байтах |
|
|
Суммарное количество записей в файлах данных |
|
|
Количество файлов equality-delete |
|
|
Суммарный размер файлов equality-delete в байтах |
|
|
Суммарное количество записей в файлах equality-delete |
|
|
Количество файлов position-delete |
|
|
Суммарный размер файлов position-delete в байтах |
|
|
Суммарное количество записей в файлах position-delete |
|
|
Количество файлов данных, к которым применяются delete-файлы (position или equality) |
|
|
Количество файлов данных, к которым применяются position-delete файлы |
|
|
Количество файлов данных, к которым применяются equality-delete файлы |
|
|
Средний размер файла данных в байтах; для пустой таблицы значение равно |
Переменные dataFilesWithDeletes, dataFilesWithPositionDeletes и dataFilesWithEqualityDeletes дополнительно требуют планирования сканирования таблицы, чтобы сопоставить delete-файлы с файлами данных, к которым они применяются. Их вычисление дороже остальных переменных этой группы.
Прочие переменные#
Переменная |
Тип |
Описание |
|---|---|---|
|
|
Количество snapshot в таблице |
|
|
Возраст самого старого snapshot в миллисекундах; для таблицы без snapshot значение равно |
|
|
Количество файлов манифестов в текущем snapshot; для таблицы без snapshot значение равно |
|
|
Время с момента последнего обновления таблицы в миллисекундах |
Операторы#
Операторы сравнения#
Оператор |
Описание |
Пример |
|---|---|---|
|
Равно |
|
|
Не равно |
|
|
Больше |
|
|
Больше или равно |
|
|
Меньше |
|
|
Меньше или равно |
|
Сравниваемые значения должны быть совместимых типов. Числовые типы LONG и DOUBLE
совместимы между собой.
Арифметические операторы#
Оператор |
Описание |
Пример |
|---|---|---|
|
Сложение |
|
|
Вычитание |
|
|
Умножение |
|
|
Деление |
|
Деление всегда возвращает DOUBLE. При делении на ноль условие считается невыполненным (false).
Унарный минус поддерживается для числовых значений: -dataFiles.
Логические операторы#
Оператор |
Описание |
Пример |
|---|---|---|
|
Логическое И |
|
|
Логическое ИЛИ |
|
|
Логическое НЕ |
|
Приоритет операторов#
От высшего к низшему:
!, унарный-*,/+,->,>=,<,<===,!=ANDOR
Скобки () позволяют изменить порядок вычисления:
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)#
Подсчитывает количество файлов данных, удовлетворяющих заданному предикату.
Аргумент — лямбда-выражение с предикатом. Параметр лямбды имеет следующие поля:
Поле |
Тип |
Описание |
|---|---|---|
|
|
Размер файла данных в байтах |
|
|
Количество записей в файле данных |
|
|
К файлу данных применяются delete-файлы (position или equality) |
|
|
К файлу данных применяются position-delete файлы |
|
|
К файлу данных применяются 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-таблицы.
Параметр |
Описание |
|---|---|
|
Имя свойства таблицы (строка) |
|
Значение по умолчанию, если свойство отсутствует или невалидно |
Тип возвращаемого значения определяется типом аргумента 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