Delta Lake: transaction log и сравнение с Iceberg snapshot model

Архитектура Delta Lake: устройство _delta_log, JSON-коммиты и checkpoint-файлы, Optimistic Concurrency Control и таксономия конфликтов, Time Travel через линейный лог, Schema Enforcement и Column Mapping, стоимость Query Planning - детальное сравнение с деревом манифестов Apache Iceberg на каждом шаге.

lakehouse storage

Девять уроков этого модуля были посвящены Apache Iceberg - архитектуре, которая решает задачу транзакционного Lakehouse через дерево независимых, иммутабельных файлов метаданных: снапшот, список манифестов, манифест, файл данных. К этому моменту понятия snapshot, manifest, partition spec и position/equality delete должны быть для читателя не абстракцией, а конкретной, разобранной по байтам механикой. Это даёт редкую возможность: изучать второй table format не с нуля, а через постоянное, явное сопоставление с уже усвоенной моделью.

Delta Lake решает ровно ту же задачу - транзакционность, версионирование, изоляция снимков поверх объектного хранилища - но выбирает физически другую структуру данных для её решения. Там, где Iceberg выбрал дерево, Delta Lake выбрал журнал. Этот урок разбирает анатомию _delta_log, механизм checkpoint-файлов, модель конкурентного доступа, Time Travel, эволюцию схемы и стоимость планирования запросов - и на каждом шаге проводит точную, конкретную параллель с уже известной архитектурой Iceberg, чтобы зафиксировать не просто «как устроен Delta Lake», а «почему два формата, решающие одну задачу, выглядят настолько разными внутри».


Философия Delta Lake: ACID-таблица как журнал транзакций поверх Parquet

Реляционная база данных, упакованная в объектное хранилище

Iceberg, как было ясно сформулировано в первом уроке модуля, спроектирован как декларативная спецификация: формат описывает структуру метаданных, и любой движок, реализующий эту спецификацию, может читать и писать таблицу независимо от других. Delta Lake исторически подходит к задаче с другим исходным предположением, которое прямо отражает происхождение проекта внутри Databricks - компании, чей основной продукт - управляемая платформа выполнения Spark-нагрузок. Архитектурная метафора Delta Lake ближе к классической реляционной СУБД, чем к открытой файловой спецификации: единый, строго последовательный журнал транзакций, к которому Spark-движок обращается как к источнику истины, очень похож на то, как PostgreSQL или MySQL ведут собственный WAL (Write-Ahead Log) для восстановления и обеспечения консистентности.

Это не оценочное утверждение «хуже» или «лучше» - это разница в исходной точке проектирования, которая будет проявляться в каждом из шести оставшихся разделов урока. Iceberg в первую очередь спроектирован для множества независимых читателей и писателей, потенциально на разных движках, которым нужно согласовываться через внешний Catalog. Delta Lake в первую очередь спроектирован для надёжной, последовательной истории изменений одной таблицы, с историческим фокусом на единственный движок - Spark - и лишь позже расширившийся на поддержку других потребителей через дополнительные библиотеки (Delta Standalone, Delta Kernel, delta-rs).

Анатомия _delta_log: скрытая директория рядом с данными

Физически Delta-таблица - это обычная директория объектного хранилища (или HDFS, или локальной файловой системы), содержащая Parquet-файлы данных и одну специальную, скрытую поддиректорию _delta_log/, в которой находится вся транзакционная история таблицы. В отличие от Iceberg, где метаданные могут логически существовать отдельно от данных (Catalog может указывать на metadata.json в произвольном месте), Delta Lake по умолчанию жёстко привязывает лог к физическому расположению самой таблицы - _delta_log всегда лежит непосредственно внутри корневой директории таблицы.

s3a://lakehouse/warehouse/orders_delta/
├── _delta_log/
│   ├── 00000000000000000000.json
│   ├── 00000000000000000001.json
│   ├── 00000000000000000002.json
│   ├── ...
│   ├── 00000000000000000010.checkpoint.parquet
│   └── _last_checkpoint
├── part-00000-9a3f...-c000.snappy.parquet
├── part-00001-2b7e...-c000.snappy.parquet
└── part-00002-7cd1...-c000.snappy.parquet

Эта структура директорий - первое наглядное архитектурное отличие от Iceberg. Во втором уроке модуля было показано, что метаданные Iceberg формируют отдельное поддерево (metadata.json → Manifest List → Manifest File), на которое Catalog хранит лишь указатель, и физическое расположение этого поддерева не обязано совпадать с расположением данных. У Delta Lake такого разделения нет: лог - неотъемлемая часть самой директории таблицы, и movement директории целиком (например, копирование в другой бакет) автоматически переносит и историю, и данные одним и тем же действием с файловой системой - что одновременно и удобство (таблица самодостаточна), и архитектурное ограничение (нельзя независимо переместить или переиспользовать метаданные без переноса самих данных).

JSON-коммит как атомарный набор инструкций

Каждая успешная транзакция Delta Lake - результат операции INSERT, UPDATE, DELETE, MERGE INTO, изменения схемы или свойств таблицы - материализуется как ровно один новый JSON-файл в _delta_log/, с именем, состоящим из монотонно растущего номера версии, дополненного нулями до 20 знаков: 00000000000000000000.json для самой первой транзакции (как правило, CREATE TABLE), 00000000000000000001.json для следующей, и так далее без пропусков. Номер версии - это и идентификатор коммита, и единица, через которую Delta реализует Time Travel (раздел 5 этого урока) и обнаружение конфликтов конкурентного доступа (раздел 4).

Каждый такой JSON-файл - это последовательность строк, каждая из которых является отдельным JSON-объектом одного из нескольких типов action (действие). Самые важные из них:

  • add - файл данных добавлен в логическое состояние таблицы этой транзакцией; содержит путь к Parquet-файлу, размер в байтах, статистику (stats - min/max значения колонок, count, null count, аналог summary-статистики Manifest File у Iceberg), и значения партиционных колонок, если таблица партиционирована.

  • remove - файл данных логически удалён из таблицы этой транзакцией; физически файл может оставаться на диске - его реальное удаление с object storage происходит позже, отдельной командой VACUUM (детально разбираемой в двенадцатом уроке модуля), полностью аналогично тому, как Iceberg откладывает физическое удаление файлов до expire_snapshots, разобранного в десятом уроке.

  • metaData - описывает схему таблицы, партиционные колонки и свойства таблицы (configuration); записывается при создании таблицы и при любом изменении схемы.

  • protocol - минимальная версия протокола чтения (minReaderVersion) и записи (minWriterVersion), которую должен поддерживать движок, чтобы безопасно работать с этой таблицей; механизм защиты от использования старым клиентом новых, несовместимых возможностей (например, Deletion Vectors требуют определённой минимальной версии протокола).

  • commitInfo - метаданные самой операции: timestamp, тип операции (operation: WRITE, MERGE, UPDATE, DELETE, OPTIMIZE...), параметры операции, идентификатор пользователя/job - именно эта запись лежит в основе команды DESCRIBE HISTORY, разбираемой в разделе про Time Travel.

{"commitInfo":{"timestamp":1718012345000,"operation":"WRITE","operationParameters":{"mode":"Append","partitionBy":"[\"order_date\"]"},"isBlindAppend":true}}
{"metaData":{"id":"a3f7...","format":{"provider":"parquet"},"schemaString":"{...}","partitionColumns":["order_date"],"configuration":{},"createdTime":1718012344000}}
{"add":{"path":"part-00000-9a3f...-c000.snappy.parquet","partitionValues":{"order_date":"2025-06-15"},"size":134821,"modificationTime":1718012345000,"dataChange":true,"stats":"{\"numRecords\":4821,\"minValues\":{\"amount\":12.4},\"maxValues\":{\"amount\":998.1}}"}}

Текущее логическое состояние таблицы на момент версии N - это результат последовательного применения всех add/remove действий из всех JSON-файлов от версии 0 до версии N включительно: каждый add добавляет файл в множество «живых» файлов таблицы, каждый remove исключает файл из этого множества. Этот процесс называется state reconstruction (восстановление состояния) или log replay (воспроизведение лога), и именно его стоимость - центральная техническая проблема, которую решает следующий раздел урока.

Диаграмма показывает механику log replay на минимальном примере: чтобы узнать, какие файлы реально составляют таблицу на версии 2, движок должен последовательно применить все три транзакции - добавить f1, f2, f3, затем добавить f4, затем убрать f2 и добавить f5 - и только после этого получить итоговое множество живых файлов {f1, f3, f4, f5}. На трёх версиях эта операция тривиальна, но при росте истории до тысяч версий её стоимость становится центральной инженерной проблемой, разбираемой в следующем разделе.

Атомарность коммита: специфика разных хранилищ

Записать новый JSON-файл с именем 00000000000000000005.json недостаточно - критически важно, чтобы при попытке двух разных писателей закоммитить версию 5 одновременно выигрывал только один из них, а второй получал явный сигнал о конфликте, а не тихую перезапись. Этот примитив называется atomic put-if-absent: «создать файл с этим именем, только если файла с таким именем ещё не существует» - и его доступность и реализация принципиально различаются между классами хранилищ.

На классическом HDFS эта примитива была изначально нативной: операция rename() атомарна на уровне самой файловой системы, и Delta Lake долгое время полагался именно на неё для гарантии последовательности номеров версий без какой-либо дополнительной координации. Объектные хранилища (S3 и совместимые с ним MinIO/Ceph RGW), однако, исторически не предоставляли подобной гарантии нативно: операция PutObject в классическом S3 API не имела условия «только если объект не существует» - что создавало реальный риск двух конкурентных писателей, каждый из которых считает, что успешно закоммитил версию 5, в то время как один из двух коммитов на самом деле тихо перезаписал другой.

Для решения этой проблемы при работе с S3-совместимым хранилищем Delta Lake исторически требовал отдельного внешнего механизма координации - LogStore с поддержкой multi-cluster записи (например, реализация на основе таблицы DynamoDB, выступающей распределённой блокировкой «кто первый запросил версию N») либо более современного commit-координатора, обслуживающего роль единой точки сериализации коммитов для конкретной таблицы. Без такого механизма Delta Lake на «голом» S3 гарантированно безопасен только для одного писателя за раз - именно эта деталь будет напрямую противопоставлена модели Iceberg в разделе про конкурентный доступ: Iceberg делегирует ту же задачу атомарного переключения указателя внешнему Catalog (Hive Metastore, JDBC, REST), который и так предоставляет compare-and-swap семантику как часть своего контракта, не требуя отдельного, специально написанного under multi-writer координатора поверх самого объектного хранилища.


Проблема роста логов и механизм контрольных точек (Checkpoints)

Почему линейный replay перестаёт масштабироваться

