Silver Layer: дедупликация, CDC и гарантии уникальности без Primary Keys
ROW_NUMBER() + Window Functions для дедупликации, обработка CDC-потоков (Insert/Update/Delete), нормализация схемы, бизнес-ключи, Soft Delete и идемпотентный MERGE INTO на Apache Iceberg
Если Bronze - это дисциплинированный «склад сырого материала», то Silver - это производственный цех. Здесь происходит главная работа: хаотичный поток событий превращается в чистые, типизированные, уникальные бизнес-сущности. Именно Silver становится Single Source of Truth (Единым источником правды) для аналитиков, ML-инженеров и всех команд, которые потребляют данные из Lakehouse.
Урок технически насыщен - мы разберём самый сложный и самый важный паттерн в Lakehouse-архитектуре: как гарантировать уникальность в распределённой системе, где нет Primary Keys.
1. Философия Silver Layer¶
1.1 Переход от событийной к сущностной модели¶
Bronze хранит события - каждая строка отвечает на вопрос «что произошло?». Silver хранит сущности - каждая строка отвечает на вопрос «каково текущее (или историческое) состояние объекта?».
На Bronze уровне пять записей про один заказ - это норма и правильное поведение (append-only журнал). На Silver уровне пять записей про один заказ - это катастрофа, потому что аналитик получит пятикратно завышенные суммы в любом отчёте. Задача Silver Pipeline - схлопнуть хронологию в актуальное состояние.
1.2 Шок для OLTP-разработчиков: никаких Primary Keys¶
Первая реакция разработчика, пришедшего из PostgreSQL или MySQL в мир Lakehouse: «Хорошо, создам таблицу и поставлю PRIMARY KEY (order_id) - система сама не даст вставить дубль». Но так не работает.
В Apache Iceberg, Delta Lake, Apache Hudi и файловых форматах Parquet/ORC не существует механизма уникальных ограничений уровня СУБД. Никаких unique constraints, никаких index lookups при вставке. Физически это просто папки с Parquet-файлами на S3, и Spark при INSERT INTO просто добавляет новый файл рядом со старыми, не проверяя уникальность.
Это не баг форматов - это осознанный дизайн. Ограничения уникальности в распределённой системе требуют глобальной координации между всеми узлами при каждой записи, что убивает производительность и масштабируемость. Вместо этого ответственность за уникальность перекладывается на ETL-логику, которую пишет инженер. И эта логика строится вокруг одного ключевого паттерна - оконных функций с ROW_NUMBER().
1.3 Концепция бизнес-ключа¶
Прежде чем писать любой код дедупликации, нужно ответить на вопрос: что делает строку уникальной с точки зрения бизнеса?
Бизнес-ключ (Business Key / Natural Key) - это атрибут или комбинация атрибутов, которые однозначно идентифицируют сущность в предметной области. Он существует независимо от Lakehouse и технических систем.
| Сущность | Бизнес-ключ | Комментарий |
|---|---|---|
| Заказ | order_id |
Назначается системой e-commerce |
| Клиент | email или phone |
Зависит от домена |
| Продукт | sku |
Stock Keeping Unit |
| Транзакция | transaction_id |
Обычно UUID |
| Сессия | session_id + device_fingerprint |
Составной ключ |
Технические поля уровня Bronze (_event_id, _batch_id) - не бизнес-ключи. Они описывают факт доставки события в озеро, а не саму бизнес-сущность. На Silver уровне мы дедуплицируем именно по бизнес-ключу.
2. Великая дедупликация: ROW_NUMBER() + Window¶
2.1 Постановка задачи¶
Представим, что в Bronze лежат следующие записи про заказы:
| order_id | updated_at | status | cdc_op | _batch_id |
|---|---|---|---|---|
| order_1 | 10:00:00 | CREATED | c | batch-A |
| order_2 | 10:05:00 | PAID | c | batch-A |
| order_2 | 10:05:00 | PAID | c | batch-A |
| order_1 | 10:15:00 | SHIPPED | u | batch-B |
| order_3 | 09:55:00 | PAID | c | batch-C |
Здесь два типа «дублей»:
- Технические дубликаты - строки 2 и 3 абсолютно идентичны. Это результат сетевого сбоя: Kafka-продюсер не получил подтверждения и повторил отправку. Для Silver нам нужна ровно одна строка.
- Бизнес-дубликаты ключа - строки 1 и 4 имеют одинаковый
order_id, но разныеupdated_atиstatus. Это не ошибка - это легитимная история изменений заказа. Для Silver нам нужна самая свежая строка.
Обе ситуации решаются одним и тем же паттерном.
2.2 Анатомия паттерна ROW_NUMBER()¶
Оконная функция ROW_NUMBER() присваивает каждой строке порядковый номер внутри «окна» (группы строк с одинаковым ключом), упорядочивая строки по заданному критерию. Строка с номером 1 - это «победитель», которого мы сохраним; все остальные отфильтруем.
Разберём каждую часть функции:
PARTITION BY order_id - говорит Spark: «разбей все строки на независимые группы по полю order_id. Нумерацию начинай заново для каждой группы». Это аналог GROUP BY, но без агрегации - каждая строка сохраняется как есть.
ORDER BY updated_at DESC - внутри каждой группы отсортируй строки от свежей к старой. Строка с самым последним updated_at получит rn = 1. Именно это поле определяет «победителя» - в нём закодировано понятие «свежести» записи.
WHERE rn = 1 - финальный фильтр, оставляющий только победителей.
2.3 Выбор поля ORDER BY: критерий свежести¶
Это самое важное дизайнерское решение при построении дедупликации. Неправильный выбор приведёт к тому, что в Silver попадёт устаревшее состояние вместо актуального.
| Поле | Когда использовать | Риски |
|---|---|---|
updated_at из источника |
Источник проставляет бизнес-время изменения | Источник может отдавать неверное время |
_ingest_ts из Bronze |
Если у источника нет надёжного updated_at |
Не отражает порядок бизнес-событий, только порядок доставки |
| Kafka offset | Для streaming CDC через Debezium | Работает только внутри одной партиции |
version / seq_num |
Источник генерирует монотонный счётчик версий | Идеальный вариант, но встречается редко |
Когда строки имеют одинаковый updated_at (технические дубликаты), нужен второй критерий сортировки. Например, ORDER BY updated_at DESC, cdc_op DESC - при равном времени приоритет отдаётся UPDATE (u) перед CREATE (c), потому что UPDATE семантически «сильнее».
2.4 Под капотом Spark: почему это дорого¶
Оконные функции - это wide transformation (широкое преобразование). Чтобы посчитать ROW_NUMBER() по order_id, Spark должен собрать все строки с одинаковым order_id на одном Executor, выполнить сортировку и только потом присвоить номера.
Каждая стрелка через SHUFFLE - это сетевая передача данных между узлами кластера. Именно поэтому дедупликация - одна из самых дорогих операций в Silver Pipeline. Практические советы:
- Установите
spark.sql.shuffle.partitionsв разумное число (по умолчанию 200, но для большой таблицы может быть нужно 2000–4000). - Дедуплицируйте инкрементально (только новые данные из Bronze с момента последнего запуска), а не всю таблицу целиком.
- Следите за Data Skew - перекос по бизнес-ключу приводит к OOM на одном Executor (подробнее в разделе «Вопрос на защите»).
2.5 dropDuplicates() vs ROW_NUMBER(): когда что использовать¶
Spark предоставляет более простой API для удаления дублей:
df.dropDuplicates(["order_id"])
Но у него есть критический изъян: вы не контролируете, какая именно строка выживет. Spark выберет произвольную строку из группы дублей - ту, которую первой прочитает в рамках своего partition scan. Это недетерминированное поведение.
| Критерий | dropDuplicates() |
ROW_NUMBER() |
|---|---|---|
| Детерминированность | ❌ - случайный выбор | ✅ - выбор по ORDER BY |
| Производительность | Чуть быстрее | Чуть медленнее (sort) |
| Контроль «победителя» | Нет | Полный |
| CDC semantics | Не работает | Работает |
| Читаемость | Проще | Сложнее, но прозрачнее |
dropDuplicates() применим только для удаления полностью идентичных строк (все поля одинаковы), когда не важно, какую оставить. Для дедупликации по бизнес-ключу с выбором актуальной версии - только ROW_NUMBER().
3. Обработка CDC-потоков¶
3.1 Что такое CDC в контексте Silver¶
Change Data Capture (CDC) - это механизм отслеживания всех изменений в OLTP-источнике (обычно через чтение бинарного лога PostgreSQL WAL или MySQL binlog) и передачи этих изменений в Lakehouse в виде потока событий.
Инструменты вроде Debezium превращают каждую операцию в источнике в событие с метаданными:
{
"order_id": "order_1",
"status": "SHIPPED",
"updated_at": "2026-05-23T10:15:00Z",
"op": "u"
}
Поле op (operation) принимает значения:
c(create) - строка была создана (INSERT в источнике)u(update) - строка была изменена (UPDATE в источнике)d(delete) - строка была удалена (DELETE в источнике)r(read/snapshot) - начальная загрузка при первом запуске CDC
Bronze честно сохраняет всё это журналом - пять op=u для одного заказа значат, что заказ менял статус пять раз. Silver должен превратить этот журнал в актуальное состояние.
3.2 Мягкое удаление (Soft Delete): почему DELETE - это тоже строка¶
Когда из Bronze приходит событие op=d (удаление), возникает соблазн физически удалить строку из Silver. Не делайте этого.
Причины использовать Soft Delete (пометка is_deleted = true) вместо физического DELETE:
- Аналитическая целостность - агрегаты по историческим периодам не должны ломаться от того, что кто-то удалил запись сегодня. Финансовый отчёт за прошлый квартал должен показывать те же цифры независимо от текущего состояния данных.
- Downstream зависимости - Gold-таблицы и ML-модели могут зависеть от этих строк. Физическое удаление нарушит их согласованность.
- Аудит и соответствие требованиям - во многих индустриях (банки, медицина, e-commerce) требуется сохранять историю даже удалённых объектов.
- Сложность MERGE - физическое удаление требует отдельной ветки
WHEN MATCHED AND source.op='d' THEN DELETE, что усложняет логику и делает pipeline менее идемпотентным.
Единственное исключение - требования GDPR «право на забвение», когда нужно физически уничтожить персональные данные. Но даже тогда правильный подход - Crypto-Shredding (удаление ключа шифрования), а не DELETE FROM, как мы разбирали в уроке о Bronze Layer.
3.3 Реализация CDC-логики с флагами операций¶
Комбинация ROW_NUMBER() с полем op позволяет одним запросом:
- выбрать последнее событие для каждого
order_id - понять, что это за событие (создание, обновление или удаление)
- превратить его в строку Silver с правильным
is_deleted
3.4 Проблема поздних событий (Late-Arriving Data)¶
В распределённых системах события могут приходить с опозданием - сеть завалила пакеты, сервис-источник был недоступен, CDC-коннектор перезапускался. Silver-pipeline должен корректно обрабатывать ситуацию, когда старое событие приходит после новых.
Пример: Silver уже знает, что order_1 имеет статус SHIPPED (обновлено в 10:15). Но в следующем батче из Bronze приходит пропущенное событие PAID с timestamp 10:05. Если мы просто сделаем MERGE INTO по updated_at DESC, пропущенное событие проиграет и не перезапишет SHIPPED - это правильное поведение. Более старая запись не должна откатывать более новое состояние.
Именно поэтому ORDER BY updated_at DESC (а не ORDER BY _ingest_ts DESC) является правильным критерием - мы доверяем бизнес-времени события, а не времени его доставки в озеро.
4. Нормализация схемы¶
4.1 Отказ от raw_payload¶
Bronze хранит raw_payload как STRING (сырой JSON). Silver - это место, где JSON разбирается, типы приводятся, поля именуются согласно корпоративному стандарту.
from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField, StringType, TimestampType, DecimalType
# Схема ожидаемого payload
payload_schema = StructType([
StructField("order_id", StringType(), nullable=False),
StructField("user_id", StringType(), nullable=False),
StructField("amount", DecimalType(18, 2), nullable=True),
StructField("currency", StringType(), nullable=True),
StructField("status", StringType(), nullable=True),
StructField("updated_at", TimestampType(), nullable=True),
])
# Разбираем JSON и явно приводим типы
parsed_df = bronze_df.select(
F.from_json(F.col("raw_payload"), payload_schema).alias("payload"),
F.col("_ingest_ts"),
F.col("_source_system"),
)
normalized_df = parsed_df.select(
F.col("payload.order_id").alias("order_id"),
F.col("payload.user_id").alias("user_id"),
F.col("payload.amount").cast("decimal(18,2)").alias("amount"),
F.coalesce(F.col("payload.currency"), F.lit("USD")).alias("currency"),
F.col("payload.status").alias("status"),
F.col("payload.updated_at").alias("updated_at"),
F.col("_ingest_ts"),
F.col("_source_system"),
)
F.from_json() принимает строку и схему, возвращает StructType. Если JSON не соответствует схеме - поле станет null, а не вызовет исключение. Это позволяет продолжить обработку батча, направив кривые записи в quarantine table.
F.coalesce() - это NULL-safe выбор первого ненулевого значения. Если currency пришёл как null, берём дефолт "USD". Это safer, чем fillna(), который работает только со строками.
4.2 Schema Evolution: что делать, когда источник меняет схему¶
Источник добавил новое поле delivery_address - что происходит с Silver?
Apache Iceberg поддерживает schema evolution «из коробки»: новые поля можно добавить командой ALTER TABLE silver.orders ADD COLUMN delivery_address STRING. Все старые Parquet-файлы, созданные до изменения схемы, при чтении вернут null для нового поля - Iceberg знает маппинг «поле X появилось в схеме версии N».
Опасный сценарий - изменение типа поля (type widening). Источник изменил amount с INT на DECIMAL(18,2). Iceberg позволяет некоторые безопасные преобразования (INT → LONG, FLOAT → DOUBLE), но опасные (STRING → INT) - нет. В таких случаях Silver Pipeline должен явно обрабатывать оба варианта через TRY_CAST или временную колонку.
4.3 Quarantine таблица для невалидных записей¶
Вместо того чтобы падать при неожиданных данных, современные Silver Pipelines направляют проблемные записи в отдельную quarantine (карантинную) таблицу:
# Определяем правило валидации: amount не может быть отрицательным
valid_df = normalized_df.filter(F.col("amount") >= 0)
invalid_df = normalized_df.filter(F.col("amount") < 0) \
.withColumn("error_reason", F.lit("negative_amount")) \
.withColumn("quarantine_ts", F.current_timestamp())
# Валидные записи идут в Silver
valid_df.writeTo("lakehouse.silver.orders").append()
# Невалидные - в карантин для ручного разбора
invalid_df.writeTo("lakehouse.silver.orders_quarantine").append()
Карантинная таблица - не мусорный бак. Её мониторят, по ней строят алерты, на неё реагируют. Стабильно ненулевой карантин сигнализирует о проблеме в источнике или в Silver-логике.
5. Практика: пишем Silver Pipeline¶
5.1 Полная архитектура пайплайна¶
5.2 Полный код Silver Pipeline¶
import uuid
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.window import Window
from pyspark.sql.types import (
StructType, StructField, StringType, TimestampType, DecimalType
)
spark = SparkSession.builder \
.appName("Silver-Orders-Pipeline") \
.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") \
.config("spark.sql.shuffle.partitions", "400") # тюним под объём данных
.getOrCreate()
# ─── 0. Параметры запуска ────────────────────────────────────────────────────
run_id = str(uuid.uuid4())
last_run_ts = spark.sql(
"SELECT COALESCE(MAX(_silver_updated_ts), TIMESTAMP '1970-01-01') "
"FROM lakehouse.silver.orders"
).collect()[0][0]
# ─── 1. Создаём целевую Silver-таблицу (идемпотентно) ───────────────────────
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.silver.orders (
order_id STRING COMMENT 'Бизнес-ключ заказа',
user_id STRING COMMENT 'Идентификатор пользователя',
amount DECIMAL(18,2),
currency STRING,
status STRING,
is_deleted BOOLEAN COMMENT 'Soft Delete флаг (op=d из CDC)',
updated_at TIMESTAMP COMMENT 'Бизнес-время последнего изменения',
_silver_updated_ts TIMESTAMP COMMENT 'Когда Silver Pipeline обновил эту строку',
_run_id STRING COMMENT 'UUID запуска Silver Pipeline'
) USING iceberg
PARTITIONED BY (days(updated_at))
""")
# ─── 2. Схема payload ────────────────────────────────────────────────────────
payload_schema = StructType([
StructField("order_id", StringType(), nullable=False),
StructField("user_id", StringType(), nullable=False),
StructField("amount", DecimalType(18, 2), nullable=True),
StructField("currency", StringType(), nullable=True),
StructField("status", StringType(), nullable=True),
StructField("updated_at", TimestampType(), nullable=True),
StructField("op", StringType(), nullable=True),
])
# ─── 3. Инкрементальное чтение из Bronze ─────────────────────────────────────
# Читаем только новые записи с момента предыдущего успешного запуска
bronze_raw = spark.table("lakehouse.bronze.orders") \
.filter(F.col("_ingest_ts") > F.lit(last_run_ts))
# ─── 4. Парсинг и нормализация ───────────────────────────────────────────────
parsed = bronze_raw.select(
F.from_json(F.col("raw_payload"), payload_schema).alias("p"),
F.col("_ingest_ts"),
)
normalized = parsed.select(
F.col("p.order_id").alias("order_id"),
F.col("p.user_id").alias("user_id"),
F.col("p.amount").cast("decimal(18,2)").alias("amount"),
F.coalesce(F.col("p.currency"), F.lit("USD")).alias("currency"),
F.col("p.status").alias("status"),
F.col("p.updated_at").alias("updated_at"),
F.col("p.op").alias("cdc_op"),
F.col("_ingest_ts"),
)
# Невалидные записи - в карантин
valid_df = normalized.filter(F.col("order_id").isNotNull())
invalid_df = normalized.filter(F.col("order_id").isNull()) \
.withColumn("error_reason", F.lit("null_order_id")) \
.withColumn("quarantine_ts", F.current_timestamp())
if invalid_df.count() > 0:
invalid_df.writeTo("lakehouse.silver.orders_quarantine").append()
# ─── 5. Дедупликация через ROW_NUMBER() ──────────────────────────────────────
# Для каждого order_id выбираем строку с самым поздним updated_at.
# При равенстве updated_at приоритет отдаётся операции 'd' (delete) > 'u' > 'c'.
window_spec = Window \
.partitionBy("order_id") \
.orderBy(F.col("updated_at").desc(), F.col("cdc_op").desc())
deduped = valid_df \
.withColumn("rn", F.row_number().over(window_spec)) \
.filter(F.col("rn") == 1) \
.drop("rn")
# ─── 6. Добавляем Silver-метаданные ──────────────────────────────────────────
silver_updates = deduped.select(
F.col("order_id"),
F.col("user_id"),
F.col("amount"),
F.col("currency"),
F.col("status"),
F.when(F.col("cdc_op") == "d", F.lit(True)).otherwise(F.lit(False))
.alias("is_deleted"),
F.col("updated_at"),
F.current_timestamp().alias("_silver_updated_ts"),
F.lit(run_id).alias("_run_id"),
)
# ─── 7. Идемпотентный MERGE INTO Silver ──────────────────────────────────────
silver_updates.createOrReplaceTempView("silver_updates_view")
spark.sql("""
MERGE INTO lakehouse.silver.orders AS target
USING silver_updates_view AS source
ON target.order_id = source.order_id
WHEN MATCHED THEN
UPDATE SET
target.user_id = source.user_id,
target.amount = source.amount,
target.currency = source.currency,
target.status = source.status,
target.is_deleted = source.is_deleted,
target.updated_at = source.updated_at,
target._silver_updated_ts = source._silver_updated_ts,
target._run_id = source._run_id
WHEN NOT MATCHED THEN
INSERT *
""")
print(f"Silver Pipeline завершён. Run ID: {run_id}")
spark.table("lakehouse.silver.orders").show(truncate=False)
5.3 Разбор плана выполнения¶
После запуска дедупликации полезно изучить физический план, чтобы убедиться, что Spark правильно организует Shuffle:
silver_updates.explain(mode="formatted")
В выводе explain() ищите следующие узлы:
Exchange hashpartitioning(order_id, N)- Shuffle-шаг. Spark перераспределяет данные поorder_id. Если N (количество партиций) слишком мало, некоторые Executors получат огромные партиции. Если слишком велико - много пустых задач.Sort [order_id ASC, updated_at DESC]- локальная сортировка внутри каждой shuffle-партиции перед расчётом окна.Window [row_number() windowspecdefinition(order_id, ...)]- сам расчётROW_NUMBER().Filter (rn = 1)- фильтрация, которую Spark пытается протолкнуть вниз по плану (predicate pushdown), но с оконными функциями это не всегда работает.
6. Проблема Data Skew и как с ней бороться¶
6.1 Что такое Data Skew¶
Data Skew (перекос данных) - ситуация, когда одно значение бизнес-ключа встречается в данных несравнимо чаще остальных. При дедупликации это приводит к тому, что один Executor получает всю «горячую» часть данных, пока остальные простаивают.
Сценарий из реальной жизни: в e-commerce системе у «технических» заказов (тестовые заказы, internal orders, возвраты через call-center) есть специальный order_id = 'INTERNAL'. Из-за особенностей интеграции 30% всех событий в Bronze имеют именно этот идентификатор.
6.2 Решения для Data Skew¶
Вариант 1: Предварительная фильтрация. Если «горячий» ключ - мусорный (например, INTERNAL заказы не нужны в Silver), просто исключите его до дедупликации:
valid_df = normalized.filter(F.col("order_id") != "INTERNAL")
Вариант 2: Salting (соление). Если горячий ключ нужен в Silver, «распределите» его искусственно:
import random
# Добавляем случайный суффикс к горячему ключу для shuffle
salt_df = normalized.withColumn(
"shuffle_key",
F.when(
F.col("order_id") == "INTERNAL",
F.concat(F.col("order_id"), F.lit("_"), (F.rand() * 10).cast("int").cast("string"))
).otherwise(F.col("order_id"))
)
# Дедупликация по shuffle_key вместо order_id
window_salted = Window.partitionBy("shuffle_key").orderBy(F.col("updated_at").desc())
deduped_salted = salt_df \
.withColumn("rn", F.row_number().over(window_salted)) \
.filter(F.col("rn") == 1) \
.drop("rn", "shuffle_key")
Вариант 3: AQE (Adaptive Query Execution). В Spark 3.x+ включите адаптивное выполнение запросов:
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
AQE автоматически обнаруживает перекошенные партиции во время выполнения и разбивает их на подпартиции, распределяя нагрузку по нескольким Executors.
7. Анти-паттерны Silver Layer¶
INSERT INTO вместо MERGE INTO - самый частый и разрушительный анти-паттерн. Каждый перезапуск pipeline добавляет новые строки без удаления старых. Через неделю в Silver будет в 7 раз больше строк, чем должно быть, и COUNT(*) vs COUNT(DISTINCT order_id) покажет разницу.
Размещение агрегатов в Silver - Silver должен быть слоем сущностей, а не витрин. Если в Silver появляется SUM(amount) GROUP BY user_id, это нужно перенести в Gold. Silver - это «что есть», Gold - это «что произошло в итоге».
Игнорирование инкрементальности - перечитывать всю Bronze при каждом запуске Silver Pipeline неэффективно. Правильно: хранить last_run_ts (как в нашем примере) и читать только новые данные.
8. Checklist: проверяем Silver Layer на качество¶
| Проверка | SQL | Ожидаемый результат |
|---|---|---|
| Нет дублей по бизнес-ключу | SELECT COUNT(*) - COUNT(DISTINCT order_id) FROM silver.orders |
0 |
| Нет строк без бизнес-ключа | SELECT COUNT(*) FROM silver.orders WHERE order_id IS NULL |
0 |
| Нет строк с некорректной суммой | SELECT COUNT(*) FROM silver.orders WHERE amount < 0 |
0 |
| Актуальность (свежесть данных) | SELECT MAX(updated_at) FROM silver.orders |
Не старше порога |
| Карантин пуст | SELECT COUNT(*) FROM silver.orders_quarantine |
0 (алерт если > 0) |
| MERGE работает идемпотентно | Повторный запуск pipeline → COUNT(*) не меняется | Одинаково |
Последняя проверка - самая важная. Запустите Silver Pipeline дважды подряд на одних и тех же данных Bronze. Количество строк и их содержимое в Silver должно остаться неизменным. Если это не так - pipeline не идемпотентен и его нельзя безопасно перезапускать при сбоях.
Итог¶
Silver Layer - самый алгоритмически нагруженный слой Medallion Architecture. Здесь инженер берёт на себя ответственность за гарантии, которые в OLTP-мире обеспечивает СУБД: уникальность строк, консистентность типов, корректная обработка удалений.
Ключевые решения, принятые на этом уровне, определяют надёжность всей аналитической платформы:
ROW_NUMBER()+ORDER BY updated_at DESC- единственный детерминированный способ дедупликации с контролем «победителя»is_deleted = trueвместо физического DELETE - сохраняет аудит и не ломает downstream зависимостиMERGE INTOвместо INSERT - гарантирует идемпотентность при повторных запусках- Инкрементальное чтение по
_ingest_ts- не читаем всю Bronze при каждом запуске - Quarantine table - не падаем при кривых данных, но и не теряем их
В следующем уроке разберём Gold Layer - как из чистых Silver-сущностей строить быстрые аналитические агрегаты и витрины для конечных потребителей.