Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 09.02 · 16 мин
Средний
Delta LakeACIDTime TravelSchema EvolutionOPTIMIZEZ-ORDERUniFormCoordinated Commits

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")
WARNING

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.

Проверка знанийKnowledge check
Почему Z-ORDER улучшает производительность запросов с фильтрацией?
ОтветAnswer
Z-ORDER пересортировывает данные внутри файлов так, чтобы связанные значения указанных колонок оказывались в одних и тех же файлах. Delta Lake хранит min/max статистику по каждому файлу. При запросе с WHERE-фильтром Spark использует data skipping -- пропускает файлы, у которых min/max значения не пересекаются с условием фильтра. Без 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:

ИзВУсловие
BYTESHORT, INT, LONGБезопасно (sign-preserving)
SHORTINT, LONGБезопасно
INTLONGБезопасно
FLOATDOUBLEБезопасно
DECIMAL(p1,s1)DECIMAL(p2,s2)p2-s2 >= p1-s1 и s2 >= s1
DATETIMESTAMP_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:

  1. DROP FEATURE deletionVectors — безопасный pathway удаления feature. Delta автоматически материализует все DV в parquet-уровне, после чего metadata-протокол можно downgrade.
  2. Streaming-aware DV reads — structured streaming в 4.0 корректно учитывает DV при micro-batch-чтении (раньше DV пропускались, что приводило к чтению удалённых записей в стриминге).
  3. 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 ... DELETERewrite N файловAppend N bitmaps
Read performanceNative parquet readRead parquet + apply DV (overhead ~5-10%)
OPTIMIZE behaviorStandardMaterializes DVs, removes bitmap files
INFO

Deletion vectors радикально улучшают latency удалений в GDPR / right-to-be-forgotten пайплайнах: вместо переписывания терабайтов файлов из-за одной row, операция становится metadata-only. После накопления критической массы DV (например, 10% rows в файле помечены deleted) запустите OPTIMIZE — он смержит изменения в новые parquet и обнулит DV.

Проверка знанийKnowledge check
Чем UniForm в Delta Lake полезен для мультиdвижковых архитектур?
ОтветAnswer
UniForm (Universal Format) автоматически генерирует metadata в формате Apache Iceberg и Apache Hudi при каждой записи в Delta-таблицу. Это позволяет другим вычислительным движкам (Trino, Presto, Snowflake) читать Delta-таблицу через привычный им Iceberg или Hudi интерфейс без конвертации данных. Решает проблему vendor lock-in -- данные физически хранятся в Delta, но доступны через любой формат.

Anti-patterns

WARNING

Частые ошибки при работе с Delta Lake

  1. Не запускать VACUUM — старые файлы данных накапливаются бесконечно, потребляя хранилище. Настройте VACUUM с retention period (минимум 7 дней для safety).

  2. Не указывать retention periodVACUUM без RETAIN N HOURS может удалить файлы, нужные для активных транзакций. Всегда указывайте retention >= времени самого длинного запроса.

  3. Слишком частый OPTIMIZE — каждый OPTIMIZE переписывает файлы данных. Для append-heavy таблиц запускайте 1 раз в день, не каждый час.

  4. 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)
TIP

Для углублённого изучения внутренней архитектуры Delta Lake (transaction log, checkpoints, data skipping) см. курс Storage Formats Deep-Dive.

Delta Lake transaction log: внутренности

Закончили урок?

Отметьте его как пройденный, чтобы отслеживать свой прогресс

Войдите чтобы оценить урок

Прогресс модуля
0 из 8