Data Vault + Medallion: Raw Vault → Business Vault → Kimball Data Marts

Синергия Data Vault 2.0 и Medallion Architecture: Raw Vault как Silver Integration Layer, Business Vault с PIT-таблицами и Bridge-таблицами, сборка Kimball Star Schema в Gold через оптимизированные JOIN с PIT, материализация витрин в Iceberg

lakehouse data-vault medallion kimball pit-tables

Мы изучили Medallion Architecture (Bronze → Silver → Gold) как процессный паттерн: как данные движутся из сырого состояния через очистку к аналитическим витринам. Мы изучили Data Vault 2.0 как модельный паттерн: как организовать интеграционный слой так, чтобы он масштабировался с любым количеством источников.

Теперь пора собрать их вместе и увидеть архитектуру, которую строят крупнейшие enterprise Data Platform'ы в мире.

Medallion и Data Vault - не альтернативы. Они созданы друг для друга. Medallion отвечает на вопрос «как организовать жизненный цикл данных», Data Vault - на вопрос «как моделировать интеграционный слой внутри этого цикла». Один описывает зоны ответственности, второй - архитектуру внутри зон.


1. Архитектурный маппинг: где Data Vault живёт внутри Medallion

1.1 Три медальона и их DV-сущности

Ключевое разделение внутри Silver Layer на Raw Vault и Business Vault - это не бюрократия, а принципиально разные зоны ответственности с разными правилами изменения.

1.2 Главный закон сквозной архитектуры

«Загружай в DV (Silver), читай из Кимбалла (Gold)»

Это не просто красивый слоган. В нём закодировано разделение двух несовместимых требований:

  • Запись требует гибкости, масштабируемости и аудита. Data Vault удовлетворяет это лучше любой другой модели.
  • Чтение требует скорости, простоты и понятности для BI-инструментов. Star Schema Кимбалла удовлетворяет это лучше любой другой модели.

Попытка объединить оба требования в одной модели неизбежно означает компромисс в обе стороны: данные хранятся недостаточно гибко для ingestion, но при этом недостаточно оптимально для BI.


2. Silver: Raw Vault - Integration Layer

2.1 Принципы Raw Vault

Raw Vault - это первый подслой Silver. Его задача: принять данные из Bronze, добавить хэш-ключи и распределить по Hub/Link/Satellite. Никакой бизнес-логики. Никаких суждений о том, «правильные» ли данные. Если в Bronze пришёл заказ со статусом UNKNOWN_STATUS - Raw Vault его примет без изменений.

Неизменяемые правила Raw Vault:

  • Append-only: строки в Hub и Link не обновляются и не удаляются
  • Нет бизнес-правил: исходные данные сохраняются в первозданном виде
  • Полный аудит: каждая строка знает, откуда и когда пришла
  • Детерминированное хэширование: sha2(trim(lower(business_key)), 256)

Именно неизменяемость Raw Vault обеспечивает replayability - возможность полностью пересобрать Business Vault и Gold-витрины, если бизнес-правила изменились. Bronze + Raw Vault = полная история всего, что когда-либо происходило в системе.

2.2 Что идёт в Raw Vault из Bronze

Все четыре таблицы загружаются параллельно из одного Staging DataFrame. Это было подробно разобрано в предыдущем уроке. Напомним суть: хэш-ключи вычисляются локально на каждом Executor без централизованной координации, что обеспечивает идеальный параллелизм.


3. Silver: Business Vault - бизнес-логика поверх Raw

3.1 Зачем Business Vault существует

После загрузки Raw Vault перед нами стоит проблема: чтобы получить «профиль клиента» на конкретную дату, нужно выполнить следующий запрос:

SELECT h.business_key, s1.full_name, s1.city, s2.segment, s2.credit_score
FROM hub_customer h
-- Первый satellite: актуальная версия на дату X
JOIN sat_customer_details s1
    ON h.hub_hk = s1.hk_customer
    AND s1.load_ts = (
        SELECT MAX(load_ts) FROM sat_customer_details
        WHERE hk_customer = h.hub_hk
        AND load_ts <= '2026-05-23'
    )