Механизм log replay, описанный в предыдущем разделе, работает безотказно на таблице с десятью или сотней версий - современный процессор перемалывает несколько сотен мелких JSON-объектов за миллисекунды. Проблема проявляется на горизонте, типичном для production-таблицы, живущей месяцы или годы под постоянной нагрузкой потоковых или батчевых джобов: таблица, в которую раз в пять минут пишет Structured Streaming job, накапливает свыше 100 000 версий в течение одного года. Чтобы прочитать такую таблицу «как она есть сейчас», движку формально нужно скачать и разобрать все 100 000 JSON-файлов, потому что состояние на версии 100 000 определено как результат применения action’ов из версий с 0 по 100 000 включительно.

На практике это означает катастрофическую деградацию Query Planning - той фазы выполнения запроса, которая предшествует собственно чтению данных и определяет, какие именно Parquet-файлы нужно прочитать. Для таблицы с десятками тысяч версий время одного только восстановления списка живых файлов может растягиваться до десятков минут - и это происходит при выполнении абсолютно любого запроса к таблице, включая тривиальный SELECT COUNT(*), потому что без восстановленного списка файлов движок не может даже начать честно вычислять агрегат. Эта деградация - системная, неизбежная плата за выбор линейного журнала вместо дерева: стоимость state reconstruction растёт пропорционально количеству исторических транзакций, а не пропорционально объёму данных, которые реально нужны запросу прямо сейчас.

Checkpoint: периодическая «свёртка» истории в Parquet

Решение, встроенное в Delta Lake с самых ранних версий протокола, - периодическая материализация уже воспроизведённого состояния в отдельный checkpoint-файл: компактный Parquet-файл, который содержит не последовательность action’ов, а уже «свёрнутый», итоговый список живых add-записей (фактически - таблицу метаданных файлов, аналог summary-таблицы) на момент конкретной версии. По умолчанию Delta Lake создаёт checkpoint после каждых десяти коммитов - это управляется свойством таблицы delta.checkpointInterval, которое при необходимости можно изменить под профиль нагрузки конкретной таблицы.

spark.sql("""
    ALTER TABLE delta.`s3a://lakehouse/warehouse/orders_delta`
    SET TBLPROPERTIES (
        'delta.checkpointInterval' = '20'
    )
""")

Имя checkpoint-файла повторяет схему именования JSON-логов, но с дополнительным суффиксом: 00000000000000000010.checkpoint.parquet означает «свёрнутое состояние таблицы по состоянию ровно на версию 10» - то есть результат применения всех add/remove action’ов из версий 0..10, уже посчитанный заранее и сохранённый в виде готовой таблицы файлов. Рядом с checkpoint-файлами Delta Lake поддерживает служебный файл _last_checkpoint - небольшой JSON-объект, содержащий номер версии последнего созданного checkpoint’а и количество файлов в нём; именно этот файл движок читает первым при открытии таблицы, чтобы сразу узнать, с какой версии начинать replay, не перебирая директорию _delta_log/ целиком.

{"version":10,"size":11,"sizeInBytes":284931}

Алгоритм восстановления состояния с учётом checkpoint

Зная номер последнего checkpoint’а из _last_checkpoint, движок Spark выполняет восстановление состояния не от версии 0, а от ближайшего предшествующего checkpoint’а: читает Parquet-файл 00000000000000000010.checkpoint.parquet как уже готовый список живых файлов, а затем применяет к нему (replay) только JSON-файлы версий, идущих после checkpoint’а - то есть 00000000000000000011.json, 00000000000000000012.json и так далее, до текущей актуальной версии. Если интервал между checkpoint’ами равен 10, то в худшем случае нужно проиграть не более 9 дополнительных JSON-файлов сверху самого checkpoint’а - независимо от того, сколько всего версий накопила таблица за всю свою историю.

Этот механизм превращает стоимость восстановления состояния из величины, растущей линейно с полной историей таблицы (O(N), где N - общее число версий с момента создания таблицы), в величину, ограниченную сверху интервалом между checkpoint’ами (O(checkpointInterval) - на практике константа, не зависящая от возраста таблицы). Именно регулярность создания checkpoint’ов - а не их наличие как таковое - определяет, удержится ли Query Planning в приемлемых рамках: если по какой-то причине автоматическое создание checkpoint’а перестаёт срабатывать (например, из-за ошибки прав на запись в _delta_log/ или из-за намеренного отключения автокомпакции на таблице с экстремально высокой частотой записи), накопленный «долг» из непросвёрнутых JSON-файлов начинает расти неограниченно, и деградация planning-фазы возвращается в полной мере - именно такой сценарий разбирается в производственном кейсе этого урока.

Очистка старых JSON и хвоста истории

Хранение и checkpoint-файлов, и всех предшествующих им JSON-логов до бесконечности постепенно превращает _delta_log/ в директорию с десятками тысяч мелких файлов - что само по себе создаёт стоимость на стороне object storage (LIST-операции, метаданные файловой системы) даже если эти файлы давно никому не нужны для чтения текущего состояния таблицы. Delta Lake управляет этим через свойство delta.logRetentionDuration (по умолчанию 30 дней): JSON-файлы и checkpoint’ы старше этого интервала становятся кандидатами на физическое удаление при следующем срабатывании внутреннего процесса очистки лога.

Критически важная деталь, прямо связывающая этот раздел с разделом про Time Travel: logRetentionDuration одновременно определяет, насколько глубоко в прошлое можно «путешествовать» через VERSION AS OF/TIMESTAMP AS OF - удаление JSON-файла версии 5 физически разрушает возможность Time Travel к этой версии, даже если соответствующие Parquet-файлы данных физически ещё существуют (их удаление управляется отдельным свойством delta.deletedFileRetentionDuration и командой VACUUM, разбираемой в двенадцатом уроке модуля). Таким образом у Delta Lake существуют две независимые точки отказа Time Travel: удаление JSON-метаданных лога и удаление физических Parquet-файлов - и обе должны быть настроены согласованно, чтобы глубина Time Travel, заявленная бизнесу, реально была доступна на практике.


Модель снимков Apache Iceberg: древовидная альтернатива

Краткое напоминание: четырёхуровневое дерево вместо журнала

Второй урок модуля детально разбирал внутреннее устройство метаданных Iceberg, и здесь стоит зафиксировать только итоговую структуру, без повторного вывода - читатель, прошедший тот урок, узнает её мгновенно. Каждая транзакция Iceberg производит новый, полностью независимый объект Snapshot, который ссылается на свой собственный Manifest List - файл, перечисляющий все Manifest File, относящиеся к этому снапшоту, с агрегированной по каждому манифесту статистикой (диапазоны значений партиционных колонок, число добавленных/удалённых файлов). Каждый Manifest File в свою очередь перечисляет конкретные Data File, физически содержащие строки таблицы, вместе с column-level статистикой по каждому файлу. Внешний Catalog (Hive Metastore, JDBC, REST Catalog) хранит единственный указатель - путь к актуальному metadata.json, внутри которого находится ссылка на текущий Snapshot.

Принципиальное свойство этой структуры, прямо контрастирующее с Delta Lake: каждый Snapshot - самодостаточный, иммутабельный корень собственного поддерева метаданных. Snapshot версии 42 не нуждается в данных снапшота версии 41 для своей интерпретации: его Manifest List перечисляет ровно те Manifest File, которые формируют полное, актуальное на момент создания этого снапшота состояние таблицы - не дельту относительно предыдущего снапшота, а полный список (с возможностью повторного использования немодифицированных Manifest File между соседними снапшотами как оптимизации, но не как логической необходимости).

«Сверху вниз» против «снизу вверх»: ключевое архитектурное различие

Это и есть фундаментальное отличие, на которое стоит обратить внимание читателя как на главный вывод первых трёх разделов урока. Delta Lake реконструирует текущее состояние таблицы «снизу вверх» (bottom-up): чтобы узнать, что представляет собой таблица сейчас, движок обязан пройти (или, при наличии checkpoint’а, частично пройти) последовательность исторических событий и свернуть их в итоговое состояние через replay. Iceberg, напротив, узнаёт состояние таблицы «сверху вниз» (top-down): Catalog мгновенно отдаёт указатель на актуальный metadata.json, в котором уже прямо записан snapshot-id текущего снапшота - и от этой точки движок спускается по дереву метаданных только в том объёме, который нужен конкретному запросу (а не всей истории таблицы).

Эта разница объясняет сразу несколько следствий, которые иначе выглядели бы как набор несвязанных фактов: почему Iceberg не нуждается во встроенном механизме checkpoint’ов в том виде, в каком он есть у Delta Lake (снапшот и так - готовая, материализованная точка состояния, а не результат накопленной дельты); почему Time Travel в Iceberg не требователен к глубине истории (переход к снапшоту - это прыжок к готовому корню дерева, а не воспроизведение пути к нему, см. раздел 5); и почему обнаружение конфликтов конкурентной записи в Iceberg технически устроено иначе, чем в Delta Lake, - потому что у Iceberg нет «текущей версии плюс N», есть только compare-and-swap указателя на снапшот в Catalog (раздел 4).

Роль внешнего Catalog как точки атомарности

Стоит явно подчеркнуть архитектурную выгоду, которую даёт Iceberg делегирование роли «единой точки правды» внешнему Catalog, а не самой файловой системе. Hive Metastore, PostgreSQL (через JDBC Catalog) или REST Catalog - все они изначально, как часть своей базовой функциональности транзакционной СУБД или сервиса, предоставляют атомарный compare-and-swap: «обновить указатель таблицы с значения X на значение Y, только если текущее значение действительно равно X» - это ровно то же семантическое требование, которое Delta Lake вынужден было решать через rename() на HDFS или через специально написанный multi-cluster LogStore с DynamoDB на S3, описанный в разделе 1.

Iceberg, иными словами, не решает проблему атомарности самостоятельно - он перекладывает её на готовую, проверенную инфраструктуру, которая и так существует в любой организации, использующей реляционные базы данных или REST-сервисы (то есть практически везде). Это объясняет, почему именно с появлением полноценного REST Catalog Protocol (стандартизованного API для управления таблицами Iceberg, не привязанного к конкретной СУБД) экосистема Iceberg получила возможность работать атомарно с любым объектным хранилищем без необходимости писать собственный multi-writer координатор для каждого нового storage backend - задача, которая для Delta Lake долгое время требовала отдельного, специфичного для каждого хранилища решения.


Конкурентный доступ и разрешение конфликтов: Optimistic Concurrency Control

Общий принцип OCC и то, в чём форматы расходятся

Оба формата заявляют поддержку ACID-транзакций через Optimistic Concurrency Control - модель, в которой ни один писатель не берёт блокировку в начале транзакции; вместо этого писатель спокойно готовит свои изменения (вычисляет новые файлы данных, формирует список action’ов или новый снапшот), и только в самый последний момент - момент коммита - проверяет, не успел ли кто-то другой закоммитить конфликтующее изменение раньше. Если конфликта нет - коммит проходит сразу. Если конфликт есть - транзакция либо автоматически разрешается (если конфликт носит безопасный характер), либо завершается ошибкой, требующей retry на уровне приложения.

