Apache Iceberg: зачем нужен table format и что он добавляет к Parquet
Архитектурные боли чистого Parquet в Data Lake, анатомия Iceberg (catalog, metadata.json, manifest list, manifest, data files), ACID через optimistic concurrency, hidden partitioning, schema evolution, partition evolution, time travel и практика на PySpark.
Жизнь до Table Formats: Data Lake как папка с файлами¶
До появления Apache Iceberg, Delta Lake и Apache Hudi типичный Data Lake выглядел предельно просто: набор Parquet-файлов, разложенных по директориям в S3, HDFS или MinIO. Инженер пишет df.write.parquet("s3a://bucket/events/"), Spark создаёт несколько файлов внутри этой директории, и с точки зрения всех участников системы это и есть «таблица events».
Проблема в том, что термин «таблица» здесь - это договорённость, а не технический факт. Сама директория ничего не знает о том, что она представляет собой логическую сущность с фиксированной схемой, историей изменений и набором гарантий целостности. Это просто путь в файловой системе, в котором лежат файлы определённого формата. Любой инструмент, который умеет читать Parquet, может прочитать эти файлы - но ни один из них не может ответить на вопрос «а что произошло с этой таблицей вчера в 14:32» или «какие именно файлы относятся к актуальной версии данных, а какие - это мусор от упавшего вчера job'а».
Чтобы превратить «папку с файлами» в нечто похожее на таблицу, в экосистему Hadoop был добавлен Hive Metastore - centralized-сервис, который хранит соответствие между именем таблицы (analytics.events) и физическим путём (s3a://bucket/events/), а также описание схемы и список партиций. Долгое время именно Hive Metastore выполнял роль «реестра таблиц» для Spark, Hive, Presto/Trino и других движков.
На диаграмме видно главное ограничение этой архитектуры: связь между Hive Metastore и реальными данными - это просто строка с путём. Metastore не хранит список файлов, их контрольные суммы, статистику по колонкам или информацию о том, какие файлы относятся к последней «правильной» версии таблицы, а какие - артефакты сбоев. Любой движок, который открывает таблицу, фактически делает LIST всей директории (а иногда и рекурсивно по всем партициям) и пытается прочитать всё, что выглядит как Parquet-файл с подходящей схемой.
Иллюзия таблицы¶
Важно зафиксировать ключевую мысль этого раздела: Hive Metastore - это каталог путей, а не табличный движок. Он отвечает на вопрос «где лежат данные таблицы X», но не отвечает на вопросы, которые умеет решать любая полноценная СУБД:
-
Какая версия данных является «текущей», если несколько job'ов писали одновременно?
-
Что делать с файлами, которые остались от прерванной записи?
-
Можно ли безопасно читать таблицу прямо сейчас, пока в неё идёт запись?
-
Как переименовать колонку без полной перезаписи всех файлов?
-
Как вернуться к состоянию данных, которое было сутки назад?
Все эти вопросы пользователю приходилось решать руками: писать обвязку поверх Spark, договариваться о конвенциях именования файлов, организовывать job'ы так, чтобы они не пересекались по времени, и вручную чистить «осиротевшие» файлы. Именно эти операционные боли и привели индустрию к идее table format - слоя, который добавляет недостающую семантику таблицы поверх обычных файлов.
Архитектурные боли чистого Parquet в Data Lake¶
Прежде чем переходить к решению, стоит подробно разобрать, почему «просто Parquet-файлы в S3» - это не временное неудобство, а фундаментальное архитектурное ограничение, которое проявляется тем сильнее, чем больше людей и job'ов работают с одной таблицей.
Боль 1: отсутствие атомарности записи¶
Когда Spark job пишет данные в Parquet-датасет, он делает это в несколько шагов: каждый task записывает свою порцию данных в файл, а после того, как все task'и успешно завершились, происходит финализация (commit) - например, перемещение файлов из временной директории _temporary в целевую. Если job упал на середине - в директории остаются частично записанные файлы и временные артефакты.
# Наивная запись большого набора партиций по одной
# Используется здесь только для иллюстрации механики сбоя -
# в реальности Spark распределяет task'и по executor'ам параллельно,
# но итоговая проблема (частично записанные данные при сбое) та же.
def naive_parquet_write_with_failure(spark, df, output_path, fail_at_partition=7):
"""
Имитация сбоя посередине записи в обычный Parquet-датасет.
При падении на партиции №7 уже записанные партиции 0-6
останутся в output_path физически, хотя job завершился с ошибкой.
"""
total_partitions = df.rdd.getNumPartitions()
for partition_id in range(total_partitions):
if partition_id == fail_at_partition:
raise RuntimeError(
f"Симулированный сбой executor'а на партиции {partition_id}"
)
part_df = df.where(f"spark_partition_id() = {partition_id}")
(
part_df.write
.mode("append")
.parquet(output_path)
)
try:
naive_parquet_write_with_failure(spark, events_df, "s3a://bucket/events_raw/")
except RuntimeError as exc:
print(f"Job упал: {exc}")
# Партиции 0-6 уже физически лежат в output_path.
# Любой, кто прочитает s3a://bucket/events_raw/ ПРЯМО СЕЙЧАС,
# увидит неполный, "оборванный" набор данных - и не узнает об этом,
# потому что для обычного Parquet-датасета это выглядит как
# "ещё одна порция файлов", а не как незавершённая транзакция.
Ключевая проблема здесь не в том, что job упал - сбои случаются всегда, это нормальная часть распределённых вычислений. Проблема в том, что частично записанные данные неотличимы от полностью записанных. Читатель, который обращается к этой директории в момент сбоя или сразу после него, не получает ни ошибки, ни предупреждения - он просто читает то подмножество файлов, которое там физически лежит, и считает это валидным состоянием таблицы.
Боль 2: проблема Eventual Consistency в облачных object storage¶
Классические объектные хранилища (в первую очередь - старые версии Amazon S3, а также многие S3-совместимые реализации вроде ранних Ceph RGW) исторически давали более слабые консистентность-гарантии, чем традиционные файловые системы. Конкретно, после PUT нового объекта операция LIST могла не сразу отразить только что записанный файл, потому что список объектов строился на отдельном индексном слое, который обновлялся с небольшой задержкой.
Для Data Lake поверх «голого» Parquet это означало конкретную и неприятную проблему: Spark job записывает 200 файлов, затем driver делает LIST, чтобы построить итоговый manifest для commit-протокола (например, в FileOutputCommitter) - и в этом списке может не оказаться части только что записанных файлов. Результат - либо потерянные данные (если committer считает, что запись завершена и опирается на неполный список), либо ошибки несоответствия количества файлов.
Современные реализации S3 API в большинстве случаев гарантируют строгую консистентность чтения после записи (strong read-after-write consistency), и многие self-hosted S3-совместимые хранилища (MinIO, современный Ceph RGW) тоже это поддерживают. Но даже при строгой консистентности отдельных операций PUT/GET/LIST остаётся более глубокая проблема: последовательность из нескольких независимых операций (LIST → GET → анализ → запись) не является атомарной транзакцией. Между шагами этой последовательности другой процесс может успеть изменить состояние директории - и именно для решения этого класса проблем нужен табличный слой с собственным протоколом коммитов, а не просто «более консистентное» хранилище.
Боль 3: дорогой листинг файлов¶
В файловой системе типа HDFS операция listStatus для директории с N файлами выполняется быстро, потому что NameNode хранит дерево файлов в памяти. В объектных хранилищах (S3, MinIO, Ceph RGW) листинг - это HTTP-запрос с пагинацией, обычно возвращающий не более 1000 ключей за один вызов.
Если таблица партиционирована по дате и содержит данные за несколько лет с гранулярностью «час», итоговое число партиций легко достигает десятков тысяч. Чтобы Spark Catalyst построил план выполнения запроса вида SELECT * FROM events WHERE event_date = '2024-06-01', движку без табличных метаданных нужно:
-
Получить список всех партиций (директорий) - потенциально тысячи
LIST-запросов с пагинацией. -
Отфильтровать те, что соответствуют партиции
event_date=2024-06-01'- obычно по самому имени пути (partition pruning на уровне путей). -
Для каждой подходящей партиции сделать ещё один
LIST, чтобы получить список Parquet-файлов внутри. -
Опционально прочитать footer каждого файла, чтобы получить статистику колонок для предиктных фильтров (row-group filtering).
Даже если итоговый результат запроса - это одна партиция с десятком файлов, движок вынужден сначала «ощупать» всё дерево директорий, чтобы понять, какие из них релевантны. На практике это превращается в десятки секунд planning time только на то, чтобы Spark решил, что читать - до того, как он прочитает хотя бы один байт данных.
Грубая оценка стоимости листинга. Объектные хранилища типа S3 возвращают не более 1000 ключей за один вызов ListObjectsV2, а типичная задержка такого вызова - 50-150 мс (зависит от региона, нагрузки на хранилище и числа объектов под префиксом). Для таблицы с тремя годами данных и партиционированием по часу это:
3 года × 365 дней × 24 часа = 26 280 партиций
Если в каждой партиции по 5-10 файлов:
26 280 партиций × ~1 LIST-вызов на партицию ≈ 26 280 HTTP-запросов
26 280 запросов × 100 мс (средняя задержка) ≈ 2 628 секунд ≈ 43 минуты
даже при разумной параллелизации (например, 50 параллельных потоков
для построения списка партиций) - это всё равно 53 секунды
ТОЛЬКО на этап "выяснить, какие партиции существуют",
до начала чтения хотя бы одного байта данных.
Эта оценка груба и не учитывает кеширование на стороне движка (Spark, например, кеширует listing партиций Hive-таблицы в driver-памяти между запросами в рамках одной сессии) - но именно отсутствие такого кеша при первом запросе, после перезапуска приложения, или при работе через движок, который не реализует собственное кеширование партиций (многие легковесные query engines), регулярно становится источником тормозящих ad-hoc запросов «к холодной» таблице. Iceberg убирает этот этап полностью: вместо построения списка партиций через LIST движок получает готовый список релевантных файлов из чтения нескольких компактных Avro-манифестов, что заменяет тысячи HTTP-запросов десятками.
Боль 4: небезопасное одновременное чтение и запись¶
Последняя и, возможно, самая опасная для production-систем проблема - отсутствие изоляции между конкурирующими операциями. Если один процесс пишет новые файлы в директорию таблицы, а другой в этот же момент читает из неё, классический Parquet-датасет не даёт никаких гарантий о том, что увидит читатель.
Эта проблема особенно болезненна для самого распространённого ETL-паттерна - overwrite, когда таблица полностью перезаписывается новой версией данных (например, ежедневный batch job, пересчитывающий агрегаты за весь период). В обычном Parquet-датасете mode("overwrite") физически означает «удалить всё в директории, затем записать новое». Любой запрос, который выполняется в этот промежуток времени, либо получит ошибку (файл не найден), либо - что хуже - тихо вернёт неполные или пустые результаты, потому что для движка чтения это просто директория с меньшим числом файлов, а не сигнал «таблица сейчас в несогласованном состоянии, отложи запрос».
Все четыре проблемы - атомарность, листинг, консистентность последовательности операций и изоляция читателей/писателей - не являются багами конкретной реализации Spark или конкретного объектного хранилища. Это структурное следствие того, что «таблица = директория с файлами» не предоставляет никакого механизма для отслеживания состояния этой таблицы как единого целого. Решение этой проблемы и есть Open Table Format.
Что такое Open Table Format¶
Open Table Format - это открытая спецификация метаданных, которая превращает набор файлов в объектном хранилище в полноценную таблицу с транзакционными свойствами, версионностью и управляемой схемой. Ключевое слово здесь - «открытая»: формат специфицирован публично, и любой движок (Spark, Trino, Flink, DuckDB, ClickHouse, Snowflake) может реализовать чтение и запись по этой спецификации независимо от других движков.
Принципиально важно разделить два понятия, которые часто путают новички в теме Lakehouse:
-
Формат файла (File Format) - то, как физически закодированы байты на диске: Parquet, ORC, Avro. Это всегда было и остаётся «нижним слоем» - именно файлы такого формата содержат фактические данные строк и колонок.
-
Формат таблицы (Table Format) - дополнительный слой метаданных поверх файлового формата, который описывает, какие именно файлы образуют таблицу в данный момент времени, какая у таблицы схема, как она партиционирована и какова история её изменений. Apache Iceberg, Delta Lake и Apache Hudi - это три наиболее известные реализации table format.
Важно понимать, что table format не заменяет Parquet - он работает вместе с ним. Данные физически продолжают лежать в Parquet-файлах (или ORC/Avro), со всеми их преимуществами: колоночным хранением, эффективным сжатием, статистикой по row group'ам для predicate pushdown. Table format добавляет отдельный, не зависящий от файлов слой, который точно знает, какой набор этих файлов представляет собой консистентную версию таблицы в любой момент времени.
Почему именно Apache Iceberg¶
На рынке существует несколько реализаций open table format, и у каждой своя история и архитектурные акценты:
-
Apache Iceberg изначально разработан в Netflix (2017) для решения именно тех проблем, которые мы разобрали выше - дорогой листинг тысяч партиций Hive-таблиц и отсутствие атомарных транзакций. Передан в Apache Software Foundation, сейчас активно развивается широким сообществом, включая инженеров из Apple, Netflix, Tabular (компания, созданная авторами Iceberg, позже приобретённая Databricks), Snowflake и других. Ключевая особенность Iceberg - дизайн, изначально ориентированный на множество движков (multi-engine): один и тот же стол читается и пишется из Spark, Trino, Flink, Snowflake, ClickHouse без переноса данных.
-
Delta Lake разработан и продвигается в первую очередь компанией Databricks, тесно интегрирован со Spark. Открыт под Apache 2.0, но историческое развитие сильнее завязано на одного вендора, хотя сообщество и Delta Lake 3.0 (формат Delta Universal Format, UniForm) активно работают над multi-engine совместимостью.
-
Apache Hudi возник в Uber, делает сильный акцент на потоковых (streaming) и инкрементальных нагрузках с частыми upsert-операциями, имеет встроенные индексы для быстрого поиска записей по ключу.
Этот курс фокусируется на Apache Iceberg как на наиболее нейтральной и быстро растущей спецификации, но архитектурные принципы (snapshot-изоляция, манифесты, ACID через атомарный commit) в большой степени общие для всех трёх форматов - сравнение Iceberg, Delta Lake и Hudi по конкретным критериям выбора будет в последнем уроке этого модуля.
Multi-engine как практическое следствие открытости спецификации¶
Слово «открытый» в Open Table Format не декларативное, а имеет конкретное архитектурное следствие: поскольку формат метаданных опубликован как спецификация (а не скрыт внутри проприетарного движка), разные вычислительные движки могут читать и писать в одну и ту же физическую таблицу, используя общий каталог как единственную точку синхронизации.
Практическая ценность этой схемы - не в абстрактной «гибкости», а в конкретном операционном сценарии, знакомом многим командам: Flink пишет события в таблицу в режиме streaming с низкой задержкой, Spark по ночам выполняет тяжёлую batch-трансформацию и компакцию той же таблицы, аналитики делают интерактивные запросы через Trino или ClickHouse напрямую к актуальным данным, и ни один из этих движков не должен «знать» о существовании остальных - вся синхронизация состояния происходит через атомарные commit'ы в общий каталог, разобранные в разделе про ACID и MVCC. Без open table format такая многодвижковая архитектура требовала бы либо привязки всех компонентов к одному движку, либо построения собственного слоя синхронизации между копиями данных для каждого движка - что само по себе создаёт ещё одну версию проблемы консистентности, которую мы пытаемся решить.
Анатомия Apache Iceberg: дерево метаданных¶
Чтобы понять, как Iceberg решает проблемы из первого раздела, нужно разобрать его внутреннее устройство. В отличие от «голого» Parquet, где состояние таблицы - это просто содержимое директории на текущий момент, Iceberg явно материализует состояние таблицы в виде дерева файлов метаданных, которое полностью описывает, какие data-файлы являются частью таблицы.
Дерево состоит из четырёх уровней, и каждый уровень решает свою конкретную задачу:
Разберём назначение каждого уровня снизу вверх, потому что именно в этом порядке движок строит план чтения:
Data Files - это обычные Parquet (или ORC, Avro) файлы с фактическими строками таблицы. Iceberg не меняет сам формат хранения данных - всё, что вы знаете про Parquet (row groups, колоночное сжатие, статистика min/max на уровне колонки), продолжает работать без изменений.
Manifest Files (.avro) - список конкретных data-файлов вместе с их метаданными: путь, размер, число строк, диапазон значений (min/max) для каждой колонки, и - что особенно важно для партиционирования - значение партиции для каждого файла, посчитанное в момент записи. Манифест - это единица, на основе которой Iceberg решает, какие файлы стоит читать для конкретного запроса, без необходимости открывать сами Parquet-файлы.
Manifest List (.avro) - список манифестов, которые в совокупности образуют один снапшот таблицы. Если манифестов много (что нормально для крупных таблиц - каждый манифест обычно описывает несколько сотен или тысяч файлов), manifest list хранит дополнительную сводную статистику по каждому манифесту (диапазон партиций, которые он покрывает), что позволяет отбросить целые манифесты без их чтения.
Metadata File (metadata.json) - корневой файл, который описывает таблицу целиком: текущую и все исторические схемы (с уникальными ID колонок - подробнее в разделе про Schema Evolution), текущую и исторические спецификации партиционирования, список всех снапшотов с их таймстампами, и указатель на manifest list, соответствующий текущему снапшоту.
Catalog - последний и критически важный элемент: внешний сервис, который хранит соответствие «имя таблицы → путь к актуальному metadata.json». Catalog - это единственное место, где происходит атомарное переключение между версиями таблицы (подробнее - в разделе про ACID и MVCC ниже).
Почему именно такая иерархия, а не один большой файл¶
Иерархическая структура (manifest list → манифесты → файлы) решает ту самую проблему «дорогого листинга», описанную в начале урока. Вместо того, чтобы делать LIST по объектному хранилищу и читать footer каждого Parquet-файла, движок выполняет план чтения целиком на основе метаданных:
-
Запрашивает у catalog путь к текущему
metadata.json- один быстрый вызов. -
Читает
metadata.json, получает путь к manifest list текущего снапшота. -
Читает manifest list - получает список манифестов и сводную статистику по каждому, что позволяет сразу отбросить манифесты, не содержащие релевантных партиций.
-
Читает только релевантные манифесты, получает точный список data-файлов с их min/max статистикой по колонкам - что позволяет дополнительно отбросить файлы, не подходящие под предикаты запроса.
-
Только теперь происходит первое обращение непосредственно к Parquet-файлам - и только к тем, которые прошли все предыдущие фильтры.
Для таблицы с миллионами файлов это кардинально меняет порядок величины операций планирования: вместо тысяч HTTP-запросов LIST/HEAD к объектному хранилищу - несколько десятков GET-запросов к компактным Avro-файлам метаданных, каждый из которых уже содержит готовую статистику.
Каталоги Iceberg: кто хранит указатель на актуальную версию¶
Слово «Catalog» в архитектуре Iceberg перегружено - это не таблица данных, а отдельный сервис, реализующий простой, но критически важный контракт: атомарно обновить указатель «имя таблицы → путь к metadata.json» и не позволить двум одновременным обновлениям конфликтовать незаметно друг для друга. Существует несколько реализаций каталога, и для self-hosted инфраструктуры (без управляемых облачных сервисов) разумны следующие варианты:
| Тип каталога | Где хранится указатель | Когда применять |
|---|---|---|
| Hive Metastore Catalog | Таблица в реляционной БД метастора (обычно MySQL/Postgres за Hive Metastore) | Уже есть Hive Metastore в инфраструктуре, нужна совместимость с legacy Hive/Trino окружением |
| JDBC Catalog | Произвольная реляционная БД (Postgres, MySQL) напрямую, без Hive Metastore как промежуточного слоя | Self-hosted сценарий без Hadoop-экосистемы: один Postgres как легковесный каталог |
| REST Catalog | Любое хранилище за HTTP API, реализующим Iceberg REST Catalog spec (например, self-hosted Project Nessie или Apache Polaris) | Нужны мультидвижковый доступ, версионирование каталога как у git (branches/tags для целых наборов таблиц) |
| Hadoop Catalog | Сам путь в файловой системе (атомарность через rename операции файловой системы) | Простейшие тестовые/учебные сценарии, без отдельного сервиса каталога |
Для практической части этого урока будет использован JDBC Catalog на базе self-hosted PostgreSQL - это минимальная по числу движущихся частей конфигурация, не требующая поднимать отдельный Hive Metastore или REST-сервис, и хорошо ложится на принцип курса «никаких внешних managed-сервисов, всё self-hosted».
Чтобы сравнение оставалось конкретным, а не только декларативным, вот как выглядит конфигурация Spark для двух других self-hosted вариантов каталога - они не используются в практическом блоке этого урока, но полезны как референс при выборе архитектуры для своего проекта.
# Вариант А: Hive Metastore Catalog
# Подходит, если в инфраструктуре уже есть Hive Metastore
# (например, общий с Trino/Presto или legacy Hive-кластером).
hive_catalog_config = {
"spark.sql.catalog.lakehouse": "org.apache.iceberg.spark.SparkCatalog",
"spark.sql.catalog.lakehouse.type": "hive",
"spark.sql.catalog.lakehouse.uri": "thrift://hive-metastore:9083",
"spark.sql.catalog.lakehouse.warehouse": "s3a://lakehouse/warehouse",
}
# Вариант Б: REST Catalog (например, self-hosted Apache Polaris
# или Project Nessie, развёрнутые как отдельный сервис в кластере).
# REST Catalog даёт версионирование самого каталога (branches/tags
# для согласованного набора таблиц), что недоступно у JDBC/Hive вариантов.
rest_catalog_config = {
"spark.sql.catalog.lakehouse": "org.apache.iceberg.spark.SparkCatalog",
"spark.sql.catalog.lakehouse.catalog-impl": "org.apache.iceberg.rest.RESTCatalog",
"spark.sql.catalog.lakehouse.uri": "http://iceberg-rest:8181",
"spark.sql.catalog.lakehouse.warehouse": "s3a://lakehouse/warehouse",
}
Принципиальной разницы в том, что умеет таблица, между этими вариантами каталога нет - набор гарантий (атомарный CAS commit, snapshot isolation) одинаков для всех реализаций, потому что это требование самой спецификации Iceberg, а не конкретного каталога. Разница - в операционных характеристиках: Hive Metastore тяжелее в поддержке, но даёт совместимость с legacy-инструментами; REST Catalog современнее и гибче (вплоть до git-like веток для каталога), но требует развёртывания и поддержки отдельного сервиса; JDBC Catalog - простейший вариант «одна таблица в существующем Postgres», без дополнительных сервисов вообще.
Table Format Spec v1 и v2: зачем нужна вторая версия¶
Сам формат Iceberg как спецификация развивается версионно, и это важно понимать отдельно от версий снапшотов конкретной таблицы. Table format spec v1 - первая стабильная версия спецификации, поддерживающая append-only нагрузки (INSERT) и полную перезапись файлов для UPDATE/DELETE (Copy-on-Write). Table format spec v2 добавляет ключевую возможность - delete files, которые позволяют логически удалять или изменять отдельные строки без переписывания всего файла, в котором они физически находятся (Merge-on-Read).
Это различие напрямую определяет, насколько эффективно таблица справляется с частыми построчными изменениями - паттерном, типичным для CDC-потоков (Change Data Capture) из транзакционных баз данных. На v1 каждый UPDATE одной строки в файле размером 256 МБ означал бы перезапись всех 256 МБ; на v2 можно записать крошечный delete-файл, ссылающийся на конкретные позиции или значения ключа, и оставить исходный файл данных нетронутым до следующей плановой компакции. Подробный разбор обоих подходов, их trade-off'ов по latency записи и чтения - тема уроков 6 («Copy-on-Write») и 7 («Merge-on-Read») этого модуля; здесь достаточно знать, что при создании таблицы стоит сразу задавать нужную версию явно:
CREATE TABLE lakehouse.analytics.events (
event_id BIGINT,
user_id BIGINT,
event_type STRING,
event_ts TIMESTAMP
)
USING iceberg
PARTITIONED BY (days(event_ts))
TBLPROPERTIES ('format-version' = '2')
Snapshot: как Iceberg фиксирует состояние таблицы¶
Центральное понятие, которое связывает все уровни дерева метаданных - снапшот (snapshot). Каждый снапшот - это неизменяемая (immutable) запись, которая указывает на один конкретный manifest list и тем самым однозначно определяет полный набор data-файлов, видимых как «таблица» в этот момент времени.
Каждая операция записи в Iceberg-таблицу - будь то INSERT, UPDATE, DELETE, MERGE или rewrite_data_files при компакции - создаёт новый снапшот. Существующие снапшоты при этом никогда не модифицируются - таблица в любой момент существует как последовательность снапшотов, где «текущим» считается только последний, на который указывает metadata.json.
Это и есть ответ на главный вопрос из первого раздела урока: «что произошло с таблицей вчера в 14:32» - в Iceberg на этот вопрос можно ответить точно, найдя снапшот с соответствующим таймстампом и заглянув в его manifest list. Полный разбор устройства снапшотов, того, как именно строится дерево манифестов при инкрементальных изменениях, и как Iceberg избегает повторной записи неизменившихся манифестов - тема следующего урока этого модуля. Здесь достаточно зафиксировать главный принцип: снапшот = неизменяемый, атомарно создаваемый снимок полного состояния таблицы.
ACID и MVCC без распределённых блокировок¶
Самый частый неверный мысленный образ, который складывается у инженеров, впервые сталкивающихся с ACID-гарантиями Iceberg: «значит, где-то есть блокировка, как в Postgres, и конкурирующие writer'ы ждут друг друга». Это не так. Iceberg (как и Delta Lake, и большинство современных табличных форматов) реализует Multi-Version Concurrency Control (MVCC) через схему optimistic concurrency control, не требующую централизованного менеджера блокировок.
Идея в следующем: каждый writer работает со своей локальной копией метаданных, не блокируя никого другого. Запись данных (загрузка Parquet-файлов в object storage) происходит полностью независимо и параллельно для разных writer'ов - конфликтов на этом этапе быть не может, потому что каждый writer создаёт файлы с уникальными именами. Единственный момент, где может возникнуть конфликт - это финальный шаг commit: атомарное переключение указателя в каталоге с старой версии metadata.json на новую.
Этот финальный шаг реализован как операция compare-and-swap (CAS): «обнови указатель таблицы на новую версию, только если текущий указатель до сих пор равен той версии, которую я читал в начале своей транзакции». Если кто-то другой уже успел сделать commit между моментом, когда writer прочитал текущую версию, и моментом, когда он пытается закоммитить свою - CAS возвращает отказ, и writer обязан повторить попытку: перечитать новую текущую версию, перестроить свои изменения метаданных на её основе и повторить commit.
Несколько важных следствий этой модели стоит проговорить отдельно, потому что они часто становятся источником путаницы:
-
Читатели никогда не блокируются писателями. Любой читающий запрос, который начался в момент, когда текущая версия была v6, продолжает читать ровно эту версию до конца своего выполнения - даже если в процессе writer B успешно закоммитил v7. Это и есть snapshot isolation: чтение видит консистентный снимок, а не «дёргающееся» состояние, которое меняется на ходу.
-
Конфликт case по умолчанию обрабатывается через retry, а не через ошибку пользователю. Iceberg-клиент (например, библиотека внутри Spark) автоматически повторяет попытку commit'а несколько раз с перестройкой метаданных - в большинстве случаев пользователь даже не замечает, что произошёл конфликт. Если конфликтов слишком много за разумное число попыток (например, десятки writer'ов одновременно делают
MERGEв одну и ту же партицию), commit в итоге завершится с ошибкой, которую приложение должно обработать явно - но это редкий сценарий high-contention нагрузки, а не норма работы. -
Атомарность относится к состоянию таблицы, а не к загрузке файлов. Если writer A успел загрузить часть data-файлов в S3, но затем его commit не прошёл (например, процесс был убит до вызова CAS) - эти файлы физически останутся в хранилище, но никогда не будут видны читателям, потому что ни один manifest не ссылается на них. Такие файлы называются orphan files, и для их периодической очистки существует maintenance-процедура
remove_orphan_files, с которой вы поработаете в одном из следующих уроков модуля.
Эта схема прекрасно масштабируется горизонтально именно потому, что не требует распределённого менеджера блокировок: единственная точка, где нужна синхронизация - это сам catalog, и для него достаточно гарантии «один atomic CAS» (что естественно поддерживается транзакциями в реляционной БД, condition-expression в DynamoDB, или версионированием объектов в Nessie/REST-каталоге).
Что Iceberg добавляет к Parquet: обзор четырёх возможностей¶
Теперь, когда понятна внутренняя механика, можно сформулировать главный практический ответ на вопрос из заголовка урока. Apache Iceberg не меняет то, как физически хранятся строки и колонки - этим продолжает заниматься Parquet. Iceberg добавляет четыре возможности, которые принципиально невозможны (или крайне болезненны) при работе с «голыми» Parquet-файлами и Hive-style партиционированием.
Hidden Partitioning: партиционирование без участия пользователя в SQL¶
В классическом Hive-подходе партиция - это явная колонка, физически закодированная в пути файла (event_date=2024-01-01/). Чтобы движок мог воспользоваться partition pruning (не читать ненужные партиции), пользователь обязан включить условие на эту колонку прямо в WHERE-выражение запроса, причём в том виде, в котором партиция была создана.
Проблема в том, что у бизнес-данных редко есть отдельная колонка «дата партиции» - чаще есть event_timestamp с точностью до секунды, а партиционирование по дате - это производная, техническая деталь физического хранения. Hive-подход вынуждает создавать дополнительную колонку специально для партиционирования (event_date = DATE(event_timestamp)), поддерживать её в актуальном состоянии при каждой записи, и - самое неприятное - помнить о ней в каждом запросе.
Iceberg решает эту проблему через partition transforms - функции, которые движок применяет к обычной колонке таблицы, чтобы вычислить значение партиции, но эта функция не создаёт отдельную видимую колонку в схеме. Партиционирование становится скрытым (hidden) - пользователь работает с обычной колонкой event_ts, а Iceberg сам, на основе сохранённой в метаданных partition spec, понимает, как сматчить предикат запроса с конкретными партициями.
Доступные partition transforms в Iceberg включают временные функции (years, months, days, hours) и функции для категориальных или числовых колонок высокой кардинальности (bucket(N, col) - хэш-распределение по N бакетов, truncate(N, col) - округление строк/чисел до N символов/единиц). Подробное практическое руководство по выбору transform под конкретный паттерн запросов - тема урока 4 этого модуля; здесь важно зафиксировать сам принцип: партиционирование становится деталью реализации, а не частью контракта между таблицей и SQL-запросами к ней.
Schema Evolution: эволюция схемы без переписывания исторических данных¶
В классическом Hive-подходе изменение схемы таблицы - рискованная операция. Переименование колонки на уровне Parquet-файла означает, что старые файлы физически содержат колонку со старым именем, а новые - с новым; для движков, которые матчат колонки по имени (а не по позиции), это означает, что после переименования старые файлы либо «теряют» переименованную колонку, либо требуют полной перезаписи всего исторического датасета.
Iceberg решает эту проблему фундаментально иначе: каждая колонка получает уникальный числовой ID в момент создания, и именно этот ID, а не имя, является источником истины при чтении данных. Имя колонки - это просто текущий «алиас» для конкретного ID, который хранится в metadata.json и может меняться сколько угодно раз без переписывания data-файлов.
Благодаря ID-based tracking Iceberg гарантирует безопасность для четырёх типов изменений схемы без необходимости трогать существующие файлы:
-
Add column - новая колонка получает новый ID; старые файлы при чтении просто возвращают
NULLдля этого ID, как если бы колонка существовала всегда, но не была заполнена. -
Drop column - ID колонки помечается как удалённый в текущей схеме; данные физически остаются в старых файлах, но больше не видны через текущую схему (и могут быть прочитаны через time travel на старую схему, если нужно восстановить историю).
-
Rename column - меняется только маппинг «имя → ID» в текущей версии схемы; ни один байт данных не переписывается.
-
Widen column type - например,
int→long, илиfloat→double(безопасное расширение типа без потери точности); чтение старых файлов сопровождается приведением типа на лету.
Это резко контрастирует с Hive-style таблицами, где аналогичные операции либо требуют полной перезаписи всех файлов, либо приводят к скрытым несоответствиям схемы между старыми и новыми партициями, которые проявляются только во время выполнения запроса непредсказуемой ошибкой парсинга.
Partition Evolution: смена стратегии партиционирования без перестройки Data Lake¶
Партиционирование Hive-таблицы выбирается один раз и фактически «вшивается» в физическую структуру директорий навсегда. Если через год выясняется, что партиционирование по дню слишком грубое (партиции по 500 ГБ, нужна большая параллельность чтения) или слишком мелкое (миллионы крошечных партиций для редко запрашиваемых исторических данных) - единственный способ изменить стратегию - переписать весь исторический Data Lake с новой схемой партиционирования.
Iceberg хранит partition spec как часть версионируемых метаданных, точно так же, как и схему колонок. Это означает, что ALTER TABLE ... ADD PARTITION FIELD создаёт новую версию partition spec, которая применяется только к новым записям - исторические data-файлы продолжают жить со старым spec'ом, а движок при планировании запроса умеет работать с несколькими спецификациями партиционирования одновременно.
Эта возможность критически важна для команд, которые принимают решение о схеме партиционирования на старте проекта с неполной информацией о реальных паттернах запросов - что почти всегда так и происходит на практике. Подробный разбор того, как Spark Catalyst строит единый план выполнения, объединяющий файлы с разными partition spec'ами, и какие есть подводные камни (например, при изменении количества бакетов в bucket(N, col)) - в уроке 5 этого модуля.
Time Travel: чтение данных на произвольный момент истории¶
Последняя из четырёх ключевых возможностей напрямую следует из snapshot-модели, разобранной выше. Поскольку каждая операция записи создаёт новый снапшот, а старые снапшоты не удаляются немедленно (пока явно не вызван expire_snapshots), у Iceberg-таблицы естественным образом есть полная история версий, доступная для чтения.
Time travel решает практические задачи, которые без table format требуют либо отдельной системы версионирования данных, либо ручного бэкапирования полных копий таблицы:
-
Расследование инцидентов - если вчера в отчёте обнаружилась аномалия, можно сравнить снапшот «до» и снапшот «после» подозрительного job'а и точно увидеть, какие строки изменились.
-
Воспроизводимость ML-моделей - модель, обученная на данных, зафиксированных конкретным snapshot ID, может быть воспроизведена точно, даже если таблица с тех пор многократно обновлялась.
-
Аудит и compliance - возможность показать регулятору точное состояние данных на конкретную дату без необходимости хранить отдельные полные копии.
-
Откат ошибочной операции - если
MERGEилиDELETEвыполнен с ошибкой в условии и затронул не те строки, откат к снапшоту перед этой операцией восстанавливает таблицу без восстановления из внешнего бэкапа.
Детальный синтаксис временных и версионных запросов, инкрементальное чтение между двумя снапшотами и механизм управляемого отката (rollback_to_snapshot) разбираются в уроке 9 - здесь будет показан только минимальный практический пример в следующем разделе.
Практический демо-блок: PySpark + Iceberg¶
Переходим от теории к практике. В этом разделе настраивается SparkSession для работы с self-hosted Iceberg-каталогом на базе PostgreSQL (в роли JDBC Catalog) и MinIO (в роли S3-совместимого хранилища), после чего разбираются три демонстрационных кейса: атомарность записи, hidden partitioning «в деле» и простейший time travel.
Настройка окружения¶
Предполагается, что в локальной/тестовой инфраструктуре уже развёрнуты self-hosted PostgreSQL (как backend для каталога) и MinIO (как S3-совместимое хранилище) - например, через docker-compose, как в лабораторных работах модулей про PostgreSQL и S3-совместимое хранилище. Для работы с Iceberg из PySpark необходимо подключить три набора jar-зависимостей: сам Iceberg runtime для конкретной версии Spark, AWS-бандл для S3FileIO, и JDBC-драйвер PostgreSQL для каталога.
from pyspark.sql import SparkSession
spark = (
SparkSession.builder
.appName("iceberg-fundamentals-demo")
.config(
"spark.jars.packages",
"org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.5.2,"
"org.apache.iceberg:iceberg-aws-bundle:1.5.2,"
"org.postgresql:postgresql:42.7.3",
)
# Подключаем Iceberg-расширения Spark SQL: без них недоступны
# CALL-процедуры (rewrite_data_files, expire_snapshots) и синтаксис
# VERSION AS OF / TIMESTAMP AS OF для time travel.
.config(
"spark.sql.extensions",
"org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
)
# Регистрируем каталог "lakehouse" как Iceberg SparkCatalog,
# а в качестве реализации самого каталога - JDBC поверх Postgres.
.config("spark.sql.catalog.lakehouse", "org.apache.iceberg.spark.SparkCatalog")
.config(
"spark.sql.catalog.lakehouse.catalog-impl",
"org.apache.iceberg.jdbc.JdbcCatalog",
)
.config(
"spark.sql.catalog.lakehouse.uri",
"jdbc:postgresql://postgres:5432/iceberg_catalog",
)
.config("spark.sql.catalog.lakehouse.jdbc.user", "iceberg")
.config("spark.sql.catalog.lakehouse.jdbc.password", "iceberg")
# warehouse - корневой путь в object storage, под которым каталог
# будет создавать директории для каждой новой таблицы.
.config("spark.sql.catalog.lakehouse.warehouse", "s3a://lakehouse/warehouse")
# io-impl: используем нативный S3FileIO Iceberg (а не Hadoop S3A),
# поэтому конфигурация endpoint'а идёт через namespace каталога,
# а не через spark.hadoop.fs.s3a.*
.config(
"spark.sql.catalog.lakehouse.io-impl",
"org.apache.iceberg.aws.s3.S3FileIO",
)
.config("spark.sql.catalog.lakehouse.s3.endpoint", "http://minio:9000")
.config("spark.sql.catalog.lakehouse.s3.path-style-access", "true")
.config("spark.sql.catalog.lakehouse.s3.access-key-id", "minioadmin")
.config("spark.sql.catalog.lakehouse.s3.secret-access-key", "minioadmin")
.config("spark.sql.catalog.lakehouse.client.region", "us-east-1")
.getOrCreate()
)
spark.sql("CREATE NAMESPACE IF NOT EXISTS lakehouse.analytics")
Обратите внимание на структуру имён конфигурации: всё, что относится конкретно к каталогу lakehouse (а не к Spark в целом), указывается через префикс spark.sql.catalog.lakehouse.*. Это позволяет в одном SparkSession одновременно работать с несколькими Iceberg-каталогами (например, production и staging), просто меняя префикс. Имя каталога lakehouse в дальнейшем используется как первая часть полностью квалифицированного имени таблицы: lakehouse.analytics.events.
Кейс 1: атомарность записи на практике¶
Демонстрация ключевого свойства из раздела про ACID и MVCC: если запись не дошла до финального commit'а, таблица остаётся ровно в том состоянии, в котором была до попытки записи - без частично примененных изменений.
def append_with_simulated_failure(spark, df, table_name, fail_before_commit=True):
"""
Имитирует сбой ДО завершения commit'а Iceberg-таблицы.
Iceberg буферизует список новых data files на стороне клиента
и отправляет один атомарный commit в конце операции записи.
Если исключение выбрасывается до вызова df.writeTo(...).append(),
ни один файл не попадает в metadata - таблица не меняется.
"""
if fail_before_commit:
raise RuntimeError("Симулированный сбой ДО вызова append()")
df.writeTo(table_name).append()
# Запоминаем количество снапшотов ДО попытки записи
snapshots_before = spark.sql(
"SELECT count(*) AS cnt FROM lakehouse.analytics.events.snapshots"
).collect()[0]["cnt"]
print(f"Снапшотов до попытки записи: {snapshots_before}")
try:
append_with_simulated_failure(
spark, new_events_df, "lakehouse.analytics.events", fail_before_commit=True
)
except RuntimeError as exc:
print(f"Job упал: {exc}")
# Проверяем: появился ли новый снапшот после неудачной попытки?
snapshots_after = spark.sql(
"SELECT count(*) AS cnt FROM lakehouse.analytics.events.snapshots"
).collect()[0]["cnt"]
print(f"Снапшотов после попытки записи: {snapshots_after}")
# Снапшотов до и после совпадает: таблица не изменилась.
# Если бы это была обычная Parquet-директория и сбой произошёл
# ПОСЛЕ записи части файлов, читатели увидели бы "оборванные" данные
# (см. демонстрацию naive_parquet_write_with_failure в начале урока).
Метаданные таблицы lakehouse.analytics.events.snapshots - это не отдельный сервис, а одна из встроенных системных таблиц метаданных (metadata tables), которые Iceberg автоматически предоставляет для любой таблицы через синтаксис <table>.snapshots, <table>.history, <table>.files, <table>.manifests. Это крайне удобный инструмент для отладки и операционного мониторинга, к которому вы вернётесь ещё несколько раз в следующих уроках модуля.
Важная оговорка, которую стоит проговорить явно, чтобы не сформировать неверную интуицию: если бы сбой произошёл после того, как Iceberg успешно загрузил Parquet-файлы в S3, но до финального commit'а в каталог, эти файлы физически остались бы в object storage как orphan files - они не нарушают консистентность таблицы (потому что ни один манифест на них не ссылается), но занимают место и со временем должны быть удалены через remove_orphan_files. Атомарность Iceberg касается видимого состояния таблицы, а не гарантии «либо все байты загружены, либо ни одного» на уровне отдельных файлов.
Кейс 2: hidden partitioning и partition pruning без участия пользователя¶
Создаём таблицу с партиционированием по дням от event_ts, но без отдельной колонки для даты - и убеждаемся, что обычный запрос с фильтром по timestamp использует partition pruning автоматически.
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.analytics.events (
event_id BIGINT,
user_id BIGINT,
event_type STRING,
event_ts TIMESTAMP
)
USING iceberg
PARTITIONED BY (days(event_ts))
""")
# Обратите внимание: в схеме таблицы нет колонки "event_date" -
# партиция вычисляется движком из event_ts на лету при записи.
spark.sql("DESCRIBE TABLE lakehouse.analytics.events").show(truncate=False)
events_batch = spark.createDataFrame(
[
(1, 101, "click", "2024-06-14 10:00:00"),
(2, 102, "click", "2024-06-14 11:30:00"),
(3, 103, "purchase", "2024-06-15 09:15:00"),
(4, 104, "click", "2024-06-15 18:45:00"),
(5, 105, "purchase", "2024-06-16 12:00:00"),
],
schema="event_id BIGINT, user_id BIGINT, event_type STRING, event_ts STRING",
).withColumn("event_ts", F.to_timestamp("event_ts"))
events_batch.writeTo("lakehouse.analytics.events").append()
Теперь выполняем запрос, который фильтрует только по event_ts - точно так, как написал бы аналитик, не знающий деталей физического партиционирования - и смотрим план выполнения:
plan_df = spark.sql("""
SELECT *
FROM lakehouse.analytics.events
WHERE event_ts >= TIMESTAMP '2024-06-15 00:00:00'
AND event_ts < TIMESTAMP '2024-06-16 00:00:00'
""")
plan_df.explain(mode="formatted")
В упрощённом виде физический план содержит характерный для Iceberg раздел Pushed Filters и сводку по числу затронутых файлов, например:
== Physical Plan ==
* Filter (event_ts >= 2024-06-15 00:00:00 AND event_ts < 2024-06-16 00:00:00)
+- * BatchScan lakehouse.analytics.events
PushedFilters: [event_ts >= 2024-06-15T00:00, event_ts < 2024-06-16T00:00]
ScanStatistics: Partitions: 1 of 3 партиций таблицы прочитано
DataFiles: 1 of 5 файлов таблицы прочитано
Запрос ни разу не упомянул партиции явно - но Iceberg, имея в metadata.json сохранённый partition spec days(event_ts), самостоятельно вычислил, что предикат event_ts >= '2024-06-15' AND event_ts < '2024-06-16' полностью попадает в партицию 2024-06-15, и отбросил манифесты, относящиеся к партициям 2024-06-14 и 2024-06-16 ещё на этапе планирования, без обращения к самим Parquet-файлам этих партиций. Аналитику, написавшему этот SQL, вообще не нужно было знать, что таблица партиционирована по дням - это и есть обещанный эффект hidden partitioning.
Кейс 3: простейший time travel¶
Делаем дополнительную запись в таблицу, после чего читаем как текущую версию, так и снапшот, существовавший до этой записи - двумя разными способами: через PySpark DataFrameReader API и через SQL-синтаксис VERSION AS OF.
# Смотрим историю снапшотов ДО второй записи
spark.sql("""
SELECT snapshot_id, committed_at, operation
FROM lakehouse.analytics.events.snapshots
ORDER BY committed_at
""").show(truncate=False)
first_snapshot_id = spark.sql(
"SELECT snapshot_id FROM lakehouse.analytics.events.snapshots "
"ORDER BY committed_at ASC LIMIT 1"
).collect()[0]["snapshot_id"]
print(f"ID первого снапшота: {first_snapshot_id}")
# Делаем вторую запись - добавляем ещё 2 события
late_events = spark.createDataFrame(
[
(6, 106, "click", "2024-06-16 14:00:00"),
(7, 107, "purchase", "2024-06-16 20:30:00"),
],
schema="event_id BIGINT, user_id BIGINT, event_type STRING, event_ts STRING",
).withColumn("event_ts", F.to_timestamp("event_ts"))
late_events.writeTo("lakehouse.analytics.events").append()
# Текущее количество строк (включая вторую запись)
current_count = spark.table("lakehouse.analytics.events").count()
print(f"Текущее количество строк: {current_count}") # 7
# Time travel через PySpark DataFrameReader API - читаем по snapshot-id
historical_df = (
spark.read
.format("iceberg")
.option("snapshot-id", first_snapshot_id)
.load("lakehouse.analytics.events")
)
print(f"Строк в первом снапшоте: {historical_df.count()}") # 5
# Тот же результат через SQL-синтаксис VERSION AS OF
historical_via_sql = spark.sql(f"""
SELECT *
FROM lakehouse.analytics.events VERSION AS OF {first_snapshot_id}
""")
print(f"Строк через VERSION AS OF: {historical_via_sql.count()}") # 5
# Time travel по времени, а не по ID снапшота
historical_by_time = spark.sql("""
SELECT *
FROM lakehouse.analytics.events
TIMESTAMP AS OF '2024-06-15 12:00:00'
""")
Обе формы синтаксиса (через option("snapshot-id", ...) в Python API и через VERSION AS OF / TIMESTAMP AS OF в SQL) обращаются к одному и тому же механизму: вместо чтения текущего указателя из metadata.json, движок напрямую переходит к manifest list конкретного исторического снапшота - и дальше план чтения строится точно так же, как для «текущей» версии, потому что для Iceberg исторический снапшот ничем структурно не отличается от текущего, кроме того, что на него больше не указывает поле current-snapshot-id.
Кейс 4: заглядываем внутрь дерева метаданных через системные таблицы¶
Последний практический кейс - не про конкретную бизнес-задачу, а про инструмент, который пригодится в каждом следующем уроке модуля: Iceberg предоставляет доступ к собственным внутренним метаданным как к обычным SQL-таблицам, без необходимости вручную парсить Avro или JSON файлы.
# Манифесты, входящие в текущий снапшот - видно, сколько файлов
# описывает каждый манифест и какой диапазон партиций он покрывает
spark.sql("""
SELECT path, length, added_data_files_count, existing_data_files_count
FROM lakehouse.analytics.events.manifests
""").show(truncate=False)
# Полный список физических data-файлов таблицы с их статистикой -
# именно эта таблица показывает то самое "DataFiles: 1 of 5",
# которое было видно в плане выполнения в Кейсе 2
spark.sql("""
SELECT file_path, file_format, record_count, file_size_in_bytes
FROM lakehouse.analytics.events.files
""").show(truncate=False)
# История операций таблицы - какая операция привела к какому снапшоту
# (append, overwrite, delete, replace) и является ли снапшот
# частью текущей линии истории (is_current_ancestor)
spark.sql("""
SELECT made_current_at, snapshot_id, parent_id, is_current_ancestor
FROM lakehouse.analytics.events.history
""").show(truncate=False)
Эти три системные таблицы (manifests, files, history) - не отдельная инфраструктура мониторинга, а прямое отражение дерева метаданных, разобранного в начале урока: files - это содержимое манифестов, manifests - это содержимое manifest list, history - это последовательность записей о снапшотах из metadata.json. На уроке про Table Maintenance (урок 10) вы будете использовать именно эти таблицы, чтобы принимать решения о том, когда таблице нужна компакция или очистка старых снапшотов - без них пришлось бы либо доверять интуиции, либо писать собственный код для парсинга Avro-файлов вручную.
Сравнительная таблица: Parquet-lake vs Iceberg-таблица¶
| Возможность | «Голый» Parquet + Hive Metastore | Apache Iceberg |
|---|---|---|
| Атомарность записи | Нет - частичные записи видимы читателям | Да - атомарный CAS commit в каталог |
| Изоляция читателей от писателей | Нет - читатель может увидеть «разорванное» состояние | Да - snapshot isolation, читатель видит консистентный снимок |
| Стоимость планирования запроса | LIST по директориям + чтение footer файлов | Чтение компактных Avro-метаданных, без LIST по object storage |
| Партиционирование | Явная колонка в пути, должна фигурировать в WHERE | Hidden partitioning - партиция вычисляется автоматически |
| Изменение схемы | Риск несовместимости старых/новых файлов, иногда полная перезапись | ID-based tracking, безопасные add/drop/rename/widen без перезаписи |
| Изменение партиционирования | Требует перестройки всего исторического датасета | Partition evolution - старые файлы остаются со старым spec |
| История изменений | Нет, если не вести отдельный лог вручную | Полная история снапшотов + time travel «из коробки» |
| Откат ошибочной операции | Восстановление из внешнего бэкапа | rollback_to_snapshot или time travel запрос |
| Работа с несколькими движками одновременно | Возможна, но без единого контракта консистентности | Возможна, контракт консистентности гарантирован спецификацией |
Типичные заблуждения про Iceberg¶
Прежде чем переходить к следующим урокам, стоит явно проговорить несколько распространённых неверных предположений, которые часто формируются у инженеров на старте работы с table format'ами:
«Iceberg - это просто ещё один формат файла, как Parquet или ORC» - неверно. Iceberg не определяет, как кодируются байты строк и колонок - этим продолжает заниматься Parquet (или ORC/Avro). Iceberg - это слой метаданных и протокол транзакций поверх файлового формата, и одна и та же Iceberg-таблица обычно состоит из Parquet-файлов без каких-либо изменений в их внутренней структуре.
«Можно просто переименовать директорию с Parquet-файлами в Iceberg-таблицу» - неверно. Iceberg-таблица требует построения полного дерева метаданных (manifest, manifest list, metadata.json) для уже существующих файлов. Для этого существует отдельная процедура миграции (add_files, migrate или snapshot процедуры в Spark SQL extension), которая регистрирует существующие файлы в новых метаданных без их физического копирования - но это явный, осознанный шаг миграции, а не переименование директории.
«Snapshot isolation работает через блокировки, как транзакции в PostgreSQL» - неверно, и это разобрано в деталях в разделе про ACID и MVCC выше. Iceberg использует optimistic concurrency через compare-and-swap на уровне каталога; конкурирующие writer'ы не блокируют друг друга, а лишь иногда вынуждены повторить commit, если кто-то их обогнал.
«Time travel хранит данные вечно бесплатно» - неверно. Старые снапшоты ссылаются на старые data-файлы, которые продолжают занимать место в object storage до тех пор, пока явно не будет вызван expire_snapshots. Time travel - это возможность, а не бесплатная функция: за глубину истории нужно платить местом на диске, и операционная практика управления retention - тема урока 10 этого модуля.
«Партиционирование больше не нужно продумывать - Iceberg всё сделает сам» - неверно. Hidden partitioning убирает необходимость вручную создавать колонку и помнить о ней в SQL, но выбор самого partition transform (days vs hours vs bucket(N, col)) и колонки для партиционирования остаётся архитектурным решением инженера, основанным на реальных паттернах запросов - это подробно разбирается в уроке 4.
Production кейс: от еженедельного инцидента к предсказуемой платформе¶
Ситуация. Команда аналитики поддерживала Data Lake на «голом» Parquet поверх self-hosted S3-совместимого хранилища: таблица analytics.orders партиционирована по order_date (явная Hive-style колонка), ежедневный батч-job перезаписывает партицию текущего дня методом mode("overwrite") с фильтром по партиции, а дашборды BI читают эту же таблицу почти непрерывно в течение рабочего дня.
Симптомы. Примерно раз в неделю BI-дашборд показывал аномально низкие цифры за текущий день - на 60-80% ниже ожидаемых - которые сами «исправлялись» в течение следующих 10-15 минут без вмешательства инженеров. Расследование показало классическую гонку чтения и записи: batch job стартовал в 09:00 по расписанию, удалял старые файлы партиции перед записью новых, и если BI-запрос приходил именно в этот промежуток (что становилось более вероятным по мере роста числа пользователей дашборда), он читал директорию партиции в момент, когда старые файлы уже удалены, а новые ещё не полностью записаны.
Первая (неудачная) попытка решения. Команда добавила «защитное окно» - запретила BI-запросы к таблице с 08:55 до 09:15 через приложенческий уровень (кэш дашборда обновлялся не чаще раза в 20 минут). Это не решило проблему по сути, а лишь снизило частоту её проявления, потому что:
-
Время выполнения batch job непредсказуемо колебалось от 8 до 22 минут в зависимости от объёма данных, и иногда «съедало» защитное окно целиком.
-
Любой ad-hoc запрос аналитика напрямую через Spark Shell (минуя BI-кэш) был всё ещё подвержен той же гонке.
-
Проблема стала реже происходить, но не была устранена архитектурно - команда тратила время на компенсацию симптома, а не на исправление причины.
Решение через миграцию на Iceberg. Таблица analytics.orders была мигрирована на Iceberg через процедуру system.migrate (регистрация существующих Parquet-файлов в новых метаданных без физического копирования), партиционирование переведено с явной колонки order_date на hidden partitioning days(order_ts), а batch job переписан с mode("overwrite") по партиции на MERGE INTO (паттерн которого подробно разбирается в уроке 8 модуля).
# Было: mode("overwrite") по партиции - удаление + запись,
# с окном гонки между этими двумя операциями
(
new_orders_df
.write
.mode("overwrite")
.option("replaceWhere", f"order_date = '{today}'")
.parquet("s3a://bucket/orders/")
)
# Стало: атомарный MERGE INTO Iceberg-таблицы -
# для читателей это одна неделимая операция, видимая
# целиком и сразу, либо не видимая вообще до commit'а
spark.sql(f"""
MERGE INTO lakehouse.analytics.orders AS target
USING new_orders_batch AS source
ON target.order_id = source.order_id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
""")
Результат через месяц эксплуатации. Инцидент «аномально низкие цифры на дашборде» не воспроизводился ни разу - не потому, что окно гонки стало уже, а потому, что для Iceberg-таблицы такого окна не существует в принципе: BI-запрос либо видит snapshot до MERGE INTO, либо snapshot после него, и ни при каком тайминге не может увидеть промежуточное, частично изменённое состояние. Дополнительным побочным эффектом стало упрощение кода job'а - больше не требовалась логика разбиения по партициям для replaceWhere, и время выполнения сократилось на 35% за счёт того, что MERGE INTO обновлял только реально изменившиеся записи, а не перезаписывал партицию целиком.
| Метрика | До (Parquet + overwrite) | После (Iceberg + MERGE INTO) |
|---|---|---|
| Частота инцидентов с «разорванным» чтением | ~1 раз в неделю | 0 за месяц наблюдения |
| Среднее время выполнения batch job | 14 минут | 9 минут |
| Необходимость защитного окна для BI | Да, искусственное ограничение | Нет |
| Сложность кода job'а (логика партиций) | Явное управление order_date |
Партиционирование скрыто от кода job'а |
| Возможность расследовать прошлые расхождения | Нет, только текущее состояние | Да, через time travel к снапшоту на момент инцидента |
Ключевой вывод. Проблема гонки чтения/записи не была багом конкретного job'а или недостаточно аккуратной инженерией - это структурное следствие отсутствия атомарности в модели «директория с файлами». Любое количество защитных окон, ретраев и кэширования на уровне приложения снижает частоту проявления проблемы, но не убирает её источник. Table format решает проблему на уровне архитектуры один раз, а не требует постоянной операционной компенсации.
Чек-лист: что получают бизнес и инженер от внедрения Iceberg¶
Подводя итог разобранным возможностям, имеет смысл сформулировать конкретную практическую пользу для двух разных ролей:
Для дата-инженера:
-
Исчезает необходимость писать и поддерживать собственную обвязку для атомарных перезаписей (
write-to-temp-then-swap), которая разобрана в начале урока и которая в реальных production-системах регулярно содержит edge-case баги. -
ETL-пайплайны могут безопасно работать параллельно с аналитическими запросами на одной и той же таблице без риска показать пользователю «разорванные» данные.
-
Schema evolution выполняется как лёгкая DDL-операция (
ALTER TABLE), а не как проект миграции с downtime. -
Отладка инцидентов через time travel и встроенные metadata tables (
snapshots,history,files) занимает минуты вместо часов раскопок в логах и резервных копиях.
Для бизнеса:
-
Снижается риск потери или искажения данных из-за сбоев ETL-процессов, потому что атомарность гарантирована на уровне платформы, а не зависит от качества кода конкретного job'а.
-
Появляется возможность аудита и compliance-отчётности «на конкретную дату» без необходимости содержать дорогую инфраструктуру полных периодических бэкапов.
-
Архитектура остаётся открытой и multi-engine: данные не привязаны к одному вычислительному движку или одному облачному вендору, что снижает риск vendor lock-in и упрощает миграцию между self-hosted и managed окружениями в будущем.
Мостик к следующим урокам модуля¶
Этот урок был намеренно широким - он вводит весь набор понятий, которые в дальнейшем разбираются глубоко и по отдельности:
-
Урок 2 («Iceberg Snapshot Model») детально разбирает, как именно строится дерево манифестов при инкрементальных изменениях, и почему Iceberg не перестраивает весь manifest list заново при каждой записи.
-
Урок 3 («Manifest Files и Manifest List») показывает, как именно движок использует статистику манифестов для построения плана сканирования без физического листинга object storage.
-
Уроки 4 и 5 («Hidden Partitioning» и «Partition Evolution») переходят от концептуального обзора к конкретным partition transform'ам, выбору стратегии под реальные паттерны запросов и деталям совместного планирования по нескольким partition spec'ам.
-
Уроки 6 и 7 («Copy-on-Write» и «Merge-on-Read») разбирают, как Iceberg v2 поддерживает построчные изменения (
UPDATE/DELETE) через delete-файлы - то, что в этом уроке было упомянуто лишь вскользь. -
Урок 8 («MERGE INTO») строит на этой основе production-паттерн upsert для CDC-потоков.
-
Урок 9 («Time Travel») детально раскрывает синтаксис, инкрементальное чтение между снапшотами и управляемый откат, которые здесь были показаны только в минимальном виде.
-
Урок 10 («Table Maintenance») закрывает тему orphan files и retention снапшотов, которая в этом уроке была упомянута как открытый вопрос.
-
Уроки 11-12 переключаются на Delta Lake, чтобы провести параллели и контрасты с архитектурой transaction log, а урок 13 завершает модуль матрицей выбора между Iceberg, Delta Lake и Hudi по конкретным критериям.
Домашнее задание¶
-
Разверните локально self-hosted связку: PostgreSQL (как backend JDBC-каталога), MinIO (как S3-совместимое хранилище) и PySpark с Iceberg-зависимостями - например, через
docker-composeс тремя сервисамиpostgres,minioиspark. Убедитесь, чтоSparkSessionиз раздела «Настройка окружения» успешно подключается иCREATE NAMESPACEвыполняется без ошибок. -
Создайте таблицу
lakehouse.analytics.ordersс hidden partitioning поdays(order_ts), запишите туда тестовый набор из 20-30 строк с разбросом по трём-четырём дням, и черезEXPLAINубедитесь, что фильтр поorder_tsв запросе действительно приводит к чтению только релевантных партиций (сравнитеDataFiles: N of Mв плане для запроса с фильтром и без него). -
Повторите демонстрацию атомарности из «Кейса 1», но усложните её: сделайте реальную (не симулированную) ошибку - например, передайте в
writeTo(...).append()DataFrame с несовместимой схемой - и убедитесь черезlakehouse.analytics.orders.snapshots, что таблица действительно не изменилась. -
Сделайте три последовательные записи в таблицу
orders(вставка, затемDELETEнескольких строк по условию, затем ещё одна вставка), запишите все триsnapshot_idизorders.history, и продемонстрируйтеVERSION AS OFдля каждого из трёх состояний, объяснив словами разницу в количестве строк между ними. -
Прочитайте код процедуры мигрции существующей Parquet-таблицы в Iceberg (
CALL ... system.snapshotилиsystem.migrateв документации Apache Iceberg) и письменно (2-3 абзаца) ответьте на вопрос: какие риски нужно учитывать перед запуском такой миграции на production-таблице, которая активно читается несколькими движками одновременно?
Итоги¶
Parquet отвечает за хранение данных, Iceberg отвечает за управление таблицей. Это главный тезис урока: Iceberg не заменяет файловый формат, а добавляет поверх него слой метаданных (catalog → metadata.json → manifest list → manifest → data files), который превращает набор файлов в полноценную ACID-таблицу.
Четыре боли «голого» Parquet решаются системно, а не точечными патчами. Отсутствие атомарности, дорогой листинг файлов, проблемы согласованности последовательности операций и небезопасный конкурентный доступ - это не баги конкретного инструмента, а структурное следствие отсутствия таблицы как единицы состояния. Iceberg вводит снапшот как такую единицу.
ACID-гарантии реализованы через optimistic concurrency, а не через блокировки. Compare-and-swap commit в каталоге позволяет множеству writer'ов работать параллельно без центрального менеджера блокировок; конфликтующие commits обрабатываются через retry, а не через ожидание.
Hidden Partitioning, Schema Evolution, Partition Evolution и Time Travel - это четыре конкретных следствия одной архитектуры, а не четыре независимых фичи. Все они становятся возможны именно благодаря тому, что схема, partition spec и история снапшотов версионируются вместе как часть единых метаданных таблицы, а не закодированы в физической структуре файлов.
Дальнейшие уроки модуля углубляются в каждый из этих элементов отдельно - от внутреннего устройства манифестов до production-паттернов MERGE INTO и операционного maintenance. Этот урок дал карту территории; следующие уроки проведут детальную экскурсию по каждому её участку.
Краткий глоссарий терминов урока¶
-
Table Format - спецификация метаданных поверх файлового формата, добавляющая таблице транзакционные свойства, версионность и управляемую схему.
-
Catalog - внешний сервис, хранящий атомарно обновляемый указатель «имя таблицы → путь к актуальному
metadata.json». -
Snapshot - неизменяемый снимок полного состояния таблицы на конкретный момент времени, создаваемый каждой операцией записи.
-
Manifest List - файл, перечисляющий манифесты одного снапшота со сводной статистикой по каждому из них.
-
Manifest File - файл со списком конкретных data-файлов и их статистикой (диапазоны значений колонок, значение партиции).
-
Optimistic Concurrency Control (MVCC) - модель конкурентного доступа без централизованных блокировок, где конфликт разрешается через retry на этапе compare-and-swap commit'а, а не через ожидание.
-
Hidden Partitioning - партиционирование через вычисляемые функции (transform) от обычных колонок, не требующее отдельной видимой колонки в схеме и явного упоминания партиции в запросах.
-
Schema Evolution - безопасное изменение схемы (add/drop/rename/widen) на основе ID-based tracking колонок, без переписывания существующих data-файлов.
-
Partition Evolution - изменение спецификации партиционирования для новых записей при сохранении старого spec'а для исторических файлов.
-
Time Travel - чтение таблицы на произвольный исторический снапшот по его ID или ближайшему по времени, через
VERSION AS OF/TIMESTAMP AS OFилиoption("snapshot-id", ...). -
Orphan Files - физически загруженные в хранилище файлы, на которые не ссылается ни один манифест текущей истории снапшотов; удаляются через
remove_orphan_files. -
Copy-on-Write / Merge-on-Read - две стратегии обработки построчных изменений (
UPDATE/DELETE) в Iceberg v2: полная перезапись затронутого файла против записи отдельного delete-файла.