Iceberg Snapshot Model: дерево файлов метаданных и атомарные коммиты
Снапшот как неизменяемый снимок таблицы, полевой разбор metadata.json / manifest list / manifest file, write-and-commit path в три стадии, разрешение конфликтов и CommitFailedException, sequence numbers в v2, инспекция метаданных через Spark и руками через Avro.
Проблема «чёрного ящика»: зачем заглядывать под капот¶
Прошлый урок объяснил, зачем нужен Apache Iceberg и какие проблемы «голого» Parquet он решает. Этого достаточно, чтобы начать пользоваться таблицей: писать INSERT, делать SELECT, доверять тому, что чтение не увидит «разорванное» состояние. Для большинства пользователей SQL-движка этого уровня абстракции достаточно навсегда.
Но дата-инженеру, который отвечает за production-таблицы, рано или поздно понадобится выйти за рамки этой абстракции. Commit падает с CommitFailedException под конкурентной нагрузкой - и нужно понимать, почему именно он падает, а не просто перезапускать job в надежде, что повезёт. Таблица «незаметно» занимает в три раза больше места на диске, чем кажется по count() - и чтобы понять, откуда взялся лишний объём, нужно знать, что именно живёт в дереве метаданных помимо текущих данных. Команда expire_snapshots из урока про maintenance удаляет «что-то» - и прежде чем доверить автоматизации право удалять файлы, стоит точно понимать, что конкретно может быть удалено безопасно, а что нет.
Все эти задачи требуют знания внутреннего устройства Iceberg - не как «чёрного ящика с гарантиями ACID», а как конкретной структуры файлов, конкретного протокола коммитов и конкретных правил, по которым один снапшот сменяется другим. Этот урок - такое «вскрытие» снапшот-модели Iceberg.
Важно сразу обозначить границу этого урока относительно соседних уроков модуля, чтобы не возникало ощущения повторения. Первый урок модуля отвечал на вопрос «зачем» - какие проблемы решает table format в принципе, и дал верхнеуровневую карту дерева метаданных. Этот урок отвечает на вопрос «как именно» - что конкретно записано в каждом файле этого дерева, какая именно последовательность сетевых вызовов и атомарных операций происходит при каждом commit, и что измеримо происходит при конфликте двух конкурирующих писателей. Третий урок модуля пойдёт ещё глубже в один конкретный аспект - как статистика внутри манифестов используется для планирования сканирования - но эта тема здесь затронута только в объёме, необходимом для понимания самой модели снапшотов.
Аналогия с Git: коммиты, указатели, деревья¶
Самая полезная мысленная модель для понимания Iceberg - это не реляционная СУБД вроде PostgreSQL, а система контроля версий вроде Git. Это не случайное сходство: оба решают одну и ту же задачу - дать возможность безопасно версионировать состояние данных, не модифицируя существующие объекты «на месте», и оба для этого используют одну и ту же базовую идею - дерево неизменяемых объектов, на которое указывает изменяемый указатель.
В Git коммит никогда не редактируется - вместо этого создаётся новый коммит, который ссылается на новое дерево файлов, а старое дерево остаётся доступным через предыдущий коммит. Ветка (branch) - это просто изменяемый указатель, который можно перемещать с одного коммита на другой; сами коммиты и деревья, на которые они ссылаются, неизменяемы.
Соответствие почти буквальное: Snapshot в Iceberg играет ту же роль, что Commit в Git - неизменяемая точка в истории, к которой всегда можно вернуться. current-snapshot-id в metadata.json играет роль HEAD - единственного изменяемого указателя, который определяет, что считается «текущей» версией. Manifest List и Manifest File - это аналог дерева Git (Tree), организующего ссылки на содержимое иерархически, а не одним плоским списком. И, наконец, Data File - это аналог Blob, конечного объекта с фактическим содержимым.
Эта аналогия не просто красивая метафора - она напрямую объясняет, откуда у Iceberg берутся возможности, которые иначе выглядели бы как магия: time travel - это то же самое, что git checkout <commit-hash>; история снапшотов - это то же самое, что git log; а будущие возможности branching и tagging (когда несколько именованных указателей могут одновременно ссылаться на разные снапшоты одной таблицы) - это прямой аналог веток и тегов Git. Если вы понимаете, почему git commit никогда не трогает прошлые коммиты, вы уже интуитивно понимаете главный архитектурный принцип Iceberg.
Главный тезис урока: snapshot как неизменяемый снимок¶
Сформулируем главную идею максимально явно, потому что весь остальной материал урока - это её развёртывание в деталях: любое состояние таблицы Iceberg в любой момент времени - это неизменяемый (immutable) снимок данных, называемый snapshot. У таблицы не существует понятия «текущие данные, которые меняются» - есть только последовательность снапшотов, и «текущим» считается тот, на который в данный момент указывает поле current-snapshot-id в metadata.json.
Это полностью переворачивает привычную модель мышления о таблице базы данных. В PostgreSQL UPDATE физически изменяет строку (или создаёт новую версию строки в рамках MVCC, но логически это воспринимается как «изменение существующих данных»). В Iceberg любая операция - INSERT, UPDATE, DELETE, MERGE, rewrite_data_files при компакции - создаёт новый снапшот, который ссылается на новый набор файлов, а предыдущий снапшот остаётся валидным и доступным для чтения до тех пор, пока его явно не «истечёт» команда expire_snapshots.
Небольшой превью того, что в деталях будет разобрано далее: после трёх операций записи (INSERT, INSERT, DELETE) системная таблица <table>.snapshots показывает не «текущее состояние таблицы», а ровно три отдельные, независимые точки истории:
+---------------------+-------------------+-----------+--------------------------+
| snapshot_id | committed_at | operation | summary['total-records'] |
+---------------------+-------------------+-----------+--------------------------+
| 5847203991111111111 | 2024-06-15 09:00 | append | 1000 |
| 5847203991222222222 | 2024-06-15 14:30 | append | 1200 |
| 5847203991234567890 | 2024-06-16 02:00 | delete | 1150 |
+---------------------+-------------------+-----------+--------------------------+
Ни одна из первых двух строк не была «перезаписана» третьей операцией - все три продолжают существовать как полноценные, независимо адресуемые снапшоты. Именно эта таблица и есть то самое «дерево коммитов», о котором шла речь в git-аналогии, только представленное в виде SQL-результата, а не графа в терминале.
Дерево метаданных под микроскопом¶
В прошлом уроке дерево метаданных было разобрано на уровне «кто на кого ссылается». Сейчас цель - разобрать каждый уровень изнутри: какие конкретно поля хранятся в каждом файле, и почему именно эти поля, а не какие-то другие.
Уровень 0: Catalog - единственная роль¶
У каталога (Hive Metastore Catalog, JDBC Catalog, REST Catalog, Hadoop Catalog - все варианты разобраны в прошлом уроке) есть ровно одна обязанность в контексте снапшот-модели: атомарно хранить и обновлять пару «имя таблицы → путь к текущему metadata.json». Каталог не знает ничего о схеме таблицы, партициях или данных - всё это содержится внутри самого metadata.json. Если представить каталог в виде SQL-таблицы (как это происходит в JDBC Catalog поверх Postgres), её структура предельно проста:
-- Упрощённая структура таблицы каталога в JDBC Catalog (Postgres)
CREATE TABLE iceberg_tables (
catalog_name VARCHAR NOT NULL,
table_namespace VARCHAR NOT NULL,
table_name VARCHAR NOT NULL,
metadata_location VARCHAR NOT NULL, -- путь к ТЕКУЩЕМУ metadata.json
previous_metadata_location VARCHAR, -- путь к предыдущему (для аудита и CAS-проверки)
PRIMARY KEY (catalog_name, table_namespace, table_name)
);
Поле metadata_location - это и есть тот самый «указатель HEAD» из git-аналогии. Любой commit в Iceberg-таблицу в конечном счёте сводится к одной SQL-операции вида UPDATE iceberg_tables SET metadata_location = '<новый путь>' WHERE metadata_location = '<путь, который writer прочитал в начале>' - классический compare-and-swap, реализованный через WHERE-условие в UPDATE. Если предыдущий путь не совпадает (потому что кто-то другой уже успел сделать commit), UPDATE затрагивает 0 строк, и клиентская библиотека Iceberg интерпретирует это как конфликт коммита.
Уровень 1: metadata.json - сердце таблицы¶
metadata.json - корневой документ, который описывает таблицу полностью: всё, что нужно знать о её структуре и истории, кроме самих данных. Ниже - упрощённое, но структурно точное представление реального содержимого такого файла после нескольких операций записи:
{
"format-version": 2,
"table-uuid": "9c2b9b1e-7e3a-4f3b-8a2e-1234567890ab",
"location": "s3a://lakehouse/warehouse/analytics/events",
"last-sequence-number": 3,
"last-updated-ms": 1718480400000,
"last-column-id": 4,
"current-schema-id": 0,
"schemas": [
{
"schema-id": 0,
"fields": [
{"id": 1, "name": "event_id", "type": "long", "required": false},
{"id": 2, "name": "user_id", "type": "long", "required": false},
{"id": 3, "name": "event_type", "type": "string", "required": false},
{"id": 4, "name": "event_ts", "type": "timestamp", "required": false}
]
}
],
"default-spec-id": 0,
"partition-specs": [
{
"spec-id": 0,
"fields": [
{"name": "event_ts_day", "transform": "day", "source-id": 4, "field-id": 1000}
]
}
],
"current-snapshot-id": 5847203991234567890,
"snapshots": [
{
"snapshot-id": 5847203991111111111,
"sequence-number": 1,
"timestamp-ms": 1718470000000,
"manifest-list": "s3a://lakehouse/warehouse/analytics/events/metadata/snap-5847203991111111111-1.avro",
"summary": {
"operation": "append",
"added-data-files": "3",
"added-records": "5",
"total-records": "5",
"total-data-files": "3"
}
},
{
"snapshot-id": 5847203991234567890,
"parent-snapshot-id": 5847203991111111111,
"sequence-number": 2,
"timestamp-ms": 1718480400000,
"manifest-list": "s3a://lakehouse/warehouse/analytics/events/metadata/snap-5847203991234567890-1.avro",
"summary": {
"operation": "append",
"added-data-files": "2",
"added-records": "2",
"total-records": "7",
"total-data-files": "5"
}
}
],
"snapshot-log": [
{"timestamp-ms": 1718470000000, "snapshot-id": 5847203991111111111},
{"timestamp-ms": 1718480400000, "snapshot-id": 5847203991234567890}
],
"metadata-log": [
{"timestamp-ms": 1718470000000, "metadata-file": "s3a://.../00000-aaa.metadata.json"},
{"timestamp-ms": 1718480400000, "metadata-file": "s3a://.../00001-bbb.metadata.json"}
]
}
Разберём смысловые блоки этого документа отдельно, потому что каждый из них отвечает за свою, отдельную часть гарантий Iceberg:
-
schemas+current-schema-id- вся история схем таблицы. Каждая колонка имеет числовойid, который никогда не переиспользуется даже после удаления колонки - это и есть фундамент ID-based schema evolution, разобранного в прошлом уроке.current-schema-idуказывает, какая запись из массиваschemasявляется «активной» прямо сейчас. -
partition-specs+default-spec-id- история спецификаций партиционирования. Массив, а не одно значение - именно это позволяет partition evolution: старые data-файлы остаются привязаны кspec-id, под которым они были записаны, и одновременно существует «текущий»default-spec-idдля новых записей. -
current-snapshot-id- это и есть тот самый указательHEADиз git-аналогии. Любой читающий запрос без явного time travel начинает планирование именно с этого ID, находит соответствующую запись в массивеsnapshots, и переходит по полюmanifest-list. -
snapshots- полная история снапшотов таблицы. Обратите внимание на полеparent-snapshot-idу второго снапшота - оно явно формирует цепочку, аналогичную истории коммитов Git, что используется, например, для инкрементального чтения «всё, что изменилось между снапшотом A и снапшотом B». -
summary- человекочитаемая (и машиночитаемая) сводка о том, что именно произошло в этом снапшоте: тип операции (append,overwrite,delete,replace), число добавленных файлов и строк, и - что особенно полезно для мониторинга - кумулятивные счётчикиtotal-recordsиtotal-data-filesна момент этого снапшота. Именно из этого поля строится системная таблица<table>.snapshots, которую вы видели в практическом блоке прошлого урока. -
properties- произвольная картаключ → значение, через которую настраивается поведение таблицы: именно в ней физически хранятся все TBLPROPERTIES, упомянутые в этом уроке (commit.manifest-merge.enabled,commit.retry.num-retries,format-versionи десятки других). Когда вы выполняетеALTER TABLE ... SET TBLPROPERTIES (...), результат этой команды - не что иное, как обновлённая запись в картеpropertiesвнутри новогоmetadata.json, созданного как обычный commit с собственным снапшотом конфигурации (хотя сама операция изменения свойств не создаёт новый snapshot данных - она лишь обновляет метаданные таблицы). -
snapshot-logиmetadata-log- два дополнительных, более «низкоуровневых» журнала.snapshot-logфиксирует, в какой момент времени каждый снапшот стал текущим (что не всегда совпадает с моментом его создания - например, приrollback_to_snapshotтекущим снова становится старый снапшот, и вsnapshot-logпоявляется новая запись с его ID и новым timestamp).metadata-logотслеживает историю самих файловmetadata.json- это нужно для очень редкого, но важного сценария восстановления после повреждения текущегоmetadata.json.
Уровень 2: Manifest List - снимок конкретного коммита¶
Почему manifest list - это отдельный файл в формате Avro, а не просто ещё один JSON-массив внутри metadata.json? Ответ напрямую связан с проблемой масштабирования метаданных: если бы список всех манифестов каждого снапшота хранился прямо в metadata.json, этот файл рос бы пропорционально числу снапшотов и числу манифестов в каждом из них, и для таблицы с тысячами снапшотов metadata.json сам по себе стал бы узким местом - его нужно перечитывать целиком при каждом обращении к таблице. Вынесение списка манифестов в отдельный Avro-файл на каждый снапшот означает, что metadata.json хранит только путь к этому файлу, а не его содержимое - сам metadata.json остаётся компактным независимо от размера таблицы.
Avro выбран не случайно - это row-based бинарный формат с компактной сериализацией и встроенной схемой, что делает его эффективным для хранения именно такого рода списков структурированных записей (в отличие от Parquet, который оптимален для аналитических сканов больших колоночных датасетов - но manifest list никогда не сканируется аналитически, его всегда читают целиком).
Ключевые поля одной записи manifest list (по одной записи на каждый manifest file, входящий в снапшот):
| Поле | Назначение |
|---|---|
manifest_path |
Путь к самому manifest-файлу (.avro) |
manifest_length |
Размер manifest-файла в байтах |
partition_spec_id |
Какой spec-id использовался при записи файлов этого манифеста |
added_snapshot_id |
В каком снапшоте этот манифест был впервые добавлен |
added_data_files_count / existing_data_files_count / deleted_data_files_count |
Сколько файлов в этом манифесте новые, сколько унаследованы от предыдущего снапшота, сколько помечены как удалённые |
partitions |
Критически важное поле: массив диапазонов (lower_bound / upper_bound) по каждому полю партиционирования, агрегированный по всем файлам внутри манифеста |
Именно поле partitions с диапазонами min/max делает возможным manifest pruning - отбрасывание целых манифестов без их открытия, если запрос фильтрует по партиции, а диапазон манифеста не пересекается с предикатом. Полный разбор того, как именно Spark Catalyst использует это поле при построении плана сканирования - тема следующего урока («Manifest Files и Manifest List: как Iceberg планирует scan без листинга»); здесь важно зафиксировать, что эта оптимизация физически возможна только потому, что manifest list заранее агрегирует статистику по партициям на уровне манифеста, а не заставляет открывать каждый манифест, чтобы это выяснить.
Уровень 3: Manifest File - непосредственный учёт файлов данных¶
Manifest file - это уровень, на котором метаданные впервые становятся привязаны к конкретным физическим data-файлам. Каждая запись (manifest entry) внутри manifest-файла описывает один data file (в v2 - также возможно delete file, к этому вернёмся в разделе про sequence numbers) и обязательно содержит поле status, которое может принимать одно из трёх значений:
Этот статус - не просто техническая деталь, а основа эффективности всей модели. Когда Iceberg создаёт новый манифест в рамках нового снапшота, он не обязан копировать туда все записи из предыдущего манифеста - вместо этого он может либо переиспользовать существующий неизменённый манифест целиком (если ни один из его файлов не затронут текущей операцией), либо создать новый манифест только для изменившейся части данных. Записи со статусом EXISTING существуют именно для тех случаев, когда новый манифест объединяет «старые» файлы (унаследованные без изменений) с «новыми» файлами одной операции - подробный разбор того, когда Iceberg решает переиспользовать манифест целиком, а когда создавать новый, и как это влияет на число «мелких» манифестов со временем - тема следующего урока.
Помимо статуса, каждая запись манифеста содержит структуру data_file со статистикой, необходимой для планирования сканов без обращения к самому Parquet-файлу:
Поле data_file |
Назначение |
|---|---|
file_path, file_format |
Путь к физическому файлу и его формат (PARQUET, ORC, AVRO) |
partition |
Значение партиции для этого конкретного файла (например, {event_ts_day: 19888} - число дней с эпохи) |
record_count |
Количество строк в файле - используется для оценки кардинальности без чтения файла |
file_size_in_bytes |
Размер файла - используется при планировании задач (bin-packing файлов в task'и) |
column_sizes, value_counts, null_value_counts |
Поколоночные карты {column_id: значение} - размер, число значений и число NULL по каждой колонке |
lower_bounds, upper_bounds |
Поколоночные карты {column_id: байтовое представление значения} - минимум и максимум по каждой колонке, на основе которых строится predicate pushdown без открытия файла |
Эти поля - буквально те же самые row-group статистики, которые Parquet хранит в своём footer, только агрегированные на уровень целого файла и продублированные в манифесте. Дублирование выглядит избыточным, но даёт принципиальный выигрыш: чтобы получить эту статистику из самого Parquet-файла, нужно сделать GET-запрос (хотя бы частичный, для footer) к каждому файлу; чтобы получить её из манифеста - достаточно одного GET-запроса к манифесту, который уже содержит статистику по сотням файлов сразу.
Уровень 4: Data Files - где заканчиваются метаданные¶
На этом уровне дерево метаданных заканчивается и начинается территория, полностью описанная спецификацией Parquet (или ORC/Avro), без какого-либо участия Iceberg. Файл data-00001.parquet ничего не знает о том, что он является частью Iceberg-таблицы - это полностью самостоятельный, валидный Parquet-файл, который можно открыть любым инструментом, понимающим Parquet, в полном отрыве от Iceberg.
Это разделение ответственности принципиально важно для эволюции экосистемы: команда Apache Parquet может продолжать улучшать формат файла (новые алгоритмы сжатия, новые типы кодирования) без необходимости менять спецификацию Iceberg, а команда Apache Iceberg может развивать модель метаданных (sequence numbers, branching, row lineage) без необходимости менять формат самого файла данных. Граница ответственности проходит ровно по этому уровню дерева.
Сравнение четырёх уровней по частоте обновления и размеру¶
Чтобы свести разобранную иерархию в одну практическую сводку, полезно явно сопоставить все четыре уровня по двум измерениям, которые определяют их роль в системе: как часто объект пересоздаётся и насколько он велик.
| Уровень | Формат | Создаётся заново | Типичный размер | Кто читает первым |
|---|---|---|---|---|
metadata.json |
JSON | При каждом commit (новая версия файла) | Килобайты - сотни килобайт (растёт со схемами/снапшотами) | Catalog (через указатель), затем сам движок |
| Manifest List | Avro | При каждом commit (новый список для нового снапшота) | Десятки - сотни килобайт на снапшот | Движок, для отбора релевантных манифестов |
| Manifest File | Avro | Только для изменившейся части данных - старые манифесты переиспользуются (EXISTING) |
Сотни килобайт - единицы мегабайт на манифест | Движок, для отбора релевантных data-файлов |
| Data File | Parquet/ORC | Только при физической записи новых строк или компакции | Десятки - сотни мегабайт (целевой размер) | Executor, для непосредственного чтения строк |
Главный паттерн, который стоит унести из этой таблицы: частота пересоздания убывает сверху вниз ровно настолько, насколько растёт стоимость пересоздания. metadata.json дешёво пересоздавать целиком при каждом commit, потому что он маленький. Manifest list тоже пересоздаётся при каждом commit, но остаётся компактным, потому что не дублирует статистику файлов - только агрегаты по манифестам. А вот сами манифесты и тем более data-файлы - дорогие в пересоздании объекты, и архитектура Iceberg специально устроена так, чтобы переиспользовать их между снапшотами максимально долго, перестраивая только то, что действительно изменилось.
Слияние манифестов: почему их число не растёт безгранично¶
Из механики Stage 2 («Metadata») может показаться, что после тысячи операций append в таблице накопится тысяча отдельных манифестов - по одному на каждый commit, переиспользуемых через EXISTING-записи последующих манифестов. Если бы это было так, manifest list со временем стал бы таким же узким местом, каким был листинг директорий для Hive - просто на другом уровне иерархии.
Iceberg решает эту проблему встроенным механизмом manifest merging: при каждом commit клиентская библиотека проверяет суммарный размер существующих манифестов снапшота, и если он превышает целевой порог, несколько мелких манифестов автоматически объединяются в один новый, переписывая их записи (включая статус EXISTING) в единый файл.
# Настройки, управляющие автоматическим слиянием манифестов
table_properties = {
# Включить (по умолчанию включено) автоматическое слияние манифестов
"commit.manifest-merge.enabled": "true",
# Целевой размер одного manifest-файла после слияния
"commit.manifest.target-size-bytes": str(8 * 1024 * 1024), # 8 МБ
# Минимальное число манифестов, при котором запускается слияние -
# защита от слияния при единичных мелких манифестах
"commit.manifest.min-count-to-merge": "100",
}
spark.sql(f"""
ALTER TABLE lakehouse.analytics.events SET TBLPROPERTIES (
'commit.manifest-merge.enabled' = 'true',
'commit.manifest.target-size-bytes' = '{8 * 1024 * 1024}',
'commit.manifest.min-count-to-merge' = '100'
)
""")
Этот механизм - прямая параллель к компакции data-файлов (разобранной в модуле про S3-совместимое хранилище), но на уровень выше: вместо объединения множества мелких Parquet-файлов в крупные, Iceberg объединяет множество мелких manifest-файлов в крупные. Без этого механизма таблица с высокой частотой записи (например, micro-batch streaming каждые 30 секунд) накопила бы тысячи манифестов за несколько дней, и каждое чтение manifest list требовало бы открытия всех их по отдельности для построения плана - именно та проблема масштабирования метаданных, которую сама иерархия дерева была призвана решить. Manifest merging - это причина, по которой Iceberg-таблица с миллионами commit'ов в истории всё равно сохраняет компактный и быстро читаемый manifest list для текущего снапшота.
Механика атомарного коммита: Write-and-Commit Path¶
Зная структуру дерева метаданных, можно детально разобрать, что именно происходит «под капотом» при выполнении, например, df.writeTo("lakehouse.analytics.events").append(). Этот процесс называется Write-and-Commit Path и состоит из трёх чётко разделённых стадий.
Stage 1: Write - запись данных независимо от других писателей¶
На этой стадии executor'ы Spark пишут новые Parquet-файлы в object storage. Это самая «долгая» по времени стадия для крупных объёмов данных, но при этом самая безопасная с точки зрения конкурентности: каждый writer генерирует файлы с гарантированно уникальными именами (обычно на основе UUID или комбинации task ID и attempt ID), поэтому даже если два независимых job'а пишут в одну и ту же таблицу одновременно, на этой стадии физически невозможен конфликт имён файлов или перезапись чужих данных. Запись на этой стадии - не транзакционная операция в смысле видимости: записанные файлы физически существуют в object storage, но ни один манифест ещё не ссылается на них, поэтому ни один читатель таблицы их не увидит, пока не произойдёт Stage 3.
Stage 2: Metadata - построение нового дерева манифестов¶
После завершения записи всех data-файлов driver (а не executor'ы) строит новые объекты метаданных: новый manifest file со списком только что записанных файлов (со статусом ADDED), и новый manifest list, который объединяет этот новый манифест с манифестами из предыдущего снапшота (для тех файлов, которые не изменились - со статусом EXISTING, как было описано в разделе про manifest file). Эта стадия происходит целиком в памяти и в виде записи новых Avro/JSON объектов в object storage - сама операция записи этих файлов метаданных также не делает таблицу видимой кому-либо, потому что metadata.json каталога ещё не обновлён.
Stage 3: Commit - атомарная попытка переключить указатель¶
Финальная и самая короткая по времени стадия: driver отправляет в catalog запрос вида «переключи metadata_location для таблицы analytics.events с пути v5.metadata.json (который я прочитал в начале операции) на путь v6.metadata.json (который я только что построил)». Это - тот самый compare-and-swap, разобранный в прошлом уроке на уровне идеи, а здесь - на уровне конкретного SQL-запроса (для JDBC Catalog):
-- Именно такой (упрощённый) запрос выполняет JDBC Catalog
-- внутри транзакции БД при попытке commit'а
UPDATE iceberg_tables
SET metadata_location = 's3a://lakehouse/warehouse/analytics/events/metadata/v6.metadata.json',
previous_metadata_location = 's3a://lakehouse/warehouse/analytics/events/metadata/v5.metadata.json'
WHERE catalog_name = 'lakehouse'
AND table_namespace = 'analytics'
AND table_name = 'events'
AND metadata_location = 's3a://lakehouse/warehouse/analytics/events/metadata/v5.metadata.json';
-- Если UPDATE затронул 1 строку - commit успешен.
-- Если UPDATE затронул 0 строк - значит, metadata_location уже
-- не равен v5 (кто-то другой успел закоммитить раньше) - конфликт.
Если результат UPDATE - одна затронутая строка, commit считается успешным, и именно с этого момента (а не с момента записи Parquet-файлов и не с момента построения манифестов) таблица становится видна читателям в новом состоянии. Если результат - ноль затронутых строк, клиентская библиотека Iceberg перехватывает это как сигнал конфликта и переходит к логике разрешения конфликтов, разобранной в следующем разделе.
Разрешение конфликтов: retry и CommitFailedException¶
Когда CAS-операция в Stage 3 завершается отказом, Iceberg-клиент не сразу сообщает об ошибке пользователю - первым делом он пытается автоматически разрешить конфликт через повторную попытку (retry). Но прежде чем повторить попытку, важно различить два принципиально разных типа конфликтов.
Когда конфликт безопасно разрешается автоматически¶
Если конфликтующая операция - это, например, два независимых append в разные партиции (или даже в одну партицию, но без перекрывающихся файлов), Iceberg может безопасно «перестроить» Stage 2 на основе новой текущей версии и повторить Stage 3. Конкретно это выглядит так: writer B обнаруживает, что текущий снапшот теперь v6 (закоммитил writer A), перечитывает v6.metadata.json, строит новый manifest list, который объединяет уже закоммиченные манифесты из v6 с собственными новыми манифестами writer B, и повторяет попытку CAS - теперь уже с условием «если текущий v6 → переключить на v7».
Ключевая деталь, которую часто упускают: retry - это не просто «повторить тот же запрос ещё раз». Это полноценная перестройка Stage 2 на основе новой базовой версии. Именно поэтому конфликт на уровне CAS-указателя в большинстве случаев не означает конфликт на уровне данных - два append в разные партиции физически не пересекаются, и Iceberg успешно объединяет результаты обоих writer'ов в одной последовательной истории снапшотов.
Когда конфликт приводит к CommitFailedException¶
Не каждый конфликт можно разрешить простым перестроением. Операции, которые требуют валидации относительно состояния таблицы на момент начала транзакции (а не только независимой записи новых файлов), могут обнаружить, что эта валидация больше не проходит на новой базовой версии. Классический пример - overwritePartitions() или операции типа MERGE INTO/DELETE, которые проверяют: «файлы, которые я собираюсь удалить/заменить, всё ещё являются частью текущего снапшота?». Если конкурирующий writer уже удалил или заменил эти же файлы (например, два параллельных MERGE INTO пытаются обработать одну и ту же партицию), валидация завершается отказом - и в этом случае retry не помогает, потому что конфликт находится не на уровне указателя, а на уровне семантики операции.
В реальности CommitFailedException (или ValidationException в более ранних версиях библиотеки) - это сигнал, что приложение должно явно решить, что делать дальше: повторить операцию MERGE INTO целиком с самого начала (включая Stage 1 - возможно, с другим набором изменений, потому что входные данные для merge могли тоже измениться), либо эскалировать ошибку оператору. Это принципиально отличается от ситуации с независимыми append-операциями, где retry полностью прозрачен для приложения.
Настройка поведения retry¶
Количество попыток и стратегия задержки между ними настраиваются на уровне каталога:
spark_config = {
# Максимальное число попыток commit перед CommitFailedException
"spark.sql.catalog.lakehouse.commit.retry.num-retries": "4",
# Начальная задержка перед первой повторной попыткой (мс)
"spark.sql.catalog.lakehouse.commit.retry.min-wait-ms": "100",
# Максимальная задержка между попытками (экспоненциальный backoff)
"spark.sql.catalog.lakehouse.commit.retry.max-wait-ms": "60000",
# Общий потолок суммарного времени всех попыток
"spark.sql.catalog.lakehouse.commit.retry.total-timeout-ms": "1800000",
}
Экспоненциальный backoff (увеличивающаяся задержка между попытками) - стандартная практика для optimistic concurrency: при высокой конкурентности (десятки writer'ов в одну партицию) немедленный повтор всеми writer'ами одновременно лишь увеличивает шанс нового конфликта, а нарастающая случайная задержка снижает вероятность того, что несколько writer'ов снова столкнутся в одно и то же окно времени.
Sequence Numbers: v1 vs v2 в контексте снапшотов¶
В прошлом уроке было упомянуто различие table format spec v1 (только append-only и Copy-on-Write) и v2 (добавляет delete files для Merge-on-Read). С точки зрения снапшот-модели это различие реализовано через концепцию sequence number - монотонно возрастающего целого числа, которое присваивается каждому снапшоту в момент его создания.
В v1 sequence number не имел практического значения для логики чтения - снапшоты упорядочивались просто по parent-snapshot-id и timestamp-ms. В v2 sequence number становится центральным механизмом, который связывает data files и delete files без необходимости физически модифицировать data file при удалении или изменении отдельных строк.
Механика работает следующим образом: каждая запись манифеста в v2 хранит два числа - sequence_number (порядковый номер снапшота, в котором эта запись данных логически появилась) и file_sequence_number (номер снапшота, в котором физический файл был фактически записан - они совпадают для большинства файлов, но могут отличаться при особых операциях типа cherry-pick). Когда движок планирует чтение, для каждого data-файла он ищет все delete-файлы, чей sequence_number больше или равен sequence number данных в этом data-файле - то есть delete-файл, созданный позже, потенциально удаляет строки из файла, созданного раньше, но не наоборот.
Это и есть техническая основа Merge-on-Read: вместо того, чтобы переписывать data-file-A целиком ради удаления нескольких строк, Iceberg v2 записывает компактный position-delete-file, который ссылается на конкретные позиции (номера строк) внутри data-file-A, помечая их как удалённые - а sequence number гарантирует, что этот delete-файл будет правильно применён при чтении любого снапшота, начиная с того, в котором он был создан, и не затронет данные, записанные позже в файлы с большим sequence number. Детальный разбор структуры delete-файлов (position deletes vs equality deletes), их формата и стоимости накопления слишком большого числа delete-файлов между компакциями - тема урока 7 этого модуля («Merge-on-Read»); здесь важно зафиксировать, что именно sequence number делает эту схему корректной без какой-либо блокировки или координации между writer'ами.
Снапшоты и Sequence Number в metadata.json¶
Поле last-sequence-number в metadata.json (которое уже встречалось в примере JSON выше) - это просто последнее выданное значение, монотонно увеличивающееся с каждым новым снапшотом. Несколько writer'ов, коммитящих конкурентно, получают разные sequence number даже если их commits происходят «почти одновременно» - именно монотонность этого счётчика, а не временные метки (которые могут быть неточными из-за рассинхронизации часов между узлами), является источником истины для порядка применения delete-файлов.
Snapshot Lineage: родительские связи и журнал переключений¶
Помимо текущего снапшота, дерево метаданных хранит полную цепочку родительских связей через поле parent-snapshot-id, что формирует линейную (а в будущем - с поддержкой branching, не линейную) историю изменений таблицы.
Этот граф родственных связей пригодится не только для отладки - он напрямую используется при инкрементальном чтении (read.option("start-snapshot-id", X).option("end-snapshot-id", Y)), когда движку нужно определить именно те файлы, которые были добавлены в диапазоне снапшотов между X и Y, без повторного чтения файлов, не изменившихся за этот период. Обратите внимание на отдельную важную деталь, изображённую на диаграмме: операция rollback_to_snapshot не удаляет снапшоты 3 и 4 из истории - она лишь переключает current-snapshot-id обратно на снапшот 2, делая его снова «текущим». Снапшоты 3 и 4 остаются полноценной частью истории snapshots и snapshot-log, и к ним всё ещё можно обратиться через time travel, пока их явно не удалит expire_snapshots. Это резко отличается от интуиции «откат как git reset --hard» - в Iceberg откат гораздо ближе к git checkout на старый коммит без потери последующей истории.
Это понятие границы между «иметь снапшоты в истории» и «быть текущим снапшотом» - прямая основа для branching и tagging, упомянутых в git-аналогии в начале урока. Это не гипотетическая возможность, а реально существующая часть спецификации Iceberg: помимо единственного безымянного указателя current-snapshot-id, таблица может хранить произвольное число дополнительных поименованных указателей в поле refs внутри metadata.json.
"refs": {
"main": {
"snapshot-id": 5847203991234567890,
"type": "branch"
},
"audit-2024-q2": {
"snapshot-id": 5847203991111111111,
"type": "tag",
"max-ref-age-ms": 7776000000
}
}
branch - изменяемая ссылка, которая продолжает двигаться вперёд с каждым новым commit'ом в эту ветку (ровно как main в Git). tag - неизменяемая, замороженная ссылка на один конкретный снапшот, обычно с опциональным сроком жизни (max-ref-age-ms), по истечении которого expire_snapshots разрешается удалить этот снапшот, даже если он не является предком текущего. Управление ветками и тегами выполняется через SQL прямо из Spark:
-- Создать тег, замороженный на конкретном снапшоте - удобно
-- для compliance-целей: "состояние таблицы на конец квартала"
ALTER TABLE lakehouse.analytics.events
CREATE TAG `audit-2024-q2` AS OF VERSION 5847203991111111111
RETAIN 90 DAYS;
-- Создать отдельную ветку для экспериментальной ETL-логики,
-- не затрагивающую основную ветку main
ALTER TABLE lakehouse.analytics.events
CREATE BRANCH `experiment-new-dedup-logic`;
-- Запись и чтение в рамках конкретной ветки
INSERT INTO lakehouse.analytics.events.branch_experiment_new_dedup_logic
SELECT * FROM staging_events;
SELECT * FROM lakehouse.analytics.events.branch_experiment_new_dedup_logic;
Практический смысл этой возможности - в том, что несколько независимых линий изменений одной физической таблицы (со своими собственными цепочками снапшотов и собственными parent-snapshot-id) могут существовать одновременно, не мешая друг другу, и в любой момент одна ветка может быть слита в другую через fast_forward (аналог git merge --ff-only) - если её история является прямым продолжением целевой ветки.
Стоит явно различать rollback_to_snapshot (упомянутый ранее в разделе про lineage) и подход через теги/ветки - у них разное назначение, хотя оба оперируют одним и тем же графом снапшотов:
| Сценарий | Инструмент | Что происходит с current-snapshot-id ветки main |
|---|---|---|
| Откатить ошибочную операцию немедленно | rollback_to_snapshot |
Переключается на старый снапшот сразу, для всех читателей main |
| Сохранить snapshot для аудита на N дней | CREATE TAG ... RETAIN N DAYS |
Не меняется - тег существует параллельно, не затрагивая main |
| Протестировать новую ETL-логику без риска для прода | CREATE BRANCH |
Не меняется - ветка живёт изолированно, пока не будет слита явно |
| Посмотреть данные на конкретный момент без изменения состояния | Time Travel (VERSION AS OF) |
Не меняется вообще - это read-only операция, не создающая commit |
Эта таблица - удобный практический ориентир: если цель - «временно посмотреть», нужен time travel; если «зафиксировать для будущего сравнения», нужен tag; если «развивать параллельно», нужна branch; и только если цель - «вернуть таблицу взад прямо сейчас для всех», нужен rollback_to_snapshot.
Практический демо-блок: инспекция метаданных через PySpark¶
Переходим к практике. В этом разделе три кейса: чтение метаданных через системные SQL-таблицы Spark (быстро и удобно для повседневной работы), прямое чтение тех же самых файлов вручную через Python и Avro-библиотеку (чтобы увидеть, что SQL-абстракция не скрывает никакой магии), и симуляция конкурентного конфликта commit'ов с воспроизведением CommitFailedException.
Конфигурация SparkSession идентична прошлому уроку - JDBC Catalog на self-hosted PostgreSQL и S3FileIO на MinIO, поэтому здесь она не повторяется; предполагается, что spark уже инициализирован, а таблица lakehouse.analytics.events создана и содержит несколько снапшотов после серии операций записи.
Кейс 1: системные таблицы метаданных под лупой¶
В прошлом уроке системные таблицы <table>.snapshots, <table>.history, <table>.files, <table>.manifests были представлены как удобный инструмент инспекции. Сейчас, имея полное понимание дерева метаданных, можно сопоставить каждую колонку этих таблиц с конкретным полем из metadata.json или manifest list, разобранными выше.
# snapshots: прямое отражение массива "snapshots" из metadata.json
spark.sql("""
SELECT
snapshot_id,
parent_id,
committed_at,
operation,
summary['added-data-files'] AS added_files,
summary['added-records'] AS added_records,
summary['total-records'] AS total_records,
summary['total-data-files'] AS total_files
FROM lakehouse.analytics.events.snapshots
ORDER BY committed_at
""").show(truncate=False)
# history: отражение snapshot-log - когда конкретный снапшот СТАЛ текущим
# (важно отличать от committed_at в snapshots - после rollback они расходятся)
spark.sql("""
SELECT made_current_at, snapshot_id, parent_id, is_current_ancestor
FROM lakehouse.analytics.events.history
ORDER BY made_current_at
""").show(truncate=False)
# manifests: прямое отражение записей manifest list текущего снапшота
spark.sql("""
SELECT
path,
length,
partition_spec_id,
added_data_files_count,
existing_data_files_count,
deleted_data_files_count
FROM lakehouse.analytics.events.manifests
""").show(truncate=False)
# entries: самый детальный уровень - по одной строке на КАЖДУЮ запись
# манифеста, включая статус ADDED/EXISTING/DELETED, разобранный выше
spark.sql("""
SELECT
status,
snapshot_id,
data_file.file_path,
data_file.record_count,
data_file.file_size_in_bytes
FROM lakehouse.analytics.events.entries
""").show(truncate=False)
Системная таблица entries - это, по сути, прямое SQL-представление содержимого manifest file, разобранного в теоретической части урока, включая то самое поле status. Если вы хотите воочию увидеть, как именно файл со статусом ADDED в одном снапшоте превращается в EXISTING в следующем (потому что унаследован без изменений), а затем в DELETED в третьем (потому что был перезаписан компакцией) - именно эта системная таблица показывает это напрямую, без необходимости лезть в сами Avro-файлы.
Кейс 2: физический скан MinIO - читаем metadata.json и manifest вручную¶
Чтобы окончательно развеять ощущение «магии», полезно один раз пройти весь путь руками: скачать metadata.json напрямую из object storage и прочитать содержимое manifest list и manifest file через стандартную Python-библиотеку для Avro, минуя Spark полностью.
import json
import boto3
import fastavro
# Подключаемся к MinIO через boto3 (S3-совместимый клиент)
s3 = boto3.client(
"s3",
endpoint_url="http://minio:9000",
aws_access_key_id="minioadmin",
aws_secret_access_key="minioadmin",
)
BUCKET = "lakehouse"
TABLE_METADATA_PREFIX = "warehouse/analytics/events/metadata/"
# Шаг 1: находим самый свежий metadata.json по времени последней модификации
# (в реальном сценарии путь к нему лучше брать из каталога, а не угадывать -
# здесь это сделано для учебной демонстрации "руками")
objects = s3.list_objects_v2(Bucket=BUCKET, Prefix=TABLE_METADATA_PREFIX)
metadata_files = [
obj for obj in objects["Contents"] if obj["Key"].endswith(".metadata.json")
]
latest_metadata_key = max(metadata_files, key=lambda o: o["LastModified"])["Key"]
print(f"Текущий metadata.json: {latest_metadata_key}")
# Шаг 2: скачиваем и парсим metadata.json - это обычный JSON,
# никакой бинарной магии на этом уровне нет
metadata_obj = s3.get_object(Bucket=BUCKET, Key=latest_metadata_key)
table_metadata = json.loads(metadata_obj["Body"].read())
print(f"current-snapshot-id: {table_metadata['current-snapshot-id']}")
print(f"Число схем в истории: {len(table_metadata['schemas'])}")
print(f"Число снапшотов: {len(table_metadata['snapshots'])}")
current_snapshot = next(
s for s in table_metadata["snapshots"]
if s["snapshot-id"] == table_metadata["current-snapshot-id"]
)
manifest_list_path = current_snapshot["manifest-list"]
print(f"Manifest List текущего снапшота: {manifest_list_path}")
# Шаг 3: скачиваем manifest list и читаем его как Avro через fastavro -
# это бинарный файл, поэтому нужна Avro-библиотека, а не json.loads
manifest_list_key = manifest_list_path.replace(f"s3a://{BUCKET}/", "")
manifest_list_obj = s3.get_object(Bucket=BUCKET, Key=manifest_list_key)
manifest_entries = list(fastavro.reader(manifest_list_obj["Body"]))
for entry in manifest_entries:
print(
f" Manifest: {entry['manifest_path']}, "
f"added={entry['added_data_files_count']}, "
f"existing={entry['existing_data_files_count']}, "
f"deleted={entry['deleted_data_files_count']}"
)
# Шаг 4: спускаемся ещё на уровень ниже - читаем первый manifest file
# и видим ровно те же поля status / data_file, которые SQL показывал
# через системную таблицу entries в Кейсе 1
first_manifest_path = manifest_entries[0]["manifest_path"]
manifest_key = first_manifest_path.replace(f"s3a://{BUCKET}/", "")
manifest_obj = s3.get_object(Bucket=BUCKET, Key=manifest_key)
for record in fastavro.reader(manifest_obj["Body"]):
status_name = {0: "EXISTING", 1: "ADDED", 2: "DELETED"}[record["status"]]
data_file = record["data_file"]
print(
f" [{status_name}] {data_file['file_path']} "
f"({data_file['record_count']} строк, "
f"{data_file['file_size_in_bytes']} байт)"
)
Результат выполнения этого скрипта построчно совпадает с тем, что показывали системные таблицы Spark SQL в Кейсе 1 - и это не случайность, а прямое доказательство главного тезиса урока: системные таблицы snapshots/manifests/entries не делают ничего сверхъестественного, они просто читают ровно те же самые JSON и Avro файлы, которые можно прочитать вручную через boto3 и fastavro, и представляют их в виде SQL-строк для удобства. Понимание этого эквивалентности - ключевой навык для самостоятельной диагностики проблем, когда привычные SQL-инструменты недоступны или дают противоречивую картину.
Кейс 3: симуляция конкурентного конфликта и CommitFailedException¶
Последний кейс - воспроизведение сценария из раздела про разрешение конфликтов: два потока выполняют конкурирующие операции overwritePartitions() на одну и ту же партицию, и один из них должен либо успешно повторить попытку (если изменения не пересекаются по файлам), либо получить ошибку (если пересекаются).
import threading
import time
from pyspark.sql import functions as F
results = {}
def run_overwrite(thread_name: str, event_type_filter: str, delay_seconds: float):
"""
Каждый поток читает текущие данные партиции 2024-06-15,
модифицирует часть строк и пытается атомарно заменить партицию.
delay_seconds сдвигает старт, чтобы оба потока с большей вероятностью
попытались закоммититься в перекрывающееся окно времени.
"""
time.sleep(delay_seconds)
try:
current_partition_df = (
spark.table("lakehouse.analytics.events")
.where("event_ts >= '2024-06-15' AND event_ts < '2024-06-16'")
)
# Оба потока модифицируют ОДНУ И ТУ ЖЕ партицию -
# это создаёт реальный, а не симулированный конфликт валидации
modified_df = current_partition_df.withColumn(
"event_type",
F.when(
F.col("event_type") == event_type_filter,
F.lit(f"updated_by_{thread_name}"),
).otherwise(F.col("event_type")),
)
modified_df.writeTo("lakehouse.analytics.events").overwritePartitions()
results[thread_name] = "SUCCESS"
except Exception as exc:
# Py4JJavaError оборачивает CommitFailedException из JVM
results[thread_name] = f"FAILED: {type(exc).__name__}: {exc}"
thread_a = threading.Thread(target=run_overwrite, args=("A", "click", 0.0))
thread_b = threading.Thread(target=run_overwrite, args=("B", "purchase", 0.05))
thread_a.start()
thread_b.start()
thread_a.join()
thread_b.join()
for name, outcome in results.items():
print(f"Поток {name}: {outcome}")
# Типичный результат:
# Поток A: SUCCESS
# Поток B: FAILED: Py4JJavaError: ...CommitFailedException:
# Found new conflicting files that can contain records
# matching partial overwrite filter...
Если оба потока успешно завершились без ошибок - значит, retry-логика, разобранная в теоретической части, успела автоматически перестроить и повторно закоммитить второй поток до того, как валидация обнаружила реальное пересечение файлов (или Iceberg определил, что модифицированные строки на самом деле не пересекаются по физическим файлам). Если один из потоков завершился с CommitFailedException - это именно тот сценарий из диаграммы «семантического конфликта», разобранной выше: валидация overwritePartitions() обнаружила, что файлы партиции, прочитанные потоком B в начале операции, больше не являются текущими после успешного commit'а потока A.
Проверить итоговое состояние таблицы и убедиться, что не произошло потери данных (тезис «либо честный commit, либо явная ошибка, но никогда тихая потеря изменений») можно тем же способом, что и в прошлом уроке:
spark.sql("""
SELECT operation, count(*) AS snapshots_count
FROM lakehouse.analytics.events.snapshots
GROUP BY operation
""").show()
# Если поток B упал с CommitFailedException, в истории снапшотов
# будет ТОЛЬКО снапшот от потока A - ни одного "наполовину применённого"
# или дублирующего снапшота от неудачной попытки B не появится,
# потому что Stage 3 (Commit) для потока B так и не завершился успехом.
Production кейс: тысячи манифестов и деградация planning time¶
Ситуация. Команда платформы данных подключила Apache Flink к Iceberg-таблице analytics.clickstream для записи событий в режиме streaming с чекпойнтом каждые 60 секунд. Каждый чекпойнт коммитил небольшой append - в среднем по 2-3 МБ новых данных. Таблица работала без видимых проблем три недели, после чего аналитики начали жаловаться, что простой запрос SELECT count(*) FROM analytics.clickstream WHERE event_date = current_date() выполняется не секунды, а десятки секунд - то есть именно то самое planning time, которое в первом уроке модуля Iceberg обещал свести почти к нулю.
Расследование. Первая гипотеза - проблема в самих данных или индексах - быстро не подтвердилась: суммарный объём таблицы был умеренным (около 200 ГБ), а партиция за текущий день - всего несколько гигабайт. Инженер применил ровно тот инструментарий, который разобран в практическом блоке этого урока:
# Шаг 1: сколько манифестов входит в текущий снапшот?
manifest_count = spark.sql(
"SELECT count(*) AS cnt FROM lakehouse.analytics.clickstream.manifests"
).collect()[0]["cnt"]
print(f"Манифестов в текущем снапшоте: {manifest_count}")
# Манифестов в текущем снапшоте: 30412
# Шаг 2: какой средний размер одного манифеста?
avg_manifest_size = spark.sql("""
SELECT avg(length) AS avg_len, sum(length) AS total_len
FROM lakehouse.analytics.clickstream.manifests
""").collect()[0]
print(f"Средний размер манифеста: {avg_manifest_size['avg_len'] / 1024:.1f} КБ")
print(f"Суммарный размер manifest list: {avg_manifest_size['total_len'] / 1024**2:.1f} МБ")
# Средний размер манифеста: 4.2 КБ
# Суммарный размер manifest list: 124.7 МБ
Причина оказалась ровно той, что разобрана в разделе про слияние манифестов: на кластере Flink свойство commit.manifest-merge.enabled было явно выставлено в false много недель назад - по легенде, временно, для отладки другой проблемы, и об этом забыли. Каждый из 30 000+ чекпойнтов создавал собственный крошечный манифест на 2-3 файла, и ни один из них никогда не объединялся с соседними. Planning time запроса включал чтение manifest list (124 МБ - уже немало для одного файла, который должен открываться при каждом запросе) и затем потенциальное обращение к десяткам тысяч отдельных мелких манифестов для построения списка релевантных файлов.
Решение.
# Шаг 1: возвращаем автоматическое слияние манифестов для будущих commit'ов
spark.sql("""
ALTER TABLE lakehouse.analytics.clickstream SET TBLPROPERTIES (
'commit.manifest-merge.enabled' = 'true',
'commit.manifest.target-size-bytes' = '8388608'
)
""")
# Шаг 2: разово принудительно переписываем существующие манифесты
# (rewrite_manifests - отдельная maintenance-процедура, не путать
# с rewrite_data_files - она трогает ТОЛЬКО уровень манифестов,
# data-файлы остаются нетронутыми)
spark.sql("""
CALL lakehouse.system.rewrite_manifests(
table => 'analytics.clickstream'
)
""")
Результат. После rewrite_manifests число манифестов текущего снапшота упало с 30 412 до 47 (каждый - близкий к целевому размеру 8 МБ, агрегирующий записи из сотен исходных мелких манифестов через статус EXISTING). Planning time запроса count() с фильтром по дате вернулся к исходным значениям в пределах одной-двух секунд. Важная деталь, которую команда специально проверила перед тем, как считать инцидент закрытым: rewrite_manifests не создаёт новые data-файлы и не запускает compaction в привычном смысле - она работает исключительно на уровне манифестов, поэтому занимает секунды-минуты, а не часы, в отличие от rewrite_data_files.
| Метрика | До инцидента (manifest-merge выключен) | После rewrite_manifests |
|---|---|---|
| Число манифестов в снапшоте | 30 412 | 47 |
| Размер manifest list | 124.7 МБ | ~190 КБ |
Planning time count() с фильтром по дате |
35-50 секунд | 1-2 секунды |
Время выполнения rewrite_manifests |
- | 4 минуты |
| Затронуты ли data-файлы | - | Нет, нулевое изменение объёма Parquet |
Ключевой вывод. Слияние манифестов - это не «фоновая магия», за которой не нужно следить, а конфигурируемый механизм, который можно выключить (намеренно или случайно) и не заметить деградацию сразу, потому что она проявляется постепенно, манифест за манифестом. Знание точного назначения поля commit.manifest-merge.enabled и умение быстро продиагностировать его состояние через <table>.manifests - это именно та компетенция, которую даёт понимание snapshot-модели на уровне внутреннего устройства, а не только на уровне SQL-интерфейса.
O(1) вместо O(N): главный архитектурный инсайт¶
Стоит явно сформулировать вывод, который незаметно пронизывает весь материал этого урока. В Hive-style подходе стоимость операции «понять состав таблицы» растёт линейно с числом партиций и файлов: O(N), где N - суммарное число объектов в object storage, потому что единственный способ узнать состав таблицы - это перелистинговать директории.
В Iceberg стоимость той же операции - чтение metadata.json (один объект) плюс чтение manifest list (один объект, размер которого растёт с числом манифестов, а не файлов, и манифесты переиспользуются между снапшотами благодаря статусу EXISTING) - то есть фактическая стоимость определения состава таблицы приближается к константной относительно общего числа data-файлов и точно не требует LIST-операций по object storage вообще. Это не означает буквально O(1) в строгом алгоритмическом смысле (manifest list всё же растёт, хоть и значительно медленнее, чем число файлов) - но по сравнению с O(N) листинга это качественный скачок на порядки, который и был продемонстрирован числовым расчётом в прошлом уроке (43 минуты листинга против десятков GET-запросов к компактным Avro-файлам).
Read Path: путь от SQL-запроса до конкретного data file¶
Весь материал урока был организован вокруг записи (write path) и структуры дерева. Завершить картину стоит обратным путём - explicit пошаговым разбором того, что происходит между моментом, когда пользователь нажимает Enter после SELECT, и моментом, когда executor открывает первый байт конкретного Parquet-файла. Этот путь - не новая механика, а просто последовательное применение всего, что было разобрано в уроке, в правильном порядке.
Несколько деталей этой последовательности стоит подчеркнуть отдельно, потому что именно они объясняют, почему этот путь принципиально быстрее эквивалентного пути для Hive-style таблицы:
-
Driver делает ровно один запрос к каталогу и один GET к
metadata.json, независимо от того, сколько партиций или файлов содержит таблица - стоимость этого шага постоянна. -
Manifest pruning (отбор релевантных манифестов по диапазонам в manifest list) происходит до открытия отдельных манифестов - это и есть прямое применение поля
partitionsиз manifest list, разобранного в разделе про Уровень 2. -
File pruning (отбор релевантных data-файлов по column min/max внутри отобранных манифестов) происходит до открытия Parquet-файлов - применение полей
lower_bounds/upper_boundsиз manifest entry, разобранных в разделе про Уровень 3. -
Только в v2-таблицах появляется дополнительный шаг сопоставления data-файлов с потенциальными delete-файлами через sequence number - в v1-таблицах этого шага просто не существует, потому что любое изменение строк там реализовано через полную перезапись файла (Copy-on-Write), и к моменту чтения delete-файлов физически нет.
-
Единственный шаг, который читает данные, а не метаданные - самый последний. Все предыдущие шаги читали только Avro/JSON объекты, общий объём которых на много порядков меньше, чем сами данные таблицы.
Это и есть полный ответ на вопрос, поставленный в начале урока про «чёрный ящик»: между SQL-запросом пользователя и конкретным открытым файлом нет ни одного шага, который нельзя было бы объяснить через дерево метаданных, разобранное в этом уроке. «Магии» не остаётся - остаётся последовательность из нескольких десятков GET-запросов к компактным объектам метаданных, каждый из которых отбрасывает часть пространства поиска для следующего шага.
Типичные заблуждения про snapshot model¶
«Снапшот - это полная копия всех данных таблицы» - неверно. Снапшот - это лёгкая структура метаданных (manifest list + манифесты), которая ссылается на данные, в основном переиспользуя файлы предыдущих снапшотов через статус EXISTING. Создание нового снапшота не означает копирование терабайт данных - оно означает запись нескольких новых, обычно небольших, Avro-объектов.
«Если commit упал с CommitFailedException, можно просто повторить запрос без изменений» - не всегда верно. Как было показано в разделе про семантические конфликты, если конфликт связан с валидацией (например, файлы, которые операция собиралась изменить, уже изменены), повтор той же самой операции с теми же входными данными снова упадёт с той же ошибкой - нужно перечитать актуальное состояние и пересчитать изменения на его основе, а не просто слепо повторить вызов.
«Старые снапшоты не занимают места, пока их явно не используют» - неверно, и это будет подробно разобрано в уроке про maintenance. Старые снапшоты держат «живыми» ссылки на свои data-файлы, и эти файлы продолжают физически занимать место в object storage до тех пор, пока снапшот не будет явно «истёкшим» через expire_snapshots.
«Sequence number и snapshot-id - это одно и то же» - неверно. snapshot-id - это уникальный (обычно случайно генерируемый) идентификатор конкретного снапшота, не несущий информации о порядке. sequence-number - монотонно возрастающий счётчик, специально предназначенный для определения порядка применения delete-файлов в v2. У снапшота, созданного раньше по времени, гарантированно меньший sequence number, но это не верно для snapshot-id, который не обязан быть монотонным.
«Manifest merging копирует физические data-файлы, поэтому работает медленно» - неверно, и это явная путаница с компакцией data-файлов. rewrite_manifests и автоматическое слияние манифестов работают исключительно с Avro-объектами метаданных - десятками-сотнями килобайт, а не с гигабайтами Parquet. Именно поэтому в production-кейсе этого урока операция заняла 4 минуты для таблицы 200 ГБ: объём данных, который физически переписывался, равнялся суммарному размеру манифестов (десятки-сотни МБ), а не размеру самой таблицы.
Мостик к следующим урокам¶
Этот урок дал детальное устройство снапшот-модели, но намеренно оставил несколько вопросов для следующих уроков модуля:
-
Урок 3 («Manifest Files и Manifest List») берёт поле
partitionsиз manifest list, разобранное здесь лишь на уровне назначения, и показывает, как именно Spark Catalyst использует эти диапазоны для отсечения манифестов при планировании сканирования - то, что в первом уроке модуля было показано только результатом (DataFiles: 1 of 5), а не механизмом. -
Урок 4 («Hidden Partitioning») возвращается к полю
partitionвнутриdata_fileи объясняет, как именно вычисляются значения transform-функций (days,bucket,truncate) и как они влияют на распределение статистики в манифестах. -
Уроки 6 и 7 («Copy-on-Write» и «Merge-on-Read») детально раскрывают delete-файлы, чей механизм применения через sequence number был показан здесь только в общих чертах.
-
Урок 9 («Time Travel») возвращается к
snapshot-logиparent-snapshot-id, разобранным в разделе про snapshot lineage, и показывает полный синтаксис инкрементального чтения между двумя снапшотами. -
Урок 10 («Table Maintenance») прямо отвечает на вопрос, поднятый в разделе про заблуждения: как именно
expire_snapshotsрешает, какие data-файлы больше не нужны ни одному живому снапшоту, и почему это не такая простая операция, как «удалить всё старше N дней».
Домашнее задание¶
-
Выполните шаги Кейса 2 самостоятельно: вручную (без Spark SQL) прочитайте
metadata.json, manifest list и хотя бы один manifest file вашей тестовой таблицы черезboto3иfastavro. Сравните результат с выводом системных таблицsnapshots/manifests/entriesи письменно зафиксируйте, какие поля совпадают буквально. -
Сделайте по таблице
orders(из домашнего задания прошлого урока) пять последовательных операций (INSERT,INSERT,DELETEчасти строк,INSERT, компакция черезrewrite_data_files), затем напишите запрос кorders.entries, который покажет полную историю статусов (ADDED→EXISTING→DELETED) для одного конкретного физического файла - выберите файл, который пережил хотя бы одну компакцию. -
Повторите Кейс 3 (симуляция конфликта), но измените сценарий так, чтобы оба потока гарантированно не конфликтовали по содержанию (например, два
appendв разные партиции вместо двухoverwritePartitions()в одну и ту же). Убедитесь, что оба потока завершаются успешно, и оба новых снапшота присутствуют в истории с корректной цепочкойparent-snapshot-id. -
Напишите Spark SQL запрос (используя
entriesиmanifests), который вычисляет: а) суммарный физический размер файлов, входящих в текущий снапшот; б) суммарный размер файлов со статусомDELETED, которые физически ещё существуют в object storage, но уже не видны ни одному текущему манифесту. Сравните результат (б) с тем, что покажетexpire_snapshotsв режиме dry-run (если ваша версия Iceberg его поддерживает) - вы только что вручную посчитали то, что эта команда автоматизирует в уроке 10. -
Найдите в документации Apache Iceberg формальное Avro-описание схемы
manifest_entry(илиmanifest_fileдля manifest list) и сравните его с упрощёнными таблицами полей из этого урока - письменно (1-2 абзаца) укажите минимум два поля, которые не были упомянуты в уроке, и предположите, для какой задачи каждое из них может использоваться. -
Создайте тег (
CREATE TAG) на одном из ранних снапшотов вашей таблицыordersс retention 1 день, затем создайте отдельную ветку (CREATE BRANCH) от текущего состояния и сделайте в неё пробныйINSERT, не затрагивающий основную веткуmain. Убедитесь черезSELECT * FROM orders.refs, что обе ссылки (branch и tag) присутствуют в метаданных, и письменно объясните, почему запись в веткуexperimentне повлияла на результатSELECT * FROM orders(без указания ветки).
Итоги¶
Snapshot - это неизменяемый снимок, а не «текущие данные». Любое чтение таблицы Iceberg без явного time travel - это чтение конкретного снапшота, на который указывает current-snapshot-id. Понимание этого снимает большую часть «магии» вокруг ACID-гарантий: ничего не блокируется и не редактируется на месте, просто появляется новый снапшот, и указатель атомарно переключается на него.
Git-аналогия - не метафора для красоты, а структурное соответствие. Снапшот ≈ коммит, manifest list/manifest ≈ дерево объектов, data file ≈ blob, current-snapshot-id ≈ HEAD. Time travel, история изменений и будущая поддержка branching/tagging - прямые следствия этой архитектуры, а не отдельные «фичи», добавленные сверху.
Каждый уровень дерева метаданных решает конкретную задачу масштабирования. metadata.json остаётся компактным, потому что список манифестов вынесен в отдельный Avro-файл на каждый снапшот. Manifest list остаётся компактным и быстро читаемым, потому что детальная статистика по файлам вынесена в манифесты, которые при этом переиспользуются между снапшотами через статус EXISTING.
Атомарный commit - это трёхстадийный Write-and-Commit Path, где только последняя стадия требует синхронизации. Write и Metadata стадии выполняются независимо и параллельно для разных writer'ов; единственная точка конкуренции - финальный CAS в каталоге, и именно поэтому модель масштабируется без распределённых блокировок.
Не все конфликты коммитов одинаковы. Конфликты на уровне указателя без пересечения по содержанию разрешаются автоматическим retry прозрачно для приложения. Конфликты на уровне семантической валидации (overwrite/delete/merge с пересекающимися файлами) приводят к CommitFailedException, который требует осознанной обработки на уровне приложения, а не бесконечного retry.
Sequence numbers в v2 - это механизм корректного применения delete-файлов без блокировок. Монотонно возрастающий номер снапшота позволяет однозначно определить, к каким data-файлам применяется конкретный delete-файл, что и делает возможным Merge-on-Read - тему следующих уроков модуля.
Read Path - это последовательное сужение пространства поиска через метаданные, а не «магия». От catalog к metadata.json, от него к manifest list с отсечением по диапазонам партиций, от выбранных манифестов к отсечению по min/max статистике колонок - и только в самом конце происходит первое обращение к фактическим байтам Parquet. Каждый шаг этой цепочки объясним через конкретное поле конкретного файла, разобранное в этом уроке.
Branching и tagging - не футуристическая возможность, а доступный сегодня механизм поверх того же графа снапшотов. CREATE TAG замораживает снапшот для аудита с заданным сроком жизни; CREATE BRANCH позволяет вести параллельную линию изменений без риска для main; оба механизма используют ровно ту же структуру parent-snapshot-id и refs, которая лежит в основе time travel и rollback.
Краткий глоссарий терминов урока¶
-
Write-and-Commit Path - трёхстадийный процесс записи: Write (данные) → Metadata (манифесты) → Commit (атомарный CAS в каталоге).
-
Compare-and-Swap (CAS) - атомарная операция «обновить указатель, только если он до сих пор равен ожидаемому значению»; единственная точка синхронизации между конкурирующими writer'ами.
-
CommitFailedException - исключение, сигнализирующее, что commit не может быть автоматически разрешён через retry, потому что конфликт обнаружен на уровне семантической валидации операции, а не только на уровне указателя.
-
Manifest Entry Status (ADDED/EXISTING/DELETED) - три состояния записи о файле внутри манифеста, отражающие, был ли файл добавлен в текущем снапшоте, унаследован без изменений, или логически удалён.
-
Sequence Number - монотонно возрастающий номер, присваиваемый каждому снапшоту; в v2 определяет, к каким data-файлам применяется конкретный delete-файл.
-
Parent Snapshot ID - ссылка на предыдущий снапшот в цепочке истории, формирующая lineage, аналогичный истории коммитов Git.
-
Snapshot Log / Metadata Log - журналы, отслеживающие, когда конкретный снапшот стал текущим, и историю самих файлов
metadata.jsonсоответственно. -
Manifest Merging - автоматическое объединение нескольких мелких манифестов в один при превышении порога
commit.manifest.target-size-bytes, предотвращающее неограниченный рост числа манифестов. -
rewrite_manifests - maintenance-процедура, принудительно объединяющая манифесты без изменения data-файлов; в отличие от
rewrite_data_files, работает только на уровне метаданных. -
Refs (branch / tag) - именованные указатели на снапшоты в
metadata.json, существующие параллельно с основнымcurrent-snapshot-id;branch- изменяемая линия истории,tag- замороженная ссылка на один снапшот с опциональным сроком жизни. -
Read Path - последовательность шагов планирования запроса: catalog →
metadata.json→ manifest list (отсечение манифестов) → manifest files (отсечение файлов по статистике) → Parquet файлы, с применением delete-файлов через sequence number в v2. -
Manifest Pruning / File Pruning - отсечение целых манифестов (по диапазонам партиций в manifest list) и отдельных data-файлов (по min/max статистике колонок в manifest entry) на этапе планирования, до открытия самих Parquet-файлов.
-
fastavro - Python-библиотека для чтения и записи Avro-файлов без JVM; используется для прямой инспекции manifest list и manifest files в обход Spark SQL.
-
Snapshot Lineage - граф родственных связей снапшотов через поле
parent-snapshot-id, лежащий в основе инкрементального чтения, time travel и операций branching/tagging. -
JDBC Catalog - реализация Iceberg-каталога поверх произвольной реляционной БД (в этом курсе - self-hosted PostgreSQL), где compare-and-swap commit реализован через условный
UPDATEс проверкой текущего значенияmetadata_location. -
Optimistic Concurrency Control (OCC) - модель конкурентного доступа, при которой writer'ы не блокируют друг друга заранее, а проверяют отсутствие конфликта только в момент финального commit'а, повторяя попытку при необходимости.