Constraints в Lakehouse: почему нет PK, идемпотентность, MERGE vs Dedup, Snapshot Rollback

Почему Iceberg не enforce Primary Keys, идемпотентность как архитектурный принцип, сравнение Write-Time Deduplication (MERGE INTO) и Read-Time Deduplication (Append + Window), Snapshot Isolation в Iceberg, атомарные коммиты и Rollback за миллисекунды

lakehouse iceberg idempotency snapshot-isolation merge

Разработчик, пришедший в мир Data Lakehouse из PostgreSQL или Oracle, первым делом задаёт вопрос: «Хорошо, где тут создать Primary Key и Foreign Key?». После нескольких экспериментов он обнаруживает, что CREATE TABLE ... PRIMARY KEY синтаксически принимается Apache Iceberg, но при двойной записи никакой ошибки не происходит. Дубли молча сохраняются.

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

В этом уроке мы разберём, почему Lakehouse-системы не могут и не должны быть клонами реляционных СУБД, как достигается целостность данных без классических constraints, и какие механизмы (идемпотентность, MERGE INTO, Snapshot Isolation) заменяют привычные гарантии базы данных.


1. Иллюзия контроля: почему Iceberg не enforce Primary Keys

1.1 Что происходит при объявлении PK в Iceberg

Apache Iceberg поддерживает синтаксис объявления Primary Key:

CREATE TABLE lakehouse.silver.orders (
    order_id STRING,
    amount   DECIMAL(10,2),
    status   STRING
) USING iceberg
TBLPROPERTIES (
    'primary-key' = 'order_id'
)

Iceberg принимает это определение и хранит его в метаданных таблицы - для интеграции с Data Catalog'ами (Apache Atlas, AWS Glue, Google Data Catalog), которые отображают схему таблицы аналитикам. Но при записи никакой проверки уникальности не происходит.

# Демонстрация: дубли записываются без ошибок
spark.sql("""
    INSERT INTO lakehouse.silver.orders
    VALUES ('order_1', 100.0, 'CREATED')
""")
spark.sql("""
    INSERT INTO lakehouse.silver.orders
    VALUES ('order_1', 100.0, 'CREATED')  -- тот же order_id!
""")

# Результат: две строки с одинаковым order_id
spark.table("lakehouse.silver.orders").show()
# | order_id | amount | status  |
# |----------|--------|---------|
# | order_1  | 100.0  | CREATED |
# | order_1  | 100.0  | CREATED |  ← дубль, никаких ошибок

Это не баг Iceberg. Это фундаментальное физическое ограничение распределённых систем.

1.2 Физика: почему enforce PK в Spark невозможно без огромной цены

Представьте ситуацию: в Silver таблице orders хранится 100 терабайт данных (несколько лет истории заказов). Пришёл новый батч из 10 ГБ. Чтобы проверить уникальность order_id во входящих данных, система должна:

Все три варианта физически неприемлемы для petabyte-scale данных. В PostgreSQL это работает потому, что B-Tree индекс помещается в RAM сервера, а записи происходят последовательно в один процесс с блокировкой. В распределённой системе с сотнями параллельных воркеров и S3 как хранилищем - нет ни глобального индекса, ни механизма блокировок на уровне строк.

1.3 Последствия молчаливых дублей для downstream-пайплайнов

Дубли в Silver таблице - не просто «лишние данные». Они ломают всю аналитику:

Двойной счёт в агрегатах - самое опасное последствие. Финансовый отчёт с удвоенной выручкой, доставленный CEO, разрушает доверие к платформе. И найти причину будет крайне сложно: аналитик не знает о дублях на уровне Silver, его запрос технически корректен.


2. Идемпотентность как архитектурный принцип

2.1 Определение и требование

Идемпотентность (Idempotency) - свойство операции, при котором её применение один раз или многократно даёт одинаковый результат.

f(f(x)) = f(x)

Для ETL-пайплайна это формулируется так:

Запуск пайплайна 1 раз и 100 раз на одних и тех же входных данных должен давать абсолютно идентичный результат в целевой таблице.

Это не опциональная хорошая практика - это обязательное требование к любому production-пайплайну, который хочет выжить в реальном мире.

2.2 Почему пайплайны перезапускаются

В распределённых системах сбои - норма, а не исключение. Причины перезапуска:

Причина Что происходит Риск без идемпотентности
OOM на Executor Задача упала на 90%, Airflow перезапускает Уже записанные 90% + повторная запись = двойник
Сетевой сбой S3 Запись прервана на середине файла Частичный файл + повторная запись = дубль
Airflow worker crash Dag перезапущен с начала Весь батч повторён, старые данные сохранились
Конфигурационная ошибка Исправили, перезапустили тот же батч Повторная запись тех же строк
Рабочий день: «запусти ещё раз для надёжности» Ручной перезапуск Двойники по всей таблице

2.3 Анти-паттерн «Слепой Append»

Самый распространённый источник дублей - наивный пайплайн с mode("append"):

# ❌ АНТИ-ПАТТЕРН: Слепой Append без идемпотентности
def process_daily_orders(date: str):
    new_orders = spark.read.parquet(f"s3://landing/{date}/orders/")
    # Если эта функция вызвана дважды с одним date - дубли неизбежны
    new_orders.write \
        .format("iceberg") \
        .mode("append") \
        .saveAsTable("lakehouse.silver.orders")

Если Airflow перезапустит задачу (из-за сбоя, таймаута, ручного запуска), те же данные будут записаны дважды. Таблица растёт, дубли накапливаются. Через месяц COUNT(*) в два раза больше COUNT(DISTINCT order_id).

2.4 Как достигается идемпотентность

Три стратегии для идемпотентного пайплайна:

  1. MERGE INTO (Write-Time Deduplication) - при записи находим совпадения по ключу и обновляем вместо добавления
  2. Overwrite партиции - полностью перезаписываем партицию за конкретную дату, не добавляя к ней
  3. Append + Downstream Deduplication - пишем всё подряд, но читатели получают данные через дедуплицирующую View

3. Битва стратегий: MERGE INTO vs Append + Dedup

3.1 Стратегия 1: Write-Time Deduplication (MERGE INTO)

При каждой записи Spark ищет совпадения по бизнес-ключу между входящим батчем и существующей таблицей. Найдено совпадение - обновляем строку. Нет совпадения - вставляем новую.

# ✅ Идемпотентный UPSERT через MERGE INTO
new_orders_df.createOrReplaceTempView("incoming_orders")

spark.sql("""
    MERGE INTO lakehouse.silver.orders AS target
    USING incoming_orders AS source
        ON target.order_id = source.order_id
    WHEN MATCHED THEN
        UPDATE SET
            target.amount = source.amount,
            target.status = source.status,
            target.updated_at = source.updated_at
    WHEN NOT MATCHED THEN
        INSERT *
""")
# Сколько бы раз ни запустили с теми же данными - результат одинаков

Плюсы MERGE INTO:

  • Таблица всегда содержит ровно одну строку на order_id
  • Аналитики не думают о дублях при запросах к таблице
  • CDC processing (insert/update/delete) естественно ложится в WHEN MATCHED/NOT MATCHED
  • Rollback и Time Travel работают корректно

Минусы MERGE INTO:

  • Дорогой JOIN между входящим батчем и всей существующей таблицей
  • Copy-on-Write: перезаписываются все Parquet-файлы, содержащие обновлённые строки
  • Write Amplification растёт с ростом таблицы
  • При частом ingestion (каждые 5 минут) создаёт огромную нагрузку

3.2 Оптимизация MERGE: сужение зоны поиска через партиционирование

Самый важный приём для масштабирования MERGE - ограничить зону поиска совпадений партицией. Если мы знаем, что обновления приходят только для заказов за последние 30 дней, незачем сканировать всю историю:

# ❌ Медленно: MERGE по всей таблице (100 ТБ)
spark.sql("""
    MERGE INTO lakehouse.silver.orders AS target
    USING incoming AS source
    ON target.order_id = source.order_id
    ...
""")

# ✅ Быстро: MERGE только по свежим партициям
spark.sql("""
    MERGE INTO lakehouse.silver.orders AS target
    USING incoming AS source
    ON  target.order_id  = source.order_id
    AND target.order_date >= current_date() - INTERVAL 30 DAYS
    ...
""")
# Spark применит Static Partition Pruning: читает только 30 Parquet-файлов
# вместо всей таблицы!

Добавление AND target.order_date >= current_date() - INTERVAL 30 DAYS сужает зону Join с «100 ТБ всей истории» до «нескольких ГБ последних 30 дней». Время выполнения MERGE падает с часов до минут.

3.3 Стратегия 2: Overwrite партиции

Для таблиц, которые хранят данные в партициях по дням, идемпотентность достигается проще: полностью перезаписывать партицию за конкретную дату:

