Silver Layer: дедупликация, CDC и гарантии уникальности без Primary Keys

ROW_NUMBER() + Window Functions для дедупликации, обработка CDC-потоков (Insert/Update/Delete), нормализация схемы, бизнес-ключи, Soft Delete и идемпотентный MERGE INTO на Apache Iceberg

lakehouse deduplication cdc 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:

  1. Аналитическая целостность - агрегаты по историческим периодам не должны ломаться от того, что кто-то удалил запись сегодня. Финансовый отчёт за прошлый квартал должен показывать те же цифры независимо от текущего состояния данных.
  2. Downstream зависимости - Gold-таблицы и ML-модели могут зависеть от этих строк. Физическое удаление нарушит их согласованность.
  3. Аудит и соответствие требованиям - во многих индустриях (банки, медицина, e-commerce) требуется сохранять историю даже удалённых объектов.
  4. Сложность 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-сущностей строить быстрые аналитические агрегаты и витрины для конечных потребителей.