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
Мы изучили 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.
Это катастрофа по трём причинам:
-
Оконные функции (
ROW_NUMBER() OVER (PARTITION BY hk_customer ORDER BY load_ts DESC)) - это Wide Transformation с Shuffle. При каждом BI-запросе Spark перемешивает всю историю Satellite'ов по сети. -
Конкурентность - если 50 аналитиков одновременно обновляют свои дашборды, запущены 50 параллельных тяжёлых Spark-заданий. Кластер перегружен.
-
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, читай из Кимбалла. Это разделение обязанностей - архитектурный принцип, а не техническая деталь.