-- Второй satellite: актуальная версия на дату X
JOIN sat_customer_scoring s2
    ON h.hub_hk = s2.hk_customer
    AND s2.load_ts = (
        SELECT MAX(load_ts) FROM sat_customer_scoring
        WHERE hk_customer = h.hub_hk
        AND load_ts <= '2026-05-23'
    )

Каждый коррелированный подзапрос (SELECT MAX(load_ts) WHERE load_ts <= X) - это потенциально тяжёлая операция в Spark, требующая агрегации всей истории Satellite. При нескольких Satellite'ах и миллионах клиентов это превращается в вычислительный ад: огромный Shuffle, невозможность использовать Broadcast Hash Join, латентность в десятки минут для построения Gold-витрины.

Business Vault решает эту проблему, предвычисляя две специализированные структуры: PIT-таблицы и Bridge-таблицы.

3.2 PIT-таблицы: временны́е индексы для Satellite'ов

PIT-таблица (Point-in-Time Table) - это предвычисленная таблица, которая для каждого Hub-ключа на каждый момент времени (снапшотную дату) хранит точные load_ts из каждого связанного Satellite. Вместо диапазонного поиска load_ts <= X - мгновенный поиск по точному значению load_ts = X (equi-join).

Структура PIT-таблицы:

Колонка Тип Описание
hk_customer STRING SHA-256 ключ Hub'а
snap_date DATE Дата снапшота (конец дня, начало месяца и т.д.)
sat_details_load_ts TIMESTAMP Точный load_ts из sat_customer_details на эту дату
sat_scoring_load_ts TIMESTAMP Точный load_ts из sat_customer_scoring на эту дату

PIT-таблица рассчитывается по расписанию (ежедневно, еженедельно) - один раз для всего периода. Результат материализуется в Iceberg. Все последующие запросы Gold-слоя используют уже готовую PIT-таблицу.

3.3 Алгоритм построения PIT-таблицы

from pyspark.sql import functions as F, Window

# Для каждого клиента и каждой снапшотной даты находим актуальный load_ts из Satellite

# Шаг 1: Генерируем сетку снапшотных дат (например, конец каждого дня за последний год)
snapshot_dates = spark.sql("""
    SELECT explode(sequence(
        DATE '2026-01-01',
        CURRENT_DATE(),
        INTERVAL 1 DAY
    )) AS snap_date
""")

# Шаг 2: Для каждого (hk_customer, snap_date) находим MAX(load_ts) из sat_customer_details
# где load_ts <= snap_date
sat_details_ts = spark.sql("""
    SELECT hk_customer, load_ts AS sat_details_ts
    FROM silver.sat_customer_details
    WHERE load_end_ts IS NULL OR load_end_ts > CURRENT_DATE()
""")

pit_details = (
    spark.table("silver.hub_customer")
    .crossJoin(snapshot_dates)  # каждый клиент × каждая дата
    .join(sat_details_ts, on="hk_customer", how="left")
    .filter(F.col("sat_details_ts") <= F.col("snap_date"))
    .groupBy("hk_customer", "snap_date")
    .agg(F.max("sat_details_ts").alias("sat_details_load_ts"))
)

# Шаг 3: Аналогично для sat_customer_scoring
sat_scoring_ts = spark.sql("""
    SELECT hk_customer, load_ts AS sat_scoring_ts
    FROM silver.sat_customer_scoring
""")

pit_scoring = (
    spark.table("silver.hub_customer")
    .crossJoin(snapshot_dates)
    .join(sat_scoring_ts, on="hk_customer", how="left")
    .filter(F.col("sat_scoring_ts") <= F.col("snap_date"))
    .groupBy("hk_customer", "snap_date")
    .agg(F.max("sat_scoring_ts").alias("sat_scoring_load_ts"))
)

# Шаг 4: Объединяем в единую PIT-таблицу
pit_customer = pit_details.join(pit_scoring, on=["hk_customer", "snap_date"], how="full")

pit_customer.writeTo("lakehouse.silver.pit_customer").overwritePartitions()