# Идемпотентная запись через overwrite конкретной партиции
daily_orders = spark.read.parquet(f"s3://landing/{processing_date}/")

daily_orders \
    .withColumn("order_date", F.lit(processing_date).cast("date")) \
    .writeTo("lakehouse.silver.orders_partitioned") \
    .overwritePartitions()
# Сколько ни перезапускай - партиция за processing_date будет содержать
# только данные из этого конкретного запуска

Ограничение: работает только если партиция полностью принадлежит одному батчу. Если за один день могут прийти несколько батчей с разными заказами - этот подход перетрёт предыдущие записи.

3.4 Стратегия 3: Append + Read-Time Deduplication

Третья стратегия разделяет запись и дедупликацию: пишем всё через APPEND (быстро и дёшево), а дедупликацию делаем при чтении через SQL View с оконной функцией.

# Запись: быстрый APPEND без проверок
new_orders_df.writeTo("lakehouse.silver.orders_raw").append()

# Дедуплицирующий View поверх raw таблицы
spark.sql("""
    CREATE OR REPLACE VIEW lakehouse.silver.orders_current AS
    SELECT *
    FROM (
        SELECT *,
            ROW_NUMBER() OVER (
                PARTITION BY order_id
                ORDER BY updated_at DESC
            ) AS rn
        FROM lakehouse.silver.orders_raw
    )
    WHERE rn = 1
""")

Аналитики работают с orders_current, а не с orders_raw. В orders_current всегда одна строка на order_id.

Плюсы Append + Read Dedup:

  • Запись максимально быстрая - просто добавить файл
  • Идеально для streaming ingestion (каждые 30 секунд)
  • orders_raw - полный аудит всей истории событий
  • Компакция и дедупликация выполняются по расписанию, не мешая записи

Минусы Append + Read Dedup:

  • Каждый SELECT к View запускает ROW_NUMBER() → Wide Transformation → Shuffle
  • При миллиардах строк в orders_raw каждый запрос аналитика дорог
  • View нельзя физически материализовать встроенными средствами (нужен отдельный scheduled pipeline)

3.5 Матрица выбора стратегии


4. Snapshot Isolation: атомарность коммитов и магия Rollback

4.1 MVCC в Iceberg: как работает изоляция снапшотов

Apache Iceberg реализует MVCC (Multi-Version Concurrency Control) - тот же механизм, что используется в PostgreSQL, Oracle и других ACID-СУБД, но адаптированный для объектного хранилища.

Каждая успешная операция записи создаёт новый снапшот - неизменяемый pointer на конкретную версию метаданных таблицы. Читатели всегда работают с конкретным снапшотом и не видят незавершённых операций.

Ключевое свойство: если Spark-задача упала после записи файлов, но до коммита метаданных - файлы существуют на S3, но они не входят ни в один снапшот. Для читателей эти файлы невидимы. Периодически вызываемая процедура expire_snapshots / remove_orphan_files подчищает такие «призрачные» файлы.

4.2 Защита от частичных записей

Это принципиальное отличие от «голого» Parquet: при записи напрямую в Parquet-файлы на S3 частичная запись (файл существует, но неполный) может привести к corrupted данным. В Iceberg метаданные и данные разделены - пока нет коммита, нет снапшота, нет видимости.

4.3 Rollback за миллисекунды

Допустим, пайплайн успешно завершился (коммит прошёл), но записал некорректные данные - отрицательные суммы, неправильный статус, случайный перезапуск со старым батчем. В СУБД пришлось бы писать DELETE или UPDATE с перечислением всех испорченных строк. В Iceberg - одна команда:

# Шаг 1: находим ID хорошего снапшота (до плохой записи)
spark.sql("""
    SELECT snapshot_id, committed_at, operation, summary
    FROM lakehouse.silver.orders.snapshots
    ORDER BY committed_at DESC
""").show(truncate=False)
# | snapshot_id        | committed_at        | operation | summary           |
# |--------------------|---------------------|-----------|-------------------|
# | 8291038471234567   | 2026-05-23 15:30:00 | append    | added 1000 files  |  ← плохой
# | 7182947362819456   | 2026-05-23 14:00:00 | merge     | changed 50 files  |  ← хороший

valid_snapshot_id = 7182947362819456

