Пять причин медленных Iceberg-запросов
Диагностика и лечение
Содержание

Дашборд грузится сорок секунд. Аналитик успевает сходить за кофе, вернуться, снова загрузить страницу и всё ещё ждать. Директор уже который раз спрашивает: «А что у нас с данными?». Данных заметно больше не стало, кластер выглядит здоровым, в приложении ничего не менялось.
Проблема не в движке запросов. Проблема — в физической структуре ваших Iceberg-таблиц.
Медленные запросы почти всегда сводятся к пяти структурным проблемам, спрятанным в метаданных и раскладке файлов: мелкие файлы, неправильный порядок сортировки, раздутые манифесты, старые снапшоты и несовпадение партиций. Оставленные без внимания, они усиливают друг друга: мелкие файлы рождают раздутые манифесты, старые снапшоты держат мёртвые данные на диске, а кривые партиции убивают файл-скиппинг. Таблица с двумя-тремя такими болезнями сразу может быть в 10–50 раз медленнее, чем те же данные, разложенные правильно.
Разберём каждую причину: как её найти, как подтвердить SQL-запросом и как починить.
Как планируется запрос в Iceberg (и где всё ломается)
Прежде чем лечить, полезно понять критический путь. Когда движок — Spark, Trino, Flink или любой другой, совместимый с Iceberg, — получает запрос, он не сканирует файлы данных сразу. Вместо этого он идёт по дереву метаданных.
Сначала движок читает текущий metadata.json, чтобы найти активный снапшот. Из снапшота читает manifest list — индекс всех манифестов таблицы. Каждая запись в manifest list содержит диапазоны значений партиций, поэтому планировщик может пропустить целые манифесты, чьи партиции не пересекаются с предикатом запроса. Это прунинг уровня 1.
Для выживших манифестов движок читает метаданные каждого файла: путь, значения партиций, min/max статистику по колонкам, количество null и строк. Он сверяет предикаты запроса с этой статистикой и отбрасывает файлы, в которых нет нужных строк. Это прунинг уровня 2 — и вот здесь решает порядок сортировки.
Планирование заканчивается, когда движок собрал список файлов для сканирования. Только после этого он шлёт GET-запросы в объектное хранилище за данными.
flowchart LR A[metadata.json] --> B[manifest list] B --> C[манифесты] C --> D[файлы данных] C -.уровень 1: партиции.-> X[(пропуск)] D -.уровень 2: min/max.-> Y[(пропуск)] D --> Z[объектное хранилище]
Каждая структурная проблема портит либо фазу планирования (слишком много манифестов и файлов для оценки), либо эффективность прунинга (неправильная сортировка, перекрывающиеся min/max, слишком широкие партиции).
Причина 1: взрыв мелких файлов
Симптомы: планирование запроса занимает 5–30 секунд до начала сканирования. Растёт число S3 GET-запросов. Счёт за облачное хранилище ползёт вверх быстрее, чем объём данных.
Почему так происходит: стриминговые джобы с короткими чекпоинтами, партицирование по высококардинальной колонке, частые MERGE INTO и конкурирующие писатели плодят недоразмеренные файлы. Flink-джоба с чекпоинтом каждые 60 секунд на 50 партициях создаёт 72 000 файлов в день — большинство меньше 5 МБ. Планировщик должен оценить каждый из них, каким бы селективным ни был запрос.
Почему это замедляет: каждый файл получает запись в манифесте. Больше файлов — больше записей, больше манифестов, больше S3 GET при планировании. Таблица на 200 000 мелких файлов планируется 15–30 секунд вместо менее чем секунды для 2 000 нормальных. Плюс мелкие файлы бьют по колоночной статистике: меньше строк в файле — шире и сильнее перекрывающиеся min/max диапазоны, хуже работает прунинг.
Диагностика
-- Spark: распределение размеров файловSELECT COUNT(*) AS total_files, ROUND(AVG(file_size_in_bytes) / 1048576, 1) AS avg_file_mb, ROUND(PERCENTILE(file_size_in_bytes, 0.5) / 1048576, 1) AS median_file_mb, ROUND(MIN(file_size_in_bytes) / 1048576, 1) AS min_file_mb, SUM(CASE WHEN file_size_in_bytes < 33554432 THEN 1 ELSE 0 END) AS files_under_32mbFROM catalog.db.my_table.files;
-- Trino: распределение размеров файловSELECT COUNT(*) AS total_files, ROUND(AVG(file_size_in_bytes) / 1048576.0, 1) AS avg_file_mb, ROUND(MIN(file_size_in_bytes) / 1048576.0, 1) AS min_file_mb, COUNT_IF(file_size_in_bytes < 33554432) AS files_under_32mbFROM "catalog"."db"."my_table$files";
-- DuckDB: плотность файлов через число записей-- (размер файла DuckDB пока не выставляет в метаданных)SELECT COUNT(*) AS live_files, ROUND(AVG(record_count), 0) AS avg_records_per_file, MIN(record_count) AS min_recordsFROM iceberg_metadata(sq.streams.yandex_metrika_101, allow_moved_paths = true)WHERE content = 'DATA' AND status = 'ADDED';
Если avg_file_mb ниже 64 или files_under_32mb больше 30% от всех файлов — у вас проблема мелких файлов. В DuckDB смотрите на avg_records_per_file: много файлов с низкой плотностью записей — та же картина мелких файлов.
Причина 2: неправильный или отсутствующий порядок сортировки
Симптомы: селективные WHERE всё равно сканируют почти всю таблицу. Объём просканированных данных велик относительно размера результата. Метрики файл-скиппинга показывают низкий прунинг.
Почему так происходит: таблица создана без порядка сортировки или с порядком, который не совпадает с тем, как её реально запрашивают. Таблица, отсортированная по created_date, бесполезна для запросов с фильтром по customer_id. Без эффективной сортировки min/max диапазон каждой колонки в каждом файле покрывает весь домен значений — на уровне 2 не отсекается ничего.
Почему это замедляет: порядок сортировки — самый сильный рычаг производительности Iceberg. Отсортированные таблицы сканируют до 51% меньше данных на запрос. Если сортировка совпадает с реальными фильтрами продакшена, файл-скиппинг отсекает 80–95% файлов ещё на этапе планирования и даёт до 12 раз быстрее запросы.
Диагностика
Проверяем текущую сортировку:
-- Spark: проверка порядка сортировки таблицыDESCRIBE EXTENDED catalog.db.my_table;-- Ищите 'Sort Columns' в выводе
Проверяем, реально ли min/max диапазоны дают прунинг. Широкие перекрывающиеся диапазоны говорят о плохой или отсутствующей сортировке:
-- Spark: диапазоны значений колонки по файламSELECT COUNT(*) AS total_files, MIN(lower_bounds['customer_id']) AS global_min, MAX(upper_bounds['customer_id']) AS global_max, COUNT(DISTINCT lower_bounds['customer_id']) AS distinct_lower_boundsFROM catalog.db.my_table.files;
-- DuckDB: диапазоны значений колонки по файламSELECT column_name, COUNT(*) AS files_with_stats, COUNT(DISTINCT lower_bound) AS distinct_lower_bounds, MIN(lower_bound) AS global_min, MAX(upper_bound) AS global_maxFROM iceberg_column_stats(sq.streams.yandex_metrika_101, allow_moved_paths = true)WHERE column_name = 'customer_id'GROUP BY column_name;
Если distinct_lower_bounds близко к 1 (у всех файлов одинаковая нижняя граница) — данные не отсортированы по этой колонке и прунинг работать не сможет. То же правило работает и для результата iceberg_column_stats.
Причина 3: раздутые манифесты
Симптомы: планирование медленное даже на маленьких результатах. Фаза планирования доминирует над общим временем запроса. «Холодные» запросы заметно медленнее повторных.
Почему так происходит: каждый коммит — каждый чекпоинт стриминга, пакетная запись, компакция или MERGE — создаёт минимум один новый манифест. Стриминговая таблица с чекпоинтом раз в 10 минут накапливает более 4 300 манифестов в месяц. Каждый манифест может отслеживать лишь горстку файлов, и слой метаданных превращается в кашу.
Почему это замедляет: при планировании движок читает все манифесты, на которые ссылается текущий снапшот. Сотни мелких манифестов — сотни S3 GET только ради построения списка файлов. Кэш манифестов помогает повторным запросам, но холодный старт — первый запрос после рестарта кластера или по давно не тронутой таблице — платит полную цену I/O. Таблицы с 500+ манифестами стабильно планируются 5–15 секунд.
Диагностика
-- Spark: количество и фрагментация манифестовSELECT COUNT(*) AS manifest_count, ROUND(AVG(added_data_files_count), 1) AS avg_files_per_manifest, SUM(CASE WHEN added_data_files_count <= 2 THEN 1 ELSE 0 END) AS tiny_manifests, ROUND(AVG(length) / 1024, 1) AS avg_manifest_size_kbFROM catalog.db.my_table.manifests;
-- Trino: количество манифестовSELECT COUNT(*) AS manifest_count, ROUND(AVG(added_data_files_count), 1) AS avg_files_per_manifestFROM "catalog"."db"."my_table$manifests";
-- DuckDB: количество и фрагментация манифестовSELECT COUNT(DISTINCT manifest_path) AS manifest_count, ROUND(COUNT(*)::DOUBLE / NULLIF(COUNT(DISTINCT manifest_path), 0), 1) AS avg_files_per_manifestFROM iceberg_metadata(sq.streams.yandex_metrika_101)WHERE content = 'DATA' AND status = 'ADDED';
Если manifest_count больше 100 или avg_files_per_manifest ниже 10 — переписывание манифестов улучшит время планирования. Если больше половины манифестов «крошечные» (следят за 1–2 файлами), эффект будет драматическим.
Причина 4: старые снапшоты держат мёртвые данные
Симптомы: объём хранилища растёт быстрее, чем скорость приёма данных. Старые файлы, которые давно должны были удалиться, всё ещё лежат на диске. S3 LIST тормозит. Файлы метаданных неудержимо толстеют.
Почему так происходит: каждая запись, обновление и удаление в Iceberg создают новый снапшот. Снапшоты ссылаются на manifest list, те — на манифесты, те — на файлы данных. Пока снапшот явно не истёк, все его файлы сохраняются — даже те, что логически заменены компакцией или переписаны MERGE. Без expire_snapshots таблица копит бесконечную историю мёртвых файлов и метаданных.
Почему это замедляет: старые снапшоты не тормозят сами запросы напрямую (они читают только текущий снапшот), но создают косвенное давление. Мёртвые файлы раздувают ответы S3 LIST и замедляют обслуживание. Растущий metadata.json дольше парсится. В крайних случаях объём мёртвых файлов вызывает троттлинг S3 (ошибки 503 SlowDown), который бьёт и по чтению того же префикса. Метаданные, которые должны жить в памяти координатора, вытесняются на диск, замедляя планирование.
Диагностика
-- Spark: накопление снапшотовSELECT COUNT(*) AS total_snapshots, MIN(committed_at) AS oldest_snapshot, MAX(committed_at) AS newest_snapshot, DATEDIFF(MAX(committed_at), MIN(committed_at)) AS snapshot_span_daysFROM catalog.db.my_table.snapshots;
-- Trino: накопление снапшотовSELECT COUNT(*) AS total_snapshots, MIN(committed_at) AS oldest_snapshot, MAX(committed_at) AS newest_snapshotFROM "catalog"."db"."my_table$snapshots";
-- DuckDB: накопление снапшотовSELECT COUNT(*) AS total_snapshots, MIN(to_timestamp(timestamp_ms / 1000)) AS oldest_snapshot, MAX(to_timestamp(timestamp_ms / 1000)) AS newest_snapshotFROM iceberg_snapshots(sq.streams.yandex_metrika_101);
Если total_snapshots больше 1 000 или snapshot_span_days больше 30 — вы держите слишком много истории. Для большинства продакшен-таблиц хватает 5–7 дней снапшотов для time-travel.
Причина 5: несовпадение партиций
Симптомы: запросы сканируют куда больше партиций, чем ожидалось. Партиционный прунинг не работает. В одних партициях миллионы строк, в других — десятки. Добавление новых фильтров не сокращает время сканирования.
Почему так происходит: таблица партицирована не по той колонке, не на том уровне гранулярности, или распределение данных изменилось с момента создания. Типичные ошибки: партицирование по дню при запросах по региону, часовые партиции на низкообъёмной таблице (тысячи крошечных партиций) или identity-партицирование по высококардинальной колонке. Partition evolution в Iceberg решает это без переписывания данных, но о нём мало кто знает.
Почему это замедляет: партиционный прунинг — фильтр уровня 1, самый грубый и дешёвый вид отбрасывания данных. Если запросы фильтруют по колонкам вне партиционного спекта, движок читает все манифесты и оценивает все файлы. Слишком мелкие партиции дают мало строк и много мелких файлов. Слишком широкие — в каждой партиции полно нерелевантных данных, которые приходится сканировать.
Диагностика
-- Spark: распределение по партициямSELECT partition, COUNT(*) AS file_count, SUM(record_count) AS total_records, ROUND(SUM(file_size_in_bytes) / 1073741824, 2) AS partition_gb, ROUND(AVG(file_size_in_bytes) / 1048576, 1) AS avg_file_mbFROM catalog.db.my_table.filesGROUP BY partitionORDER BY file_count DESCLIMIT 20;
-- DuckDB: профиль партиционных полейSELECT partition_field_name, partition_source_columns, partition_field_transform, MIN(lower_bound) AS global_min, MAX(upper_bound) AS global_maxFROM iceberg_partition_stats(sq.streams.yandex_metrika_101)GROUP BY partition_field_name, partition_source_columns, partition_field_transform;
Ищите партиции с очень большим числом файлов (тысячи файлов на партицию) и партиции с крошечными partition_gb (меньше 100 МБ). И то и другое говорит о несовпадении. Проверьте и то, совпадает ли колонка партиции с самыми частыми предикатами: если всегда фильтруете по region, а таблица партицирована по event_date, прунинг помогает только запросам по дате.
Полный диагностический чек-лист
Прогоните эти запросы, чтобы собрать полный профиль здоровья любой медленной таблицы:
-- Spark: комплексная проверка здоровья таблицы -- 1. Статистика файловSELECT COUNT(*) AS total_files, ROUND(AVG(file_size_in_bytes) / 1048576, 1) AS avg_file_mb, SUM(CASE WHEN file_size_in_bytes < 33554432 THEN 1 ELSE 0 END) AS small_files, ROUND(SUM(file_size_in_bytes) / 1073741824, 2) AS total_gbFROM catalog.db.my_table.files; -- 2. Здоровье манифестовSELECT COUNT(*) AS manifest_count, ROUND(AVG(added_data_files_count), 1) AS avg_files_per_manifestFROM catalog.db.my_table.manifests; -- 3. Накопление снапшотовSELECT COUNT(*) AS total_snapshots, MIN(committed_at) AS oldest, MAX(committed_at) AS newestFROM catalog.db.my_table.snapshots; -- 4. Перекос партицийSELECT partition, COUNT(*) AS file_count, ROUND(AVG(file_size_in_bytes) / 1048576, 1) AS avg_file_mbFROM catalog.db.my_table.filesGROUP BY partitionORDER BY file_count DESCLIMIT 10;
Это ваш скоринг-кард здоровья: размеры файлов говорят о необходимости компакции, число манифестов — о накладных расходах планирования, число снапшотов — о раздувании метаданных, распределение партиций — о согласованности раскладки.
Ручное лечение: все пять причин одним скриптом
Если хотите управлять оптимизацией сами, вот продакшен-скрипт обслуживания, закрывающий все пять причин. Планируйте его через Airflow, Dagster или любой оркестратор:
-- Шаг 1: компакция и сортировка мелких файлов (ежедневно для стриминговых таблиц)CALL catalog.system.rewrite_data_files( table => 'db.my_table', strategy => 'sort', sort_order => 'customer_id ASC NULLS LAST, event_date ASC NULLS LAST', options => map( 'target-file-size-bytes', '536870912', 'min-file-size-bytes', '67108864', 'min-input-files', '5', 'partial-progress.enabled', 'true', 'partial-progress.max-commits', '10' )); -- Шаг 2: консолидация манифестов (еженедельно)CALL catalog.system.rewrite_manifests( table => 'db.my_table'); -- Шаг 3: удаление старых снапшотов (ежедневно)CALL catalog.system.expire_snapshots( table => 'db.my_table', older_than => TIMESTAMP '2026-07-26 00:00:00', retain_last => 5, stream_results => true); -- Шаг 4: очистка осиротевших файлов (еженедельно)CALL catalog.system.remove_orphan_files( table => 'db.my_table', older_than => TIMESTAMP '2026-07-29 00:00:00');
В Trino компакция файлов использует процедуру optimize:
-- Trino: компакция файлов по партицииALTER TABLE "catalog"."db"."my_table" EXECUTE optimize WHERE partition_col = 'value';
Trino не умеет rewrite_manifests и expire_snapshots нативно — эти операции требуют Spark. Если основной движок у вас Trino, понадобится Spark как sidecar для обслуживания метаданных.
Закодируйте политики ретеншена и размера файлов прямо в свойствах таблицы, чтобы любой движок и планировщик их уважал:
-- Политики ретеншена и размера файлов как свойства таблицыALTER TABLE catalog.db.my_table SET TBLPROPERTIES ( 'history.expire.max-snapshot-age-ms' = '604800000', 'history.expire.min-snapshots-to-keep' = '5', 'write.target-file-size-bytes' = '268435456');
Для эволюции партиций Iceberg поддерживает недеструктивные изменения спекта:
-- Spark: эволюция схемы партицированияALTER TABLE catalog.db.my_tableADD PARTITION FIELD bucket(16, customer_id); -- Замена слишком гранулярной партицииALTER TABLE catalog.db.my_tableREPLACE PARTITION FIELD hour(event_ts)WITH day(event_ts) AS event_day;
Эволюция партиций — операция только над метаданными: существующие данные сохраняют исходную раскладку, новые записи идут по новому спекту. Чтобы применить новую раскладку к старым данным, после эволюции запустите sort-компакцию.
Ручной путь работает для нескольких таблиц. Проблемы начинаются в масштабе: каким таблицам компакция нужна прямо сегодня? Какая сортировка нужна каждой? Манифесты копятся быстрее, чем еженедельная перезапись их чистит? Схема партиций оторвалась от реальных паттернов запросов? Эти вопросы превращают простой cron-джоб в систему оптимизации на множество таблиц и движков.
Замер результата
После применения фиксов — любым путём — перепроверьте улучшение теми же диагностическими запросами и следите за метриками:
- Средний размер файла: должен быть 256–512 МБ после компакции. Следите за откатом после стриминговых коммитов.
- Число манифестов: должно упасть в 5–10 раз после перезаписи. Отслеживайте недельный прирост.
- Время планирования: мерьте от отправки запроса до первой строки. Ожидайте ускорения в 2–5 раз от компакции и перезаписи манифестов вместе.
- Объём просканированных данных: сравнивайте байты на запрос до и после оптимизации сортировки. Ожидайте сокращения на 50–90% для запросов, совпадающих с сортировкой.
- Число снапшотов: должно стабилизироваться на целевом ретеншене (5–10 снапшотов для большинства таблиц).
Вывод
Медленные Iceberg-запросы — не мистика. Это структурная проблема с пятью конкретными причинами, каждая из которых диагностируется SQL-запросом и чинится штатными процедурами Iceberg. Сложность не в том, что делать, а в том, чтобы делать это непрерывно и для каждой таблицы, пока данные и паттерны запросов меняются.
Диагностические запросы дают мгновенную видимость того, какие таблицы деградировали и почему. Ручные процедуры — инструменты для починки. Компакция, порядок сортировки, перезапись манифестов, жизненный цикл снапшотов и согласованность партиций — это не пять отдельных проблем, а пять граней здоровья таблицы. Лечите их вместе, измеряйте непрерывно — и ваши Iceberg-запросы станут такими же быстрыми, как задумано форматом.