Стоимость этой операции: crossJoin между Hub (N клиентов) и снапшотными датами (M дней) создаёт N×M строк. Для 1 млн клиентов и 365 дней - это 365 млн строк. Это тяжёлый расчёт, поэтому PIT перестраивается инкрементально: только для новых дат и только для изменившихся клиентов.

3.4 Bridge-таблицы: плоские иерархии связей

Bridge-таблица решает другую проблему: сложные цепочки Hub-Link-Hub при анализе иерархических или многоуровневых бизнес-связей.

Например, бизнес-вопрос: «Какие продукты входят в каждый заказ каждого клиента?» требует JOIN'а: hub_customer → lnk_order_customer → hub_order → lnk_order_product → hub_product. Пять таблиц, четыре JOIN'а, неизбежный тяжёлый SortMergeJoin.

Bridge-таблица «схлопывает» эту цепочку в плоскую таблицу хэш-ключей:

# Построение Bridge-таблицы: все комбинации (клиент, заказ, продукт)
bridge_customer_product = spark.sql("""
    SELECT
        loc.hk_customer,
        loc.hk_order,
        lop.hk_product,
        loc.load_ts AS lnk_load_ts
    FROM lakehouse.silver.lnk_order_customer loc
    JOIN lakehouse.silver.lnk_order_product  lop
        ON loc.hk_order = lop.hk_order
""")

bridge_customer_product.writeTo("lakehouse.silver.bridge_customer_product") \
    .overwritePartitions()

После материализации Bridge-таблицы запрос «продукты каждого клиента» превращается в:

SELECT hk_customer, hk_product FROM silver.bridge_customer_product

Один APPEND-only scan без единого JOIN.

3.5 Вычисляемые сущности (Computed Satellites)

Business Vault также содержит computed satellites - сателлиты, которые не приходят из источника, а вычисляются по бизнес-правилам на основе Raw Vault данных.

Пример: sat_customer_segmentation - расчётный сателлит, который присваивает клиенту сегмент на основе его истории заказов:

# Computed Satellite: сегментация клиентов на основе статистики заказов
customer_stats = spark.sql("""
    SELECT
        loc.hk_customer,
        COUNT(DISTINCT loc.hk_order)    AS total_orders,
        SUM(sos.amount)                 AS lifetime_value,
        MAX(sos.load_ts)                AS last_order_ts
    FROM lakehouse.silver.lnk_order_customer loc
    JOIN lakehouse.silver.sat_order_status sos
        ON loc.hk_order = sos.hk_order
        AND sos.load_end_ts IS NULL
    GROUP BY loc.hk_customer
""")

computed_segmentation = customer_stats.withColumn(
    "segment",
    F.when(F.col("lifetime_value") > 10_000, F.lit("PLATINUM"))
     .when(F.col("lifetime_value") > 2_000,  F.lit("GOLD"))
     .when(F.col("lifetime_value") > 500,    F.lit("SILVER"))
     .otherwise(F.lit("BRONZE"))
).withColumn(
    "hash_diff",
    F.sha2(F.concat_ws("||", F.col("segment"), F.col("total_orders").cast("string")), 256)
).withColumn("load_ts", F.current_timestamp()) \
 .withColumn("load_end_ts", F.lit(None).cast("timestamp")) \
 .withColumn("record_source", F.lit("bv_segmentation_pipeline"))

computed_segmentation.writeTo("lakehouse.silver.sat_customer_segmentation") \
    .overwritePartitions()

Критически важно: computed satellites хранятся отдельно от Raw Vault сателлитов. Если бизнес-правило сегментации изменится - мы пересчитываем sat_customer_segmentation, не трогая sat_customer_details из источника. Разделение ответственности соблюдено.


4. Gold: Кимбалл поверх Business Vault

4.1 Полная схема сборки Gold-слоя

4.2 Почему Gold должен быть физически материализован

Альтернатива материализации - создать VIEW поверх Data Vault с оконными функциями для нахождения актуальной версии каждого Satellite. Каждый раз, когда аналитик открывает дашборд, Spark запускал бы этот VIEW.