# Шаг 2: откатываемся к хорошему снапшоту
spark.sql(f"""
    CALL lakehouse.system.rollback_to_snapshot(
        'silver.orders',
        {valid_snapshot_id}
    )
""")
# Это НЕ удаляет файлы. Это просто переключает текущий указатель снапшота.
# Время выполнения: ~100 мс независимо от размера таблицы

Что происходит физически: Iceberg создаёт новый снапшот, который указывает на ту же коллекцию файлов, что и valid_snapshot_id. «Плохой» снапшот перестаёт быть текущим, но его файлы остаются на S3 - они будут очищены expire_snapshots через настроенный retention period (обычно 7–30 дней).

4.4 Time Travel: диагностика через прошлое

Snapshot Isolation открывает ещё одну мощную возможность: Time Travel - запросы к прошлым состояниям таблицы.

# Запрос к таблице на конкретный момент времени (до плохой записи)
spark.sql("""
    SELECT COUNT(*), SUM(amount)
    FROM lakehouse.silver.orders
    TIMESTAMP AS OF '2026-05-23 14:00:00'
""").show()

# Запрос к конкретному снапшоту
spark.sql(f"""
    SELECT COUNT(*), SUM(amount)
    FROM lakehouse.silver.orders
    VERSION AS OF {valid_snapshot_id}
""").show()

# Сравниваем текущее состояние с состоянием вчера
yesterday_snapshot = spark.sql("""
    SELECT snapshot_id
    FROM lakehouse.silver.orders.snapshots
    WHERE committed_at < current_timestamp() - INTERVAL 1 DAY
    ORDER BY committed_at DESC
    LIMIT 1
""").collect()[0]["snapshot_id"]

current_count  = spark.table("lakehouse.silver.orders").count()
yesterday_count = spark.sql(f"""
    SELECT COUNT(*) AS cnt
    FROM lakehouse.silver.orders VERSION AS OF {yesterday_snapshot}
""").collect()[0]["cnt"]

print(f"Сегодня: {current_count} строк, вчера: {yesterday_count} строк, дельта: {current_count - yesterday_count}")

Time Travel - незаменимый инструмент при дебаггинге пайплайнов: «когда появились эти странные строки», «каков был COUNT(*) вчера в 10 утра», «что изменилось между двумя запусками».


5. Практика: полный сценарий с сбоем и восстановлением

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

