Medallion Architecture: зачем три слоя, design-принципы каждого, anti-patterns - почему 'папки в S3 без метаслоя' уже legacy
От Data Swamp к Data Lakehouse: архитектурные принципы Bronze, Silver и Gold слоёв, идемпотентные ETL-пайплайны на PySpark + Iceberg и роль метаслоя в управлении качеством данных
Medallion Architecture: три слоя, которые превращают Data Swamp в управляемый Lakehouse¶
Прежде в этом курсе мы учились хранить данные быстро - партиционирование, бакеты, компакцию, Z-ordering. Но всё это бессмысленно, если данные сначала неправильно организованы структурно. Медальонная архитектура (Medallion Architecture) - это архитектурный шаблон, который отвечает на вопрос «как правильно организовать данные в Lakehouse», а не «как сделать их быстрыми».
Без чёткого разделения зон ответственности (слоёв) даже самый быстрый Spark-кластер будет обслуживать хаотичное, неуправляемое «болото данных» (Data Swamp). Именно этот урок отделяет разработчика, умеющего писать PySpark-код, от архитектора, умеющего строить Data Platform.
1. Эволюция парадигм: как Data Lake превратился в Data Swamp¶
Чтобы понять, почему медальонная архитектура стала стандартом индустрии, нужно пройти путь от первых Hadoop-кластеров до современных Lakehouse.
Эпоха Hadoop (2006–2013): большие данные, простая архитектура¶
Первые Data-платформы строились на Hadoop: HDFS хранил файлы, MapReduce обрабатывал их батчами. Архитектура была простой: есть сырые данные в HDFS, есть результаты обработки в другой директории HDFS. Schema была жёстко фиксирована в схемах Hive. Все данные принадлежали одной команде из 3–5 человек, которая знала каждую папку наизусть.
Это работало, пока компании были небольшими и данные принадлежали одной команде.
Эпоха Data Lake (2014–2018): «складывайте всё в S3»¶
С появлением дешёвого объектного хранилища S3 и колоночного формата Parquet возникла идеология Data Lake: «Храните все данные как есть, а схему определите потом (schema-on-read)». Это позволило копить данные без предварительного проектирования схем.
Компании дрейфовали к этой модели, потому что она казалась гибкой и дешёвой. На старте - да. Через три года - катастрофа.
Эпоха Data Swamp (2018–2021): когда гибкость превратилась в хаос¶
Когда Data Lake рос - к нему подключались десятки команд, каждая писала данные в «свои» папки, придумывала «свои» схемы, применяла «свои» правила обработки. Без централизованного управления начался распад:
- Schema Drift: команда А добавила колонку
user_uuidвместоuser_id- всё, что использовалоuser_id, упало - Параллельные записи: две джобы одновременно пишут в одну папку - половина файлов пустые, половина битые
- Нет атомарности: батч упал на середине записи - часть данных на диске, часть нет, Lake испорчен
- Нет lineage: откуда взялась таблица
final_customers_v7_FINAL_REAL? Никто не знает
Диаграмма: эволюция парадигм хранения данных. Путь от Hadoop (одна команда, простая архитектура) через Data Lake (гибкость, schema-on-read) к Data Swamp (хаос из-за отсутствия governance) и наконец к Data Lakehouse с ACID-метаслоем и медальонной архитектурой. Каждый переход был вынужденным ответом на проблемы предыдущей парадигмы.
Спасение: Data Lakehouse и метаслой¶
В 2020–2021 годах Apache Iceberg и Delta Lake достигли production-зрелости. Они добавили к объектному хранилищу то, чего не хватало Data Lake:
- ACID-транзакции: атомарная запись файлов - либо всё записалось, либо ничего
- Schema Evolution: добавление и переименование колонок без поломки существующих читателей
- Time Travel:
SELECT * FROM orders FOR SYSTEM_TIME AS OF '2024-01-01'- запрос к историческому состоянию - Snapshot Isolation: читатели видят консистентный срез данных во время параллельной записи
Метаслой - это фундамент, на котором строится медальонная архитектура. Без него она бессмысленна: «папки в S3» без ACID превращают любую архитектуру в Data Swamp при первом сбое.
2. Философия и цели медальонной архитектуры¶
Medallion Architecture - это не просто «три папки» (bronze, silver, gold). Это философия разделения ответственности в Data Pipeline: каждый слой имеет чёткое назначение, гарантии качества и контракт с потребителями.
Главный принцип: повышение качества, структуры и бизнес-ценности данных от слоя к слою. Сырые данные на входе - бизнес-инсайты на выходе.
Главное преимущество: возможность полностью пересобрать (replay) любую витрину с нуля. Если в Silver-логике нашлась ошибка - её можно исправить и перезапустить весь пайплайн без повторного похода в транзакционную СУБД источника. Bronze хранит все первоначальные данные вечно.
Диаграмма: полная архитектура трёх слоёв. Источники данных (PostgreSQL, Kafka, REST API) пишут в Bronze без трансформаций. Bronze преобразуется в структурированный Silver с типизированными колонками. Silver агрегируется в Gold-витрины под конкретные бизнес-нужды. Потребители (BI, ML, API) работают только с Gold, никогда не обращаясь напрямую к Bronze или Silver. Схемы в узлах показывают, как качество и структурированность данных нарастают от слоя к слою.
Разделение ответственности - ключевое преимущество архитектуры:
- Команда Data Engineering отвечает за Bronze и Silver-пайплайны
- Команда Analytics Engineering отвечает за Gold (часто через dbt)
- Команда Data Science работает с Gold и иногда с Silver
Это разделение позволяет каждой команде развиваться независимо: аналитики могут добавлять новые Gold-витрины, не трогая Bronze-пайплайны.
3. Bronze Layer: Landing Zone¶
Назначение¶
Bronze - это зеркало источника данных. Его задача - максимально точно и полно воспроизвести то, что пришло из внешней системы, не утеряв ни одной детали. Bronze - это страховка: если через полгода обнаружится, что Silver-логика содержала ошибку, Bronze позволит пересчитать всё с нуля без повторного обращения к продакшн-базе источника.
Дизайн-принципы Bronze¶
Принцип 1: Append-Only - никаких UPDATE и OVERWRITE
Bronze должен быть накопительным (append-only) логом. Даже если одни и те же данные придут дважды (дубликаты из источника) - Bronze сохраняет их оба раза. Дедупликация - ответственность Silver, не Bronze. Даже «плохие» записи (NULL в важных полях, невалидные суммы) попадают в Bronze. Фильтрация по качеству - тоже ответственность Silver.
Если вы делаете OVERWRITE в Bronze - вы уничтожаете возможность Replay.
Принцип 2: Минимум трансформаций
Единственные допустимые добавления в Bronze:
ingested_at = current_timestamp()- технический штамп времени загрузкиsource = 'kafka.orders_topic'- откуда пришли данныеbatch_id = uuid()- идентификатор батча для отладки
Никакой бизнес-логики, никаких JOIN, никаких фильтров, никакого приведения типов. Если из Kafka пришёл JSON - Bronze хранит его как строку. Если из источника пришла цифра в виде текста "42" - Bronze хранит "42", не 42.
Принцип 3: Flexible Schema
Источники данных меняются: добавляют поля, переименовывают, удаляют. Bronze должен быть готов к этим изменениям. Лучшая стратегия: хранить весь payload как STRING или BINARY. Iceberg поддерживает добавление колонок без перезаписи данных.
Принцип 4: Долгосрочное хранение
Bronze - самый дешёвый в хранении слой (сырые данные часто избыточны по объёму, но S3 дёшев). Типичный горизонт хранения: 3–7 лет. На Bronze лежит ответственность за регуляторные требования (GDPR, HIPAA) - данные должны быть доступны для аудита.
Диаграмма: процесс записи в Bronze. Spark-джоба читает Kafka-топик или batch-файлы, добавляет минимальные метаданные (timestamp, источник, batch_id) и записывает в Iceberg в режиме APPEND - только дописывание, никогда OVERWRITE. Такой подход гарантирует три ключевых свойства: неизменность (immutability), отсутствие потерь (даже дублей) и возможность Replay для перестройки вышестоящих слоёв.
Практика: запись в Bronze¶
from pyspark.sql import SparkSession
import pyspark.sql.functions as F
import uuid
spark = SparkSession.builder \
.appName("Medallion-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-таблицу: raw_payload как STRING (весь JSON как строка)
# Schema максимально гибкая: если источник добавит поле, ничего не сломается
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.bronze.orders (
raw_payload STRING COMMENT 'JSON payload from source as-is',
ingested_at TIMESTAMP COMMENT 'When we received this record',
source STRING COMMENT 'Source system identifier',
batch_id STRING COMMENT 'ETL batch ID for debugging'
) USING iceberg
""")
# Имитируем сырые данные из Kafka
batch_id = str(uuid.uuid4())
raw_events = [
('{"order_id": 1001, "user_id": 42, "amount": "99.99", "status": "PAID"}',),
('{"order_id": 1002, "user_id": 15, "amount": "149.00", "status": "PENDING"}',),
('{"order_id": 1001, "user_id": 42, "amount": "99.99", "status": "PAID"}',), # дубль из источника!
('{"order_id": 1003, "user_id": 42, "amount": "bad_data", "status": null}',), # плохие данные
]
raw_df = spark.createDataFrame(raw_events, ["raw_payload"])
# Добавляем только технические метаданные - никаких бизнес-трансформаций
bronze_df = raw_df \
.withColumn("ingested_at", F.current_timestamp()) \
.withColumn("source", F.lit("kafka.orders_raw")) \
.withColumn("batch_id", F.lit(batch_id))
# ТОЛЬКО APPEND - никогда не используем mode("overwrite") в Bronze!
bronze_df.writeTo("lakehouse.bronze.orders").append()
print(f"Bronze: записано {bronze_df.count()} строк (включая дубли и битые данные)")
Обратите внимание: дубль (order_id=1001) и «плохая» запись (amount='bad_data') записываются в Bronze без каких-либо фильтров. Это правильно. Silver разберётся с ними позже.
4. Silver Layer: единый источник правды¶
Назначение¶
Silver - это Single Source of Truth всей компании. Если Bronze - это «архив первоисточника», то Silver - это «нормализованная правда о реальности». Каждая сущность (пользователь, заказ, продукт) имеет ровно одну запись в Silver, отражающую её актуальное состояние.
Аналитики и дата-саентисты не должны видеть Bronze. Они работают с Silver. Если аналитику нужна история изменений - в Silver должна быть таблица с историей (SCD Type 2), а не доступ к Bronze.
Дизайн-принципы Silver¶
Принцип 1: Schema Enforcement - строгие типы данных
В Silver нет STRING для числовых значений. amount = "99.99" из Bronze превращается в amount DECIMAL(10,2) в Silver. user_id = "42" превращается в user_id BIGINT. Временные метки приводятся к единому формату и часовому поясу (UTC). Это делает Silver предсказуемым для всех потребителей.
Принцип 2: Дедупликация
Silver содержит каждую уникальную сущность ровно один раз. Дубли из Bronze (из-за at-least-once семантики Kafka) удаляются по бизнес-ключу (order_id). Это реализуется через MERGE INTO в Iceberg/Delta - идемпотентный UPSERT.
Принцип 3: Идемпотентность пайплайна
Одна и та же Silver-джоба, запущенная N раз с теми же данными, должна давать одинаковый результат. Это достигается через MERGE INTO (не INSERT). Если джоба упала и перезапустилась - никаких дублей. Это критически важно в продакшне, где рестарты неизбежны.
Принцип 4: Валидация качества
Silver отклоняет или маркирует плохие записи. Отклонённые записи уходят в отдельную таблицу silver.quarantine - они не потеряны, но помечены для анализа:
WHERE user_id IS NOT NULL AND user_id > 0- отбрасываем невалидные user_idWHERE amount IS NOT NULL AND amount > 0- отбрасываем нулевые суммыWHERE status IN ('PAID', 'PENDING', 'CANCELLED')- только известные статусы
Принцип 5: Conformed Dimensions
Silver унифицирует данные между источниками. Если user_id в одном источнике называется customer_id, а в другом - uid - Silver приводит всё к единому user_id. Валюты конвертируются в одну (например, USD). Временные зоны стандартизируются (UTC). Silver строит единый язык для всей компании.
Диаграмма: Silver-пайплайн от Bronze до структурированных данных. Bronze поставляет сырые JSON-строки. Silver-пайплайн проходит четыре шага: парсинг JSON, приведение типов, валидация качества (с отправкой плохих записей в карантин), дедупликация. Финальная запись идёт через MERGE INTO - идемпотентный UPSERT, который гарантирует отсутствие дублей при повторных запусках пайплайна.
Практика: Silver-пайплайн с MERGE INTO¶
from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, LongType, StringType, DecimalType
import pyspark.sql.functions as F
spark = SparkSession.builder \
.appName("Medallion-Silver-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") \
.getOrCreate()
# Создаём Silver-таблицу с явной схемой и строгими типами
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.silver.orders (
order_id BIGINT NOT NULL COMMENT 'Business key',
user_id BIGINT NOT NULL,
amount DECIMAL(12,2),
status STRING,
order_ts TIMESTAMP,
updated_at TIMESTAMP COMMENT 'When Silver was last updated'
) USING iceberg
""")
# Схема для парсинга JSON из Bronze
order_schema = StructType([
StructField("order_id", StringType()),
StructField("user_id", StringType()),
StructField("amount", StringType()),
StructField("status", StringType()),
])
# Читаем из Bronze - сырые JSON-строки
bronze_df = spark.table("lakehouse.bronze.orders")
# Шаг 1: Парсим JSON из raw_payload
parsed_df = bronze_df.withColumn(
"parsed", F.from_json(F.col("raw_payload"), order_schema)
).select("parsed.*", "ingested_at")
# Шаг 2: Приводим типы и добавляем технические поля
typed_df = parsed_df \
.withColumn("order_id", F.col("order_id").cast("bigint")) \
.withColumn("user_id", F.col("user_id").cast("bigint")) \
.withColumn("amount", F.col("amount").cast("decimal(12,2)")) \
.withColumn("order_ts", F.col("ingested_at")) \
.withColumn("updated_at", F.current_timestamp())
# Шаг 3: Хорошие записи - в Silver, плохие - в Quarantine
good_filter = (
F.col("order_id").isNotNull() &
F.col("user_id").isNotNull() &
F.col("amount").isNotNull() &
(F.col("amount") > 0) &
F.col("status").isin("PAID", "PENDING", "CANCELLED")
)
valid_df = typed_df.filter(good_filter) \
.select("order_id", "user_id", "amount", "status", "order_ts", "updated_at")
bad_df = typed_df.filter(~good_filter)
if bad_df.count() > 0:
bad_df.writeTo("lakehouse.silver.quarantine").append()
print(f"Quarantine: {bad_df.count()} плохих записей отправлено на анализ")
# Шаг 4: Дедупликация внутри батча (по бизнес-ключу)
deduped_df = valid_df.dropDuplicates(["order_id"])
# Шаг 5: Идемпотентный MERGE INTO Silver
# Этот запрос можно запустить 100 раз - результат будет одинаковым
deduped_df.createOrReplaceTempView("silver_updates")
spark.sql("""
MERGE INTO lakehouse.silver.orders AS target
USING silver_updates AS source
ON target.order_id = source.order_id
WHEN MATCHED THEN
UPDATE SET
target.status = source.status,
target.amount = source.amount,
target.updated_at = source.updated_at
WHEN NOT MATCHED THEN
INSERT (order_id, user_id, amount, status, order_ts, updated_at)
VALUES (source.order_id, source.user_id, source.amount,
source.status, source.order_ts, source.updated_at)
""")
print(f"Silver: обновлено/добавлено {deduped_df.count()} записей")
spark.table("lakehouse.silver.orders").show()
Ключевой момент этого кода - MERGE INTO. Он одновременно обновляет существующие записи (если order_id уже есть в Silver) и вставляет новые. Повторный запуск этой джобы с теми же данными не создаст дублей - Silver останется идемпотентным.
Обогащение данных в Silver: JOIN со справочниками¶
Очистка и дедупликация - первый шаг Silver. Второй шаг - обогащение (enrichment): дополнение базовых сущностей данными из справочных таблиц (reference tables / dimension tables), чтобы создать полноценную, самодостаточную запись.
Почему обогащение происходит в Silver, а не в Gold? Потому что Gold-пайплайны разные - каждый строит свои агрегаты под конкретных потребителей. Если не обогащать в Silver, каждый Gold-пайплайн будет дублировать один и тот же JOIN со справочником. Если справочная таблица изменится (например, пользователь сменил страну) - нужно будет обновить все Gold-пайплайны. Silver делает этот JOIN один раз централизованно, и все Gold-пайплайны автоматически получают актуальные данные.
Типичные сценарии обогащения в Silver:
- Добавление страны, города и сегмента пользователя из
silver.usersк таблице заказовsilver.orders - Добавление категории, бренда и поставщика из
silver.productsк таблице транзакций - Конвертация валюты по таблице обменных курсов
silver.exchange_rates - Геокодирование: IP-адрес источника → страна/регион через GeoIP-справочник
# Пример: обогащение silver.orders данными из справочников
# silver.users уже существует с актуальным country, segment
# silver.exchange_rates содержит курсы на каждый день
orders_df = spark.table("lakehouse.silver.orders")
users_ref_df = spark.table("lakehouse.silver.users") \
.select("user_id", "country", "user_segment")
rates_ref_df = spark.table("lakehouse.silver.exchange_rates") \
.filter(F.col("rate_date") == F.current_date()) \
.select("from_currency", "rate_to_usd")
# LEFT JOIN: заказы без пользователя не теряем (пользователь мог удалить аккаунт)
enriched_df = orders_df \
.join(users_ref_df, on="user_id", how="left") \
.join(
rates_ref_df,
on=orders_df["currency"] == rates_ref_df["from_currency"],
how="left"
) \
.withColumn(
"amount_usd",
F.when(F.col("currency") == "USD", F.col("amount"))
.otherwise(F.col("amount") * F.col("rate_to_usd"))
) \
.select(
"order_id", "user_id",
"country", # из справочника users
"user_segment", # из справочника users
"amount", "amount_usd", "currency",
"status", "order_ts", "updated_at"
)
# Обогащённую сущность сохраняем как отдельную Silver-таблицу
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.silver.orders_enriched (
order_id BIGINT,
user_id BIGINT,
country STRING COMMENT 'Из справочника silver.users',
user_segment STRING COMMENT 'Из справочника silver.users',
amount DECIMAL(12,2),
amount_usd DECIMAL(12,2) COMMENT 'Конвертированная сумма в USD',
currency STRING,
status STRING,
order_ts TIMESTAMP,
updated_at TIMESTAMP
) USING iceberg
PARTITIONED BY (days(order_ts))
""")
enriched_df.createOrReplaceTempView("enriched_updates")
# MERGE INTO обогащённой таблицы - тоже идемпотентный
spark.sql("""
MERGE INTO lakehouse.silver.orders_enriched AS target
USING enriched_updates AS source
ON target.order_id = source.order_id
WHEN MATCHED THEN
UPDATE SET
target.country = source.country,
target.user_segment = source.user_segment,
target.amount_usd = source.amount_usd,
target.status = source.status,
target.updated_at = source.updated_at
WHEN NOT MATCHED THEN
INSERT *
""")
Обогащённая silver.orders_enriched становится единым источником для нескольких Gold-пайплайнов. gold.revenue_daily читает из неё и сразу имеет country для группировки - без дополнительного JOIN. gold.user_kpi тоже читает из неё и получает user_segment для сегментации LTV.
Важный нюанс: обогащение в Silver добавляет базовые атрибуты (country, segment, нормализованную валюту), но не вычисляет бизнес-метрики (LTV, churn_score, revenue_30d). Метрики - это задача Gold. Silver отвечает на вопрос «что это за объект», Gold - на вопрос «каковы его бизнес-показатели».
CDC и SCD: как Silver работает с изменяющимися данными¶
Реальные источники присылают не только новые записи, но и изменения существующих (Change Data Capture - CDC). Например, заказ перешёл из PENDING в PAID. Silver должен это отразить.
SCD Type 1 (полная перезапись текущего состояния): Silver хранит только последнее состояние сущности. При изменении статуса заказа - MERGE INTO обновляет только изменённые поля. История не сохраняется. Подходит для сущностей, где история не важна (например, текущий адрес доставки).
SCD Type 2 (историческое хранение версий): Silver хранит все версии сущности с колонками valid_from, valid_to, is_current. При изменении - старая запись закрывается (valid_to = now(), is_current = false), создаётся новая (valid_from = now(), is_current = true). Подходит для аудита и аналитики исторических трендов.
-- SCD Type 2: пример схемы Silver для изменяемых сущностей
CREATE TABLE lakehouse.silver.orders_history (
order_id BIGINT,
user_id BIGINT,
status STRING,
valid_from TIMESTAMP,
valid_to TIMESTAMP, -- NULL означает "текущая версия"
is_current BOOLEAN
) USING iceberg
PARTITIONED BY (is_current, days(valid_from));
5. Gold Layer: бизнес-витрины¶
Назначение¶
Gold - это язык бизнеса. Если Silver говорит на языке технических сущностей (заказы, пользователи, продукты), то Gold говорит на языке KPI (выручка за месяц, количество активных пользователей, churn rate). Аналитики, BI-инструменты и ML-модели работают исключительно с Gold.
Gold-таблицы проектируются под конкретные запросы - это принципиально отличает их от Silver. Silver оптимизирован для обновления и хранения полного состояния, Gold оптимизирован для скорости чтения и удобства аналитики.
Дизайн-принципы Gold¶
Принцип 1: Денормализация для скорости чтения
Gold не нормализован - это не OLTP. Вместо JOIN-ов в момент запроса, Gold заранее преджойнивает нужные таблицы. gold.user_kpi может содержать user_name и user_country прямо в таблице агрегатов - чтобы аналитику не нужно было джойнить её с silver.users.
Принцип 2: Агрегации по бизнес-логике
Gold считает то, что спрашивает бизнес: LTV (пожизненная ценность клиента), Churn Rate, DAU/MAU, средний чек. Бизнес-логика вычисления этих метрик живёт в Gold-пайплайнах централизованно. Если определение LTV изменится - его меняют в одном месте, а не в 20 отчётах аналитиков.
Принцип 3: Оптимизация чтения
Gold-таблицы - главные кандидаты на Z-ordering и Bucketing (которые мы изучали в модуле 04). Именно здесь имеет смысл OPTIMIZE ZORDER BY (user_id, country) для BI-запросов. Bronze и Silver часто обновляются - Z-ordering их замедлит. Gold относительно стабилен - Z-ordering окупается многократно.
Диаграмма: Gold Layer - от Silver к бизнес-витринам. Три Silver-таблицы (заказы, пользователи, продукты) служат источником для трёх независимых Gold-пайплайнов: KPI пользователей, ежедневная выручка и ML-признаки. Каждый пайплайн оптимизирован под своих потребителей: BI-таблицы с Z-ordering для быстрых аналитических запросов, ML-таблицы с оптимальной структурой для batch feature inference. SLA на свежесть задаётся отдельно для каждого Gold-пайплайна.
Практика: Gold-пайплайн агрегации выручки¶
# Gold-пайплайн: ежедневная агрегация выручки по стране
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.gold.revenue_daily (
report_date DATE NOT NULL,
country STRING,
revenue DECIMAL(15,2),
orders_count BIGINT,
avg_order DECIMAL(10,2),
computed_at TIMESTAMP
) USING iceberg
PARTITIONED BY (report_date)
""")
# Агрегация с бизнес-логикой: только PAID заказы входят в выручку
# Денормализация: JOIN с silver.users прямо в Gold-запросе
spark.sql("""
MERGE INTO lakehouse.gold.revenue_daily AS target
USING (
SELECT
DATE(o.order_ts) AS report_date,
u.country AS country,
SUM(o.amount) AS revenue,
COUNT(*) AS orders_count,
AVG(o.amount) AS avg_order,
current_timestamp() AS computed_at
FROM lakehouse.silver.orders o
JOIN lakehouse.silver.users u ON o.user_id = u.user_id
WHERE o.status = 'PAID'
AND DATE(o.order_ts) = CURRENT_DATE - INTERVAL 1 DAY
GROUP BY DATE(o.order_ts), u.country
) AS source
ON target.report_date = source.report_date
AND target.country = source.country
WHEN MATCHED THEN
UPDATE SET
revenue = source.revenue,
orders_count = source.orders_count,
avg_order = source.avg_order,
computed_at = source.computed_at
WHEN NOT MATCHED THEN
INSERT *
""")
# Gold оптимизируем под BI-запросы - Z-ORDER для скорости
# Запускаем еженедельно (не каждый день - Z-order дорогой)
spark.sql("""
OPTIMIZE lakehouse.gold.revenue_daily
ZORDER BY (country, report_date)
""")
print("Gold: revenue_daily обновлён")
spark.table("lakehouse.gold.revenue_daily").show()
6. Анти-паттерны медальонной архитектуры¶
Четыре самых распространённых ошибки, которые превращают Medallion в анти-паттерн:
Диаграмма: четыре анти-паттерна медальонной архитектуры. Каждый приводит к реальной производственной проблеме. The Shortcut (прямая ссылка Bronze→Gold) - к дублированию логики и расхождению KPI между витринами. Бизнес-логика в Bronze - к невозможности аудита исторических данных. Монолитный Silver - к единой точке отказа всей платформы. Отсутствие идемпотентности - к задвоенным данным в отчётах, которые замечает CEO в понедельник утром.
Анти-паттерн 1: The Shortcut - пропуск Silver¶
Ситуация: дедлайн через два дня, аналитик хочет данные прямо сейчас. Инженер читает из Bronze напрямую в Gold, добавляя дедупликацию прямо в Gold-запрос.
Через три месяца аналитик создаёт вторую Gold-витрину с похожей логикой - снова из Bronze, снова со своей дедупликацией. Теперь в двух витринах - две чуть разные версии логики очистки. Расхождение составляет 0.3%. Никто не знает, какая правильная.
Правило: Gold всегда читает только из Silver, никогда напрямую из Bronze.
Анти-паттерн 2: бизнес-логика в Bronze¶
Ситуация: инженер добавляет в Bronze-джобу WHERE status = 'OK', чтобы «сразу почистить мусор». Кажется логичным - зачем хранить ERROR-записи?
Через полгода регулятор требует аудит всех ERROR-запросов за прошедший год. Данных нет - они были отфильтрованы на Bronze-шаге. Replay невозможен: исходные данные уничтожены.
Правило: Bronze никогда не применяет фильтры по бизнес-критериям. Всё попадает в Bronze, Silver решает, что валидно.
Анти-паттерн 3: монолитный Silver¶
Ситуация: команда создаёт «универсальную» Silver-таблицу silver.everything с 500 колонками и 30 JOIN-ами. Логика «всё в одном».
Пересчёт занимает 8 часов. Если джоба падает на 7-м часу - нужно перезапускать с нуля. Все downstream Gold-пайплайны блокируются и опаздывают. SLA нарушен.
Правило: Silver - набор тематических таблиц (orders, users, products), а не одна таблица на все случаи жизни. Пайплайны должны быть независимыми.
Анти-паттерн 4: игнорирование идемпотентности¶
Ситуация: Silver-джоба использует INSERT INTO вместо MERGE INTO. Оркестратор (Airflow) перезапустил задачу из-за временного сбоя. Та же джоба выполнилась дважды. Silver теперь содержит дубли. Gold агрегирует Silver и возвращает удвоенную выручку.
Правило: любая трансформация в Silver и Gold должна быть идемпотентной - используйте MERGE INTO, INSERT OVERWRITE по партиции или CREATE OR REPLACE TABLE AS SELECT.
7. Data Catalog как клей медальонной архитектуры¶
Медальонная архитектура без Data Catalog - это три слоя папок, которые только выглядят как архитектура. Data Catalog (Unity Catalog, Apache Atlas, AWS Glue, Hive Metastore) выполняет три ключевые функции, без которых Medallion невозможен:
Диаграмма: Data Catalog как управляющий слой медальонной архитектуры. Catalog не просто каталогизирует таблицы - он реализует RBAC (кто имеет доступ к каждому слою), Data Lineage (откуда пришли данные в Gold-таблице и через какие джобы) и Quality Rules с SLA на свежесть данных. Без Catalog Bronze/Silver/Gold - это просто папки. С Catalog - это управляемая Data Platform.
RBAC (Role-Based Access Control):
- Bronze: только служебные системные аккаунты имеют право записи. Ни один аналитик не должен видеть Bronze напрямую
- Silver: чтение доступно командам Data Engineering и Analytics Engineering. Прямой доступ BI-инструментов запрещён
- Gold: полный READ-доступ для аналитиков, BI-инструментов и ML-платформ
Data Lineage: Catalog должен знать, что gold.revenue_daily вычисляется из silver.orders и silver.users через ETL-джобу версии 1.3. Если схема silver.orders изменится - Catalog автоматически предупредит, что gold.revenue_daily может сломаться.
Кейс для обсуждения: бизнес-пользователь просит прямой доступ к Bronze, потому что «ему срочно нужны данные, а Gold-витрина будет готова только через неделю».
Правильный ответ архитектора: отказать. Допуск в Bronze создаёт зависимость от сырых, нестабильных схем. Когда источник изменит структуру JSON (что неизбежно), у пользователя сломаются все его отчёты. Вместо этого: ускорить доставку Silver-данных или предоставить доступ к Silver с чётким предупреждением об SLA.
8. Итоги: чек-лист готовности архитектуры¶
Контрольные вопросы для каждого слоя¶
| Вопрос | Bronze | Silver | Gold |
|---|---|---|---|
| Данные append-only? | Да | Нет (MERGE) | Нет (агрегации) |
| Схема фиксированная? | Нет (flexible) | Да | Да |
| Дедупликация? | Нет | Да | Да |
| Идемпотентен ли пайплайн? | Да (append) | Да (MERGE) | Да (MERGE/OVERWRITE) |
| Кто имеет READ-доступ? | Только системы | Инженеры | Все |
| Горизонт хранения | 3–7 лет | 2–3 года | 1–2 года |
| Формат хранения | Iceberg / Delta | Iceberg / Delta | Iceberg / Delta |
Пять признаков здорового Medallion¶
- Bronze никогда не теряет данные: дубли, ошибки, «мусор» - всё хранится
- Silver - единственный источник правды: никто не читает из Bronze для аналитики
- Gold говорит на языке бизнеса: в именах колонок нет технических артефактов типа
user_id_fk - Любой Gold-пайплайн можно пересобрать с нуля: запустил Replay из Bronze - получил тот же Gold через Silver
- Catalog знает lineage каждой таблицы: откуда данные, кто написал, когда последний раз обновились
Медальонная архитектура - это не просто три папки. Это инженерная дисциплина, которая превращает Data Lake из склада в Data Platform. В следующих уроках этого модуля мы рассмотрим конкретные паттерны моделирования данных внутри Silver и Gold слоёв: Slowly Changing Dimensions, CDC-паттерны и Domain-Driven Design в контексте Lakehouse.