Это катастрофа по трём причинам:

  1. Оконные функции (ROW_NUMBER() OVER (PARTITION BY hk_customer ORDER BY load_ts DESC)) - это Wide Transformation с Shuffle. При каждом BI-запросе Spark перемешивает всю историю Satellite'ов по сети.

  2. Конкурентность - если 50 аналитиков одновременно обновляют свои дашборды, запущены 50 параллельных тяжёлых Spark-заданий. Кластер перегружен.

  3. Latency - пользователь BI ожидает ответа за секунды. Пересчёт через Raw Vault займёт минуты.

Физическая материализация Gold-витрин в Iceberg переносит тяжёлый пересчёт на scheduled pipeline (раз в ночь или раз в час), а аналитику отдаёт быстрый scan уже готовой таблицы. Это точно та же логика, что и у PIT-таблиц: вычисляем один раз, читаем быстро много раз.

4.3 Практика: полный пайплайн DV → Gold

from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder \
    .appName("DataVault-to-Kimball-Gold") \
    .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()

# ─── 1. Имитируем Raw Vault данные ────────────────────────────────────────────
# hub_customer
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.silver.hub_customer (
        hub_hk STRING, business_key STRING, load_ts TIMESTAMP, record_source STRING
    ) USING iceberg
""")
spark.sql("""
    INSERT INTO lakehouse.silver.hub_customer VALUES
    ('hk_cust_001', 'CUST-001', TIMESTAMP '2026-01-01', 'crm'),
    ('hk_cust_002', 'CUST-002', TIMESTAMP '2026-01-01', 'crm'),
    ('hk_cust_003', 'CUST-003', TIMESTAMP '2026-02-15', 'ecom')
""")

# sat_customer_details - персональные данные (обновлялись дважды для CUST-001)
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.silver.sat_customer_details (
        hk_customer STRING, load_ts TIMESTAMP, load_end_ts TIMESTAMP,
        hash_diff STRING, full_name STRING, city STRING, email STRING
    ) USING iceberg
""")
spark.sql("""
    INSERT INTO lakehouse.silver.sat_customer_details VALUES
    ('hk_cust_001', TIMESTAMP '2026-01-01', TIMESTAMP '2026-03-01', 'hd_1a', 'Иван Иванов',  'Казань',  'ivan@m.ru'),
    ('hk_cust_001', TIMESTAMP '2026-03-01', NULL,                   'hd_1b', 'Иван Иванов',  'Москва',  'ivan@m.ru'),
    ('hk_cust_002', TIMESTAMP '2026-01-01', NULL,                   'hd_2a', 'Alice Smith',  'London',  'alice@m.com'),
    ('hk_cust_003', TIMESTAMP '2026-02-15', NULL,                   'hd_3a', 'Bob Johnson',  'Berlin',  'bob@m.de')
""")

# sat_customer_scoring - скоринг (обновлялся один раз для CUST-001)
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.silver.sat_customer_scoring (
        hk_customer STRING, load_ts TIMESTAMP, load_end_ts TIMESTAMP,
        hash_diff STRING, credit_score INT, risk_class STRING
    ) USING iceberg
""")
spark.sql("""
    INSERT INTO lakehouse.silver.sat_customer_scoring VALUES
    ('hk_cust_001', TIMESTAMP '2026-01-01', TIMESTAMP '2026-05-01', 'hs_1a', 720, 'LOW'),
    ('hk_cust_001', TIMESTAMP '2026-05-01', NULL,                   'hs_1b', 780, 'VERY_LOW'),
    ('hk_cust_002', TIMESTAMP '2026-01-01', NULL,                   'hs_2a', 650, 'MEDIUM'),
    ('hk_cust_003', TIMESTAMP '2026-02-15', NULL,                   'hs_3a', 590, 'MEDIUM')
