Manifest Files и Manifest List: как Iceberg планирует scan без листинга
Алгоритм Scan Planning в Apache Iceberg: manifest-level pruning через partition summary bounds, file-level pruning через InclusiveMetricsEvaluator, residual-фильтры, bin-packing в CombinedScanTask, компромисс write.metadata.metrics.default и инспекция через системную таблицу manifests.
От «что это» к «как это работает»¶
Прошлый урок разобрал дерево метаданных Iceberg на уровне структуры: какие поля хранятся в metadata.json, что такое manifest list и manifest file, как устроен атомарный commit. Этого достаточно, чтобы понимать, почему Iceberg даёт ACID-гарантии. Но остался открытым отдельный, не менее важный вопрос: как именно движок использует это дерево, чтобы за миллисекунды решить, какие из миллионов файлов таблицы стоит читать для конкретного запроса - без единого обращения к листингу object storage.
Этот урок - детальный разбор именно этого алгоритма, который называется Scan Planning. Если прошлый урок объяснял архитектуру хранения метаданных, то этот объясняет архитектуру их использования. Это разделение специально проведено так, чтобы не дублировать материал: там, где прошлый урок останавливался на уровне «вот это поле существует и хранит статистику», этот урок идёт дальше - «вот конкретный алгоритм, который использует это поле, чтобы принять решение читать файл или нет, и вот почему это решение всегда безопасно».
Практическая мотивация изучать именно этот уровень детализации простая: знание архитектуры дерева метаданных само по себе не объясняет, почему конкретный production-запрос вдруг стал планироваться 10 секунд вместо 100 миллисекунд. Ответ на такой вопрос всегда лежит в одном из четырёх этапов алгоритма, разбираемых в этом уроке - и умение быстро определить, в каком именно этапе кроется проблема (через EXPLAIN, Spark UI или прямой запрос к системным таблицам), напрямую конвертируется в способность диагностировать и исправлять реальные инциденты производительности, а не гадать вслепую.
Проклятие listStatus: как ищут файлы традиционные движки¶
Чтобы оценить масштаб решения, которое даёт Iceberg, нужно явно зафиксировать, как работает планирование в традиционной Hive-style архитектуре - той самой, что разобрана в первом уроке модуля как источник архитектурных болей.
Hive Metastore хранит метаданные о партициях таблицы (путь к директории каждой партиции, её схему), но принципиально не хранит список файлов внутри этих партиций. Это сознательное архитектурное решение, появившееся в эпоху, когда таблицы были на порядки меньше, а HDFS NameNode дешево обслуживал листинг директорий в памяти. Когда движок (Hive, Spark, Trino) планирует запрос к Hive-таблице, он обязан выполнить отдельный шаг - получить список файлов внутри каждой релевантной партиции через вызов listStatus к файловой системе или ListObjectsV2 к object storage.
Этот шаг - не теоретическая деталь, а конкретная операционная боль, знакомая каждому, кто администрировал Hive-кластер: команда MSCK REPAIR TABLE (синхронизация метастора с реальным содержимым файловой системы после добавления новых партиций вручную) в больших таблицах с десятками тысяч партиций может выполняться часами именно потому, что под капотом она выполняет листинг каждой партиции по отдельности. Даже без MSCK REPAIR, обычное планирование запроса к таблице с тысячами партиций вынуждено повторять этот листинг при каждом обращении (если результат не закеширован движком между запросами в рамках одной долгоживущей сессии).
Почему Object Storage - это не файловая система¶
Листинг в HDFS и листинг в S3-совместимом object storage (MinIO, Ceph RGW, AWS S3) - принципиально разные по стоимости операции, и это различие лежит в основе всей проблемы.
В HDFS NameNode хранит дерево директорий целиком в памяти - listStatus для директории с тысячей файлов выполняется за единицы миллисекунд, потому что это просто обращение к структуре данных в RAM, без сетевого round-trip к каждому файлу.
В object storage не существует понятия «директория» в файловом смысле - есть только плоское пространство ключей (bucket/key), и операция ListObjectsV2 эмулирует иерархию через префиксный поиск. Важное технической ограничение этой операции: она возвращает не более 1000 ключей за один вызов, независимо от того, сколько объектов реально подходит под префикс. Если под префиксом country=DE/ лежит 50 000 файлов, движку придётся сделать 50 последовательных (или умеренно параллельных) вызовов ListObjectsV2, каждый раз передавая токен пагинации (continuation-token) от предыдущего вызова.
Помимо самой пагинации, у S3-совместимых хранилищ есть практический предел пропускной способности запросов на один префикс - например, опубликованная для AWS S3 ориентировочная норма около 5500 GET/HEAD запросов в секунду на партиционированный префикс (после внутренней балансировки нагрузки по диапазонам ключей). LIST-запросы по стоимости ближе к GET, и при конкурентной нагрузке (несколько аналитиков одновременно запускают тяжёлые запросы к одной таблице с тысячами партиций) суммарное число листинговых запросов может приблизиться к этому потолку, что добавляет дополнительную задержку (throttling, retry с экспоненциальной задержкой) сверх базовой латентности отдельного вызова.
Тяжёлое планирование: когда planning дороже самого scan¶
Соединим два фактора - множество партиций и пагинированный листинг - в конкретный числовой пример, который объясняет фразу «10 минут планирования ради 10 секунд исполнения», знакомую многим инженерам, работавшим с крупными Hive-таблицами.
Таблица: 5 лет данных, партиционирование по (country, event_date)
50 стран × 1825 дней = 91 250 партиций
В среднем 8 файлов на партицию = 730 000 файлов всего
Запрос: SELECT * FROM sales WHERE country = 'DE'
Релевантно: 1 страна × 1825 дней = 1825 партиций (после partition pruning
на уровне ПУТИ - это Hive ещё может сделать без листинга,
потому что страна явно в пути)
Но чтобы узнать ФАЙЛЫ внутри каждой из 1825 партиций:
1825 партиций × 1 LIST-вызов (если файлов <1000 на партицию) = 1825 запросов
1825 запросов × ~80 мс средней латентности = 146 000 мс ≈ 146 секунд
Даже с параллелизацией в 30 потоков на стороне Driver:
146 секунд / 30 ≈ 5 секунд - уже немало для ТОЛЬКО листинга,
а реальные кластеры часто хуже параллелизируют листинг,
чем кажется на бумаге, из-за throttling и retry на холодном старте.
Этот пример соответствует расчёту из первого урока модуля, но смотрит на проблему с другой стороны: там оценивался листинг всех партиций таблицы для построения первичного списка; здесь - листинг даже уже отфильтрованного подмножества партиций, которое Hive смог определить по пути без обращения к файлам. Проблема не исчезает с улучшением partition pruning на уровне путей - она просто становится на порядок меньше, но не равной нулю, потому что листинг файлов внутри партиций остаётся неизбежным шагом архитектуры «метаданные о партициях есть, метаданных о файлах нет».
Главная цель урока: из I/O-bound в CPU-bound планирование¶
Сформулируем цель, к которой ведёт весь оставшийся материал: Iceberg должен превратить операцию «определить, какие файлы читать» из операции, доминирующей по стоимости сетевыми round-trip'ами к object storage (I/O-bound), в операцию, доминирующую по стоимости вычислениями над уже загруженными в память данными (CPU-bound). Дерево метаданных, разобранное в прошлом уроке, - это структура данных, которая делает такое превращение возможным; алгоритм Scan Planning, разбираемый в этом уроке, - конкретный способ её использовать для этого превращения.
Разделение обязанностей: Manifest List vs Manifest Files¶
Прежде чем переходить к алгоритму, стоит ещё раз чётко развести роли двух уровней метаданных - не с точки зрения «что в них хранится» (это разобрано в прошлом уроке), а с точки зрения «какую конкретную задачу пруннинга каждый из них решает».
Manifest List как карта снапшота: zoom на partition summary¶
Поле, которое делает manifest-level pruning возможным - partitions (в спецификации называется partition_field_summary) внутри каждой записи manifest list. Это не статистика по отдельным файлам - это агрегированная статистика по партиционным полям для всех файлов, перечисленных в данном конкретном манифесте.
// Одна запись manifest list, поле "partitions" развёрнуто подробно
{
"manifest_path": "s3a://lakehouse/.../manifest-aa.avro",
"manifest_length": 184320,
"partition_spec_id": 0,
"added_snapshot_id": 5847203991234567890,
"added_data_files_count": 312,
"existing_data_files_count": 688,
"deleted_data_files_count": 0,
"partitions": [
{
"contains_null": false,
"contains_nan": false,
"lower_bound": 19888,
"upper_bound": 19891
}
]
}
Поле partitions здесь - массив из одного элемента, потому что партиция в этом примере - одно поле (days(event_ts)). lower_bound: 19888 и upper_bound: 19891 - это число дней с эпохи Unix, то есть диапазон значений партиции по всем файлам, входящим в этот манифест, составляет всего 4 дня. Если запрос фильтрует по event_ts >= '2024-08-01', а 2024-08-01 соответствует дню 19932 (что больше upper_bound: 19891), движок может сделать однозначный вывод: ни один файл в этом манифесте не может содержать релевантных строк, и манифест отбрасывается без открытия - то есть без единого GET-запроса к самому манифесту.
Важная деталь, объясняющая значение слова contains_null и contains_nan отдельно от диапазона: некоторые предикаты (IS NULL, IS NOT NULL) не могут быть оценены через числовой диапазон вообще - для них движку нужно явно знать, встречаются ли NULL-значения партиции внутри манифеста. Это отдельный флаг, а не часть числового диапазона, именно потому что NULL логически не входит ни в какой числовой интервал.
Manifest File как индекс data files: что используется для file-level pruning¶
На уровне manifest file (детально разобранном в прошлом уроке - статус ADDED/EXISTING/DELETED, поля data_file) для file-level pruning используются конкретно три группы полей: точное значение partition файла (не диапазон, а конкретное значение - например, {event_ts_day: 19889} для одного файла), и поколоночные карты lower_bounds/upper_bounds/null_value_counts/nan_value_counts/value_counts. Принципиальное отличие от manifest list: здесь статистика не агрегирована по множеству файлов - она относится к одному конкретному физическому файлу, что позволяет на следующем шаге отбрасывать отдельные файлы, а не только целые группы.
Partition Pruning vs Predicate Pushdown: две вещи, которые часто путают¶
Прежде чем переходить к алгоритму планирования, стоит явно развести два термина, которые в обсуждениях Spark и Iceberg часто используются взаимозаменяемо, хотя означают разные механизмы на разных уровнях стека.
Partition Pruning - это специфичный для table format механизм, разобранный в этом уроке: отбрасывание целых файлов или манифестов на основе того, какой партиции они принадлежат, без открытия самих файлов. Это происходит на уровне Iceberg (Этапы 2 и 3 алгоритма), до того, как Spark вообще получает список файлов на чтение.
Predicate Pushdown - более общий термин, означающий передачу предиката запроса как можно ближе к источнику данных, чтобы избежать чтения и последующей фильтрации ненужных строк в памяти Spark. Это происходит на нескольких независимых уровнях одновременно: Catalyst передаёт предикат источнику данных (Iceberg Scan) как Expression; Iceberg использует часть этого предиката для partition/metrics pruning (то, что разобрано в этом уроке); а та часть, что не была разрешена через метаданные (residual expression), передаётся дальше вниз - к Parquet reader, который использует её для Row Group filtering (пропуск целых row group по их собственной min/max статистике) и Dictionary filtering (для строковых/категориальных колонок со словарным кодированием).
Эта цепочка объясняет, почему «predicate pushdown» как термин в обсуждениях Spark обычно шире, чем то, что разбирается в этом уроке: Iceberg отвечает только за первый сегмент цепочки (отбор файлов/манифестов), а дальнейшая фильтрация внутри отобранных файлов - заслуга формата Parquet и движка выполнения Spark, уже не связанная с table format напрямую. Путать эти уровни - частая причина неверной диагностики: если запрос всё равно читает «слишком много данных», стоит сначала проверить Manifests/Data Files skipped в EXPLAIN (зона ответственности Iceberg), а затем отдельно - метрики Parquet reader в Spark UI (зона ответственности файлового формата).
Особый случай: bucket() и truncate() трансформы при pruning¶
Не все partition transform одинаково хорошо работают с описанным алгоритмом manifest-level pruning. Разница принципиальна и часто становится сюрпризом для инженеров, которые интуитивно ожидают одинаковое поведение от любого partition transform.
Для временных трансформов (years, months, days, hours) и для truncate(N) на числовых/строковых колонках сохраняется монотонность: если исходное значение колонки больше, то и значение трансформа не меньше. Это значит, что диапазон lower_bound/upper_bound в partition summary манифеста по-прежнему осмысленно сравнивается с диапазоном предиката через операторы >, <, >=, <= - то есть range-предикаты (event_ts >= '2024-06-01') эффективно проходят manifest-level pruning.
Для bucket(N, col) ситуация другая: бакет вычисляется как hash(col) % N, и хэш-функция не сохраняет порядок - значения col=5 и col=1005 могут попасть в один и тот же бакет, а соседние по значению col=5 и col=6 - в совершенно разные бакеты. Из-за этого диапазон lower_bound/upper_bound бакета в partition summary малополезен для range-предикатов (user_id > 1000) - такой предикат почти никогда не позволяет отбросить манифест, потому что числа бакетов разбросаны без видимой связи с диапазоном исходной колонки. Зато bucket-партиционирование отлично работает для equality-предикатов (user_id = 12345) - движок детерминированно вычисляет hash(12345) % N, получает точный номер бакета и может отбросить все манифесты/файлы с другими номерами бакетов, что эквивалентно прямому попаданию в правильный «ящик».
Этот нюанс - прямое практическое следствие алгоритма, разобранного в этом уроке, и одна из причин, почему выбор partition transform (тема следующего урока модуля) нельзя делать механически: тип трансформа должен соответствовать тому, какие именно предикаты (range или equality) реально доминируют в нагрузке запросов к конкретной таблице.
Anatomy of Scan Planning: алгоритм от SQL до Parquet¶
Это центральная часть урока - детальный, шаг за шагом, разбор того, что происходит между моментом, когда Spark Catalyst строит логический план для df.filter(col("country") == "DE"), и моментом, когда конкретные Parquet-файлы передаются executor'ам на чтение.
Этап 1: точка входа¶
Этот этап целиком разобран в прошлом уроке как часть read path - Driver получает от каталога путь к текущему metadata.json, читает его, находит снапшот по current-snapshot-id и его поле manifest-list. Единственное, что важно подчеркнуть здесь отдельно: к этому моменту Driver сделал ровно два сетевых запроса (catalog + GET metadata.json), независимо от того, сколько партиций или файлов содержит таблица - оба запроса имеют постоянную стоимость.
Дополнительная деталь, которая становится важна именно в контексте Scan Planning (а не только commit-протокола, как в прошлом уроке): на этом же шаге Driver получает из metadata.json текущую партиционную спецификацию (partition-specs + default-spec-id) и схему (schemas + current-schema-id). Эти два объекта - не просто справочная информация, а необходимый входной материал для следующего этапа: чтобы понять, какие именно поля предиката являются партиционными (а значит, могут быть проверены через manifest-level pruning) и каким ID колонки соответствует каждое имя в условии запроса (вспомните ID-based schema tracking из первого урока модуля), движку нужна именно эта пара объектов, прочитанная на Этапе 1.
Этап 2: Manifest-level Pruning¶
Driver читает manifest list (один GET-запрос - но это уже не O(N) от числа файлов, а O(1) относительно их числа, потому что manifest list - один файл) и применяет к каждой записи manifest-level фильтр, построенный на основе предиката запроса, спроецированного на партиционные поля.
# Концептуальный (упрощённый) псевдокод того, что происходит внутри
# движка на этапе manifest-level pruning - НЕ реальный API,
# а иллюстрация логики принятия решения
def manifest_might_match(manifest_entry, partition_predicate):
"""
Консервативная проверка: может ли манифест содержать
хотя бы один файл, удовлетворяющий предикату.
Важнейшее свойство: эта функция НИКОГДА не возвращает False,
если манифест на самом деле может содержать совпадения -
ложноотрицательный результат здесь недопустим (это бы означало
потерю данных в результате запроса). Ложноположительные
результаты (манифест оставлен, хотя совпадений в итоге нет)
допустимы - они просто означают, что File-level Pruning
на следующем шаге сделает дополнительную работу.
"""
for field_summary, predicate_bound in zip(
manifest_entry.partitions, partition_predicate.bounds
):
if field_summary.upper_bound < predicate_bound.lower_required:
return False # ROWS_CANNOT_MATCH - однозначно отбросить
if field_summary.lower_bound > predicate_bound.upper_required:
return False # ROWS_CANNOT_MATCH - однозначно отбросить
return True # ROWS_MIGHT_MATCH - оставить для следующего этапа
Ключевое архитектурное свойство этой проверки, которое стоит выделить отдельно: она консервативна (conservative). Это означает не «иногда угадывает», а строгую гарантию направления ошибки: функция может ошибочно оставить манифест, который в итоге не содержит совпадений (это просто потерянная эффективность, но не потерянная корректность), но никогда не может ошибочно отбросить манифест, который реально содержит совпадения (это была бы потеря данных в результате запроса - недопустимая ошибка). Эта же гарантия применяется на всех последующих этапах пруннинга, и именно она позволяет агрессивно оптимизировать planning без риска получить неполный результат.
Манифесты, не прошедшие эту проверку, отбрасываются без единого обращения к самому manifest-файлу - их пути просто не попадают в список «манифестов для чтения на Этапе 3». Это и есть буквальная реализация фразы из плана урока: «манифесты, не прошедшие фильтр, даже не скачиваются с S3».
Этап 3: File-level Pruning¶
Для манифестов, прошедших Этап 2, Driver выполняет GET-запросы (параллельно, по числу прошедших манифестов - что, для типичной таблицы, число на порядки меньше, чем число файлов) и применяет уже знакомую логику, но на уровне отдельных записей data_file.
Здесь происходит две независимые проверки для каждой записи манифеста, и обе должны пройти, чтобы файл остался в плане сканирования:
-
Точная проверка по значению партиции - значение
partitionконкретного файла сравнивается с предикатом напрямую (не через диапазон, как на Этапе 2, а через точное равенство/сравнение), потому что на уровне отдельного файла партиция - это конкретное значение, а не диапазон. -
Проверка по статистике колонок (Metrics Pruning) - для предикатов на любых колонках (не только партиционных - и это принципиальное расширение по сравнению с Этапом 2, который работает только с партиционными полями), движок сравнивает
lower_bounds/upper_boundsфайла с границами предиката, используя ту же консервативную логику ROWS_MIGHT_MATCH / ROWS_CANNOT_MATCH.
# Концептуальный псевдокод file-level metrics pruning -
# работает для ЛЮБОЙ колонки, не только партиционной
def file_might_match_predicate(data_file, predicate):
column_id = predicate.column_id
if predicate.is_null_check:
# Для IS NULL: если null_value_counts == 0, файл точно
# не содержит NULL в этой колонке - можно отбросить
if predicate.kind == "IS_NULL":
return data_file.null_value_counts.get(column_id, 0) > 0
if predicate.kind == "IS_NOT_NULL":
total_rows = data_file.value_counts.get(column_id, 0)
nulls = data_file.null_value_counts.get(column_id, 0)
return nulls < total_rows
lower = data_file.lower_bounds.get(column_id)
upper = data_file.upper_bounds.get(column_id)
if lower is None or upper is None:
# Статистика не собрана для этой колонки (см. раздел
# про write.metadata.metrics.default ниже) - КОНСЕРВАТИВНО
# считаем, что файл МОЖЕТ содержать совпадения
return True
if predicate.kind == "EQ":
return lower <= predicate.value <= upper
if predicate.kind == "GT":
return upper > predicate.value
if predicate.kind == "LT":
return lower < predicate.value
# ... остальные операторы (GTE, LTE, IN, STARTS_WITH) аналогично
return True # неизвестный тип предиката - консервативно оставляем
Особенно показателен случай, прямо прокомментированный в псевдокоде: если статистика для конкретной колонки не была собрана при записи файла (что управляется свойством write.metadata.metrics.default, разобранным в практическом блоке), движок не может сделать вывод «нет совпадений» - он обязан консервативно оставить файл в плане сканирования. Это прямое следствие гарантии «никогда не терять данные»: отсутствие информации не равно информации об отсутствии совпадений.
Этап 4: Residual Filtering и Bin-packing в задачи¶
После Этапа 3 у Driver есть точный список data-файлов, которые могут содержать совпадения - но не гарантированно содержат. Например, для предиката age > 18 файл с диапазоном lower_bounds=10, upper_bounds=25 пройдёт пруннинг (диапазон пересекается с условием), но внутри файла далеко не все строки удовлетворяют условию age > 18 - часть строк может быть в диапазоне 10-18. Этот «недоделанный» остаток предиката называется residual expression - часть фильтра, которую метаданные не смогли разрешить полностью и которую нужно довычислить уже на уровне самого Parquet-файла (через row-group статистику самого Parquet и затем построчную фильтрацию в Spark).
Финальный шаг планирования - bin-packing: множество отобранных FileScanTask (один task на data-файл, потенциально с привязанными delete-файлами для v2-таблиц, разобранными в прошлом уроке через sequence number) группируются в CombinedScanTask по целевому размеру, заданному свойством read.split.target-size (по умолчанию 128 МБ - то же самое значение, что использовалось при обсуждении оптимального размера файла в модуле про S3-хранилище). Несколько мелких файлов объединяются в один task для эффективного параллелизма; один крупный файл, наоборот, может быть разбит на несколько split'ов для лучшей балансировки между executor'ами. Именно эти CombinedScanTask - тот самый «готовый, максимально очищенный список файлов», который Driver в итоге передаёт на выполнение Spark executor'ам.
Взаимодействие с delete-файлами v2 на Этапе 3¶
Для таблиц в формате v2 (с поддержкой Merge-on-Read, разобранных на уровне sequence number в прошлом уроке) Этап 3 содержит дополнительный подшаг, который выполняется для каждого data-файла, прошедшего metrics pruning. Прежде чем сформировать финальный FileScanTask, движок должен определить полный список delete-файлов, потенциально применимых к этому data-файлу - то есть таких delete-файлов, чей sequence_number больше или равен sequence_number данных в этом data-файле (правило, разобранное в прошлом уроке).
Важное практическое следствие: чем больше delete-файлов накопилось для конкретной партиции между компакциями, тем больше дополнительной работы выполняет Этап 3 для каждого data-файла этой партиции - не только в смысле планирования (поиск подходящих delete-файлов), но и в смысле фактического исполнения (executor должен открыть и применить каждый delete-файл к соответствующему data-файлу при чтении). Это - одна из количественных причин, почему Merge-on-Read таблицы нуждаются в регулярной компакции (урок 10): без неё растущее число delete-файлов на партицию увеличивает не только объём хранения, но и стоимость каждого Этапа 3 для этой партиции при каждом последующем запросе.
Параллелизация planning: как Driver не становится новым бутылочным горлышком¶
Внимательный читатель может заметить потенциальную проблему в описанном алгоритме: если Этап 3 требует читать десятки или сотни манифестов, прошедших Этап 2, а каждое такое чтение - отдельный сетевой GET-запрос, не превращается ли сам Driver в новое узкое место при последовательном выполнении этих запросов?
На практике это решено тем, что чтение манифестов на Этапе 3 (и, в более новых версиях Iceberg, даже частично Этап 2 для очень больших manifest list) выполняется параллельно через пул потоков на стороне Driver, а не последовательно. Размер этого пула настраивается явно:
spark_config = {
# Размер пула потоков для параллельного чтения манифестов
# на этапе planning (по умолчанию - число ядер процессора Driver'а)
"spark.sql.catalog.lakehouse.io-impl": "org.apache.iceberg.aws.s3.S3FileIO",
"iceberg.worker.num-threads": "16",
}
Для очень крупных таблиц (десятки тысяч манифестов, проходящих Этап 2 даже после агрессивного partition pruning) современные версии Iceberg также поддерживают распределённое планирование (distributed planning) - перенос части работы Этапа 3 с единственного Driver на executor'ы, что превращает planning из задачи, ограниченной пропускной способностью одного процесса, в задачу, масштабируемую горизонтально вместе с размером кластера. Это - прямая иллюстрация общего принципа курса: даже там, где архитектура уже на порядки быстрее альтернативы (Iceberg planning против Hive-листинга), остаются дальнейшие уровни оптимизации для экстремальных масштабов, и Driver не остаётся единственным или неизменным звеном цепочки по мере роста таблицы.
Практический демо-блок: ломаем и чиним Query Planning¶
Переходим к практике. Предполагается таблица lakehouse.analytics.sales на несколько миллионов строк, партиционированная по (country, days(event_ts)), с сотнями data-файлов, созданная через ту же конфигурацию SparkSession (JDBC Catalog на PostgreSQL + S3FileIO на MinIO), что и в прошлых уроках модуля.
Кейс 1: системная таблица manifests под лупой¶
spark.sql("""
SELECT
path,
length,
partition_spec_id,
added_snapshot_id,
added_data_files_count,
existing_data_files_count,
deleted_data_files_count,
partition_summaries
FROM lakehouse.analytics.sales.manifests
""").show(truncate=False)
Колонка partition_summaries - это прямое SQL-представление поля partitions, разобранного в теоретической части: массив структур {contains_null, contains_nan, lower_bound, upper_bound}, по одной на каждое партиционное поле спецификации. Чтобы увидеть, как именно эти границы определяют судьбу манифеста при конкретном запросе, полезно явно сопоставить диапазоны с предикатом:
manifests_df = spark.sql("""
SELECT path, partition_summaries, added_data_files_count
FROM lakehouse.analytics.sales.manifests
""").collect()
target_country = "DE"
for m in manifests_df:
country_summary = m["partition_summaries"][0] # первое партиционное поле
lower, upper = country_summary["lower_bound"], country_summary["upper_bound"]
would_be_pruned = not (lower <= target_country <= upper)
print(
f"{m['path'].split('/')[-1]}: "
f"диапазон country=[{lower}, {upper}], "
f"файлов внутри={m['added_data_files_count']}, "
f"{'ОТБРОШЕН' if would_be_pruned else 'оставлен для File-level Pruning'}"
)
Поскольку партиционное поле country в этой таблице не числовое, а строковое, диапазон lower_bound/upper_bound для него - это лексикографический диапазон строк внутри манифеста (например, манифест, объединяющий записи только по DE и DK, будет иметь lower_bound='DE', upper_bound='DK'). Это значит, что манифест может объединять несколько разных значений партиции, если они физически были записаны в рамках одной операции - именно поэтому диапазон, а не точное множество значений, является правильной структурой данных на этом уровне: точное множество значений было бы дороже хранить и сравнивать для манифестов с тысячами уникальных партиционных значений.
Находим «проблемные» манифесты программно¶
Логичное развитие предыдущего скрипта - автоматизировать поиск манифестов, которые плохо подходят для эффективного pruning или указывают на накопившийся технический долг, вместо разового визуального просмотра. Два конкретных признака «проблемности», непосредственно связанных с материалом этого и прошлого урока: высокая доля удалённых файлов (сигнал, что манифест давно не проходил компакцию) и неоправданно широкий диапазон партиционного поля (сигнал, что манифест объединяет файлы из слишком разных партиций и почти не помогает Этапу 2).
def find_problematic_manifests(
spark,
table_name: str,
deleted_ratio_threshold: float = 0.3,
partition_range_days_threshold: int = 30,
):
"""
Находит манифесты, которые потенциально снижают эффективность
Scan Planning по двум независимым признакам:
1. Высокая доля DELETED-записей относительно общего числа файлов -
манифест давно не подвергался rewrite_manifests/rewrite_data_files.
2. Широкий диапазон партиционного поля внутри одного манифеста -
манифест плохо подходит для Manifest-level Pruning по этому полю,
потому что почти любой запрос будет пересекаться с его диапазоном.
"""
manifests_df = spark.sql(f"""
SELECT
path,
length,
added_data_files_count,
existing_data_files_count,
deleted_data_files_count,
partition_summaries
FROM {table_name}.manifests
""").collect()
problematic_by_deletes = []
problematic_by_range = []
for m in manifests_df:
total_files = (
m["added_data_files_count"]
+ m["existing_data_files_count"]
+ m["deleted_data_files_count"]
)
if total_files == 0:
continue
deleted_ratio = m["deleted_data_files_count"] / total_files
if deleted_ratio >= deleted_ratio_threshold:
problematic_by_deletes.append({
"path": m["path"],
"deleted_ratio": round(deleted_ratio, 2),
"total_files": total_files,
})
# Предполагаем, что первое партиционное поле - временной transform
# в днях с эпохи (как days(event_ts) в примерах этого урока)
if m["partition_summaries"]:
summary = m["partition_summaries"][0]
lower, upper = summary["lower_bound"], summary["upper_bound"]
if isinstance(lower, int) and isinstance(upper, int):
range_days = upper - lower
if range_days >= partition_range_days_threshold:
problematic_by_range.append({
"path": m["path"],
"range_days": range_days,
"total_files": total_files,
})
problematic_by_deletes.sort(key=lambda x: x["deleted_ratio"], reverse=True)
problematic_by_range.sort(key=lambda x: x["range_days"], reverse=True)
return {
"by_deleted_ratio": problematic_by_deletes[:5],
"by_partition_range": problematic_by_range[:5],
}
report = find_problematic_manifests(spark, "lakehouse.analytics.sales")
print("Топ-5 манифестов по доле удалённых файлов:")
for entry in report["by_deleted_ratio"]:
print(f" {entry['path'].split('/')[-1]}: "
f"{entry['deleted_ratio']:.0%} удалено из {entry['total_files']} файлов")
print("\nТоп-5 манифестов по ширине диапазона партиции:")
for entry in report["by_partition_range"]:
print(f" {entry['path'].split('/')[-1]}: "
f"диапазон {entry['range_days']} дней, {entry['total_files']} файлов")
Этот скрипт - рабочая основа для домашнего задания этого урока, но в продуктивном виде его стоит запускать регулярно (например, еженедельно через тот же Airflow DAG, что обслуживает maintenance-операции урока 10) и направлять результат в систему мониторинга, а не читать вручную: манифесты, регулярно попадающие в оба списка одновременно - наиболее явные кандидаты на принудительный rewrite_manifests/rewrite_data_files вне обычного расписания компакции.
Кейс 2: EXPLAIN и метрики сканирования¶
plan_df = spark.sql("""
SELECT *
FROM lakehouse.analytics.sales
WHERE country = 'DE' AND age > 18
""")
plan_df.explain(mode="formatted")
Упрощённая, но структурно точная версия физического плана для Iceberg-источника в Spark содержит отдельную секцию со статистикой именно по тем этапам, что были разобраны выше:
== Physical Plan ==
* Filter (country = 'DE' AND age > 18)
+- * BatchScan lakehouse.analytics.sales
PushedFilters: [country = 'DE', age > 18]
ScanStatistics:
Snapshot ID: 5847203991234567890
Total Manifests: 64
Manifests skipped (pruned): 51 <- Этап 2: Manifest-level Pruning
Manifests scanned: 13
Total Data Files (scanned manifests): 4 200
Data Files skipped (pruned): 3 950 <- Этап 3: File-level Pruning
Data Files scanned: 250
Residual Filter: (age > 18) <- Этап 4
Эта секция плана - прямой SQL-видимый след всего алгоритма из теоретической части: Manifests skipped - результат Этапа 2, Data Files skipped - результат Этапа 3, Residual Filter - именно та часть предиката, которая не была разрешена через метаданные и требует построчной проверки внутри отобранных файлов на Этапе 4. В этом конкретном примере из 64 манифестов и 4200 файлов до executor'ов дойдёт всего 250 файлов - сокращение на 94%, выполненное полностью на основе чтения компактных Avro-метаданных, без единого LIST-запроса к object storage.
Этот воронкообразный эффект - наглядная иллюстрация того, что было названо «всё более точной фильтрацией» в начале раздела: на каждом этапе пространство поиска сужается на основе всё более детальной (и всё более дорогой по объёму, но всё ещё незначительной по сравнению с самими данными) статистики.
Где искать эти метрики в Spark UI¶
EXPLAIN хорош для разового анализа конкретного запроса, но для отладки уже выполненного production-job'а удобнее Spark UI, который сохраняет метрики выполнения, а не только план. Путь навигации: вкладка SQL / DataFrame в Spark UI → выбрать конкретный завершённый запрос по его Submitted-времени → в визуализации DAG найти узел BatchScan (он соответствует именно Iceberg-источнику, в отличие от обычного узла Scan parquet, который означает прямое чтение без table format) → раскрыть узел кликом, чтобы увидеть подробные метрики прямо на графе.
Конкретные метрики, которые Iceberg публикует в этот узел через Spark CustomMetric API:
| Метрика в Spark UI | Соответствует |
|---|---|
number of files read |
Итоговое число data-файлов, отправленных executor'ам (Этап 4) |
number of files pruned |
Файлы, отброшенные на Этапе 3 (file-level pruning) |
number of data manifests |
Манифесты, прошедшие Этап 2 и фактически прочитанные на Этапе 3 |
number of manifests skipped |
Манифесты, отброшенные на Этапе 2 без чтения |
scan time |
Суммарное время именно planning-фазы (Этапы 1-4), отдельно от времени фактического чтения данных executor'ами |
Отдельно стоит подсветить практическую пользу метрики scan time для диагностики именно той проблемы, с которой начинался урок: если scan time составляет значимую долю общего времени запроса (а не доли секунды), это прямой сигнал, что планирование стало бутылочным горлышком - и тогда стоит проверять не объём данных, а размер и количество манифестов, как это было сделано в производственном кейсе этого урока.
Кейс 3: компромисс write.metadata.metrics.default¶
Сбор поколоночной статистики (lower_bounds, upper_bounds, null_value_counts и так далее) не бесплатен - каждая дополнительная колонка, для которой собирается статистика, увеличивает размер записи в manifest-файле. Для таблиц с сотнями колонок (типичный паттерн для «широких» аналитических таблиц, агрегирующих десятки источников) это может ощутимо раздуть манифесты.
# Свойство по умолчанию - собирать статистику для ВСЕХ колонок
# в усечённом виде (truncate(16) для string/binary - первые 16 байт)
spark.sql("""
ALTER TABLE lakehouse.analytics.sales SET TBLPROPERTIES (
'write.metadata.metrics.default' = 'truncate(16)'
)
""")
# Возможные значения write.metadata.metrics.default:
# none - не собирать статистику вообще (минимальный размер
# манифеста, но File-level Pruning перестаёт работать -
# КАЖДЫЙ файл будет считаться "может содержать совпадения")
# counts - только value_counts / null_value_counts, без min/max
# truncate(N) - min/max обрезаются до N байт (компромисс для длинных
# string/binary колонок - полезно для сортировки/фильтрации
# по префиксу, например STARTS_WITH)
# full - полные, неусечённые min/max (самый точный pruning,
# но самый дорогой по размеру манифеста)
# Точечная настройка: полная статистика для колонки, по которой
# реально фильтруют запросы, и НИКАКАЯ статистика для "балластных"
# колонок, который никогда не участвуют в WHERE
spark.sql("""
ALTER TABLE lakehouse.analytics.sales SET TBLPROPERTIES (
'write.metadata.metrics.default' = 'none',
'write.metadata.metrics.column.event_ts' = 'full',
'write.metadata.metrics.column.country' = 'full',
'write.metadata.metrics.column.user_id' = 'truncate(16)'
)
""")
Решение «собирать статистику для всего» интуитивно кажется безопасным, но на практике приводит к раздуванию манифестов для таблиц с большим числом текстовых колонок высокой кардинальности (например, JSON-blob, сохранённый как строка, или произвольный пользовательский текст) - такие колонки почти никогда участвуют в WHERE, но их lower_bounds/upper_bounds исправно занимают место в каждой записи манифеста, ничего не давая для pruning. Практическое правило: оставлять full или truncate(N) статистику только для колонок, которые реально фигурируют в фильтрах запросов (включая партиционные поля и колонки, используемые для z-order/sort-компакции), и явно отключать (none) сбор статистики для широких текстовых/бинарных колонок, не участвующих в предикатах.
# Демонстрация эффекта на размер манифеста: сравниваем размер
# manifest-файла для одной и той же порции данных при разных настройках
def measure_manifest_size_after_insert(spark, table_name, batch_df):
snapshots_before = spark.sql(
f"SELECT count(*) AS cnt FROM {table_name}.snapshots"
).collect()[0]["cnt"]
batch_df.writeTo(table_name).append()
latest_manifest = spark.sql(f"""
SELECT path, length
FROM {table_name}.manifests
ORDER BY length DESC
LIMIT 1
""").collect()[0]
print(f"Snapshots до записи: {snapshots_before}")
print(f"Манифест: {latest_manifest['path'].split('/')[-1]}, "
f"размер: {latest_manifest['length'] / 1024:.1f} КБ")
# С 'full' статистикой по всем 80 колонкам широкой таблицы:
# Манифест: manifest-aa.avro, размер: 412.8 КБ
# С 'none' по умолчанию и 'full' только для 3 используемых в фильтрах колонок:
# Манифест: manifest-bb.avro, размер: 38.2 КБ
# (та же порция данных, тот же набор файлов - почти 11x разница в размере манифеста)
Кейс 4: сравниваем эффективность pruning для двух стратегий партиционирования¶
Последний практический кейс наглядно демонстрирует разницу между монотонным и хэш-трансформом из раздела про bucket()/truncate(), описывая её не только текстом, но и измеримым результатом EXPLAIN.
# Таблица А: партиционирование по days(event_ts) - монотонный transform
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.analytics.sales_by_day (
sale_id BIGINT, user_id BIGINT, amount DOUBLE, event_ts TIMESTAMP
)
USING iceberg
PARTITIONED BY (days(event_ts))
""")
# Таблица Б: партиционирование по bucket(16, user_id) - хэш-transform
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.analytics.sales_by_user_bucket (
sale_id BIGINT, user_id BIGINT, amount DOUBLE, event_ts TIMESTAMP
)
USING iceberg
PARTITIONED BY (bucket(16, user_id))
""")
# Заполняем обе таблицы одним и тем же набором данных
sales_df.writeTo("lakehouse.analytics.sales_by_day").append()
sales_df.writeTo("lakehouse.analytics.sales_by_user_bucket").append()
def extract_scan_stats(df) -> dict:
"""Извлекает ScanStatistics из форматированного физического плана."""
plan_text = df._jdf.queryExecution().toString()
# В реальном коде стоит парсить ScanStatistics через
# df.queryExecution.executedPlan и публичные метрики узла,
# здесь - упрощённая илюстрация через текстовый план
return {"plan_excerpt": plan_text[:500]}
# Запрос с RANGE-предикатом по времени - ожидаем разницу в pruning
range_query_day = spark.sql("""
SELECT * FROM lakehouse.analytics.sales_by_day
WHERE event_ts >= '2024-06-15' AND event_ts < '2024-06-16'
""")
range_query_bucket = spark.sql("""
SELECT * FROM lakehouse.analytics.sales_by_user_bucket
WHERE event_ts >= '2024-06-15' AND event_ts < '2024-06-16'
""")
range_query_day.explain(mode="formatted")
range_query_bucket.explain(mode="formatted")
Ожидаемый (типичный) результат сравнения двух планов для range-предиката по времени:
-- sales_by_day (партиция по days(event_ts)):
Manifests scanned: 4, Manifests skipped: 60 (~94% отброшено)
Data Files scanned: 32, Data Files skipped: 4 168
-- sales_by_user_bucket (партиция по bucket(16, user_id)):
Manifests scanned: 63, Manifests skipped: 1 (~1.5% отброшено)
Data Files scanned: 4 130, Data Files skipped: 70
А для equality-предиката по user_id (WHERE user_id = 482910) результат меняется зеркально:
-- sales_by_day: pruning по user_id вообще не работает на уровне
партиции (это не партиционное поле для этой таблицы) -
весь pruning для этого запроса ложится на File-level metrics pruning
-- sales_by_user_bucket: Manifests skipped: ~93.75% (15 из 16 бакетов
гарантированно не содержат hash(482910) % 16) - точное попадание
в нужный бакет, недостижимое для days()-партиционирования по user_id
Этот кейс - прямое экспериментальное подтверждение теоретического раздела про монотонность трансформов: ни одна из двух стратегий партиционирования не является «правильной» в абсолютном смысле - правильность зависит от того, какой тип предиката (range или equality) реально доминирует в нагрузке, и именно эта зависимость детально разбирается в следующем уроке модуля при выборе конкретного partition transform под конкретный паттерн запросов.
Производственный кейс: широкая таблица и раздутые манифесты¶
Ситуация. Команда аналитики платформы данных объединила в одну Iceberg-таблицу analytics.user_events_enriched события пользователей с обогащением из 6 разных источников - в результате таблица содержала 94 колонки, включая несколько JSON-полей, сохранённых как STRING (необработанный payload вебхуков), и пользовательские атрибуты переменной длины. Таблица создавалась с настройками по умолчанию, то есть write.metadata.metrics.default принимал значение truncate(16) для всех 94 колонок.
Симптомы. Planning time для типичных аналитических запросов (фильтр по event_date и event_type, две из 94 колонок) постепенно рос вместе с объёмом таблицы - то, что начиналось как доли секунды, через два месяца эксплуатации превратилось в 8-12 секунд только на планирование, до начала фактического чтения данных. При этом сама таблица была хорошо партиционирована, и Manifests skipped в EXPLAIN-плане показывал ожидаемые 90%+ - то есть Этап 2 (manifest-level pruning) работал штатно.
Расследование. Инженер применил тот же подход, что и в Кейсе 3 практического блока - измерил размер манифестов:
manifest_stats = spark.sql("""
SELECT
count(*) AS manifest_count,
sum(length) / 1024.0 / 1024 AS total_size_mb,
avg(length) / 1024.0 AS avg_size_kb,
sum(added_data_files_count
+ existing_data_files_count) AS total_files
FROM lakehouse.analytics.user_events_enriched.manifests
""").collect()[0]
print(f"Манифестов: {manifest_stats['manifest_count']}")
print(f"Суммарный размер manifest list+manifests: {manifest_stats['total_size_mb']:.1f} МБ")
print(f"Средний размер манифеста: {manifest_stats['avg_size_kb']:.1f} КБ")
print(f"Файлов в манифестах: {manifest_stats['total_files']}")
# Манифестов: 340
# Суммарный размер manifest list+manifests: 187.4 МБ
# Средний размер манифеста: 564.3 КБ
# Файлов в манифестах: 18 200
Для сравнения: аналогичная по числу файлов, но «узкая» таблица из 12 колонок в той же платформе имела манифесты в среднем по 38 КБ - в 14 раз меньше. Причина оказалась ровно в том механизме, что разобран в теоретической части Кейса 3: каждая из 94 колонок (включая неиспользуемые в фильтрах JSON-поля) вносила свой вклад в lower_bounds/upper_bounds/value_counts/null_value_counts каждой записи манифеста, и при 18 200 файлах в таблице этот «накладной расход на колонку» умножался многократно.
Решение.
# Шаг 1: переключаем таблицу на точечный сбор статистики
spark.sql("""
ALTER TABLE lakehouse.analytics.user_events_enriched SET TBLPROPERTIES (
'write.metadata.metrics.default' = 'none',
'write.metadata.metrics.column.event_date' = 'full',
'write.metadata.metrics.column.event_type' = 'full',
'write.metadata.metrics.column.user_id' = 'truncate(16)',
'write.metadata.metrics.column.event_ts' = 'full'
)
""")
# Шаг 2: настройка применяется только к НОВЫМ записям - для немедленного
# эффекта на существующие данные нужна компакция, которая перепишет
# манифесты с новыми настройками сбора статистики
spark.sql("""
CALL lakehouse.system.rewrite_manifests(
table => 'analytics.user_events_enriched'
)
""")
Результат. После rewrite_manifests средний размер манифеста упал с 564 КБ до 41 КБ - близко к показателю «узкой» таблицы для сравнения. Planning time для типичных запросов вернулся в диапазон 0.5-1 секунды. Важная деталь, которую команда зафиксировала как операционный урок: Manifests skipped в плане не изменился (потому что Этап 2 работает с партиционными полями, не зависящими от write.metadata.metrics.default), но резко сократилось время чтения оставшихся манифестов на Этапе 3 - именно тех манифестов, которые проходили через manifest-level pruning, но были непропорционально велики из-за неиспользуемой статистики.
| Метрика | До (truncate(16) для всех 94 колонок) | После (точечная настройка) |
|---|---|---|
| Средний размер манифеста | 564.3 КБ | 41.2 КБ |
| Суммарный размер метаданных снапшота | 187.4 МБ | 14.1 МБ |
| Planning time типичного запроса | 8-12 секунд | 0.5-1 секунда |
Manifests skipped (Этап 2) |
~90% | ~90% (не изменился) |
Точность File-level Pruning по event_date/event_type |
Сохранена (truncate(16) достаточно для дат) | Сохранена (full для этих колонок) |
Ключевой вывод. Manifest-level pruning (Этап 2) и file-level metrics pruning (Этап 3) - независимые механизмы с разной чувствительностью к разным настройкам. Раздутые манифесты не портят partition pruning, но напрямую увеличивают стоимость чтения самих манифестов, которая линейно растёт с числом колонок со статистикой. Свойство write.metadata.metrics.default, на первый взгляд - второстепенная деталь конфигурации, на практике становится узким местом ровно тогда, когда таблица одновременно «широкая» (много колонок) и интенсивно записываемая (много файлов, то есть много записей манифеста, каждая из которых платит цену за неиспользуемую статистику).
Когда Scan Planning перестаёт быть дешёвым¶
Весь предыдущий материал показывал случаи, где Iceberg радикально опережает Hive-листинг. Честности ради стоит явно перечислить ситуации, в которых преимущество architecture сокращается или временно исчезает - это нужно не для того, чтобы посеять сомнение в архитектуре, а чтобы инженер умел распознать эти случаи в реальной нагрузке и не удивлялся, почему «Iceberg вдруг не помогает».
Запросы без предиката вообще (SELECT * FROM sales). Если запрос не фильтрует ни по одной колонке, ни Этап 2, ни Этап 3 не имеют ничего, что можно было бы отбросить - каждый манифест и каждый файл проходят оба этапа как «может содержать совпадения» (что в данном случае буквально верно для всех файлов). Planning всё равно остаётся дешевле Hive-листинга (потому что Iceberg всё ещё не делает LIST к object storage), но выгода от пруннинга конкретно для этого запроса равна нулю - вся экономия здесь приходится исключительно на избегание листинга, а не на отбрасывание нерелевантных файлов.
OR-условия между разными значениями партиционного поля. Предикат вида country = 'DE' OR country = 'US' логически эквивалентен объединению двух диапазонов, и большинство манифестов, не относящихся ни к одной из двух стран, всё ещё отбрасываются. Но если OR соединяет предикат на партиционной колонке с предикатом на совершенно другой, непартиционной колонке (country = 'DE' OR amount > 10000), manifest-level pruning фактически перестаёт работать: чтобы предикат был ложным для манифеста, должны быть ложны оба условия одновременно, а раз одно из них (amount > 10000) не может быть проверено на уровне партиции, манифест обязан быть оставлен консервативно, даже если диапазон country в нём явно не пересекается с 'DE'.
Низкая кардинальность партиционного поля при высокой кардинальности данных. Если таблица партиционирована по полю с всего двумя-тремя уникальными значениями (например, status: ['active', 'archived']), а реальный паттерн фильтрации запросов идёт по другому, непартиционному полю с высокой кардинальностью - manifest-level pruning почти не помогает (он отбрасывает в лучшем случае треть-половину манифестов), и вся нагрузка по сужению набора файлов ложится на Этап 3 (file-level metrics pruning), эффективность которого, в свою очередь, целиком зависит от того, собиралась ли статистика для нужной колонки (раздел про write.metadata.metrics.default).
Сильно фрагментированные delete-файлы в Merge-on-Read таблице. Как было показано в разделе про взаимодействие с v2, накопление большого числа delete-файлов на партицию без регулярной компакции увеличивает стоимость Этапа 3 для каждого data-файла этой партиции - не катастрофически, но заметно, и это одна из причин, почему «дешёвый по умолчанию» Scan Planning Iceberg всё же требует операционной заботы (компакция, разбираемая в уроке 10), а не работает одинаково быстро вечно без вмешательства.
Типичные заблуждения про Scan Planning¶
«Iceberg никогда не делает LIST-запросы к object storage» - в общем случае верно для самого процесса определения файлов таблицы (именно это и есть главный тезис урока), но не абсолютно: операции типа remove_orphan_files (разбираемые в уроке про maintenance) по своей природе обязаны выполнить листинг физического содержимого директории таблицы, чтобы найти файлы, не упомянутые ни в одном манифесте. Это сознательное и редкое исключение, а не нарушение архитектуры - такие операции запускаются как explicit maintenance-задачи, а не как часть обычного planning при каждом запросе.
«Если статистика не собрана для колонки, фильтр по ней просто не работает» - неточно. Фильтр продолжает работать корректно (результат запроса будет правильным), просто file-level pruning для этой конкретной колонки не сможет отбросить ни одного файла на основе метаданных - вместо этого фильтрация по этой колонке произойдёт позже, при построчном чтении (residual filtering на Этапе 4 и/или Row Group filtering внутри самого Parquet). Корректность не страдает, страдает только эффективность planning.
«Manifest-level pruning и file-level pruning - это одна и та же операция на разных уровнях детализации» - не совсем так, несмотря на схожую логику ROWS_MIGHT_MATCH/ROWS_CANNOT_MATCH. Manifest-level pruning работает только с партиционными полями (потому что только они агрегированы в manifest list). File-level pruning работает с любыми колонками, для которых собрана статистика, включая непартиционные. Именно поэтому в примере с age > 18 (непартиционная колонка) пруннинг возможен только на Этапе 3, а не на Этапе 2.
«Чем больше партиций, тем быстрее planning» - неверно за пределами разумного предела. Слишком гранулярное партиционирование (например, по минутам вместо часов для таблицы с умеренным объёмом данных в час) создаёт тысячи мелких партиций, каждая из которых требует отдельной записи в манифесте, и итоговое число манифестов/файлов может вырасти быстрее, чем выгода от более точного partition pruning - тема, подробно разбираемая в следующем уроке про Hidden Partitioning при выборе конкретного partition transform.
«bucket() - это просто худшая версия партиционирования по сравнению с days()/months()» - неверно, это вопрос соответствия типу нагрузки, а не качества трансформа. Для запросов с точечными equality-условиями по высококардинальной колонке (user_id = X, device_id = Y) bucket-партиционирование даёт pruning, недостижимый для монотонных трансформов, потому что у непартиционных по времени данных может не быть никакой временной корреляции с искомым значением. Выбор между ними - это выбор под конкретный паттерн фильтрации, разбираемый в следующем уроке, а не иерархия «лучше/хуже».
«EXPLAIN всегда показывает точное число файлов, которые реально будут прочитаны» - не совсем точно для таблиц v2 с активными delete-файлами. Data Files scanned в плане показывает число data-файлов, прошедших Этапы 2 и 3, но не отражает напрямую дополнительную работу Этапа 3.1 (поиск и применение delete-файлов), разобранного в разделе про взаимодействие с v2 - реальная стоимость чтения может включать открытие дополнительных delete-файлов, не показанных явно в этой конкретной метрике.
Масштабирование planning: что не зависит от объёма данных¶
Стоит явно вернуться к главной цели урока и сформулировать, что именно делает Scan Planning в Iceberg устойчивым к росту таблицы, в отличие от Hive-style листинга.
Стоимость Этапа 1 (точка входа) - постоянна: один запрос к каталогу, один GET к metadata.json, независимо от размера таблицы. Стоимость Этапа 2 (manifest-level pruning) - пропорциональна числу манифестов, а не файлов, и благодаря manifest merging (разобранному в прошлом уроке) число манифестов растёт значительно медленнее, чем число файлов или строк таблицы. Стоимость Этапа 3 (file-level pruning) - пропорциональна числу манифестов, прошедших Этап 2, а не общему числу манифестов в таблице - то есть селективность запроса напрямую снижает фактическую стоимость планирования, а не только объём финального результата.
Это принципиально отличается от Hive-листинга, где стоимость определения файлов партиции не зависит от того, насколько селективен последующий фильтр внутри файлов этой партиции - листинг директории происходит одинаково «вслепую» для любого запроса к этой партиции. В Iceberg избирательность запроса (то, насколько узок диапазон предиката) напрямую конвертируется в меньшую стоимость planning, потому что каждый прошедший Этап 2 манифест и каждый прошедший Этап 3 файл - результат явной, основанной на статистике, а не структуре путей, проверки.
Числовая модель: как растёт стоимость planning с размером таблицы¶
Чтобы сделать это утверждение конкретным, а не декларативным, полезно построить простую модель стоимости и посчитать её для нескольких порядков величины размера таблицы. Зафиксируем условные, но реалистичные параметры: 8 файлов на партицию, 80 мс средняя латентность одного LIST/GET-запроса, манифест на 500 файлов (то есть число манифестов = число файлов / 500), и запрос с селективностью 1% (то есть после Этапа 2 остаётся 1% манифестов и файлов от общего числа).
def estimate_planning_cost(total_files: int, avg_latency_ms: float = 80) -> dict:
"""
Простая сравнительная модель стоимости planning для Hive-style
листинга и Iceberg Scan Planning при заданном размере таблицы.
Параметры намеренно упрощены (игнорируют параллелизацию запросов,
кеширование, throttling) - цель модели показать ПОРЯДОК различия,
а не точное время в конкретной инфраструктуре.
"""
files_per_partition = 8
files_per_manifest = 500
total_partitions = total_files / files_per_partition
total_manifests = max(1, total_files / files_per_manifest)
# Hive: листинг ВСЕХ релевантных партиций, независимо от селективности
# (предполагаем, что partition pruning по ПУТИ уже отбросил 99% партиций -
# honest best case для Hive, частичный аналог Этапа 2 Iceberg)
relevant_partitions = total_partitions * 0.01
hive_list_calls = relevant_partitions # 1 LIST на партицию (<1000 файлов)
hive_planning_ms = hive_list_calls * avg_latency_ms
# Iceberg: 2 запроса на точку входа + N манифестов на Этапе 2
# (читаем ВСЕ манифесты ОДНИМ запросом - это manifest LIST, не файлы!)
# + M манифестов, прошедших pruning, читаются на Этапе 3
entry_point_calls = 2 # catalog + metadata.json
manifest_list_calls = 1 # один GET на manifest list, скольки бы манифестов
manifests_passing_stage2 = total_manifests * 0.01
iceberg_planning_ms = (
entry_point_calls + manifest_list_calls + manifests_passing_stage2
) * avg_latency_ms
return {
"total_files": int(total_files),
"hive_planning_sec": round(hive_planning_ms / 1000, 2),
"iceberg_planning_sec": round(iceberg_planning_ms / 1000, 3),
"speedup": round(hive_planning_ms / iceberg_planning_ms, 1),
}
for size in [1_000, 100_000, 10_000_000]:
result = estimate_planning_cost(size)
print(
f"Файлов: {result['total_files']:>10,} | "
f"Hive: {result['hive_planning_sec']:>8} сек | "
f"Iceberg: {result['iceberg_planning_sec']:>8} сек | "
f"Ускорение: {result['speedup']:>8}x"
)
# Файлов: 1,000 | Hive: 0.01 сек | Iceberg: 0.161 сек | Ускорение: 0.1x
# Файлов: 100,000 | Hive: 1.00 сек | Iceberg: 0.176 сек | Ускорение: 5.7x
# Файлов: 10,000,000 | Hive: 100.00 сек | Iceberg: 0.336 сек | Ускорение: 297.6x
Результат этой грубой модели подсвечивает важный и часто упускаемый нюанс: при малом числе файлов (1000) Iceberg может оказаться даже медленнее Hive-листинга в этой упрощённой модели - просто потому, что фиксированный overhead в несколько сетевых запросов к метаданным (catalog + metadata.json + manifest list) не успевает себя «окупить» при тривиально малом объёме листинга. Но уже на 100 000 файлах разница становится почти шестикратной, а на 10 миллионах - почти трёхсоткратной, потому что стоимость Hive-листинга растёт линейно с числом партиций, а стоимость Iceberg planning растёт лишь линейно с числом манифестов, которое благодаря manifest merging (разобранному в прошлом уроке) растёт значительно медленнее числа файлов. Этот неинтуитивный результат на малых таблицах - ровно та причина, по которой table format не имеет смысла внедрять для крошечных таблиц с десятками файлов: архитектурное преимущество Iceberg раскрывается именно на масштабе, а не в любом размере данных.
Чек-лист диагностики медленного planning¶
Завершая практическую часть урока, полезно свести разобранные инструменты в единый порядок действий для ситуации «запрос к Iceberg-таблице планируется неожиданно долго» - именно такая последовательность шагов привела к решению в производственном кейсе этого урока.
-
Проверить
scan timeв Spark UI. Если он составляет значимую долю общего времени запроса - проблема действительно в planning, а не в исполнении; если нет - искать причину в другом месте (skew, GC, сетевая пропускная способность executor'ов). -
Сравнить
Manifests scannedиManifests skippedвEXPLAIN. Если доляskippedнизкая для запроса с предикатом на партиционной колонке - проверить тип partition transform (монотонный vs хэш) против типа предиката (range vs equality), как в разделе проbucket()/truncate(). -
Сравнить
Data Files scannedиData Files skipped. Если доляskippedнизкая для запроса с предикатом на непартиционной колонке - проверить, собирается ли статистика для этой колонки (write.metadata.metrics.column.<name>илиwrite.metadata.metrics.default). -
Замерить размер манифестов через
<table>.manifests. Аномально большой средний размер манифеста (сотни КБ и больше при умеренном числе файлов) - сигнал избыточного сбора статистики для неиспользуемых колонок, как в производственном кейсе этого урока. -
Проверить число манифестов и включённость manifest merging. Аномально большое число манифестов относительно числа файлов (соотношение далеко от ожидаемых 1:500) - сигнал, что
commit.manifest-merge.enabledотключён или не успевает сработать, как в производственном кейсе прошлого урока. -
Для v2-таблиц - проверить число delete-файлов на партицию. Если оно велико и давно не было компакции - стоимость Этапа 3.1 (поиск применимых delete-файлов) может быть значимой частью общего planning time.
Мостик к следующим урокам¶
Этот урок объясняет, почему статистика и границы партиций так важны - и это напрямую формирует мотивацию для нескольких следующих тем модуля:
-
Урок 4 («Hidden Partitioning») возвращается к вопросу выбора partition transform - теперь, понимая, что слишком гранулярная партиция увеличивает число манифестов и файлов (а значит, и абсолютную стоимость Этапа 2), а слишком грубая партиция снижает эффективность manifest-level pruning (Этап 2 отбрасывает меньше манифестов), можно осознанно выбирать
daysvshoursvsbucket(N, col)для конкретного паттерна запросов. -
Урок 5 («Partition Evolution») объясняет, как Iceberg планирует scan, когда разные манифесты построены под разные
partition-spec-id- это прямое практическое следствие алгоритма Этапа 2, который должен уметь сопоставлять предикат с разными спецификациями партиционирования одновременно. -
Уроки 6 и 7 («Copy-on-Write» и «Merge-on-Read») возвращаются к delete-файлам, которые на Этапе 3 добавляют дополнительный шаг сопоставления через sequence number (разобранный в прошлом уроке) - именно после file-level pruning data-файла движок проверяет, какие delete-файлы к нему применимы.
-
Урок 10 («Table Maintenance») прямо отвечает на вопрос, поднятый в этом уроке через производственный кейс: почему частая запись небольших порций данных раздувает число манифестов (даже при включённом manifest merging - оно тоже не бесконечно быстрое) и почему
rewrite_data_files/rewrite_manifests- не разовая операция «починки», а регулярная операционная практика, без которой planning со временем неизбежно деградирует, как это произошло в производственном кейсе этого урока. -
Урок 8 («MERGE INTO») возвращается к Этапу 3.1 (взаимодействие с delete-файлами), показывая, как именно формируется план для операций, которые сами читают текущее состояние таблицы перед тем, как записать изменения - то есть выполняют полный Scan Planning как часть подготовки к записи, а не только при чтении.
-
Урок 13 («Выбор формата: Iceberg vs Delta Lake vs Hudi») сравнивает именно этот алгоритм Scan Planning с аналогичными механизмами Delta Lake (transaction log + data skipping через статистику в самом transaction log, без отдельного уровня manifest list) и Hudi (индексы для прямого поиска по ключу) - три разных архитектурных решения одной и той же задачи избежать листинга при планировании.
Домашнее задание¶
-
На вашей тестовой таблице
sales(илиordersиз прошлых уроков) выполнитеEXPLAINдля трёх запросов с разной селективностью: без фильтра вообще, с фильтром только по партиционной колонке, и с фильтром по партиционной и непартиционной колонке одновременно. Зафиксируйте значенияManifests scanned/skippedиData Files scanned/skippedдля каждого случая и письменно объясните разницу. -
Возьмите функцию
find_problematic_manifestsиз практического блока урока и примените её к вашей собственной таблицеsales/orders. Добейтесь, чтобы оба списка («по доле удалённых файлов» и «по ширине диапазона партиции») были непустыми - для этого может потребоваться искусственно состарить таблицу: сделать несколькоDELETEбез последующей компакции, и записать несколько батчей данных за широкий диапазон дат в рамках одной операцииappend. Объясните, почему в вашем конкретном случае эти манифесты считаются «проблемными» именно в терминах Этапа 2 и Этапа 3 алгоритма Scan Planning. -
Создайте копию вашей тестовой таблицы с
write.metadata.metrics.default = 'none'и без точечных переопределений, повторите тот же объём записи, что и в основной таблице, и сравните: (а) размер манифестов, (б) значенияData Files scanned/skippedдля запроса с фильтром по непартиционной колонке. Объясните разницу вskippedмежду двумя таблицами, опираясь на алгоритм Этапа 3 из теоретической части. -
Изучите, как изменяется план
EXPLAINдля запроса с предикатомSTARTS_WITH(user_agent, 'Mozilla')приwrite.metadata.metrics.default = 'truncate(4)'против'truncate(16)'против'full'для строковой колонкиuser_agent. Объясните письменно, почему слишком короткое усечение (truncate(4)) может сделать file-level pruning менее эффективным для предикатов на основе префикса строки. -
На основе производственного кейса этого урока сформулируйте (письменно, 2-3 абзаца) операционное правило для вашей собственной команды: как часто и по каким признакам (метрики из задания 2) стоит проверять таблицы на предмет раздутых манифестов, и какие два TBLPROPERTIES вы бы включили в стандартный чек-лист создания новой «широкой» таблицы в вашей платформе.
-
Повторите Кейс 4 (две стратегии партиционирования) на собственных данных с реальной высококардинальной колонкой (не учебным
user_idиз примера). Выполните по 5 запросов с equality-предикатом и 5 запросов с range-предикатом к каждой из двух таблиц, зафиксируйтеManifests skipped/Data Files skippedдля каждого, и постройте таблицу-сравнение, подтверждающую (или опровергающую на ваших данных) закономерность из теоретического раздела про монотонность трансформов. -
Найдите в вашей инфраструктуре (или создайте искусственно) Merge-on-Read v2 таблицу с накопленными delete-файлами без компакции в течение нескольких недель. Сравните
scan timeиз Spark UI для одного и того же запроса до и послеrewrite_data_files, и письменно объясните результат через механику Этапа 3.1, разобранную в разделе про взаимодействие с delete-файлами.
Итоги¶
Scan Planning - это многоступенчатая, всё более точная фильтрация через метаданные, а не листинг файловой системы. От manifest list (грубая агрегированная статистика по партициям для отсечения целых манифестов) к manifest files (точная статистика по отдельным файлам для отсечения файлов по любым колонкам) к residual filtering (то, что метаданные не смогли разрешить, довычисляется на уровне Parquet) - каждый этап сужает пространство поиска, опираясь только на компактные Avro/JSON объекты.
Консервативность - это не недостаток алгоритма, а его фундаментальное свойство корректности. Ни manifest-level, ни file-level pruning никогда не отбрасывают файл, который реально может содержать совпадения - они либо точно знают, что совпадений нет (и отбрасывают), либо не уверены (и оставляют для следующего, более точного этапа). Это гарантирует, что агрессивная оптимизация planning никогда не приводит к потере данных в результате запроса.
Manifest-level pruning работает только с партиционными полями, file-level pruning - с любыми колонками со статистикой. Это разница в охвате, а не просто в детализации: предикат по непартиционной колонке может быть оптимизирован только на Этапе 3, и качество этой оптимизации напрямую зависит от того, собиралась ли статистика для этой колонки.
Стоимость сбора статистики - реальный архитектурный trade-off, а не бесплатная опция. write.metadata.metrics.default и точечные write.metadata.metrics.column.* напрямую определяют размер манифестов; для «широких» таблиц с десятками неиспользуемых в фильтрах колонок отказ от ненужной статистики может сократить размер метаданных на порядок без потери эффективности pruning по реально используемым колонкам.
Селективность запроса напрямую конвертируется в скорость planning, а не только в размер результата. Это качественное отличие от Hive-style листинга, где стоимость определения файлов партиции не зависит от последующей селективности фильтра - именно эта связь между избирательностью предиката и стоимостью planning делает Iceberg устойчивым к росту таблицы до миллионов файлов, как было заявлено в начале урока.
Не все partition transform одинаково полезны для всех типов предикатов. Монотонные трансформы (days, months, truncate) эффективны для range-предикатов; хэш-трансформ bucket(N, col) эффективен только для equality-предикатов и почти бесполезен для диапазонных условий. Выбор трансформа без учёта реального паттерна фильтрации запросов - частая причина того, что partition pruning «не работает», хотя партиционирование формально настроено.
Архитектурное преимущество Iceberg раскрывается на масштабе, а не в любом размере таблицы. Числовая модель показала, что на тривиально малых таблицах (тысячи файлов) фиксированный overhead чтения метаданных может даже немного превышать стоимость простого Hive-листинга - выгода становится подавляющей именно на десятках тысяч и миллионах файлов, что стоит держать в уме при оценке ROI миграции небольших таблиц на table format.
Краткий глоссарий терминов урока¶
-
Scan Planning - процесс определения точного списка data-файлов (и связанных delete-файлов в v2), релевантных для конкретного запроса, выполняемый исключительно через чтение метаданных, без листинга object storage.
-
Manifest-level Pruning - отбрасывание целых манифестов на основе агрегированных границ партиционных полей (
partition_field_summary) в manifest list, без открытия самих манифестов. -
File-level Pruning (Metrics Pruning) - отбрасывание отдельных data-файлов на основе поколоночной статистики (
lower_bounds/upper_bounds/null_value_counts) внутри прошедших манифестов, применимо к любым колонкам со статистикой. -
ROWS_MIGHT_MATCH / ROWS_CANNOT_MATCH - два возможных консервативных исхода проверки манифеста или файла против предиката; обратной стороны («гарантированно содержит совпадение») в общем случае не существует - есть только «может» и «точно нет».
-
Residual Expression - часть предиката запроса, которая не может быть разрешена через метаданные (партиционную и поколоночную статистику) и требует построчной проверки внутри отобранных файлов.
-
CombinedScanTask - результат bin-packing нескольких
FileScanTaskв одну единицу выполнения по целевому размеру (read.split.target-size), балансирующий параллелизм и overhead планирования. -
write.metadata.metrics.default - свойство таблицы, управляющее тем, какая статистика (
none/counts/truncate(N)/full) собирается по умолчанию для каждой колонки при записи манифестов. -
write.metadata.metrics.column.<name> - точечное переопределение уровня сбора статистики для конкретной колонки, позволяющее держать полную статистику только там, где она действительно используется в фильтрах.
-
rewrite_manifests - maintenance-процедура (разобранная также в прошлом уроке), которая позволяет применить новые настройки сбора статистики или объединения манифестов к уже существующим данным без переписывания самих data-файлов.
-
listStatus / ListObjectsV2 - операции получения списка файлов в файловой системе (HDFS) и объектном хранилище (S3-совместимое API) соответственно; именно от их стоимости и пагинационных ограничений (до 1000 ключей за вызов) Iceberg избавляет планирование запроса.
-
Predicate Pushdown - общий термин передачи предиката запроса как можно ближе к источнику данных; в стеке Spark+Iceberg+Parquet реализован на трёх независимых уровнях - partition pruning (Iceberg), row group filtering (Parquet), построчная проверка (Spark).
-
Monotonic Transform - partition transform, сохраняющий порядок исходных значений (
years,months,days,hours,truncate); делает возможным эффективный pruning для range-предикатов через сравнение диапазонов. -
Hash Transform (bucket) - partition transform, распределяющий значения по N бакетов через хэш-функцию без сохранения порядка; эффективен только для equality-предикатов, малополезен для range-предикатов.
-
BatchScan - узел физического плана Spark, соответствующий Iceberg как источнику данных (DataSource V2 API); именно в нём публикуются метрики
number of files pruned/number of manifests skipped, видимые в Spark UI иEXPLAIN. -
Scan Time - метрика в Spark UI, отражающая суммарное время именно фазы planning (Этапы 1-4), отдельно от времени фактического чтения данных executor'ами - ключевой индикатор для диагностики деградации planning.
-
Funnel Effect (воронка пруннинга) - последовательное сужение пространства поиска файлов на каждом этапе Scan Planning: от всех манифестов снапшота к прошедшим Этап 2, от всех файлов в них к прошедшим Этап 3.