То, в чём Delta Lake и Iceberg расходятся - не сам принцип OCC, а гранулярность, на которой проверяется наличие конфликта, и она прямо производна от различия, разобранного в разделе 3: у Delta Lake единица проверки - физический файл данных (объект add/remove в JSON-логе), у Iceberg - партиция и диапазон значений, описанные в манифестах. Эта разница в гранулярности определяет, в каких именно сценариях два конкурентных писателя смогут безопасно ужиться, а в каких - один из них обязательно получит ошибку, даже если по бизнес-логике их изменения никак друг другу не противоречат.

Delta Lake: проверка перекрытия файлов на уровне JSON

Представим двух писателей, Job A и Job B, которые одновременно читают таблицу на версии 10 и оба готовятся закоммитить версию 11. Job A выполняет INSERT нескольких новых строк (чистый append, без удаления существующих файлов). Job B выполняет UPDATE строк, физически удовлетворяющих условию WHERE region = 'EU', что для Delta Lake (в режиме Copy-on-Write, который используется по умолчанию) означает: прочитать все файлы, содержащие хотя бы одну подходящую строку, переписать их полностью без обновлённых строк и добавить новый файл с обновлёнными версиями строк - то есть транзакция Job B готовит набор remove для затронутых старых файлов и add для новых.

Допустим, Job A успевает закоммитить версию 11 первым - его JSON-файл 00000000000000000011.json содержит только add-записи новых файлов, без единого remove. Job B, пытаясь закоммитить следом, обнаруживает, что версия 11 уже занята, и обязан перед попыткой коммита версии 12 проверить: пересекается ли множество файлов, которые я собираюсь удалить (remove), с множеством файлов, которые были изменены конкурентной транзакцией, прошедшей между моментом, когда я начал читать таблицу, и текущим моментом? Поскольку Job A выполнил чистый append и не тронул ни одного существующего файла, пересечения нет - и Delta Lake выполняет то, что называется автоматическим разрешением конфликта (conflict resolution): Job B просто перенумеровывает свою подготовленную транзакцию в версию 12 и коммитит её без повторного выполнения вычислений и без ошибки на уровне пользовательского кода.

Если бы вместо чистого append Job A выполнил DELETE или UPDATE, затронувший хотя бы один из тех же физических файлов, что и Job B, проверка перекрытия обнаружила бы конфликт, и Job B получил бы исключение конкретного типа - в зависимости от характера операций это может быть ConcurrentAppendException (конкурентная транзакция добавила файлы в тот же диапазон, что нарушает предположения о изоляции для операций, читающих по предикату), ConcurrentDeleteReadException (конкурентная транзакция удалила файл, который текущая транзакция прочитала и попыталась изменить) или MetadataChangedException (конкурентная транзакция изменила схему или свойства таблицы). Ответственность за обработку такого исключения лежит на вызывающем коде: Spark не делает retry автоматически - приложение (Airflow-таск, Structured Streaming micro-batch с собственной retry-логикой) должно перехватить исключение и повторить всю транзакцию с начала, заново прочитав актуальное состояние таблицы.

Iceberg: проверка на уровне партиций и предикатов через манифесты

Второй урок модуля уже подробно разбирал Write-and-Commit Path Iceberg - три стадии (Write, Metadata, Commit) - и поведение при конфликте: если конкурентная транзакция успела обновить указатель в Catalog между тем, как текущая транзакция прочитала базовый снапшот, и моментом, когда она попыталась закоммитить свой собственный, Iceberg бросает CommitFailedException, и runtime (либо явно написанный retry-цикл, либо встроенный механизм Spark-процедур) должен заново прочитать актуальный снапшот, пересчитать манифесты и повторить попытку коммита. Здесь важно зафиксировать конкретное отличие от Delta Lake в том, что именно проверяется при определении конфликта.

Iceberg сверяет не точное множество затронутых физических файлов, а более широкую, основанную на статистике манифестов информацию: диапазоны значений партиционных колонок и - для операций MERGE INTO/UPDATE/DELETE, реализованных через row-level операции - конкретные предикаты, по которым шла операция. Если Job B выполнял UPDATE ... WHERE region = 'EU', Iceberg проверяет, не модифицировала ли конкурентная транзакция файлы, принадлежащие тем же партициям (например, партиции region=EU), которые затрагивает предикат Job B, - и в случае Copy-on-Write делает это даже строже: проверяет, не изменился ли набор файлов, прочитанных как базис для операции, независимо от того, в какой партиции они физически лежат. Эта проверка в целом строже файловой проверки Delta Lake: даже частичное пересечение диапазонов партиций двух предикатов может быть достаточным основанием для конфликта, тогда как Delta Lake пропустил бы такую ситуацию как «безопасную», если бы физические файлы не пересекались буквально.

Сводная таблица различий в модели конфликтов

Аспект Delta Lake Apache Iceberg
Единица проверки конфликта Физический файл (add/remove в JSON) Партиция / диапазон значений в манифесте, предикат операции
Точка атомарности коммита put-if-absent на следующий номер версии в _delta_log compare-and-swap указателя snapshot-id во внешнем Catalog
Поведение при безопасном непересечении Автоматическое разрешение конфликта, без ошибки в коде приложения Зависит от движка: часто требуется явный retry даже при отсутствии реального пересечения данных
Исключение при конфликте ConcurrentAppendException, ConcurrentDeleteReadException, MetadataChangedException CommitFailedException
Кто выполняет retry Явно вызывающий код (Delta сам не повторяет коммит) Явно вызывающий код, либо retry-логика конкретного клиента/процедуры
Строгость проверки Точное файловое пересечение Более широкая проверка по статистике партиций/предикатов - может быть строже даже без реального пересечения файлов

Практический вывод для инженера, проектирующего конкурентные пайплайны (несколько параллельных Airflow DAG, пишущих в одну таблицу, или Structured Streaming job, работающий параллельно с батчевым backfill): обе модели требуют, чтобы код приложения был готов к retry, и ни одна из них не превращает конкурентную запись в полностью бесплатную операцию без проектирования партиционирования с учётом паттерна записи (если разные джобы стабильно пишут в разные партиции, конфликты в обоих форматах редки; если джобы случайным образом «размазывают» запись по всей таблице, конфликты становятся систематической проблемой производительности, а не разовой случайностью).


Реализация Time Travel и аудит изменений

Delta Lake: Time Travel как ограниченный replay

Девятый урок модуля подробно разбирал Time Travel в Iceberg через три режима адресации - snapshot-id, timestamp и именованные ссылки (tag/branch). У Delta Lake синтаксис на поверхности выглядит почти идентично - те же ключевые слова VERSION AS OF и TIMESTAMP AS OF в Spark SQL, тот же доступ через DataFrame API - но механика, скрытая за этим синтаксисом, принципиально другая, и она прямо производна от bottom-up модели, разобранной в разделе 3.

-- Time Travel по номеру версии
SELECT * FROM delta.`s3a://lakehouse/warehouse/orders_delta` VERSION AS OF 42;

-- Time Travel по timestamp
SELECT * FROM delta.`s3a://lakehouse/warehouse/orders_delta`
TIMESTAMP AS OF '2025-06-15 10:00:00';
df = (
    spark.read
    .format("delta")
    .option("versionAsOf", 42)
    .load("s3a://lakehouse/warehouse/orders_delta")
)

df_by_time = (
    spark.read
    .format("delta")
    .option("timestampAsOf", "2025-06-15 10:00:00")
    .load("s3a://lakehouse/warehouse/orders_delta")
)

Когда движок получает запрос VERSION AS OF 42, он выполняет ровно ту же процедуру восстановления состояния, что описана в разделах 1 и 2, но останавливает replay на версии 42 включительно, полностью игнорируя все JSON-файлы с номером больше 42 - они просто не читаются и не учитываются. Это означает, что стоимость Time Travel-запроса в Delta Lake не константна, а зависит от расстояния между запрошенной версией и ближайшим предшествующим ей checkpoint’ом: если checkpoint существует на версии 40, движок читает checkpoint и проигрывает всего две дополнительные JSON-транзакции (41 и 42); если же запрошена версия 42, а ближайший checkpoint - на версии 0 (потому что checkpoint’ы давно перестали создаваться по какой-то операционной причине), движок будет вынужден проиграть все 42 транзакции с нуля.

Iceberg: мгновенный прыжок к корню дерева

Iceberg реализует Time Travel принципиально иначе - не как ограниченный по дальности replay, а как прямую адресацию к готовому, уже материализованному корню поддерева метаданных. Каждый Snapshot в Iceberg имеет собственный уникальный snapshot-id и timestamp-ms, явно записанные в списке snapshots внутри metadata.json. Запрос VERSION AS OF <snapshot-id> или TIMESTAMP AS OF <timestamp> в Iceberg означает: найти соответствующую запись в списке снапшотов и перейти прямо к Manifest List этого снапшота - минуя необходимость воспроизводить путь к нему через предшествующие снапшоты, потому что, как было зафиксировано в разделе 3, каждый снапшот самодостаточен.

-- Iceberg: Time Travel по snapshot-id (синтаксис, разобранный в девятом уроке)
SELECT * FROM lakehouse.analytics.orders_iceberg VERSION AS OF 8073050618693726890;

SELECT * FROM lakehouse.analytics.orders_iceberg
TIMESTAMP AS OF '2025-06-15 10:00:00';

Эта разница - прямое следствие top-down vs bottom-up архитектуры из раздела 3, и именно здесь она проявляется наиболее ощутимо на практике: стоимость Time Travel-запроса в Iceberg практически не зависит от того, насколько глубоко в историю таблицы лежит запрошенный снапшот - снапшот версии 42 и снапшот версии 4 200 читаются за сопоставимое время, потому что оба - готовые, независимые корни, а не точки на пути последовательного replay. Девятый урок формулировал это как «O(1) вместо O(N)» применительно к самому Iceberg; в контексте сравнения с Delta Lake эта формула приобретает более острый смысл: для Delta Lake Time Travel - это в буквальном смысле O(N) относительно расстояния до checkpoint’а, тогда как для Iceberg - O(1) относительно глубины всей истории таблицы.

Системные таблицы для аудита: DESCRIBE HISTORY против .history/.snapshots

Оба формата предоставляют способ просмотреть историю изменений таблицы без выполнения самого Time Travel - то есть способ ответить на вопрос «что вообще происходило с этой таблицей» прежде, чем решать, к какой именно точке времени откатиться. В Delta Lake это команда DESCRIBE HISTORY, которая читает поля commitInfo из каждого JSON-файла лога и представляет их в виде таблицы:

DESCRIBE HISTORY delta.`s3a://lakehouse/warehouse/orders_delta`;
+-------+-------------------+---------+-------------+--------------------------------+
|version|timestamp          |operation|isBlindAppend|operationParameters            |
+-------+-------------------+---------+-------------+--------------------------------+
|42     |2025-06-15 10:03:21|MERGE    |false        |{predicate -> ...}             |
|41     |2025-06-15 09:58:04|UPDATE   |false        |{predicate -> region = 'EU'}   |
|40     |2025-06-15 09:40:11|WRITE    |true         |{mode -> Append}               |
+-------+-------------------+---------+-------------+--------------------------------+