""")

# ─── 2. Строим PIT-таблицу (Business Vault) ───────────────────────────────────
# Снапшотные даты: несколько ключевых дат для анализа
snapshot_dates = spark.createDataFrame([
    ("2026-02-01",),
    ("2026-04-01",),
    ("2026-05-23",),  # сегодня
], ["snap_date"]).withColumn("snap_date", F.to_date(F.col("snap_date")))

hub_customer = spark.table("lakehouse.silver.hub_customer")

# Для каждого (customer, snap_date) находим MAX(load_ts) <= snap_date из каждого Satellite
sat_details_hist = spark.table("lakehouse.silver.sat_customer_details") \
    .select("hk_customer", "load_ts")

sat_scoring_hist = spark.table("lakehouse.silver.sat_customer_scoring") \
    .select("hk_customer", "load_ts")

# Строим PIT через LEFT JOIN с фильтром по дате и группировкой
pit_with_details = (
    hub_customer.select("hub_hk").withColumnRenamed("hub_hk", "hk_customer")
    .crossJoin(snapshot_dates)
    .join(
        sat_details_hist.withColumnRenamed("load_ts", "sat_details_ts"),
        on="hk_customer", how="left"
    )
    .filter(F.col("sat_details_ts") <= F.col("snap_date").cast("timestamp"))
    .groupBy("hk_customer", "snap_date")
    .agg(F.max("sat_details_ts").alias("sat_details_load_ts"))
)

pit_full = (
    pit_with_details
    .join(
        hub_customer.select("hub_hk").withColumnRenamed("hub_hk", "hk_customer")
        .crossJoin(snapshot_dates)
        .join(
            sat_scoring_hist.withColumnRenamed("load_ts", "sat_scoring_ts"),
            on="hk_customer", how="left"
        )
        .filter(F.col("sat_scoring_ts") <= F.col("snap_date").cast("timestamp"))
        .groupBy("hk_customer", "snap_date")
        .agg(F.max("sat_scoring_ts").alias("sat_scoring_load_ts")),
        on=["hk_customer", "snap_date"],
        how="left"
    )
)

spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.silver.pit_customer (
        hk_customer          STRING,
        snap_date            DATE,
        sat_details_load_ts  TIMESTAMP,
        sat_scoring_load_ts  TIMESTAMP
    ) USING iceberg
    PARTITIONED BY (snap_date)
""")

pit_full.writeTo("lakehouse.silver.pit_customer").overwritePartitions()

print("PIT-таблица построена:")
spark.table("lakehouse.silver.pit_customer").orderBy("hk_customer", "snap_date").show(truncate=False)

# ─── 3. Строим Gold dim_customer через PIT + equi-join ────────────────────────
# Ключевое: используем pit.sat_details_load_ts = sat.load_ts (equi-join!)
# Spark может применить BroadcastHashJoin для PIT-таблицы - она небольшая
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.gold.dim_customer (
        customer_sk      STRING   COMMENT 'hk_customer - становится PK Gold-измерения',
        snap_date        DATE,
        full_name        STRING,
        city             STRING,
        email            STRING,
        credit_score     INT,
        risk_class       STRING
    ) USING iceberg
    PARTITIONED BY (snap_date)
""")

spark.sql("""
    INSERT OVERWRITE lakehouse.gold.dim_customer
    SELECT
        pit.hk_customer       AS customer_sk,
        pit.snap_date,
        det.full_name,
        det.city,
        det.email,
        scr.credit_score,
        scr.risk_class
    FROM lakehouse.silver.pit_customer pit

    -- equi-join по точному load_ts из PIT: O(1) lookup, возможен BroadcastHashJoin
    LEFT JOIN /*+ BROADCAST(det) */ lakehouse.silver.sat_customer_details det
        ON  pit.hk_customer = det.hk_customer
        AND pit.sat_details_load_ts = det.load_ts

    LEFT JOIN /*+ BROADCAST(scr) */ lakehouse.silver.sat_customer_scoring scr
        ON  pit.hk_customer = scr.hk_customer
        AND pit.sat_scoring_load_ts = scr.load_ts
