Medallion Architecture: зачем три слоя, design-принципы каждого, anti-patterns - почему 'папки в S3 без метаслоя' уже legacy

От Data Swamp к Data Lakehouse: архитектурные принципы Bronze, Silver и Gold слоёв, идемпотентные ETL-пайплайны на PySpark + Iceberg и роль метаслоя в управлении качеством данных

lakehouse medallion architecture iceberg etl

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_id
  • WHERE 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.