Iceberg предоставляет эквивалентную информацию не через одну команду, а через несколько системных таблиц метаданных, каждая отвечающая за свой срез истории - что само по себе отражает более модульную, дерево-ориентированную природу метаданных формата. Системная таблица history показывает последовательность снапшотов во времени (аналог DESCRIBE HISTORY), snapshots - детальную информацию о каждом снапшоте (включая summary-статистику добавленных/удалённых файлов и записей), а files/manifests - текущий физический состав таблицы на уровне файлов и манифестов соответственно:

SELECT * FROM lakehouse.analytics.orders_iceberg.history;
SELECT * FROM lakehouse.analytics.orders_iceberg.snapshots;

Стоимость хранения метаданных для долгосрочного Time Travel

Финальный практический аспект, который стоит сопоставить напрямую: стоимость хранения, необходимая для поддержания глубокого Time Travel на горизонте месяцев. У Delta Lake эта стоимость определяется количеством непросвёрнутых JSON-файлов и checkpoint’ов, которые нужно сохранять (управляется delta.logRetentionDuration, разобранным в разделе 2) - сами по себе JSON-файлы малы, но при высокой частоте записи (потоковая нагрузка с коммитом каждые несколько секунд) их совокупный объём и количество объектов в _delta_log/ растут быстро, увеличивая не только объём хранения, но и стоимость LIST-операций при открытии лога.

У Iceberg стоимость долгосрочного Time Travel определяется количеством сохранённых снапшотов и - что важнее с точки зрения объёма - количеством физических Data File, на которые ссылаются эти старые снапшоты и которые поэтому не могут быть удалены через expire_snapshots (детально разобранный в десятом уроке). Поскольку каждый снапшот Iceberg - это полноценное, независимое поддерево манифестов, глубокая история Time Travel в Iceberg обычно «весит» в метаданных больше на один снапшот, чем у Delta Lake на одну версию, но даёт несравнимо более дешёвый и предсказуемый по времени доступ к каждой исторической точке - классический инженерный trade-off между стоимостью хранения метаданных и стоимостью их последующего чтения, который проявляется одинаково последовательно во всех разделах этого урока.


Эволюция схем данных (Schema Evolution): строгость против гибкости

Историческая боль: почему «эволюция схемы» вообще стала отдельной темой

До появления транзакционных table format’ов таблицы на Hive, состоящие из голых Parquet/ORC-файлов под одной директорией, эволюционировали мучительно и опасно: переименование колонки, смена её типа или удаление часто требовали либо полного переписывания всех исторических файлов (дорого и долго на таблицах в десятки терабайт), либо - что хуже - молчаливого расхождения между объявленной схемой в Hive Metastore и фактической физической схемой старых файлов, из-за которого запрос к историческим партициям мог либо упасть с ошибкой типов, либо, незаметно для аналитика, вернуть NULL вместо реальных значений из-за позиционного маппинга колонок по их порядковому номеру, а не по имени. Оба table format’а в этом уроке решают эту проблему, но выбирают разные точки на шкале между строгостью (защита от ошибок) и гибкостью (свобода менять схему без оверхеда).

Delta Lake: Schema Enforcement по умолчанию

Базовое поведение Delta Lake при записи - Schema Enforcement (принудительная проверка схемы): если структура входящего DataFrame не совпадает точно со схемой, зафиксированной в последней metaData-записи _delta_log, Spark немедленно бросает исключение AnalysisException ещё до начала физической записи данных, отказываясь молча принять несовместимые данные. Это касается и добавления новых колонок, и попытки записать данные с отсутствующей обязательной колонкой, и попытки изменить тип существующей колонки на несовместимый.

# Существующая схема таблицы: order_id, amount, region
new_batch = spark.createDataFrame(
    [(1, 99.5, "EU", "promo_code_applied")],
    ["order_id", "amount", "region", "promo_flag"],
)

new_batch.write.format("delta").mode("append").save(
    "s3a://lakehouse/warehouse/orders_delta"
)
# AnalysisException: A schema mismatch detected when writing to the Delta table.
# Table schema: order_id, amount, region
# Data schema:  order_id, amount, region, promo_flag

Чтобы осознанно расширить схему, инженер должен явно указать намерение через опцию mergeSchema - в этом случае Delta Lake добавляет новую колонку в metaData-запись очередной транзакции (со значением NULL для всех ранее записанных строк, потому что они физически не содержат эту колонку), и с этого момента колонка становится частью официальной схемы таблицы:

new_batch.write.format("delta").mode("append") \
    .option("mergeSchema", "true") \
    .save("s3a://lakehouse/warehouse/orders_delta")

Аналогичная защита есть и в обратную сторону - опция overwriteSchema требуется, если операция overwrite должна полностью заменить схему таблицы (а не просто данные), что чаще встречается при пересоздании таблицы с нуля, чем при штатной эволюции.

Delta Column Mapping: ответ на главное ограничение

Исходная модель Delta Lake идентифицирует колонки по имени в физическом Parquet-файле - то есть колонка region в логической схеме таблицы напрямую соответствует физической колонке region в каждом Parquet-файле. Это устроено просто, но создаёт конкретное ограничение: переименование колонки (regioncountry_region) в этой модели логически означает, что все физические Parquet-файлы должны быть переписаны с новым именем колонки - дорогая операция на больших таблицах, и до появления Column Mapping Delta Lake не предлагал способа избежать полного rewrite при переименовании.

Современные версии Delta Lake (начиная с протокола, поддерживающего delta.columnMapping.mode) добавили режим name (а также промежуточный режим id), который вводит уровень индирекции, концептуально близкий к подходу Iceberg: каждой логической колонке присваивается стабильный внутренний идентификатор, и физические Parquet-файлы продолжают использовать старые физические имена колонок, в то время как логическая схема таблицы хранит соответствие «логическое имя → физический идентификатор колонки» в metaData. Включение этого режима на существующей таблице:

spark.sql("""
    ALTER TABLE delta.`s3a://lakehouse/warehouse/orders_delta`
    SET TBLPROPERTIES (
        'delta.columnMapping.mode' = 'name',
        'delta.minReaderVersion' = '2',
        'delta.minWriterVersion' = '5'
    )
""")

spark.sql("""
    ALTER TABLE delta.`s3a://lakehouse/warehouse/orders_delta`
    RENAME COLUMN region TO country_region
""")

С Column Mapping переименование становится чисто метаданной операцией (изменяется только metaData-запись, физические файлы не трогаются) - это значимо сократило разрыв между Delta Lake и Iceberg на этом конкретном сценарии, хотя само включение режима требует поднятия минимальной версии протокола чтения/записи, что может потребовать обновления клиентских библиотек, читающих таблицу.

Iceberg: id-by-name mapping как архитектурная норма, а не опция

Первый и пятый уроки модуля уже разбирали, что Iceberg с самого начала спроектирован вокруг стабильных, никогда не переиспользуемых field-id - каждая колонка (и каждое поле внутри вложенной структуры) получает уникальный численный идентификатор в момент создания, и именно этот идентификатор, а не имя колонки, записывается в физический footer каждого Parquet-файла как метаданные Iceberg-уровня. Переименование колонки в Iceberg - всегда чисто метаданная операция: ALTER TABLE ... RENAME COLUMN обновляет только текстовое имя, привязанное к существующему field-id, в новой версии схемы внутри metadata.json, не трогая ни одного физического файла.

ALTER TABLE lakehouse.analytics.orders_iceberg
RENAME COLUMN region TO country_region;

ALTER TABLE lakehouse.analytics.orders_iceberg
DROP COLUMN promo_flag;

ALTER TABLE lakehouse.analytics.orders_iceberg
ALTER COLUMN amount TYPE double;

Что особенно важно - и что в этой модели не требует отдельного «режима совместимости», в отличие от Column Mapping у Delta Lake, - удаление колонки в Iceberg не разрушает возможность Time Travel к более старым снапшотам, где эта колонка ещё существовала и была заполнена данными: чтение старого снапшота просто использует версию схемы, актуальную на момент создания этого снапшота, потому что Iceberg хранит не одну глобальную схему, а историю версий схемы, привязанных к конкретным снапшотам через их собственные field-id. Удалённый field-id никогда не переиспользуется для новой колонки даже спустя много версий схемы - что исключает категорию ошибок, в которой старые Parquet-файлы случайно «совпадают» по позиции колонки с новой, не связанной с ними по смыслу колонкой.

Сводная таблица и практический эксперимент

Аспект Delta Lake (по умолчанию) Delta Lake (Column Mapping) Apache Iceberg
Идентификация колонки По имени в физическом файле По стабильному внутреннему id, имя - алиас По стабильному field-id, имя - алиас
Добавление колонки Требует mergeSchema Требует mergeSchema Чистая метаданная операция
Переименование колонки Полный rewrite физических файлов Чистая метаданная операция Чистая метаданная операция
Удаление колонки и Time Travel к старым данным Безопасно, но без Column Mapping переименования вперемешку с удалением рискованны Безопасно Безопасно по построению, без отдельного режима
Поведение по умолчанию при несовпадении схемы Жёсткое исключение (AnalysisException) Жёсткое исключение Жёсткое исключение (через валидацию схемы при записи)

Практический эксперимент, который стоит проделать самостоятельно на тестовых таблицах: создать одинаковую таблицу в обоих форматах, выполнить переименование одной колонки в каждой, и сравнить, что физически изменилось. Для Delta Lake без Column Mapping понадобится DESCRIBE HISTORY и сравнение размера и количества физических Parquet-файлов до и после операции - переименование без Column Mapping в старых версиях Delta Lake требовало полного rewrite, что видно по новому набору файлов с более поздними modificationTime. Для Iceberg и Delta Lake с включённым Column Mapping RENAME COLUMN не создаёт ни одного нового физического файла - только новую запись в метаданных, что легко проверить, сравнив листинг директории данных таблицы до и после команды.


Сравнение производительности планирования запросов (Query Planning)

Из чего складывается стоимость планирования у Delta Lake

Третий урок модуля вводил различие между I/O-bound и CPU-bound планированием, разбирая, как Iceberg переносит стоимость планирования с дорогих обращений к файловой системе на дешёвые операции в памяти. Имеет смысл применить ту же рамку к Delta Lake. Чтобы спланировать выполнение запроса (определить, какие именно Parquet-файлы реально нужно прочитать), Spark должен сначала восстановить полный список живых файлов таблицы - тот самый state reconstruction из разделов 1 и 2 - и только после этого применить partition pruning и data skipping по статистике файлов, перечисленной в add-записях.