""")

print("\nGold dim_customer:")
spark.table("lakehouse.gold.dim_customer").orderBy("customer_sk", "snap_date").show(truncate=False)

# ─── 4. Проверяем point-in-time корректность ──────────────────────────────────
print("\n=== CUST-001 на 01 февраля: должен быть в Казани с credit_score=720 ===")
spark.sql("""
    SELECT snap_date, city, credit_score
    FROM lakehouse.gold.dim_customer
    WHERE customer_sk = 'hk_cust_001'
    ORDER BY snap_date
""").show()
# Ожидаем: 2026-02-01 → Казань, 720 | 2026-04-01 → Москва, 720 | 2026-05-23 → Москва, 780

5. Жизненный цикл схемы: Schema Evolution сквозь все слои

Одно из главных преимуществ архитектуры DV + Medallion - контролируемая эволюция схемы. Когда источник добавляет новое поле loyalty_points к клиентской сущности, изменения распространяются чётко по слоям:

Принцип распространения изменений: каждый слой добавляет только то, что относится к его зоне ответственности. Bronze не знает о семантике поля. Raw Vault добавляет его в Satellite как данность. Business Vault решает, влияет ли оно на бизнес-правила. Gold добавляет его в измерение, если оно нужно аналитикам.

Без Data Vault в Silver: изменение источника немедленно требует изменения fact_sales или dim_customer - то есть затрагивает Gold напрямую. С DV: изменение источника затрагивает только Satellite в Raw Vault, а Gold пересчитывается позже в плановом порядке.


6. Операционные паттерны: мониторинг и обслуживание

6.1 Checklist ежедневного обслуживания

Операция Слой Частота Описание
Инжект из Bronze Raw Vault Непрерывно / hourly Hub + Link + Satellite APPEND
Пересчёт PIT Business Vault Ежедневно Добавляем новую снапшотную дату
Пересчёт Computed Satellites Business Vault Ежедневно Обновляем сегментацию, скоринг
Пересчёт Bridge Business Vault При изменении Links Полный пересчёт или инкремент
Пересчёт Gold-витрин Gold Ежедневно DV → Kimball материализация
Compaction Raw Vault + Gold Еженедельно rewrite_data_files в Iceberg
Expire Snapshots Все слои Еженедельно Очистка старых metadata

6.2 Мониторинг качества данных в DV

# Проверяем, что каждый hk_customer в Link существует в Hub
orphan_links = spark.sql("""
    SELECT l.hk_customer
    FROM lakehouse.silver.lnk_order_customer l
    LEFT JOIN lakehouse.silver.hub_customer h ON l.hk_customer = h.hub_hk
    WHERE h.hub_hk IS NULL
""")

if orphan_links.count() > 0:
    print(f"⚠️  Orphan Links: {orphan_links.count()} записей без соответствующего Hub")

# Проверяем, что PIT-таблица покрывает все активные клиентов
pit_coverage = spark.sql("""
    SELECT
        COUNT(DISTINCT hk_customer)               AS pit_customers,
        (SELECT COUNT(*) FROM silver.hub_customer) AS total_customers
    FROM lakehouse.silver.pit_customer
    WHERE snap_date = CURRENT_DATE() - 1
""")
pit_coverage.show()

7. Когда применять полную архитектуру


Итог

Data Vault + Medallion - это не просто сложение двух методологий. Это синергия: каждый паттерн закрывает слабые места другого.

Medallion без Data Vault имеет проблему: при большом количестве источников Silver-слой деградирует в хаос разношёрстных таблиц без единой модели.

Data Vault без Medallion имеет проблему: BI-инструменты не могут эффективно работать с высоконормализованной DV-моделью - слишком много JOIN'ов, слишком медленные запросы.

Вместе они образуют полноценную enterprise-архитектуру:

  • Bronze - иммутабельный журнал всего, что пришло из источников
  • Raw Vault - параллельный, масштабируемый, аудируемый Integration Layer
  • Business Vault (PIT + Bridge + Computed Satellites) - предвычисленные структуры, превращающие дорогие диапазонные JOIN'ы в дешёвые equi-join'ы
  • Gold (Kimball) - физически материализованные витрины для мгновенного BI-доступа

Главный закон: загружай в DV, читай из Кимбалла. Это разделение обязанностей - архитектурный принцип, а не техническая деталь.