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 за миллисекунды
Разработчик, пришедший в мир 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 Как достигается идемпотентность¶
Три стратегии для идемпотентного пайплайна:
- MERGE INTO (Write-Time Deduplication) - при записи находим совпадения по ключу и обновляем вместо добавления
- Overwrite партиции - полностью перезаписываем партицию за конкретную дату, не добавляя к ней
- 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 - не пускаем плохие данные в витрины через явные проверки с исключениями
Правило-маяк: «База данных не защитит. Твой код должен быть идемпотентным».