Проблема в том, что на таблице с большим количеством живых файлов сам по себе checkpoint-файл, содержащий список этих файлов, может разрастись до сотен мегабайт или нескольких гигабайт - на таблице с миллионами файлов данных каждая add-запись с её статистикой (stats - min/max по каждой колонке) занимает заметный объём, и checkpoint, будучи обычным Parquet-файлом, требует полноценного чтения и фильтрации в памяти, прежде чем Spark сможет понять, какие именно строки checkpoint’а (то есть какие физические файлы) релевантны текущему запросу. Хотя сам checkpoint читается распределённо (это обычный Parquet-файл, и Spark может распараллелить его чтение по executor’ам), результат этого чтения - список релевантных файлов - должен быть собран и согласован, прежде чем можно перейти к следующей фазе выполнения запроса, что создаёт менее гибкую, более «один большой проход» структуру планирования по сравнению с Iceberg.

Из чего складывается стоимость планирования у Iceberg

Третий урок детально разбирал четырёхстадийный алгоритм Scan Planning у Iceberg: точка входа через metadata.json, manifest-level pruning (отсечение целых манифестов по агрегированной статистике в Manifest List без открытия самих манифестов), file-level pruning (отсечение конкретных файлов внутри оставшихся манифестов по их собственной статистике) и residual filtering с bin-packing для финального формирования сплитов задачи. Ключевое свойство этой архитектуры, прямо релевантное сравнению с Delta Lake: каждый Manifest File - небольшой, независимый файл, который можно обрабатывать отдельно от других манифестов, и эту обработку Iceberg распределяет по executor’ам как полноценную параллельную Spark-задачу, а не как единый монолитный проход по одному большому файлу состояния.

На таблице с сотнями тысяч партиций и миллионами файлов это различие становится практически ощутимым: Iceberg может пропустить (skip) целые манифесты, основываясь только на сводной статистике в Manifest List, без необходимости даже открывать соответствующие Manifest File - то есть существенная доля файлов метаданных вообще не читается при типичном запросе с выборочным предикатом по партиционной колонке. У Delta Lake аналогичная по духу оптимизация существует (data skipping по статистике в checkpoint’е), но она происходит после того, как сам checkpoint уже свёрнут в единую структуру состояния, а не до того, на уровне отдельных, независимо пропускаемых блоков метаданных.

Метрики Spark UI: на что смотреть при диагностике

При диагностике медленного запроса к Delta-таблице через вкладку SQL/DataFrame в Spark UI инженер должен обращать внимание на фазу, которую сообщество условно называет Log Replay (или Delta Scan в плане выполнения) - время, затраченное на восстановление состояния таблицы из checkpoint’а и хвоста JSON-файлов, прежде чем начнётся фактическое чтение данных. Если эта фаза занимает заметную долю общего времени выполнения запроса (что особенно заметно на коротких, избирательных запросах, где собственно чтение данных должно быть быстрым), это прямой сигнал, что checkpoint давно не создавался или разросся избыточно - типичный повод проверить delta.checkpointInterval и историю фактических checkpoint’ов в _delta_log/.

Для Iceberg эквивалентная фаза в Spark UI и в EXPLAIN-плане отображается как Scan Planning (или конкретнее - время построения BatchScan/TableScan на основе разобранных манифестов); третий урок учил интерпретировать связанные с ней метрики (number of files read, number of files pruned, время на каждую из четырёх стадий алгоритма). Сопоставление двух планов выполнения - EXPLAIN EXTENDED для Delta-таблицы и для Iceberg-таблицы с тем же запросом и сопоставимым объёмом данных - наглядно показывает разницу в количестве физических объектов метаданных, которые движок обязан прочитать перед тем, как приступить к собственно чтению данных, и именно эта разница лежит в основе всех практических рекомендаций по выбору формата для таблиц с экстремальным числом партиций и файлов, обсуждаемых в финальном уроке модуля.

Итоговый вывод раздела

Различие в стоимости планирования между двумя форматами - не абстрактная теоретическая деталь, а прямое, измеримое следствие архитектурного решения, сделанного в самом начале урока: журнал против дерева. Линейный журнал по своей природе вынуждает движок в какой-то момент свернуть всю историю в единую структуру состояния, размер которой растёт с числом живых файлов независимо от того, насколько избирателен конкретный запрос. Дерево манифестов позволяет отбрасывать целые ветви метаданных на основе сводной статистики, не материализуя единое глобальное состояние вообще - и именно эта возможность отбрасывать целые ветви, а не каждый файл по отдельности, даёт Iceberg структурное преимущество на таблицах с действительно большим числом партиций и файлов, тогда как на небольших и средних таблицах с регулярными checkpoint’ами разница между форматами по стоимости планирования обычно практически не ощущается.


Практический демо-блок: PySpark + Delta Lake + Apache Iceberg

Семь теоретических разделов выше построены на постоянном сопоставлении двух архитектур. Эта часть урока переводит то же сопоставление в воспроизводимые команды: один и тот же SparkSession настраивается на работу одновременно с Delta Lake (через spark_catalog, путь-based таблицы на MinIO) и с Iceberg (через каталог lakehouse, уже знакомый из предыдущих уроков модуля), что позволяет выполнять параллельные эксперименты над двумя форматами в одном ноутбуке/скрипте без переключения окружения.

Настройка окружения

Расширения Delta Lake и Iceberg для Spark SQL не конфликтуют друг с другом и могут быть подключены в одном SparkSession одновременно - единственное отличие в том, что Delta Lake в этой конфигурации регистрируется как реализация default-каталога (spark_catalog), в то время как Iceberg, как и в предыдущих уроках, регистрируется как отдельный, явно именованный каталог lakehouse.

from pyspark.sql import SparkSession

spark = (
    SparkSession.builder
    .appName("delta-vs-iceberg-internals")
    .config(
        "spark.jars.packages",
        "io.delta:delta-spark_2.12:3.1.0,"
        "org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.5.2,"
        "org.apache.iceberg:iceberg-aws-bundle:1.5.2",
    )
    # Delta Lake расширения и default-каталог
    .config(
        "spark.sql.extensions",
        "io.delta.sql.DeltaSparkSessionExtension,"
        "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
    )
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")
    # Iceberg каталог "lakehouse", как в первом уроке модуля
    .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")
    .config("spark.sql.catalog.lakehouse.warehouse", "s3a://lakehouse/warehouse")
    .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.access-key-id", "minioadmin")
    .config("spark.sql.catalog.lakehouse.s3.secret-access-key", "minioadmin")
    # Доступ к MinIO для Delta-таблиц через Hadoop S3A (Delta не использует
    # отдельный io-impl namespace, как Iceberg, а работает через fs.s3a.*)
    .config("spark.hadoop.fs.s3a.endpoint", "http://minio:9000")
    .config("spark.hadoop.fs.s3a.access.key", "minioadmin")
    .config("spark.hadoop.fs.s3a.secret.key", "minioadmin")
    .config("spark.hadoop.fs.s3a.path.style.access", "true")
    .getOrCreate()
)

DELTA_PATH = "s3a://lakehouse/warehouse/orders_delta"
spark.sql("CREATE NAMESPACE IF NOT EXISTS lakehouse.analytics")

Кейс 1: создание Delta-таблицы и инспекция _delta_log напрямую

Первый эксперимент - создать минимальную таблицу, выполнить несколько последовательных операций и заглянуть внутрь _delta_log/ напрямую через файловую систему, минуя Spark SQL, чтобы увидеть JSON-коммиты ровно в том виде, в каком они физически лежат на MinIO.

df0 = spark.createDataFrame(
    [(1, 49.9, "EU"), (2, 120.0, "US"), (3, 17.5, "EU")],
    ["order_id", "amount", "region"],
)
df0.write.format("delta").mode("overwrite").save(DELTA_PATH)

df1 = spark.createDataFrame([(4, 88.0, "APAC")], ["order_id", "amount", "region"])
df1.write.format("delta").mode("append").save(DELTA_PATH)

spark.sql(f"UPDATE delta.`{DELTA_PATH}` SET amount = amount * 1.1 WHERE region = 'EU'")
import boto3
import json

s3 = boto3.client(
    "s3",
    endpoint_url="http://minio:9000",
    aws_access_key_id="minioadmin",
    aws_secret_access_key="minioadmin",
)

objects = s3.list_objects_v2(Bucket="lakehouse", Prefix="warehouse/orders_delta/_delta_log/")
for obj in sorted(objects["Contents"], key=lambda o: o["Key"]):
    print(obj["Key"], obj["Size"])
warehouse/orders_delta/_delta_log/00000000000000000000.json   1842
warehouse/orders_delta/_delta_log/00000000000000000001.json    734
warehouse/orders_delta/_delta_log/00000000000000000002.json   1103
body = s3.get_object(
    Bucket="lakehouse",
    Key="warehouse/orders_delta/_delta_log/00000000000000000002.json",
)["Body"].read().decode("utf-8")

for line in body.strip().split("\n"):
    action = json.loads(line)
    print(list(action.keys())[0], "->", action)

Результат второго блока должен показать ровно ту структуру, что разбиралась в разделе 1: версия 2 (результат UPDATE) содержит remove-записи для исходных файлов, физически содержавших строки region = 'EU', и add-записи для новых файлов с уже умноженным значением amount, плюс одну commitInfo-запись с "operation": "UPDATE" и предикатом операции в operationParameters.

Кейс 2: наблюдение за созданием checkpoint

Чтобы увидеть механизм checkpoint’а из раздела 2 в действии, нужно довести таблицу до десятой транзакции (или явно уменьшить delta.checkpointInterval, чтобы не выполнять десять реальных операций для одного учебного эксперимента):

spark.sql(f"""
    ALTER TABLE delta.`{DELTA_PATH}`
    SET TBLPROPERTIES ('delta.checkpointInterval' = '3')
""")

for i in range(5):
    spark.createDataFrame(
        [(100 + i, 10.0 + i, "US")], ["order_id", "amount", "region"]
    ).write.format("delta").mode("append").save(DELTA_PATH)
objects = s3.list_objects_v2(Bucket="lakehouse", Prefix="warehouse/orders_delta/_delta_log/")
for obj in sorted(objects["Contents"], key=lambda o: o["Key"]):
    print(obj["Key"])
warehouse/orders_delta/_delta_log/00000000000000000000.json
warehouse/orders_delta/_delta_log/00000000000000000001.json
warehouse/orders_delta/_delta_log/00000000000000000002.json
warehouse/orders_delta/_delta_log/00000000000000000003.checkpoint.parquet
warehouse/orders_delta/_delta_log/00000000000000000003.json
warehouse/orders_delta/_delta_log/00000000000000000004.json
warehouse/orders_delta/_delta_log/00000000000000000005.json
warehouse/orders_delta/_delta_log/00000000000000000006.checkpoint.parquet
warehouse/orders_delta/_delta_log/00000000000000000006.json
warehouse/orders_delta/_delta_log/_last_checkpoint