spark = SparkSession.builder \
    .appName("Lakehouse-Constraints-Demo") \
    .config("spark.sql.extensions",
            "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
    .config("spark.sql.catalog.lakehouse",
            "org.apache.iceberg.spark.SparkCatalog") \
    .config("spark.sql.catalog.lakehouse.type", "hadoop") \
    .config("spark.sql.catalog.lakehouse.warehouse", "/tmp/warehouse") \
    .getOrCreate()

# ─── 1. Создаём таблицу с «информационным» PK ────────────────────────────────
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.silver.orders (
        order_id   STRING,
        amount     DECIMAL(10,2),
        status     STRING,
        updated_at TIMESTAMP
    ) USING iceberg
""")

# ─── 2. Начальная загрузка (идемпотентный MERGE) ──────────────────────────────
initial_data = [
    ("order-001", 500.0,  "CREATED",  "2026-05-23 10:00:00"),
    ("order-002", 1200.0, "PAID",     "2026-05-23 10:05:00"),
    ("order-003", 75.0,   "SHIPPED",  "2026-05-23 10:10:00"),
]
spark.createDataFrame(initial_data, ["order_id", "amount", "status", "updated_at"]) \
    .createOrReplaceTempView("initial_batch")

spark.sql("""
    MERGE INTO lakehouse.silver.orders AS target
    USING initial_batch AS source ON target.order_id = source.order_id
    WHEN MATCHED THEN UPDATE SET *
    WHEN NOT MATCHED THEN INSERT *
""")
print("После начальной загрузки:")
spark.table("lakehouse.silver.orders").show()

# Запоминаем ID хорошего снапшота
good_snapshot_id = spark.sql("""
    SELECT snapshot_id FROM lakehouse.silver.orders.snapshots
    ORDER BY committed_at DESC LIMIT 1
""").collect()[0]["snapshot_id"]
print(f"✅ Хороший снапшот: {good_snapshot_id}")

# ─── 3. Демонстрация: MERGE идемпотентен ──────────────────────────────────────
print("\nПовторный запуск MERGE с теми же данными:")
spark.sql("""
    MERGE INTO lakehouse.silver.orders AS target
    USING initial_batch AS source ON target.order_id = source.order_id
    WHEN MATCHED THEN UPDATE SET *
    WHEN NOT MATCHED THEN INSERT *
""")
count_after_retry = spark.table("lakehouse.silver.orders").count()
print(f"Количество строк после повторного запуска: {count_after_retry}")
# Должно остаться 3 - MERGE не создал дублей

# ─── 4. Демонстрация: «плохой» батч ───────────────────────────────────────────
print("\n⚠️  Имитируем плохой батч (случайные значения):")
spark.sql("""
    INSERT INTO lakehouse.silver.orders VALUES
    ('CORRUPTED-999', -9999.99, 'FATAL_ERROR', CURRENT_TIMESTAMP()),
    ('order-001', -1.0, 'BAD_OVERRIDE', CURRENT_TIMESTAMP())
""")
print("Таблица испорчена:")
spark.table("lakehouse.silver.orders").show()
print(f"Строк теперь: {spark.table('lakehouse.silver.orders').count()}")

# ─── 5. ROLLBACK к хорошему снапшоту ──────────────────────────────────────────
print(f"\n🔄 Rollback к снапшоту {good_snapshot_id}...")
spark.sql(f"""
    CALL lakehouse.system.rollback_to_snapshot(
        'silver.orders',
        {good_snapshot_id}
    )
""")
print("✅ Таблица восстановлена:")
spark.table("lakehouse.silver.orders").show()
print(f"Строк после rollback: {spark.table('lakehouse.silver.orders').count()}")
# Должно вернуться 3 строки без CORRUPTED и без BAD_OVERRIDE

# ─── 6. Time Travel: смотрим историю снапшотов ────────────────────────────────
print("\n📜 История снапшотов:")
spark.sql("""
    SELECT snapshot_id, committed_at, operation,
           summary['added-records']   AS added_records,
           summary['deleted-records'] AS deleted_records
    FROM lakehouse.silver.orders.snapshots
    ORDER BY committed_at
""").show(truncate=False)

6. Data Quality вместо Hard Constraints

6.1 Мониторинг уникальности через проверки

Поскольку Iceberg не enforce PK, инженер должен проверять уникальность самостоятельно через scheduled DQ-проверки:

def check_uniqueness(table_name: str, key_columns: list[str]) -> int:
    """Возвращает количество дублей по бизнес-ключу. 0 = норма."""
    key_expr = ", ".join(key_columns)
    result = spark.sql(f"""
        SELECT COUNT(*) - COUNT(DISTINCT {key_expr}) AS duplicate_count
        FROM {table_name}
    """).collect()[0]["duplicate_count"]
    return result

dups = check_uniqueness("lakehouse.silver.orders", ["order_id"])
if dups > 0:
    raise RuntimeError(f"❌ DQ FAIL: {dups} дублей по order_id в silver.orders!")
print("✅ DQ PASS: дублей нет")

6.2 Checklist дата-инженера по гарантии уникальности

Проверка Инструмент Когда
Дубли по бизнес-ключу COUNT(*) vs COUNT(DISTINCT key) После каждого батча
Нет NULL в ключевых колонках WHERE key IS NULL После парсинга
Pipeline идемпотентен Повторный запуск + сравнение COUNT При развёртывании
Rollback проверен rollback_to_snapshot в dev Перед production
Метрики свежести данных MAX(updated_at) > threshold Мониторинг SLA
Orphan файлы очищены remove_orphan_files Еженедельно
Старые снапшоты очищены expire_snapshots Еженедельно

7. Lakehouse vs OLTP: философия целостности данных

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


Итог

Lakehouse-архитектура работает не вопреки отсутствию классических constraints, а благодаря осознанному отказу от них. Это разумный инженерный компромисс: жертвуем автоматической проверкой уникальности ради горизонтального масштабирования на петабайтах.

Ответственность за целостность данных перекладывается на инженера:

  • Идемпотентные пайплайны вместо «система сама защитит» - MERGE INTO, Overwrite Partitions, или Append + Downstream Dedup
  • Snapshot Isolation вместо row-level locks - атомарные коммиты, Rollback за миллисекунды, Time Travel для диагностики
  • Scheduled DQ-проверки вместо database constraints - мониторим уникальность, свежесть, NULL-ы после каждого батча
  • Circuit Breaker перед Gold - не пускаем плохие данные в витрины через явные проверки с исключениями

Правило-маяк: «База данных не защитит. Твой код должен быть идемпотентным».