Medallion Architecture: bronze, silver, gold - принципы и границы слоёв
Medallion Architecture: bronze, silver, gold - принципы и границы слоёв
Каждый Data Engineer на определённом этапе карьеры сталкивается с одной и той же проблемой: данные хранятся, но никто не знает, можно ли им доверять. Таблицы дублируются, схемы разъезжаются, аналитики используют разные «правильные» версии одного и того же показателя. Озеро данных (Data Lake), которое должно было стать спасением от ограничений Data Warehouse, превращается в Data Swamp - непроходимое болото, где никто не может найти нужные данные, а те, что находят, вызывают сомнения.
Medallion Architecture - это архитектурный ответ на эту проблему. Она не меняет технологии, не требует перехода на новую платформу. Она вводит дисциплину: каждый байт данных проходит через три чётко разграниченных слоя с понятными правилами, обязанностями и контрактами качества. Именно эта дисциплина отличает промышленный Data Lakehouse от хаоса набора Parquet-файлов на S3.
1. Почему классические подходы перестали работать¶
1.1 Эволюция от ETL к ELT и Lakehouse¶
Data Warehouse (1990-е - 2010-е) был построен вокруг концепции ETL: данные трансформировались до загрузки в хранилище. Это обеспечивало контроль качества, но создавало критические ограничения: схема должна быть известна заранее, изменения дорогостоящи, хранилище не масштабируется горизонтально, а хранить «сырые» данные слишком дорого.
Data Lake появился как ответ: сначала храним всё (ELT - Extract, Load, Transform), трансформируем потом. S3 и HDFS сделали хранение дешёвым. Но без архитектурной дисциплины Data Lake превращался в кладбище Parquet-файлов. Никаких транзакций, никакого управления схемами, никакого контроля качества.
Data Lakehouse объединил лучшее из обоих миров:
- Дешёвое объектное хранилище (как Data Lake)
- ACID-транзакции, управление схемами, Time Travel (как Data Warehouse)
- Горизонтальное масштабирование Spark для обработки (как Data Lake)
Apache Iceberg, Delta Lake и Apache Hudi - это table formats, которые добавляют транзакционный слой поверх обычных Parquet-файлов на S3 или HDFS.
Диаграмма показывает эволюцию: Data Warehouse обеспечивал качество, но не масштабировался. Data Lake масштабировался, но терял качество. Data Lakehouse объединяет обе характеристики за счёт table formats нового поколения.
1.2 Почему Data Lake становится Data Swamp¶
Без архитектурной дисциплины происходит следующее:
- Нет единого источника истины: команды создают «свои» версии таблицы
orders_final,orders_clean_v2,orders_new_20240315 - Нет контроля схем: поставщик изменил формат JSON, и пайплайн молча пишет
nullво все поля - Нет идемпотентности: повторный запуск пайплайна дублирует строки
- Нет lineage: непонятно, откуда взялось значение в конкретной колонке
- Нет разделения ответственности: трансформации перемешаны с ingestion, бизнес-логика живёт в скрипте загрузки
Medallion Architecture решает каждую из этих проблем через разделение слоёв с чёткими контрактами.
1.3 Три слоя, три контракта¶
Контракт Bronze: «Я храню данные точно такими, какими они пришли из источника. Я никогда ничего не меняю и не удаляю. Если данные пришли сломанными - я их всё равно сохраняю.»
Контракт Silver: «Я предоставляю чистые, дедуплицированные, типизированные данные. Аналитик может доверять каждой строке. Я не занимаюсь бизнес-агрегатами.»
Контракт Gold: «Я оптимизирован под конкретный use-case: BI-дашборд, ML-модель, или API. Я содержу предвычисленные метрики и агрегаты.»
2. Bronze Layer: принцип «храни всё, меняй ничего»¶
2.1 Назначение и философия Bronze¶
Bronze - это иммутабельный архив источника. Его единственная задача - сохранить данные как можно ближе к оригиналу, максимально быстро и без потерь. Слой Bronze отвечает на вопрос: «Что именно пришло из источника в момент времени T?»
Ключевые принципы Bronze:
- Append-only: данные только добавляются, никогда не удаляются и не изменяются
- No business logic: никаких бизнес-правил, никаких JOIN-ов с другими таблицами
- Technical metadata: к каждой записи добавляются технические поля - время загрузки, имя файла-источника, ID запуска пайплайна
- Schema preservation: сохраняем оригинальную схему источника, даже если она «плохая»
- Reproducibility: из Bronze всегда можно воспроизвести любой последующий слой
Стоимость хранения на S3 или MinIO составляет около $23 за терабайт в месяц. Это несравнимо дешевле стоимости потери данных или невозможности пересчитать Silver после обнаружения бага в трансформации. Именно поэтому хранить всё в Bronze навсегда (или с длинным TTL) - это не расточительство, а инженерная осторожность.
2.2 Что добавляется при записи в Bronze¶
Технические метаданные - это «паспорт» каждой строки в Bronze. Они позволяют ответить на вопросы:
_ingested_at- когда строка попала в систему (время обработки, а не время события)_source_file- из какого файла или Kafka-топика пришла запись_pipeline_run_id- какой запуск пайплайна её создал (для отладки и повторной обработки)_source_system- из какой исходной системы (PostgreSQL, Salesforce, Kafka)
Эти поля добавляются автоматически в момент записи и никогда не приходят из источника.
2.3 PySpark: чтение из источника и запись в Bronze¶
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = SparkSession.builder \
.appName("Bronze-Ingestion") \
.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()
# Создаём Bronze-таблицу с Iceberg
# PARTITIONED BY (ingestion_date) - партиционирование по дате загрузки
# для быстрого чтения инкрементальных обновлений
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.bronze.orders (
order_id STRING,
customer_id STRING,
status STRING,
amount STRING, -- намеренно STRING: Bronze хранит оригинал
created_at STRING, -- тоже STRING: не делаем cast на Bronze
raw_payload STRING, -- полный оригинальный JSON (для отладки)
_ingested_at TIMESTAMP COMMENT 'Время загрузки в Bronze',
_source_file STRING COMMENT 'Путь к исходному файлу',
_pipeline_run_id STRING COMMENT 'ID запуска пайплайна',
_source_system STRING COMMENT 'Имя системы-источника',
ingestion_date DATE COMMENT 'Партиция по дате загрузки'
)
USING iceberg
PARTITIONED BY (ingestion_date)
""")
# Имитируем «грязные» сырые данные из источника
# Обратите внимание: amount - строка, дата в нестандартном формате,
# есть дубликат (order_id=1001 пришёл дважды), одна запись явно битая
raw_orders_data = [
("1001", "cust_A", "confirmed", "1500.00", "2026-05-23 10:00:00", None),
("1001", "cust_A", "confirmed", "1500.00", "2026-05-23 10:00:00", None), # дубликат
("1002", "cust_B", "pending", "750.50", "23-05-2026 11:30:00", None), # другой формат даты
("1003", None, "shipped", "N/A", "2026-05-23 12:00:00", None), # битая запись
("1004", "cust_C", "cancelled", "-200.00", "2026-05-23 13:00:00", None),
]
raw_df = spark.createDataFrame(
raw_orders_data,
["order_id", "customer_id", "status", "amount", "created_at", "raw_payload"]
)
# Добавляем технические метаданные - ЕДИНСТВЕННАЯ трансформация на Bronze
bronze_df = raw_df \
.withColumn("raw_payload",
F.to_json(F.struct([F.col(c) for c in raw_df.columns]))) \
.withColumn("_ingested_at", F.current_timestamp()) \
.withColumn("_source_file", F.lit("s3://raw/orders/2026-05-23.json")) \
.withColumn("_pipeline_run_id", F.lit("run-20260523-001")) \
.withColumn("_source_system", F.lit("orders_db")) \
.withColumn("ingestion_date", F.current_date())
# APPEND - единственный режим записи на Bronze
# Дубликаты, битые строки - всё пишем как есть
bronze_df.writeTo("lakehouse.bronze.orders").append()
print(f"Записано в Bronze: {bronze_df.count()} строк (включая дубликаты и битые)")
spark.table("lakehouse.bronze.orders").show(truncate=False)
Обратите внимание: мы намеренно оставляем amount как STRING, дату в нестандартном формате 23-05-2026, дубликат и битую запись с customer_id=null. Bronze не занимается исправлением. Если мы уберём битые строки на Bronze, мы навсегда потеряем информацию о том, что источник когда-то присылал некорректные данные. Это важно для аудита и отладки.
2.4 Schema Drift: что делать, когда источник меняет схему¶
Одна из самых частых проблем в продакшн-пайплайнах - schema drift: поставщик без предупреждения добавляет новые поля, переименовывает колонки или меняет типы. Bronze должен справляться с этим автоматически.
Iceberg поддерживает Schema Evolution: добавление новых колонок, изменение типов (с ограничениями), переименование. При чтении старых данных новые колонки возвращают null - обратная совместимость сохраняется.
# Пример: источник добавил новое поле discount_code
# Bronze просто принимает расширенную схему
new_data_with_extra_field = spark.createDataFrame([
("1005", "cust_D", "confirmed", "2000.00", "2026-05-24 09:00:00", "SALE20"),
], ["order_id", "customer_id", "status", "amount", "created_at", "discount_code"])
# Iceberg автоматически добавит колонку discount_code в таблицу
new_bronze_df = new_data_with_extra_field \
.withColumn("raw_payload", F.to_json(F.struct([F.col(c)
for c in new_data_with_extra_field.columns]))) \
.withColumn("_ingested_at", F.current_timestamp()) \
.withColumn("_source_file", F.lit("s3://raw/orders/2026-05-24.json")) \
.withColumn("_pipeline_run_id", F.lit("run-20260524-001")) \
.withColumn("_source_system", F.lit("orders_db")) \
.withColumn("ingestion_date", F.current_date())
# mergeSchema=True позволяет добавлять новые колонки автоматически
new_bronze_df.writeTo("lakehouse.bronze.orders") \
.option("mergeSchema", "true") \
.append()
После этой операции Bronze-таблица содержит колонку discount_code. Для старых строк (записанных до появления этого поля) она будет null. Silver-слой, читая Bronze, увидит это изменение и сможет обработать его в контролируемый момент - не при ingestion, а при трансформации.
2.5 Что нельзя делать на Bronze¶
Почему запрещено удалять «плохие» строки на Bronze? Если запись с customer_id=null пришла из источника - это важная информация. Может оказаться, что через месяц будет обнаружен баг в источнике, и эти строки нужно будет переобработать с правильными ключами. Если мы удалили их при ingestion, эта возможность потеряна навсегда.
Почему запрещены агрегации? Агрегация разрушает грануляцию данных. Если мы агрегируем заказы по дням на Bronze, мы теряем способность воспроизвести почасовую аналитику или расследовать конкретный заказ.
3. Silver Layer: слой качества и доверия¶
3.1 Назначение Silver: от «сырых фактов» к «доверенным сущностям»¶
Silver - это центральный слой обработки в Medallion Architecture. Именно здесь данные трансформируются из технического артефакта (Bronze) в бизнес-сущность, которой доверяют аналитики и инженеры.
Аналогия: если Bronze - это необработанный металл из шахты, то Silver - это очищенный слиток, который проверен на чистоту и соответствует стандартам. Из него можно делать продукты (Gold), но сам по себе он уже имеет ценность.
Ключевые задачи Silver:
- Дедупликация: устранение дублей, неизбежных при at-least-once delivery
- Типизация:
amount: STRING "1500.00"→amount: DECIMAL(10,2) 1500.00 - Стандартизация: единые форматы дат, валют, кодировок
- Null-handling: стратегия обработки пропущенных значений (fill, drop, quarantine)
- Enrichment: обогащение данными из справочников (например, добавление имени города по
city_id) - SCD-логика: обновление медленно изменяющихся измерений
- Quarantine: помещение некорректных записей в карантинную зону вместо удаления
3.2 Дедупликация: два подхода¶
На Silver данные из Bronze могут содержать дубликаты - как из-за повторной отправки источником, так и из-за нескольких запусков пайплайна. Дедупликация - первая операция в Silver-пайплайне.
from pyspark.sql.window import Window
# Читаем из Bronze - берём только новые данные за нужную дату
bronze_df = spark.table("lakehouse.bronze.orders") \
.filter(F.col("ingestion_date") == F.current_date())
# Подход 1: dropDuplicates - простой, но не всегда точный
# Хорошо работает, когда дубликат - это точная копия строки
dedup_simple = bronze_df.dropDuplicates(["order_id"])
# Подход 2: ROW_NUMBER по бизнес-ключу с сортировкой по времени
# Лучше: выбираем «самую позднюю» версию по _ingested_at
window_spec = Window.partitionBy("order_id").orderBy(F.col("_ingested_at").desc())
dedup_advanced = (
bronze_df
.withColumn("rn", F.row_number().over(window_spec))
.filter(F.col("rn") == 1)
.drop("rn")
)
print("Строк в Bronze (с дубликатами):", bronze_df.count())
print("Строк после дедупликации:", dedup_advanced.count())
Второй подход предпочтительнее в production: если источник прислал исправленную версию той же записи (update, а не точный дубликат), ROW_NUMBER выберет более позднюю версию.
3.3 Типизация, стандартизация и Quarantine¶
# Разделяем записи на валидные и невалидные ПЕРЕД трансформацией
# Это паттерн "fail soft" - не падаем на ошибке, а помещаем в карантин
# Условия валидности бизнес-правил Silver-слоя
valid_condition = (
F.col("order_id").isNotNull() &
F.col("customer_id").isNotNull() &
F.col("amount").rlike(r"^-?\d+(\.\d+)?$") # Число (можно отрицательное)
)
valid_records = dedup_advanced.filter(valid_condition)
invalid_records = dedup_advanced.filter(~valid_condition)
# Карантинная зона - сохраняем для расследования, не удаляем
invalid_records \
.withColumn("rejection_reason",
F.when(F.col("order_id").isNull(), F.lit("null_order_id"))
.when(F.col("customer_id").isNull(), F.lit("null_customer_id"))
.when(~F.col("amount").rlike(r"^-?\d+(\.\d+)?$"),
F.lit("invalid_amount"))
.otherwise(F.lit("unknown"))) \
.withColumn("quarantined_at", F.current_timestamp()) \
.writeTo("lakehouse.silver.orders_quarantine") \
.append()
# Типизация и стандартизация валидных записей
silver_df = (
valid_records
.withColumn("order_id", F.col("order_id").cast("integer"))
.withColumn("amount", F.col("amount").cast("decimal(10,2)"))
# Парсим оба формата даты: "2026-05-23 10:00:00" и "23-05-2026 11:30:00"
.withColumn("created_at",
F.coalesce(
F.to_timestamp("created_at", "yyyy-MM-dd HH:mm:ss"),
F.to_timestamp("created_at", "dd-MM-yyyy HH:mm:ss"),
)
)
# Стандартизация статусов: приводим к нижнему регистру
.withColumn("status", F.lower(F.trim(F.col("status"))))
# Отрицательный amount для cancelled-заказов - бизнес-аномалия, ставим 0
.withColumn("amount",
F.when(
(F.col("status") == "cancelled") & (F.col("amount") < 0),
F.lit(0).cast("decimal(10,2)")
).otherwise(F.col("amount"))
)
# Добавляем метаданные обработки Silver
.withColumn("_silver_processed_at", F.current_timestamp())
.select(
"order_id", "customer_id", "status", "amount", "created_at",
"_silver_processed_at"
)
)
print(f"Валидных строк: {silver_df.count()}, в карантине: {invalid_records.count()}")
Паттерн Quarantine критически важен. Вместо того чтобы падать с исключением или молча удалять проблемные строки, мы помещаем их в отдельную таблицу с причиной отклонения. Data Engineer может периодически проверять карантин, исправлять проблемы на источнике и принимать решения о повторной обработке.
3.4 Запись в Silver через MERGE INTO: идемпотентность¶
Запись в Silver должна быть идемпотентной - повторный запуск пайплайна не должен дублировать данные. Для этого используется MERGE INTO (UPSERT): если строка с таким order_id уже существует - обновляем, если нет - вставляем.
# Регистрируем Silver-данные как временное представление для Spark SQL
silver_df.createOrReplaceTempView("incoming_silver_orders")
# Создаём Silver-таблицу если не существует
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.silver.orders (
order_id INTEGER,
customer_id STRING,
status STRING,
amount DECIMAL(10,2),
created_at TIMESTAMP,
_silver_processed_at TIMESTAMP
)
USING iceberg
PARTITIONED BY (month(created_at))
""")
# MERGE INTO: идемпотентный UPSERT
# При повторном запуске order_id=1001 обновится (не продублируется)
spark.sql("""
MERGE INTO lakehouse.silver.orders AS target
USING incoming_silver_orders AS source
ON target.order_id = source.order_id
WHEN MATCHED AND (
target.status != source.status OR
target.amount != source.amount
) THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
""")
print("Silver обновлён. Финальное состояние:")
spark.table("lakehouse.silver.orders").show()
MERGE INTO в Iceberg (и Delta Lake) реализован через Copy-on-Write: Spark находит все data files, содержащие затронутые строки, перечитывает их, применяет изменения и записывает новые файлы. Старые файлы помечаются как удалённые (но физически не удаляются - это обеспечивает Time Travel).
3.5 Enrichment: обогащение данными из справочников¶
Обогащение (Enrichment) - это добавление атрибутов из справочных таблиц, которые не пришли из источника. Например, к заказу добавляем имя клиента из Silver-таблицы клиентов:
# Silver-таблица клиентов (предположим, что уже обработана)
customers_df = spark.table("lakehouse.silver.customers") \
.select("customer_id", "customer_name", "city", "segment")
# Enrichment через LEFT JOIN
# LEFT JOIN: заказы без совпадающего клиента всё равно остаются (не теряются)
enriched_orders = (
spark.table("lakehouse.silver.orders")
.join(
F.broadcast(customers_df), # BROADCAST hint: customers_df маленький
on="customer_id",
how="left"
)
)
# Результат: каждый заказ теперь содержит customer_name, city, segment
F.broadcast() - оптимизационный hint для Spark: вместо SortMergeJoin (требует Shuffle) Spark будет использовать BroadcastHashJoin (таблица клиентов рассылается на все executors). Это в разы быстрее для небольших справочников (< 200 МБ).
4. Gold Layer: слой потребления и бизнес-ценности¶
4.1 Назначение Gold: данные, готовые к принятию решений¶
Gold - это финальный слой Medallion Architecture. Если Silver - это «чистый металл», то Gold - это «ювелирное изделие», созданное под конкретную задачу. Каждая Gold-таблица оптимизирована под один конкретный use-case: BI-дашборд, ML-модель, API-ответ, executive-отчёт.
Ключевые характеристики Gold:
- Высокая степень агрегации: не отдельные события, а метрики и показатели
- Бизнес-логика: расчёт KPI, применение бизнес-правил компании
- Оптимизация чтения: партиционирование и Z-Ordering под паттерн запросов
- Денормализация: wide tables для быстрых аналитических запросов без дополнительных JOIN
- Обновляемость: Gold-таблицы регулярно пересчитываются (ежедневно, ежечасно)
4.2 Построение аналитической витрины¶
# Gold-витрина: ежедневная выручка по сегменту клиентов
# Потребитель: BI-дашборд для коммерческого директора
daily_revenue_mart = (
enriched_orders
.filter(F.col("status") == "confirmed")
.withColumn("order_date", F.to_date("created_at"))
.groupBy("order_date", "segment")
.agg(
F.sum("amount").alias("total_revenue"),
F.count("order_id").alias("orders_count"),
F.avg("amount").alias("avg_order_value"),
F.countDistinct("customer_id").alias("unique_customers"),
)
.withColumn("revenue_per_customer",
F.col("total_revenue") / F.col("unique_customers"))
.orderBy("order_date", "segment")
)
# Создаём Gold-таблицу
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.gold.daily_revenue_by_segment (
order_date DATE,
segment STRING,
total_revenue DECIMAL(15,2),
orders_count BIGINT,
avg_order_value DECIMAL(10,2),
unique_customers BIGINT,
revenue_per_customer DECIMAL(10,2),
_updated_at TIMESTAMP
)
USING iceberg
PARTITIONED BY (order_date)
""")
# Overwrite Partitions: для Gold часто используем перезапись партиций
# (вместо MERGE INTO) - проще и быстрее для агрегатов
(
daily_revenue_mart
.withColumn("_updated_at", F.current_timestamp())
.writeTo("lakehouse.gold.daily_revenue_by_segment")
.overwritePartitions()
)
print("Gold витрина обновлена:")
spark.table("lakehouse.gold.daily_revenue_by_segment").show()
4.3 ROLLUP и CUBE для многомерных агрегатов¶
Часто BI-дашборды требуют агрегации на разных уровнях иерархии одновременно: по дню, по месяцу, по году, итого. Вместо отдельных запросов для каждого уровня Spark SQL поддерживает ROLLUP и CUBE:
# ROLLUP: иерархическая агрегация (order_date → segment → итого)
rollup_df = (
enriched_orders
.filter(F.col("status") == "confirmed")
.withColumn("order_date", F.to_date("created_at"))
.rollup("order_date", "segment")
.agg(
F.sum("amount").alias("total_revenue"),
F.count("order_id").alias("orders_count")
)
# grouping() == 1 означает: это итоговая строка для данного уровня
.withColumn("level",
F.when(F.grouping("order_date") == 1, F.lit("GRAND_TOTAL"))
.when(F.grouping("segment") == 1, F.lit("DATE_TOTAL"))
.otherwise(F.lit("DETAIL"))
)
)
rollup_df.show(20)
# Результат содержит:
# DETAIL: конкретная дата + сегмент → их выручка
# DATE_TOTAL: конкретная дата + null → итог за день
# GRAND_TOTAL: null + null → итог за всё время
4.4 Сквозная схема трёх слоёв¶
Диаграмма показывает полный поток данных: три разных источника (PostgreSQL, Kafka, CSV) загружаются в три отдельных Bronze-таблицы с сохранением исходной структуры. Silver трансформирует их в бизнес-сущности с разными стратегиями - UPSERT для заказов, SCD Type 2 для клиентов, карантин для проблемных записей. Gold строит тематические витрины под конкретных потребителей.
5. Управление данными: идемпотентность, партиционирование и оптимизация¶
5.1 Идемпотентность: почему перезапуск не должен создавать дубликаты¶
Идемпотентность - это свойство операции: её повторное выполнение с теми же данными даёт тот же результат. В контексте data pipelines: если пайплайн упал на полпути и был перезапущен, итоговые данные должны быть теми же, что и при первом успешном запуске.
Три паттерна идемпотентной записи:
| Паттерн | Уровень | Механизм | Когда использовать |
|---|---|---|---|
| MERGE INTO (UPSERT) | Silver | Iceberg MERGE | Точечные обновления сущностей |
| Overwrite Partitions | Gold | overwritePartitions() |
Пересчёт агрегатов за дату |
| Append + Dedup | Bronze | dropDuplicates() при чтении Silver |
Append-only источники |
# Паттерн "Overwrite Partitions" для Gold
# При повторном запуске партиция за дату просто перезаписывается
# Нет дублей, нет MERGE overhead
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
(
daily_revenue_mart
.withColumn("_updated_at", F.current_timestamp())
.write
.format("iceberg")
.mode("overwrite") # Перезаписываем только тронутые партиции
.option("partitionOverwriteMode", "dynamic")
.saveAsTable("lakehouse.gold.daily_revenue_by_segment")
)
5.2 Партиционирование: стратегии для каждого слоя¶
Партиционирование - это физическое разделение данных по ключу на отдельные директории. Правильная стратегия позволяет Spark читать только нужные партиции, игнорируя остальные (Partition Pruning).
Правило партиционирования Bronze: партиционируем по ingestion_date (дате загрузки), а не по дате события. Инкрементальный пайплайн читает Bronze за «сегодня» - и читает только одну партицию.
Правило партиционирования Silver: партиционируем по month(created_at) - по бизнес-времени события. Аналитические запросы «за апрель» читают одну партицию.
Правило партиционирования Gold: партиционируем по дате отчёта order_date. Ежедневный пайплайн перезаписывает только текущую партицию.
5.3 Оптимизация Iceberg: OPTIMIZE и EXPIRE SNAPSHOTS¶
После интенсивных MERGE INTO и APPEND операций в Iceberg накапливаются мелкие файлы, что замедляет чтение. Регулярные операции обслуживания поддерживают производительность:
# Компакция мелких файлов в Silver-таблице
# Рекомендуемый целевой размер: 128-512 МБ на файл
spark.sql("""
CALL lakehouse.system.rewrite_data_files(
table => 'silver.orders',
strategy => 'binpack',
options => map('target-file-size-bytes', '134217728')
)
""")
# Удаление устаревших снапшотов Iceberg
# Оставляем последние 7 дней для Time Travel
spark.sql("""
CALL lakehouse.system.expire_snapshots(
table => 'silver.orders',
older_than => TIMESTAMP '2026-05-17 00:00:00',
retain_last => 10
)
""")
# Обновление статистики для Cost-Based Optimizer
spark.sql("ANALYZE TABLE lakehouse.silver.orders COMPUTE STATISTICS FOR ALL COLUMNS")
Когда запускать: компакцию - еженощно после завершения всех пайплайнов. Очистку снапшотов - еженедельно. Это задачи для orchestrator (Airflow, Dagster) - отдельные DAG-задачи после основных шагов пайплайна.
5.4 Time Travel: отладка и воспроизводимость¶
Одно из ключевых преимуществ Iceberg - возможность читать исторические снапшоты таблицы. Это критически важно для отладки и аудита:
# Смотрим на Silver-таблицу, какой она была вчера
spark.sql("""
SELECT * FROM lakehouse.silver.orders
TIMESTAMP AS OF '2026-05-22 23:59:59'
""")
# Сравниваем текущее состояние с прошлым
spark.sql("""
SELECT
current.order_id,
current.status AS status_now,
past.status AS status_yesterday,
current.amount AS amount_now,
past.amount AS amount_yesterday
FROM lakehouse.silver.orders current
LEFT JOIN (
SELECT * FROM lakehouse.silver.orders
TIMESTAMP AS OF '2026-05-22 23:59:59'
) past USING (order_id)
WHERE current.status != past.status
""")
6. End-to-End кейс: от «грязного события» до BI-метрики¶
6.1 Полный pipeline в одном скрипте¶
import datetime
from pyspark.sql import SparkSession, functions as F
from pyspark.sql.window import Window
spark = SparkSession.builder \
.appName("Medallion-E2E") \
.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()
TODAY = datetime.date.today().isoformat()
PIPELINE_RUN_ID = f"run-{TODAY}-001"
# ─── ШАГ 1: BRONZE INGESTION ─────────────────────────────────────────────────
print("=== ШАГ 1: BRONZE ===")
raw_data = [
("1001","cust_A","confirmed","1500.00","2026-05-23 10:00:00"),
("1001","cust_A","confirmed","1500.00","2026-05-23 10:00:00"), # дубликат
("1002","cust_B","pending", "750.50", "2026-05-23 11:30:00"),
("1003",None, "shipped", "N/A", "2026-05-23 12:00:00"), # битая
("1004","cust_C","confirmed","200.00", "2026-05-23 13:00:00"),
]
raw_df = spark.createDataFrame(
raw_data, ["order_id","customer_id","status","amount","created_at"]
)
bronze_df = (
raw_df
.withColumn("_ingested_at", F.current_timestamp())
.withColumn("_pipeline_run_id", F.lit(PIPELINE_RUN_ID))
.withColumn("_source_system", F.lit("orders_db"))
.withColumn("ingestion_date", F.current_date())
)
spark.sql("""CREATE TABLE IF NOT EXISTS lakehouse.bronze.orders (
order_id STRING, customer_id STRING, status STRING,
amount STRING, created_at STRING,
_ingested_at TIMESTAMP, _pipeline_run_id STRING,
_source_system STRING, ingestion_date DATE
) USING iceberg PARTITIONED BY (ingestion_date)""")
bronze_df.writeTo("lakehouse.bronze.orders").append()
print(f"Bronze: {bronze_df.count()} строк записано (включая дубли и битые)")
# ─── ШАГ 2: SILVER TRANSFORMATION ────────────────────────────────────────────
print("\n=== ШАГ 2: SILVER ===")
bronze_today = spark.table("lakehouse.bronze.orders") \
.filter(F.col("ingestion_date") == F.current_date())
# Дедупликация: берём последнюю версию по _ingested_at
w = Window.partitionBy("order_id").orderBy(F.col("_ingested_at").desc())
deduped = bronze_today \
.withColumn("rn", F.row_number().over(w)) \
.filter(F.col("rn") == 1).drop("rn")
# Валидация и карантин
valid_cond = (
F.col("order_id").isNotNull() &
F.col("customer_id").isNotNull() &
F.col("amount").rlike(r"^\d+(\.\d+)?$")
)
valid = deduped.filter(valid_cond)
invalid = deduped.filter(~valid_cond)
if invalid.count() > 0:
spark.sql("""CREATE TABLE IF NOT EXISTS lakehouse.silver.orders_quarantine (
order_id STRING, customer_id STRING, status STRING,
amount STRING, rejection_reason STRING, quarantined_at TIMESTAMP
) USING iceberg""")
invalid.withColumn("rejection_reason", F.lit("validation_failed")) \
.withColumn("quarantined_at", F.current_timestamp()) \
.select("order_id","customer_id","status","amount",
"rejection_reason","quarantined_at") \
.writeTo("lakehouse.silver.orders_quarantine").append()
print(f"В карантин отправлено: {invalid.count()} строк")
# Типизация
silver_df = (
valid
.withColumn("order_id", F.col("order_id").cast("integer"))
.withColumn("amount", F.col("amount").cast("decimal(10,2)"))
.withColumn("created_at",
F.coalesce(
F.to_timestamp("created_at","yyyy-MM-dd HH:mm:ss"),
F.to_timestamp("created_at","dd-MM-yyyy HH:mm:ss"),
))
.withColumn("status", F.lower(F.trim("status")))
.withColumn("_silver_processed_at", F.current_timestamp())
.select("order_id","customer_id","status","amount",
"created_at","_silver_processed_at")
)
spark.sql("""CREATE TABLE IF NOT EXISTS lakehouse.silver.orders (
order_id INTEGER, customer_id STRING, status STRING,
amount DECIMAL(10,2), created_at TIMESTAMP, _silver_processed_at TIMESTAMP
) USING iceberg PARTITIONED BY (month(created_at))""")
silver_df.createOrReplaceTempView("incoming_silver")
spark.sql("""
MERGE INTO lakehouse.silver.orders t
USING incoming_silver s ON t.order_id = s.order_id
WHEN MATCHED AND (t.status != s.status OR t.amount != s.amount)
THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
""")
print(f"Silver обновлён: {silver_df.count()} строк")
# ─── ШАГ 3: GOLD AGGREGATION ──────────────────────────────────────────────────
print("\n=== ШАГ 3: GOLD ===")
gold_df = (
spark.table("lakehouse.silver.orders")
.filter(F.col("status") == "confirmed")
.withColumn("order_date", F.to_date("created_at"))
.groupBy("order_date")
.agg(
F.sum("amount").alias("total_revenue"),
F.count("order_id").alias("orders_count"),
F.avg("amount").alias("avg_order_value"),
F.countDistinct("customer_id").alias("unique_customers"),
)
.withColumn("_updated_at", F.current_timestamp())
)
spark.sql("""CREATE TABLE IF NOT EXISTS lakehouse.gold.daily_revenue (
order_date DATE, total_revenue DECIMAL(15,2),
orders_count BIGINT, avg_order_value DECIMAL(10,2),
unique_customers BIGINT, _updated_at TIMESTAMP
) USING iceberg PARTITIONED BY (order_date)""")
gold_df.writeTo("lakehouse.gold.daily_revenue").overwritePartitions()
print("=== ИТОГ ПАЙПЛАЙНА ===")
print("Gold витрина:")
spark.table("lakehouse.gold.daily_revenue").show()
6.2 Результат выполнения пайплайна¶
=== ШАГ 1: BRONZE ===
Bronze: 5 строк записано (включая дубли и битые)
=== ШАГ 2: SILVER ===
В карантин отправлено: 1 строк (order_id=1003, customer_id=null, amount="N/A")
Silver обновлён: 3 строк (1001 deduped, 1002, 1004)
=== ШАГ 3: GOLD ===
Gold витрина:
order_date | total_revenue | orders_count | avg_order_value | unique_customers
2026-05-23 | 1700.00 | 2 | 850.00 | 2
Из пяти исходных строк: одна отфильтрована как дубликат (ROW_NUMBER), одна помещена в карантин (null customer_id). Итого две уникальные confirmed-записи: 1500.00 (order_id=1001) + 200.00 (order_id=1004) = 1700.00.
7. Чек-лист: как определить принадлежность данных к слою¶
7.1 Чек-лист для Bronze¶
- [ ] Данные пишутся только в режиме APPEND
- [ ] Схема содержит все оригинальные поля источника без изменений типов
- [ ] Добавлены технические метаданные:
_ingested_at,_source_file,_pipeline_run_id - [ ] Таблица партиционирована по дате загрузки (
ingestion_date) - [ ] Нет JOIN-ов с другими таблицами
- [ ] Нет агрегаций и GROUP BY
- [ ] Битые и дублирующиеся строки присутствуют (не удалены)
7.2 Чек-лист для Silver¶
- [ ] Все поля приведены к правильным типам (DECIMAL, TIMESTAMP, INTEGER)
- [ ] Дубликаты устранены по бизнес-ключу
- [ ] Добавлена карантинная стратегия для невалидных записей
- [ ] Запись идемпотентна (MERGE INTO или Overwrite Partitions)
- [ ] Нет бизнес-агрегатов (GROUP BY, SUM, COUNT)
- [ ] Нет KPI-расчётов и метрик
- [ ] Партиционирование по бизнес-дате события (не по дате загрузки)
7.3 Чек-лист для Gold¶
- [ ] Таблица создана под конкретный use-case (BI, ML, API)
- [ ] Данные денормализованы и предагрегированы
- [ ] Нет логики очистки и валидации
- [ ] Партиционирование под паттерн запросов потребителей
- [ ] Имя таблицы отражает бизнес-сущность, а не технический процесс
7.4 Типичные ошибки распределения логики¶
| Ошибка | Почему это проблема | Правильный слой |
|---|---|---|
| Агрегация выручки на Bronze | Потеря возможности replay | Gold |
| Расчёт LTV на Silver | Слишком специфично для одного use-case | Gold |
| Дедупликация отсутствует на Silver | Дубли в Gold и аналитике | Silver |
| Бизнес-правила на Bronze | Хрупкость: изменение правил = потеря истории | Silver |
| Удаление «плохих» строк на Bronze | Невозможно пересчитать после фикса | Silver (quarantine) |
8. Итоги: Medallion Architecture как инженерная дисциплина¶
Medallion Architecture - это не технология и не продукт. Это организационная дисциплина, которая заставляет команду думать о данных в терминах качества, ответственности и жизненного цикла. Каждый слой имеет один чёткий контракт, и нарушение этого контракта (агрегация на Bronze, бизнес-логика в ingestion-скрипте) сразу становится видимым - оно нарушает принципы слоя.
Три главных вывода урока:
-
Bronze - это архив, а не источник для аналитики. Его ценность - в воспроизводимости и полноте. Если Bronze правильно спроектирован, любой последующий слой можно пересчитать заново после исправления ошибок.
-
Silver - это контракт качества между Engineering и Analytics. Аналитик, читающий Silver, должен доверять каждой строке. Задача инженера - обеспечить это доверие через дедупликацию, типизацию и карантин.
-
Gold - это язык бизнеса. Разные потребители получают разные Gold-витрины, оптимизированные именно под их паттерны использования. Один размер не подходит всем.
Следующие уроки модуля разберут, как строить инкрементальные пайплайны между слоями с использованием Structured Streaming, как организовывать оркестрацию через Airflow/Dagster, и как настраивать мониторинг качества данных на Silver-слое.