Появление 00000000000000000003.checkpoint.parquet и 00000000000000000006.checkpoint.parquet ровно после третьего и шестого коммита подтверждает работу настроенного checkpointInterval = 3. Прочитать содержимое checkpoint’а как обычный Parquet-файл и убедиться, что он содержит уже свёрнутый список живых add-записей, можно напрямую через Spark:

checkpoint_df = spark.read.parquet(
    f"{DELTA_PATH}/_delta_log/00000000000000000006.checkpoint.parquet"
)
checkpoint_df.select("add.path", "add.size", "remove.path").show(truncate=False)

Кейс 3: конкурентный конфликт в Delta Lake - автоматическое разрешение и явный сбой

Раздел 4 описывал два сценария: безопасное автоматическое разрешение конфликта при непересекающихся файлах и явный сбой при пересечении. Воспроизвести оба можно, не поднимая реальные параллельные процессы, а явно эмулируя ситуацию через два независимых объекта транзакции, читающих таблицу на одной и той же версии:

from delta.tables import DeltaTable

# "Job A" читает текущую версию и готовит чистый append
delta_table = DeltaTable.forPath(spark, DELTA_PATH)
version_before = delta_table.history(1).select("version").first()["version"]

# "Job A" коммитит первым - чистый append в новую партицию
spark.createDataFrame(
    [(200, 5.0, "LATAM")], ["order_id", "amount", "region"]
).write.format("delta").mode("append").save(DELTA_PATH)

# "Job B" пытается выполнить UPDATE, как будто всё ещё стоит на version_before -
# в реальном конкурентном сценарии это два параллельных процесса;
# здесь мы наблюдаем итоговый эффект через DESCRIBE HISTORY
spark.sql(f"UPDATE delta.`{DELTA_PATH}` SET amount = amount + 1 WHERE region = 'US'")

spark.sql(f"DESCRIBE HISTORY delta.`{DELTA_PATH}`").select(
    "version", "operation", "operationParameters"
).show(truncate=False)

В реальном продакшен-сценарии симуляции настоящего файлового пересечения (когда Job B обновляет ту же партицию LATAM, которую только что создал Job A) проявляются исключения ConcurrentAppendException - воспроизвести его честно можно только запуском двух параллельных Spark-приложений (например, через два отдельных spark-submit с искусственной задержкой между чтением и записью у одного из них), что выходит за рамки одного интерактивного ноутбука, но является рекомендованным самостоятельным упражнением (см. домашнее задание).

Кейс 4: то же самое для Iceberg - CommitFailedException

Для контраста повторим аналогичный сценарий на Iceberg-таблице, уже зная из второго урока модуля, как там устроен CommitFailedException:

spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.analytics.orders_iceberg (
        order_id BIGINT, amount DOUBLE, region STRING
    ) USING iceberg
    PARTITIONED BY (region)
""")

spark.sql("""
    INSERT INTO lakehouse.analytics.orders_iceberg VALUES
    (1, 49.9, 'EU'), (2, 120.0, 'US'), (3, 17.5, 'EU')
""")

spark.sql("""
    INSERT INTO lakehouse.analytics.orders_iceberg VALUES (200, 5.0, 'LATAM')
""")

spark.sql("""
    UPDATE lakehouse.analytics.orders_iceberg
    SET amount = amount + 1 WHERE region = 'US'
""")

spark.sql("SELECT * FROM lakehouse.analytics.orders_iceberg.history").show(truncate=False)
spark.sql("SELECT * FROM lakehouse.analytics.orders_iceberg.snapshots").show(truncate=False)

Системная таблица snapshots сразу показывает количество добавленных/удалённых файлов на каждом снапшоте (поле summary) - удобный способ убедиться, что INSERT в партицию LATAM и UPDATE по предикату region = 'US' действительно затронули разные, непересекающиеся партиции, и в реальном конкурентном запуске такая пара операций безопасно прошла бы без CommitFailedException даже при строгой, партиционно-ориентированной проверке конфликтов Iceberg.

Кейс 5: Time Travel - сравнение синтаксиса и стоимости

# Delta Lake: Time Travel по версии и по timestamp
spark.read.format("delta").option("versionAsOf", 1).load(DELTA_PATH).show()

spark.read.format("delta") \
    .option("timestampAsOf", "2025-06-15 10:00:00") \
    .load(DELTA_PATH).show()

# Iceberg: Time Travel по snapshot-id, полученному из системной таблицы
snapshot_id = (
    spark.sql("SELECT snapshot_id FROM lakehouse.analytics.orders_iceberg.snapshots ORDER BY committed_at LIMIT 1")
    .first()["snapshot_id"]
)
spark.read.format("iceberg") \
    .option("snapshot-id", snapshot_id) \
    .load("lakehouse.analytics.orders_iceberg").show()

Для измерения практической стоимости каждого запроса в учебных целях достаточно обернуть оба вызова в time.perf_counter() и сравнить - на таблице из нескольких десятков версий разница будет незаметна, но именно это и есть ожидаемый результат: раздел 5 объяснял, что разница проявляется не на маленькой таблице, а на таблице с большим расстоянием от запрошенной версии до ближайшего checkpoint’а (для Delta Lake) или просто отсутствует вовсе (для Iceberg) - поэтому осмысленный замер требует искусственно «состарить» Delta-таблицу до сотен версий без checkpoint’ов, что и предлагается сделать в домашнем задании.

Кейс 6: Schema Evolution - переименование колонки в обоих форматах

# Delta Lake без Column Mapping: попытка переименования потребует rewrite
spark.sql(f"""
    ALTER TABLE delta.`{DELTA_PATH}`
    SET TBLPROPERTIES (
        'delta.columnMapping.mode' = 'name',
        'delta.minReaderVersion' = '2',
        'delta.minWriterVersion' = '5'
    )
""")
spark.sql(f"ALTER TABLE delta.`{DELTA_PATH}` RENAME COLUMN region TO country_region")

# Iceberg: переименование - всегда чистая метаданная операция
spark.sql("""
    ALTER TABLE lakehouse.analytics.orders_iceberg
    RENAME COLUMN region TO country_region
