Архитектура Transaction Log
Delta Lake — это открытый протокол хранения таблиц поверх Parquet-файлов на object storage (S3, ADLS, GCS, локальная FS). Ключевое отличие от голого Parquet — наличие transaction log (журнала транзакций), который превращает набор файлов в ACID-таблицу.
В Модуле 10 мы видели, как Delta Lake решает schema evolution. Этот урок — глубже: мы разберём побайтово структуру transaction log, каждый тип action, механику checkpoints и систему версионирования протокола.
Этот курс использует engine-agnostic подход. Примеры кода — на Python с библиотекой deltalake (delta-rs, версия 1.5.0), без привязки к Apache Spark. Delta Lake — это протокол, а не движок.
Структура _delta_log/
Каждая Delta-таблица содержит директорию _delta_log/ рядом с Parquet-файлами данных. Эта директория — единственный источник истины о состоянии таблицы:
my_table/
Корневая директория Delta-таблицы. Содержит Parquet-файлы с данными и служебную директорию _delta_log/ с журналом транзакций.Имена JSON-файлов — 20-значные zero-padded номера версий. Версия 0 — это первый коммит при создании таблицы. Каждая последующая операция (INSERT, UPDATE, DELETE, ALTER TABLE) создаёт новый JSON-файл с инкрементированным номером.
_delta_log/
├── 00000000000000000000.json ← версия 0 (CREATE TABLE)
├── 00000000000000000001.json ← версия 1 (INSERT)
├── 00000000000000000002.json ← версия 2 (UPDATE)
├── ...
├── 00000000000000000010.checkpoint.parquet ← чекпоинт версии 10
├── 00000000000000000011.json ← версия 11
├── _last_checkpoint ← указатель на последний чекпоинт
Анатомия JSON-коммита
Каждый JSON-файл содержит одну или несколько строк (newline-delimited JSON — ndjson). Каждая строка — один action. Действия применяются к состоянию таблицы в порядке записи.
Шесть типов action:
add — добавление файла
Action add — самый частый. Каждый INSERT, UPDATE, MERGE создаёт новые Parquet-файлы и записывает add action для каждого:
{
"add": {
"path": "part-00000-1a2b3c4d-5e6f-7890-abcd-ef1234567890-c000.snappy.parquet",
"size": 1048576,
"modificationTime": 1700000000000,
"dataChange": true,
"partitionValues": { "date": "2025-01-15" },
"stats": "{\"numRecords\":50000,\"minValues\":{\"id\":1,\"amount\":0.01},\"maxValues\":{\"id\":50000,\"amount\":9999.99},\"nullCount\":{\"id\":0,\"amount\":12}}",
"tags": null
}
}
Ключевые поля:
| Поле | Тип | Описание |
|---|---|---|
path | string | Относительный путь к Parquet-файлу от корня таблицы |
size | long | Размер файла в байтах |
modificationTime | long | Epoch milliseconds — время последней модификации файла |
dataChange | boolean | true если данные изменились (INSERT/UPDATE/DELETE). false для compaction/OPTIMIZE |
partitionValues | map | Значения partition-колонок для этого файла. Пустой map для non-partitioned таблиц |
stats | string | JSON-строка со статистикой: numRecords, minValues, maxValues, nullCount по каждой колонке |
tags | map | Произвольные метки. Используется для deletion vectors и других расширений |
Поле stats — это строка, а не объект. Парсинг двухуровневый: сначала JSON коммита, потом JSON статистики внутри. Движки используют эту статистику для data skipping — пропуска файлов, которые не могут содержать нужные данные. Подробнее — в уроке 04.
remove — удаление файла
{
"remove": {
"path": "part-00000-old-file.snappy.parquet",
"deletionTimestamp": 1700000100000,
"dataChange": true
}
}
remove не удаляет файл физически — только помечает его как невидимый для текущей версии таблицы. Файл остаётся на диске до VACUUM. Это обеспечивает time travel: старые версии таблицы ещё могут ссылаться на этот файл.
dataChange: false означает compaction — файл заменён оптимизированной версией без изменения данных.
metaData — метаданные таблицы
Присутствует в первом коммите (версия 0) и при ALTER TABLE:
{
"metaData": {
"id": "af23c9d7-fff1-4a5a-a2c8-55c59bd782aa",
"name": "users",
"format": { "provider": "parquet", "options": {} },
"schemaString": "{\"type\":\"struct\",\"fields\":[{\"name\":\"id\",\"type\":\"long\",\"nullable\":false,\"metadata\":{}},{\"name\":\"name\",\"type\":\"string\",\"nullable\":true,\"metadata\":{}}]}",
"partitionColumns": ["date"],
"configuration": {
"delta.enableChangeDataFeed": "true",
"delta.minReaderVersion": "1",
"delta.minWriterVersion": "4"
}
}
}
schemaString — JSON-сериализация Spark StructType. Это единственный источник истины для схемы таблицы. Отдельные Parquet-файлы могут иметь подмножество колонок (из-за schema evolution), но metaData.schemaString всегда содержит полную текущую схему.
protocol — версионирование протокола
{
"protocol": {
"minReaderVersion": 2,
"minWriterVersion": 7,
"readerFeatures": ["columnMapping", "deletionVectors", "timestampNtz", "v2Checkpoint"],
"writerFeatures": ["columnMapping", "deletionVectors", "timestampNtz", "v2Checkpoint", "coordinatedCommits", "clusteringTable"]
}
}
Этот action определяет, какие версии клиентов могут читать и писать в таблицу. Подробнее о протоколе — в следующем разделе.
txn — идемпотентность
{
"txn": {
"appId": "streaming-job-42",
"version": 157,
"lastUpdated": 1700000200000
}
}
Используется streaming-приложениями для exactly-once semantics. Если приложение streaming-job-42 уже записало версию 157 — повторная попытка будет отклонена. Это не влияет на обычные batch-записи.
commitInfo — аудит
{
"commitInfo": {
"timestamp": 1700000300000,
"operation": "WRITE",
"operationParameters": { "mode": "Append", "partitionBy": "[date]" },
"engineInfo": "Apache-Spark/4.0.0 Delta-Lake/4.0.0",
"isBlindAppend": true
}
}
Необязательный action. Не влияет на состояние таблицы, но полезен для аудита: кто, когда и какую операцию выполнил. Поле operation содержит тип операции: WRITE, MERGE, DELETE, UPDATE, OPTIMIZE, VACUUM, CREATE TABLE и другие.
Чтение transaction log через delta-rs
from deltalake import DeltaTable
import json
# Открываем таблицу
dt = DeltaTable("./my_table")
# Текущая версия
print(f"Версия: {dt.version()}")
# Метаданные
meta = dt.metadata()
print(f"ID таблицы: {meta.id}")
print(f"Partition columns: {meta.partition_columns}")
# Схема таблицы
schema = dt.schema()
for field in schema.fields:
print(f" {field.name}: {field.type} (nullable={field.nullable})")
# История коммитов
history = dt.history()
for entry in history[:5]:
print(f" v{entry['version']}: {entry['operation']} @ {entry['timestamp']}")
# Активные файлы с статистикой
actions = dt.get_add_actions(flatten=True)
df = actions.to_pandas()
print(f"\nАктивных файлов: {len(df)}")
print(f"Общий размер: {df['size_bytes'].sum() / 1024**2:.1f} MB")
Для низкоуровневого доступа к сырым JSON-коммитам можно прочитать файлы напрямую:
import json
from pathlib import Path
log_dir = Path("./my_table/_delta_log")
# Читаем конкретный коммит
commit_file = log_dir / "00000000000000000005.json"
with open(commit_file) as f:
for line in f:
action = json.loads(line)
action_type = list(action.keys())[0]
print(f"Action: {action_type}")
print(json.dumps(action[action_type], indent=2))
Protocol Versioning и Table Features
Delta Lake использует двухуровневую систему совместимости: версии протокола (грубые) и table features (гранулярные).
Protocol = Reader Version + Writer Version
Версии Reader и Writer — числовые значения, которые клиент должен понимать, чтобы корректно читать или писать в таблицу. Если клиент не поддерживает нужную версию — операция отклоняется.Table Features (Writer V7 / Reader V2)
С Writer V7 Delta Lake перешёл на гранулярную систему table features. Вместо одного числа версии — явный список feature-флагов:
# Проверка протокола через delta-rs
dt = DeltaTable("./my_table")
protocol = dt.protocol()
print(f"Reader: V{protocol.min_reader_version}")
print(f"Writer: V{protocol.min_writer_version}")
# Table features доступны при Writer V7
if hasattr(protocol, 'reader_features') and protocol.reader_features:
print(f"Reader features: {protocol.reader_features}")
if hasattr(protocol, 'writer_features') and protocol.writer_features:
print(f"Writer features: {protocol.writer_features}")
Повышение версии протокола — необратимая операция. Если вы установили Writer V7 для Coordinated Commits, клиенты с поддержкой только Writer V5 больше не смогут писать в эту таблицу. В Delta Lake 4.0 появилась возможность понижать table features без потери истории, но версия протокола не понижается.
Checkpoints: Parquet-снимки состояния
По мере накопления коммитов чтение лога замедляется: для восстановления текущего состояния нужно прочитать все JSON-файлы с самого начала. Checkpoints решают эту проблему.
Механика checkpoint
Checkpoint — это Parquet-файл, содержащий полное накопленное состояние таблицы на момент определённой версии. Вместо replay сотен JSON-коммитов — одно чтение Parquet.
Коммиты: v0.json → v1.json → … → v9.json
Клиент выполняет операции записи. Каждая операция создаёт новый JSON-файл в _delta_log/ с инкрементированным номером версии.Создание checkpoint для v10
Движок создаёт checkpoint: читает все JSON до v10, накапливает состояние (все активные add, текущие metaData и protocol), и записывает результат как Parquet-файл.Чтение v15: checkpoint(v10) + replay v11..v15
Для восстановления текущего состояния (например, v15): прочитать checkpoint v10 + replay JSON v11, v12, v13, v14, v15. Пять файлов вместо шестнадцати.Формат _last_checkpoint
Single-part checkpoint:
{"version": 10, "size": 42}
Multi-part checkpoint (для очень больших таблиц):
{"version": 100, "size": 15000, "parts": 4}
Здесь size — количество action-записей в checkpoint, а parts — количество Parquet-файлов, на которые разбит checkpoint. Multi-part checkpoints появляются при сотнях тысяч активных файлов.
V2 Checkpoints (Writer V7)
Delta Lake 4.0 ввёл v2 checkpoints с улучшениями:
- Sidecar files: вспомогательные файлы с дополнительной информацией (file actions выносятся в sidecar, основной файл — только metaData + protocol)
- Version checksums: хеш-сумма для быстрой проверки целостности
- File counts и table size: метрики прямо в checkpoint — не нужен полный scan для базовых вопросов
# Чтение _last_checkpoint
import json
from pathlib import Path
checkpoint_file = Path("./my_table/_delta_log/_last_checkpoint")
with open(checkpoint_file) as f:
info = json.loads(f.read())
print(f"Checkpoint version: {info['version']}")
print(f"Actions count: {info['size']}")
if 'parts' in info:
print(f"Multi-part: {info['parts']} файлов")
Алгоритм Log Replay
Когда клиент открывает Delta-таблицу, он выполняет log replay — восстановление текущего состояния из журнала. Алгоритм:
- Прочитать _last_checkpoint
- Загрузить checkpoint.parquet
- Replay JSON-коммитов после checkpoint
- Reconcile: add − remove = активные файлы
Текущее состояние таблицы готово к запросу
Результат: клиент знает текущую схему (metaData.schemaString), протокол (protocol), и точный список активных Parquet-файлов с их статистикой (stats). Можно выполнять запрос.# delta-rs делает log replay автоматически при открытии таблицы
dt = DeltaTable("./my_table")
# Результат replay: список активных файлов
active_files = dt.file_uris()
print(f"Активных файлов: {len(active_files)}")
# Можно загрузить конкретную версию (time travel — подробнее в уроке 03)
dt_v5 = DeltaTable("./my_table", version=5)
print(f"Файлов в версии 5: {len(dt_v5.file_uris())}")
Без checkpoint log replay для таблицы с 10 000 коммитов означает чтение 10 000 JSON-файлов. С checkpoint каждые 10 коммитов — максимум 10 JSON-файлов после последнего checkpoint. Конфигурация: delta.checkpointInterval (по умолчанию 10).
Собираем всё вместе: жизненный цикл коммита
Что происходит при INSERT INTO my_table VALUES (...):
- Запись данных: движок создаёт новый Parquet-файл с данными
- Подготовка action: формирует
addaction с путём, размером, stats, иcommitInfo - Атомарный коммит: записывает JSON-файл
{version}.jsonв_delta_log/. Атомарность зависит от storage:
- HDFS: atomic rename (
temp.json→00000000000000000042.json) - S3: conditional put (put-if-absent) — через DynamoDB или S3 conditional write
- ADLS: conditional write (If-None-Match ETag)
- Checkpoint (опционально): если
version % checkpointInterval == 0, создаёт checkpoint
Если два клиента пытаются записать одну и ту же версию одновременно — один из них получит conflict. Механизм разрешения конфликтов — в следующем уроке.
Итоги
| Компонент | Роль |
|---|---|
_delta_log/*.json | Версионированные коммиты — ndjson с action |
add / remove | Файловые операции: добавление и логическое удаление |
metaData | Схема, partition columns, конфигурация таблицы |
protocol | Версионирование: minReader/Writer Version + table features |
txn | Идемпотентность для streaming |
commitInfo | Аудит: кто, когда, какая операция |
checkpoint.parquet | Parquet-снимок накопленного состояния |
_last_checkpoint | Указатель на последний checkpoint |
| Log Replay | Checkpoint + JSON replay → текущее состояние |
В следующем уроке мы разберём, как Delta Lake обеспечивает ACID-транзакции при конкурентной записи — optimistic concurrency control, conflict resolution, и Coordinated Commits.
Delta Lake в Spark — production ClickHouse + Delta Lake — архитектуры интеграций Lakehouse architecture — system design выбор формата