Delta Lake: глубокое погружение
Архитектура
Delta Lake — lakehouse-формат, разработанный Databricks в 2019 году. Архитектура построена на transaction log (_delta_log/), который записывает каждую атомарную операцию над таблицей.
Delta таблица на диске:
/data/events/
├── _delta_log/
│ ├── 00000000000000000000.json ← версия 0: initial write
│ ├── 00000000000000000001.json ← версия 1: append
│ ├── 00000000000000000002.json ← версия 2: update
│ └── 00000000000000000010.checkpoint.parquet ← checkpoint
├── part-00000-...snappy.parquet
├── part-00001-...snappy.parquet
└── part-00002-...snappy.parquet
Каждый JSON-файл в _delta_log/ содержит список actions: добавленные файлы (add), удалённые файлы (remove), обновления metadata (metaData). Checkpoint-файлы в Parquet создаются каждые 10 версий для ускорения чтения.
ACID-гарантии: optimistic concurrency control. Каждая транзакция атомарно записывает новый JSON в _delta_log/. При конфликте (два writer одновременно) один из них получает ошибку и перезапускается.
Настройка SparkSession
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("DeltaLakePipeline") \
.config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
.config("spark.sql.catalog.spark_catalog",
"org.apache.spark.sql.delta.catalog.DeltaCatalog") \
.config("spark.jars.packages", "io.delta:delta-spark_2.13:4.0.0") \
.getOrCreate()
Ключевые возможности
Time Travel
Delta Lake сохраняет полную историю изменений. Вы можете читать данные на любую предыдущую версию или timestamp:
from delta.tables import DeltaTable
# Создание таблицы
df.write.format("delta").save("/data/events")
# Чтение конкретной версии
df_v0 = spark.read.format("delta") \
.option("versionAsOf", 0) \
.load("/data/events")
# Чтение по timestamp
df_yesterday = spark.read.format("delta") \
.option("timestampAsOf", "2024-01-15T10:00:00") \
.load("/data/events")
# Просмотр истории
spark.sql("DESCRIBE HISTORY delta.`/data/events`").show()
Восстановление данных — откат таблицы к предыдущей версии:
dt = DeltaTable.forPath(spark, "/data/events")
dt.restoreToVersion(0) # Откат к версии 0
# или
dt.restoreToTimestamp("2024-01-15T10:00:00")
Anti-pattern: не запускать VACUUM
Каждая версия таблицы хранит полные файлы данных. Без VACUUM старые файлы накапливаются и потребляют хранилище. Настройте регулярный VACUUM с retention period (по умолчанию 7 дней):
# Удаление файлов старше 7 дней
spark.sql("VACUUM delta.`/data/events` RETAIN 168 HOURS")
# Важно: после VACUUM time travel к удалённым версиям невозможен!Schema Evolution
Delta Lake поддерживает эволюцию схемы при записи:
# Добавление новых колонок (mergeSchema)
new_data.write.format("delta") \
.option("mergeSchema", "true") \
.mode("append") \
.save("/data/events")
# Полная замена схемы (overwriteSchema)
reshaped_data.write.format("delta") \
.option("overwriteSchema", "true") \
.mode("overwrite") \
.save("/data/events")
Поддерживаемые операции: добавление колонок, переименование (через ALTER TABLE), удаление (через ALTER TABLE DROP COLUMN).
OPTIMIZE и Z-ORDER
Малые файлы — главный враг производительности Delta Lake. OPTIMIZE объединяет файлы, Z-ORDER оптимизирует их для фильтрации по определённым колонкам:
# Compaction: объединение мелких файлов
spark.sql("OPTIMIZE delta.`/data/events`")
# Z-ORDER: оптимизация data skipping по колонкам
spark.sql("""
OPTIMIZE delta.`/data/events`
ZORDER BY (event_date, user_id)
""")
Data skipping: Delta Lake хранит min/max статистику по каждому файлу. При фильтрации WHERE event_date = '2024-01-15' Spark пропускает файлы, где min > ‘2024-01-15’ или max < ‘2024-01-15’. Z-ORDER группирует связанные значения в одних файлах, увеличивая эффективность skipping.
Spark 4.0: новые возможности
UniForm Spark4.0
UniForm (Universal Format) позволяет читать Delta-таблицу как Apache Iceberg или Apache Hudi без конвертации данных. Delta автоматически генерирует metadata в формате Iceberg/Hudi при каждой записи.
# Создание Delta-таблицы с UniForm
spark.sql("""
CREATE TABLE events (
event_id STRING,
event_time TIMESTAMP,
payload STRING
) USING delta
TBLPROPERTIES (
'delta.universalFormat.enabledFormats' = 'iceberg,hudi'
)
""")
# Теперь таблицу можно читать из Trino/Presto как Iceberg!
UniForm решает проблему vendor lock-in: вы пишете в Delta, а другие движки (Trino, Presto, Snowflake) читают через Iceberg-интерфейс.
Coordinated Commits Spark4.0
В Spark 4.0 Delta Lake получил Coordinated Commits — внешний координатор транзакций для multi-cluster сценариев. До этого optimistic concurrency работал только в рамках одного кластера. С Coordinated Commits несколько кластеров могут безопасно писать в одну таблицу, используя внешний commit service (DynamoDB, Azure CosmosDB).
Variant Type Spark4.0
Новый тип данных Variant для хранения semi-structured данных (JSON) внутри Delta-таблицы без фиксированной схемы:
spark.sql("""
CREATE TABLE logs (
log_id BIGINT,
timestamp TIMESTAMP,
payload VARIANT -- semi-structured JSON
) USING delta
""")
Type widening Spark4.0
Type widening — возможность изменить тип колонки на более широкий (INT -> BIGINT, FLOAT -> DOUBLE, DECIMAL(10,2) -> DECIMAL(20,2)) без переписывания parquet-файлов. Доступно в preview с Delta 3.2 и полностью production-ready в Delta 4.0.
Принципиальное отличие от классического ALTER TABLE ALTER COLUMN: широкий тип записывается в schema metadata, а старые parquet-файлы продолжают использовать узкий тип на физическом уровне. При чтении Delta автоматически промотирует значения. Никакого rewrite — мутация почти бесплатная.
# Включение type widening
spark.sql("""
ALTER TABLE events
SET TBLPROPERTIES ('delta.enableTypeWidening' = 'true')
""")
# Расширение типа колонки -- мгновенная metadata-операция
spark.sql("ALTER TABLE events ALTER COLUMN amount TYPE BIGINT")
spark.sql("ALTER TABLE events ALTER COLUMN ratio TYPE DOUBLE")
Поддерживаемые widenings:
| Из | В | Условие |
|---|---|---|
BYTE | SHORT, INT, LONG | Безопасно (sign-preserving) |
SHORT | INT, LONG | Безопасно |
INT | LONG | Безопасно |
FLOAT | DOUBLE | Безопасно |
DECIMAL(p1,s1) | DECIMAL(p2,s2) | p2-s2 >= p1-s1 и s2 >= s1 |
DATE | TIMESTAMP_NTZ | Сохраняется wall-clock |
Schema evolution на write автоматически использует widening, если включён mergeSchema:
# Автоматическое расширение типа при write новых данных
new_data_with_bigger_int.write.format("delta") \
.option("mergeSchema", "true") \
.mode("append") \
.save("/data/events")
# amount был INT -> Delta автоматически продвинет к BIGINT
Drop column / Drop feature Spark4.0
Delta 4.0 формализует operation DROP FEATURE — удаление table feature вместе с возможностью снизить protocol version для совместимости с старыми клиентами.
# Удалить колонку (column mapping должен быть включён)
spark.sql("ALTER TABLE events DROP COLUMN deprecated_field")
-- Удалить целую table feature (например, deletion vectors)
ALTER TABLE events DROP FEATURE deletionVectors;
-- Если активные deletion vectors есть, Delta переписывает parquet-файлы,
-- материализуя удаления. После этого feature можно безопасно убрать.
Drop feature workflow в 4.0 умнее, чем в 3.x: сохраняется set of “safe checkpoints”, доступных и старым (downgrade-target), и новым клиентам, чтобы во время drop-операции table оставалась читаемой обоими протокольными версиями.
-- Полный workflow drop type widening:
ALTER TABLE events DROP FEATURE typeWidening;
-- Delta перезаписывает parquet-файлы, where применённое расширение типа
-- было только в metadata. После завершения протокол can be downgraded.
Deletion vectors evolution Spark4.0
Deletion vectors — storage optimization, появившаяся ещё до 4.0. Без них удаление одной строки из parquet-файла требует rewrite целого файла. С DV: Delta создаёт bitmap “deleted rows” в отдельном файле, parquet остаётся неизменным, а при чтении DV применяется поверх данных.
# Включить deletion vectors на новой таблице
spark.sql("""
CREATE TABLE events (
event_id STRING,
timestamp TIMESTAMP,
payload STRING
) USING delta
TBLPROPERTIES ('delta.enableDeletionVectors' = 'true')
""")
# DELETE становится почти мгновенным
spark.sql("DELETE FROM events WHERE event_id = 'abc-123'")
# Delta пишет DV-файл вместо rewrite parquet
Что нового в 4.0 для deletion vectors:
DROP FEATURE deletionVectors— безопасный pathway удаления feature. Delta автоматически материализует все DV в parquet-уровне, после чего metadata-протокол можно downgrade.- Streaming-aware DV reads — structured streaming в 4.0 корректно учитывает DV при micro-batch-чтении (раньше DV пропускались, что приводило к чтению удалённых записей в стриминге).
- DV + Z-ORDER coexistence — Z-ORDER теперь корректно re-applies DV при rewrite (раньше Z-ORDER обнулял DV, потеряв удаления).
| Операция | Без DV | С DV |
|---|---|---|
DELETE WHERE id = X (1 row) | Rewrite целого файла (~128 MB) | Append bitmap (~10 KB) |
MERGE INTO ... DELETE | Rewrite N файлов | Append N bitmaps |
| Read performance | Native parquet read | Read parquet + apply DV (overhead ~5-10%) |
| OPTIMIZE behavior | Standard | Materializes DVs, removes bitmap files |
Deletion vectors радикально улучшают latency удалений в GDPR / right-to-be-forgotten пайплайнах: вместо переписывания терабайтов файлов из-за одной row, операция становится metadata-only. После накопления критической массы DV (например, 10% rows в файле помечены deleted) запустите OPTIMIZE — он смержит изменения в новые parquet и обнулит DV.
Anti-patterns
Частые ошибки при работе с Delta Lake
-
Не запускать VACUUM — старые файлы данных накапливаются бесконечно, потребляя хранилище. Настройте VACUUM с retention period (минимум 7 дней для safety).
-
Не указывать retention period —
VACUUMбезRETAIN N HOURSможет удалить файлы, нужные для активных транзакций. Всегда указывайте retention >= времени самого длинного запроса. -
Слишком частый OPTIMIZE — каждый OPTIMIZE переписывает файлы данных. Для append-heavy таблиц запускайте 1 раз в день, не каждый час.
-
Z-ORDER по high-cardinality колонкам — Z-ORDER по
user_idс миллионами уникальных значений неэффективен. Используйте колонки с medium cardinality (date, region, category).
Когда использовать Delta Lake
Лучший выбор, когда:
- Вы работаете в Databricks экосистеме
- Вам нужна лучшая интеграция со Spark (Delta — default lakehouse-формат Spark 4.0)
- Хотите UniForm для мультиdвижкового доступа
- Нужны простые в настройке ACID-транзакции
Не лучший выбор, когда:
- Вам критична vendor-нейтральность (рассмотрите Iceberg)
- Основная нагрузка — upserts с incremental pull (рассмотрите Hudi)
- Streaming-first архитектура с changelog CDC (рассмотрите Paimon)
Для углублённого изучения внутренней архитектуры Delta Lake (transaction log, checkpoints, data skipping) см. курс Storage Formats Deep-Dive.