""")
# Сравнение: листинг физических файлов до и после операции
before = set(s3.list_objects_v2(Bucket="lakehouse", Prefix="warehouse/orders_delta/")["Contents"])
# ... выполнить RENAME COLUMN ...
after = set(s3.list_objects_v2(Bucket="lakehouse", Prefix="warehouse/orders_delta/")["Contents"])
print("Новые физические файлы данных:", after - before)

При включённом Column Mapping и на Delta, и на Iceberg список физических файлов данных (part-...parquet) должен остаться без изменений - появится только новая запись в логе/метаданных (_delta_log/0000...00N.json с обновлённой metaData у Delta, новая версия схемы в metadata.json у Iceberg).

Кейс 7: Query Planning - чтение плана выполнения

spark.sql(f"SELECT * FROM delta.`{DELTA_PATH}` WHERE region = 'EU'").explain(True)
spark.sql("SELECT * FROM lakehouse.analytics.orders_iceberg WHERE country_region = 'EU'").explain(True)

В физическом плане Delta-запроса нужно искать узел FileSourceScan, в текстовом представлении которого Spark указывает число файлов после data skipping; в плане Iceberg-запроса - узел BatchScan, в Pushed Filters/PartitionFilters которого видно, какие именно фильтры были применены на уровне манифестов ещё до перехода к чтению самих данных. Сопоставление содержимого вкладки SQL / DataFrame в Spark UI (раздел с разбивкой времени по стадиям выполнения) для обоих запросов на таблице, искусственно доведённой до нескольких тысяч файлов (см. домашнее задание), даёт наглядное, измеримое подтверждение раздела 7: доля времени, уходящая на фазу планирования относительно фазы непосредственного чтения данных, у Delta Lake растёт быстрее с числом файлов, чем у Iceberg.


Производственный кейс

Команда платформы данных среднего fintech-стартапа перевела основную таблицу транзакций (payments.transactions) на Delta Lake почти три года назад - на момент миграции выбор в пользу Delta Lake объяснялся тем, что вся аналитика уже была написана под Databricks-нотбуки, а Iceberg на тот момент не имел зрелого self-hosted REST Catalog. Таблица получала запись от Structured Streaming job, обрабатывающего события платёжного шлюза с частотой коммита раз в 15 секунд - то есть около 5 760 новых версий _delta_log в сутки, или свыше двух миллионов версий за три года эксплуатации.

Проблема начала проявляться постепенно: дашборд для финансового отдела, построенный на простом агрегатном запросе SELECT SUM(amount) FROM transactions WHERE date = current_date(), изначально выполнялся за 4-6 секунд, но за последние полгода время выполнения выросло до 3-4 минут - при том, что объём данных, реально нужных для ответа на этот конкретный запрос (одна календарная дата), оставался стабильным. Команда сначала подозревала проблему в самом Spark-кластере дашборда (нехватка executor’ов, неправильный размер партиций), но профилирование через Spark UI показало, что свыше 90% времени запроса уходило в фазу Delta Scan - то есть в восстановление состояния таблицы из _delta_log, ещё до того, как Spark успевал применить фильтр по дате.

Расследование вскрыло конкретную причину: при настройке Structured Streaming job три года назад инженер по ошибке указал слишком большой интервал между checkpoint’ами (delta.checkpointInterval был установлен в 1000 вместо стандартных 10 - предположительно, в попытке «снизить нагрузку на запись метаданных», без понимания, что плата за это - линейный рост стоимости каждого последующего чтения). За три года это означало, что между соседними checkpoint’ами накапливалось до тысячи невоспроизведённых JSON-транзакций, и при каждом открытии таблицы Spark был вынужден проигрывать сотни мелких файлов сверху последнего checkpoint’а - а поскольку поток коммитов не прекращался, «расстояние от последнего checkpoint’а» постоянно росло между плановыми перерасчётами, и среднее время planning-фазы дрейфовало вверх вместе с накопленным долгом.

Исправление состояло из двух частей: немедленного снижения delta.checkpointInterval обратно до 10 и принудительного запуска ручного чекпойнтинга (DeltaTable.forPath(spark, path).history() плюс служебная команда генерации checkpoint’а через io.delta.tables.DeltaTable, доступную в Delta Lake API) для немедленного создания нового checkpoint’а без ожидания следующих десяти органических коммитов. После исправления время выполнения дашборда вернулось к исходным 4-6 секундам в течение суток, как только Spark начал использовать свежий, близкий checkpoint вместо replay многотысячного хвоста. Команда добавила отдельный alert в систему мониторинга, отслеживающий «возраст» последнего checkpoint’а в версиях (current_version - last_checkpoint_version), чтобы аналогичная конфигурационная ошибка обнаруживалась автоматически, а не через жалобу финансового отдела на медленный дашборд спустя полгода деградации.

Дополнительный вывод, который команда зафиксировала во внутреннем post-mortem: при параллельном пилотном проекте на Iceberg для новой таблицы аналогичного профиля нагрузки (частая потоковая запись) аналогичной проблемы структурно не возникло бы даже при ошибке конфигурации, потому что у Iceberg нет эквивалента «интервала между checkpoint’ами», за настройку которого можно было ошибочно отвечать - манифесты создаются на каждый коммит независимо от частоты, а стоимость их чтения определяется числом партиций, релевантных конкретному запросу, а не общим числом исторических транзакций таблицы. Это не означает, что у Iceberg нет своих аналогичных по духу операционных рисков (избыточное количество мелких манифестов при очень частых коммитах без регулярного rewrite_manifests создаёт похожую по характеру, хотя не идентичную по механике проблему, разобранную в третьем уроке), но конкретно класс ошибки «забыли настроить checkpoint interval» для Iceberg структурно не существует.


Типичные заблуждения

«Delta Lake и Iceberg хранят данные принципиально по-разному». Неверно: оба формата хранят сами строки таблицы в обычных Parquet-файлах, абсолютно идентичных по внутреннему устройству. Различие касается исключительно слоя метаданных, описывающего, какие именно файлы и в каком виде составляют логическое состояние таблицы - это и есть центральная тема всего урока.

«Checkpoint в Delta Lake - это то же самое, что снапшот в Iceberg». Неверно, и это одна из самых частых путаниц при первом знакомстве с обоими форматами. Checkpoint Delta Lake - это оптимизация скорости restore состояния (свёрнутая версия предшествующей истории, не создающая нового логического состояния), тогда как снапшот Iceberg - это и есть само логическое состояние таблицы на конкретный момент, независимо от того, существует ли для него отдельная «свёртка» предыдущей истории (свёртывать нечего - снапшот уже полон).

«Если конфликт записи не возникает в Delta Lake, значит, конфликта не возникнет и в Iceberg при тех же операциях». Неверно: раздел 4 показал, что гранулярность проверки у Iceberg обычно строже (партиция/предикат), чем у Delta Lake (точное файловое пересечение) - сценарий, безопасно разрешившийся автоматически в Delta Lake, в Iceberg при тех же входных данных вполне может потребовать явного retry с CommitFailedException.

«Time Travel в Delta Lake работает медленнее, чем в Iceberg, всегда и при любых обстоятельствах». Неточно: на таблице с регулярными checkpoint’ами и небольшим расстоянием до запрошенной версии разница может быть незаметна на практике. Замедление становится систематическим только при накоплении большого «долга» непросвёрнутых JSON-файлов - что, как показывает производственный кейс, чаще является следствием конфигурационной ошибки или сбоя автоматического checkpoint’а, чем неизбежным свойством формата при разумной настройке.

«Schema Enforcement в Delta Lake - это недостаток по сравнению с гибкостью Iceberg». Контекстно зависит: для таблиц с нестабильной, часто меняющейся внешней схемой (например, входящие JSON-логи от множества разных источников) строгая проверка Delta Lake по умолчанию действительно требует больше явных действий программиста. Но для таблиц, где случайное расширение схемы - признак ошибки в восходящем пайплайне, а не ожидаемое поведение, та же строгость - осознанное защитное свойство, а не недосмотр архитектуры.

«Поскольку Delta Lake появился в экосистеме Databricks, его нельзя использовать в полностью self-hosted инфраструктуре». Неверно: Delta Lake - открытый протокол с открытой референсной реализацией (delta-spark), полностью работающей на собственном Spark-кластере и собственном S3-совместимом хранилище без какой-либо зависимости от управляемого сервиса Databricks - что демонстрирует весь практический демо-блок этого урока, выполненный исключительно на self-hosted MinIO и self-hosted Spark.

«Delta Lake не поддерживает партиционирование так же гибко, как Iceberg». Частично верно, но требует уточнения: модель Hidden Partitioning Iceberg (разобранная в четвёртом уроке модуля) действительно избавляет аналитика от необходимости явно указывать партиционные предикаты в запросе. У Delta Lake партиционирование ближе к классической Hive-модели (явные партиционные колонки, видимые в предикате запроса), что менее эргономично, но не делает партиционирование Delta Lake принципиально нерабочим или несовместимым с эффективным data skipping - просто механизм отсечения данных у Delta Lake опирается в первую очередь на статистику файлов в checkpoint’е, а не на структуру партиционного дерева.

«checkpoint и Vacuum - это одна и та же операция обслуживания Delta-таблицы». Неверно: checkpoint - это операция над метаданными (свёртка JSON-истории в Parquet), выполняемая автоматически на каждый N-й коммит; Vacuum - операция над физическими данными (удаление файлов, помеченных remove и старше порога retention), требующая явного вызова и разбираемая отдельно в двенадцатом уроке модуля. Обе операции снижают накопленный «мусор», но в совершенно разных слоях системы и с разной частотой и механикой запуска.


Производственный чек-лист: что нельзя забыть до того, как Delta-таблица уйдёт в production

  • delta.checkpointInterval явно проверен на каждой таблице с интенсивной записью, а не оставлен на дефолтное значение «потому что никто не трогал» - производственный кейс этого урока показал, что одна забытая настройка способна за несколько месяцев превратить простой SELECT COUNT(*) в трёхминутный запрос.

  • «Возраст» последнего checkpoint’а (current_version - last_checkpoint_version) вынесен в мониторинг с алертом на разумном пороге - это превращает деградацию Query Planning из тихого, обнаруживаемого только жалобами аналитиков процесса в наблюдаемую, проактивно отслеживаемую метрику.

  • delta.logRetentionDuration и delta.deletedFileRetentionDuration настроены согласованно друг с другом и с реальным горизонтом Time Travel, заявленным бизнесу или аудиту - раздел 2 показал, что это две независимые точки отказа, и рассинхрон между ними создаёт ситуацию, в которой Time Travel «формально доступен», но фактически упирается в более раннее из двух ограничений.

  • Команда понимает разницу между файловой проверкой конфликтов Delta Lake и партиционной проверкой Iceberg прежде, чем переносить retry-логику из одного формата в другой при миграции пайплайна - сценарий, безопасно проходящий в одном формате, не гарантированно безопасен в другом при тех же входных операциях (раздел 4).

  • Конкурентные пайплайны, пишущие в одну таблицу, спроектированы с учётом партиционирования по паттерну записи - если разные джобы стабильно затрагивают разные партиции, конфликты в обоих форматах остаются редкими исключениями, а не систематической причиной retry-штормов под нагрузкой.

  • Переход на Column Mapping (delta.columnMapping.mode) выполнен осознанно, с обновлением минимальной версии протокола, а не реактивно, в момент, когда понадобилось впервые переименовать колонку - повышение minReaderVersion/minWriterVersion может потребовать синхронного обновления всех клиентов, читающих таблицу, и должно быть запланированной миграцией, а не точечным патчем.

  • Профиль нагрузки таблицы (частота записи, число конкурентных писателей, типичная избирательность запросов) задокументирован и пересмотрен при значимом росте объёма - выбор, оптимальный на старте проекта при сотне версий в день, не обязан оставаться оптимальным при достижении сотен тысяч версий, и периодическая ревизия конфигурации checkpoint’ов и retention - часть нормальной эксплуатации, а не разовая настройка «один раз и навсегда».


Мостик к следующему уроку

Этот урок дал детальную, конкретную модель того, как Delta Lake устроен внутри, и зафиксировал, в чём именно она отличается от уже изученной модели Iceberg - на уровне физических файлов журнала, механизма checkpoint’ов, конфликтов конкурентного доступа, Time Travel, эволюции схемы и стоимости планирования. Этого достаточно, чтобы прочитать _delta_log глазами, понимающими, что там физически лежит, и чтобы объяснить коллеге, почему два «одинаковых на первый взгляд» ACID-table format’а на практике ведут себя по-разному в конкретных операционных сценариях.

Но знание устройства лога и checkpoint’ов - это только половина практической компетенции, нужной для эксплуатации Delta-таблицы в production. Производственный кейс этого урока уже показал, к чему приводит забытая или неправильно настроенная checkpoint-стратегия; то же самое в полный рост касается физического удаления устаревших файлов данных, борьбы с фрагментацией мелких файлов после интенсивной потоковой записи, и целенаправленной сортировки данных для ускорения последующих запросов. Двенадцатый урок модуля («Delta Lake: OPTIMIZE, VACUUM, Z-ordering и auto-optimize») разбирает именно эти три процедуры обслуживания - прямой аналог rewrite_data_files, expire_snapshots и remove_orphan_files, уже изученных для Iceberg в десятом уроке, но устроенный по-своему под линейную, журнальную архитектуру Delta Lake. Имея сегодняшнюю модель _delta_log и checkpoint’ов в голове, механика OPTIMIZE и VACUUM в следующем уроке будет читаться не как новый набор команд для запоминания, а как прямое, ожидаемое продолжение того, что уже было разобрано здесь.


Домашнее задание

  1. Разверните self-hosted Delta Lake поверх локального MinIO (или существующей инфраструктуры из предыдущих уроков модуля) и создайте тестовую таблицу с десятью последовательными INSERT-операциями. После каждой операции выводите листинг _delta_log/ и зафиксируйте, на какой именно операции появляется первый checkpoint-файл при настройке по умолчанию (checkpointInterval = 10).

  2. Прочитайте содержимое любого JSON-файла из _delta_log/ вашей тестовой таблицы напрямую (через boto3/файловую систему, без Spark) и для каждой найденной add-записи вручную сопоставьте поле stats с реальными значениями колонок в соответствующем Parquet-файле - убедитесь, что min/max действительно совпадают с физическим содержимым файла.

  3. Искусственно создайте таблицу с большим расстоянием от последней версии до последнего checkpoint’а (установите delta.checkpointInterval в большое значение, например 500, и выполните 100+ операций append по одной строке). Замерьте время выполнения простого SELECT COUNT(*) до и после принудительного запуска создания checkpoint’а и зафиксируйте разницу в секундах.

  4. Создайте идентичную по схеме и объёму данных таблицу в Iceberg и в Delta Lake. Выполните на обеих 50 операций append небольшими батчами (по 1-5 строк), затем сравните: (а) суммарный объём всех файлов метаданных (_delta_log/ против суммы Manifest List + Manifest File), (б) количество отдельных объектов метаданных, (в) время выполнения одного и того же selective-запроса с фильтром по партиционной колонке.

  5. Воспроизведите конкурентный конфликт в Delta Lake по-настоящему - запустите два независимых процесса spark-submit, оба читающих таблицу на одной версии: первый выполняет UPDATE по предикату region = 'EU', второй - DELETE по предикату region = 'EU' AND amount > 100. Добавьте в один из процессов искусственную задержку (time.sleep) между чтением состояния и началом записи, чтобы гарантированно создать пересечение. Зафиксируйте полученное исключение и опишите, как его обработка должна выглядеть в production-коде (retry-цикл с ограничением попыток).

  6. Повторите тот же эксперимент на Iceberg-таблице с аналогичными операциями и партиционированием по region. Сравните тип и текст полученного исключения, а также то, насколько строго или гибко каждый формат отнёсся к частично пересекающимся, но не идентичным предикатам двух конкурентных транзакций.

  7. Включите delta.columnMapping.mode = 'name' на тестовой Delta-таблице, переименуйте одну колонку, затем выполните VERSION AS OF к версии, существовавшей до переименования. Убедитесь, что историческое чтение продолжает работать корректно, и объясните своими словами, через какой механизм это обеспечивается (стабильный внутренний идентификатор колонки, а не текстовое имя).

  8. Постройте сравнительную таблицу (на основе разделов 1-7 этого урока) из 7 строк - по одной на каждый раздел - и для каждой строки сформулируйте одним предложением: «При каком профиле нагрузки разница между Delta Lake и Iceberg в этом аспекте становится критичной, а при каком - не имеет практического значения?». Эта таблица должна стать личным чек-листом для будущего выбора формата под конкретный production-сценарий.

  9. Найдите в документации Delta Lake (Delta Lake Protocol Specification) перечень всех типов исключений конкурентного доступа (ConcurrentAppendException, ConcurrentDeleteReadException, ConcurrentDeleteDeleteException, MetadataChangedException, ProtocolChangedException) и для каждого опишите конкретную пару операций, которая может его спровоцировать - используя как образец сценарии, разобранные в разделе 4 этого урока.

  10. На таблице, искусственно доведённой до нескольких тысяч мелких файлов (множественные append по одной-двум строкам без промежуточной компакции), сравните вывод EXPLAIN для одного и того же selective-запроса в Delta Lake и в Iceberg, и зафиксируйте конкретные числовые метрики planning-фазы из Spark UI (число файлов до/после pruning, время фазы scan/planning) - это прямая практическая проверка выводов раздела 7.


Полная картина

Диаграмма ниже сводит весь урок в одну схему: путь одной записи от момента коммита до момента, когда аналитик читает таблицу - параллельно для Delta Lake и для Iceberg, с явным указанием, на каком шаге каждая архитектура платит свою цену и получает свою выгоду.

Обе ветви диаграммы начинаются с одного и того же шага - Spark вычисляет, какие физические файлы нужно добавить или удалить - и заканчиваются одним и тем же результатом: аналитик получает корректный, изолированный, версионированный ответ на свой запрос. Всё содержательное различие сосредоточено в middle-секции: способ зафиксировать коммит атомарно, способ проверить конфликт, способ ускорить будущее чтение и способ, которым это чтение в итоге считывает состояние. Эта диаграмма - не новая информация, а компактная карта уже разобранного материала, на которую можно опираться при объяснении различий коллеге за пять минут вместо часовой лекции.


Итоги

Delta Lake реализует ACID-транзакционность через единый, строго последовательный журнал коммитов (_delta_log), в котором каждая транзакция - это ровно один новый JSON-файл с монотонно растущим номером версии, содержащий action’ы add/remove/metaData/protocol/commitInfo. Текущее состояние таблицы определяется как результат последовательного применения (replay) всех action’ов от версии 0 до текущей.

Рост числа версий без оптимизации делает восстановление состояния всё более дорогим, поэтому Delta Lake автоматически создаёт checkpoint-файлы - свёрнутые Parquet-снимки состояния каждые delta.checkpointInterval коммитов (по умолчанию 10) - которые ограничивают стоимость replay сверху интервалом между checkpoint’ами, а не полной длиной истории таблицы.

Iceberg решает ту же задачу принципиально иначе - через дерево независимых, иммутабельных снапшотов, каждый из которых самодостаточен и не требует воспроизведения предыдущей истории для своей интерпретации. Delta Lake реконструирует состояние «снизу вверх» через replay, Iceberg узнаёт его «сверху вниз» через мгновенный переход к готовому корню метаданных конкретного снапшота.

Оба формата используют Optimistic Concurrency Control, но проверяют конфликты на разной гранулярности: Delta Lake - на уровне точного пересечения физических файлов (с конкретной таксономией исключений: ConcurrentAppendException, ConcurrentDeleteReadException, MetadataChangedException), Iceberg - на уровне партиций и предикатов операции через CommitFailedException, обычно строже, чем файловая проверка Delta Lake.

Time Travel у Delta Lake - это ограниченный по дальности replay (стоимость зависит от расстояния между запрошенной версией и ближайшим checkpoint’ом), тогда как у Iceberg - это прямой переход к снапшоту по его уникальному идентификатору, практически не зависящий от глубины истории таблицы.

Schema Enforcement Delta Lake по умолчанию строго проверяет соответствие схемы при записи, требуя явного mergeSchema для расширения; полноценная безболезненная эволюция (переименование без rewrite) доступна только при включении Column Mapping. Iceberg реализует id-by-name field mapping как архитектурную норму с самого начала, без отдельного «режима совместимости».

Стоимость Query Planning у Delta Lake определяется размером checkpoint-файла, который должен быть загружен и профильтрован как единая структура состояния, тогда как у Iceberg планирование распределяется по независимым манифестам и может параллельно отбрасывать целые ветви метаданных по сводной статистике - разница, ощутимая прежде всего на таблицах с сотнями тысяч партиций и миллионами файлов.

Ни один из двух форматов не является «более правильным» в абсолютном смысле - оба решают одну транзакционную задачу разными структурами данных, и выбор между ними определяется конкретным профилем нагрузки, существующей экосистемой инструментов и операционной зрелостью команды, что детально разбирается в финальном уроке модуля.


Краткий глоссарий терминов урока

_delta_log - скрытая директория внутри Delta-таблицы, содержащая весь транзакционный журнал: последовательность JSON-коммитов и периодические checkpoint-файлы.

Action - отдельная JSON-запись внутри коммита Delta Lake, описывающая одно элементарное изменение состояния: add, remove, metaData, protocol или commitInfo.

add / remove - типы action, отвечающие за логическое добавление и логическое удаление конкретного физического файла данных из состояния таблицы.

State reconstruction (log replay) - процесс восстановления текущего списка живых файлов таблицы путём последовательного применения всех action’ов от начала истории (или от последнего checkpoint’а) до целевой версии.

Checkpoint - периодически создаваемый Parquet-файл, содержащий уже свёрнутое состояние таблицы на конкретную версию, избавляющий движок от необходимости проигрывать всю историю с нуля.

delta.checkpointInterval - свойство таблицы, определяющее частоту автоматического создания checkpoint-файлов (по умолчанию каждые 10 коммитов).

delta.logRetentionDuration - свойство таблицы, определяющее, сколько времени хранятся старые JSON-логи и checkpoint’ы перед физическим удалением; напрямую ограничивает глубину доступного Time Travel.

Atomic put-if-absent - примитив, гарантирующий, что файл с конкретным именем (номером версии) будет создан только если он ещё не существует; основа атомарности коммита Delta Lake.

LogStore / commit-координатор - механизм координации многопользовательской записи в Delta-таблицу на хранилищах без нативной поддержки атомарного rename (исторически - на S3).

Optimistic Concurrency Control (OCC) - модель конкурентного доступа, при которой писатель не берёт блокировку заранее, а проверяет наличие конфликта только в момент коммита.

ConcurrentAppendException / ConcurrentDeleteReadException / MetadataChangedException - таксономия исключений Delta Lake, бросаемых при обнаружении конкретного типа конфликта конкурентной записи.

CommitFailedException - исключение Iceberg, бросаемое при неудачной попытке compare-and-swap указателя snapshot-id в Catalog из-за конкурентного изменения.

Snapshot (Iceberg) - иммутабельный, самодостаточный объект, описывающий полное состояние таблицы на конкретный момент через собственный Manifest List.

Bottom-up / Top-down реконструкция - два архитектурных подхода к определению текущего состояния таблицы: через накопительный replay истории (Delta Lake) или через прямой переход к материализованному корню метаданных (Iceberg).

VERSION AS OF / TIMESTAMP AS OF - синтаксис Spark SQL для Time Travel, общий для обоих форматов на уровне поверхности, но различающийся по внутренней механике выполнения.

DESCRIBE HISTORY - команда Delta Lake для просмотра истории коммитов таблицы на основе записей commitInfo.

.history / .snapshots (Iceberg) - системные таблицы метаданных Iceberg, предоставляющие аналогичный DESCRIBE HISTORY обзор истории снапшотов и детальную статистику по каждому из них.

Schema Enforcement - поведение Delta Lake по умолчанию, при котором запись с несовпадающей схемой завершается исключением, а не молчаливым приведением типов или пропуском данных.

mergeSchema - опция записи Delta Lake, явно разрешающая расширение схемы таблицы новыми колонками из входящего DataFrame.

Column Mapping (delta.columnMapping.mode) - механизм Delta Lake, вводящий стабильный внутренний идентификатор колонки, отделённый от её текстового имени, что делает переименование чистой метаданной операцией.

Field-id - стабильный, никогда не переиспользуемый числовой идентификатор колонки в Iceberg, записываемый в footer физических Parquet-файлов и обеспечивающий безопасную эволюцию схемы без рерайта данных.

Query Planning - фаза выполнения запроса, предшествующая фактическому чтению данных, на которой движок определяет точный список физических файлов, релевантных запросу.

Manifest-level pruning / file-level pruning - две последовательные стадии отсечения нерелевантных файлов в Iceberg, основанные на сводной статистике манифестов и на собственной статистике отдельных файлов соответственно.

Serializable Isolation - уровень изоляции транзакций, гарантирующий, что результат выполнения конкурентных транзакций эквивалентен какому-то их последовательному порядку выполнения; заявляемая гарантия Delta Lake для конкурентных коммитов в одну таблицу.

protocol (action) - запись в _delta_log, фиксирующая минимальную версию протокола чтения (minReaderVersion) и записи (minWriterVersion), необходимую клиенту для безопасной работы с таблицей; растёт при включении новых возможностей формата, например Column Mapping или Deletion Vectors.

_last_checkpoint - небольшой служебный JSON-файл внутри _delta_log/, хранящий номер версии последнего созданного checkpoint’а; читается движком первым при открытии таблицы, чтобы не перебирать директорию лога целиком.

Write Amplification (в контексте Delta Lake) - дополнительный объём перезаписанных данных, возникающий из-за Copy-on-Write поведения операций UPDATE/DELETE/MERGE, при котором затронутый файл переписывается целиком; подробно разобран применительно к Iceberg в шестом уроке модуля и сохраняет ту же природу у Delta Lake.

Delta Standalone / Delta Kernel / delta-rs - библиотеки, позволяющие читать и писать Delta-таблицы без полноценного Spark-движка (из Java/Scala, из унифицированного Kernel API, или из Rust/Python-инструментов соответственно); расширяют изначально Spark-центричную модель Delta Lake на другие движки и языки.

REST Catalog Protocol - стандартизованный, не привязанный к конкретной СУБД API для управления Iceberg-таблицами, упомянутый в разделе 3 как ключевая инфраструктура, благодаря которой Iceberg получает атомарность коммита без необходимости писать собственный multi-writer координатор для каждого нового storage backend.