Data Vault 2.0 на Spark: Hubs, Links, Satellites и параллельный инжект
Философия Data Vault 2.0 Дэна Линстедта, анатомия Hub/Link/Satellite, детерминированные SHA-256 хэш-ключи без центрального генератора, Independent Load Patterns и параллельная запись в Iceberg, PIT-таблицы и мост к Gold-витринам
Представьте enterprise-компанию с двадцатью разнородными источниками данных: CRM-система Salesforce, ERP-система SAP, собственная e-commerce платформа, call-center система, мобильное приложение, сторонние маркетплейсы. Каждый источник имеет свою схему, свои идентификаторы, свои форматы дат. Источники меняются: в CRM добавили новое поле, ERP поменяла формат даты, e-commerce запустила новую категорию товаров.
В этой ситуации Star Schema по Кимбаллу начинает трещать по швам. Каждое изменение источника требует переосмысления и переписывания центральной таблицы фактов. Добавление нового источника означает каскадные изменения во всех таблицах измерений. Загрузить данные параллельно невозможно: нужно сначала загрузить Dimension (чтобы получить Surrogate Key), только потом - Fact.
Data Vault 2.0 (автор - Дэн Линстедт, 2000–2013) создавался именно для этого сценария. Его архитектурный манифест: система никогда не должна меняться при добавлении нового источника - только расширяться. А физическая конструкция из отдельных независимых таблиц обеспечивает полный параллелизм загрузки, что идеально ложится на архитектуру Apache Spark.
1. Почему Кимбалл и 3NF не справляются в Enterprise¶
1.1 Проблема гибкости Star Schema¶
Рассмотрим конкретный сценарий: в e-commerce появилась новая бизнес-модель - совместные покупки, когда один заказ оформляется от нескольких клиентов. В Кимбалловской Star Schema fact_orders содержит customer_sk как единственный Foreign Key к dim_customer. Один заказ - один клиент.
Новая бизнес-логика означает: нужно добавить co_customer_sk в fact_orders? Или создать отдельную таблицу fact_order_customers? И то и другое требует изменения существующей схемы, перезаписи исторических данных и обновления всех BI-запросов, которые уже работают с fact_orders.
В Data Vault 2.0 добавление новой связи «заказ - совладелец» - это создание нового Link (таблицы связей) между существующими hub_order и hub_customer. Существующие таблицы не меняются. Исторические данные сохраняются. BI-запросы продолжают работать. Новая связь просто начинает наполняться с этого момента.
1.2 Проблема последовательной загрузки¶
В Кимбалловском подходе существует жёсткая зависимость при загрузке: чтобы записать строку в fact_sales, нужно знать customer_sk. Чтобы узнать customer_sk, нужно сначала убедиться, что клиент существует в dim_customer. Это означает последовательную загрузку: Dimension → Fact.
В Data Vault 2.0 все четыре таблицы (hub_customer, hub_order, lnk_order_customer, sat_customer_details) вычисляют свои ключи локально - это SHA-256 хэш бизнес-ключа. Каждый Executor вычисляет хэш независимо и параллельно. Никаких JOIN'ов для получения Surrogate Key из уже существующих таблиц.
2. Три кита Data Vault 2.0¶
2.1 Hub: бизнес-идентичность без контекста¶
Hub - самая простая и самая важная таблица в Data Vault. Она хранит одну вещь: уникальный бизнес-ключ сущности. Ничего лишнего.
Структура Hub предельно минималистична:
| Колонка | Тип | Описание |
|---|---|---|
hub_hash_key |
STRING (SHA-256) | Суррогатный ключ - хэш от бизнес-ключа |
business_key |
STRING | Оригинальный бизнес-ключ из источника |
load_ts |
TIMESTAMP | Когда запись впервые появилась в DV |
record_source |
STRING | Из какой системы пришёл бизнес-ключ |
Всё. В Hub нет имени клиента, нет его города, нет категории продукта. Только идентификатор и метаданные о его происхождении.
Зачем такая строгость? Hub отвечает на единственный вопрос: «Существует ли эта бизнес-сущность в системе?». Он неизменен по природе: если клиент CUST-99 существует, его Hub-запись создаётся один раз и никогда не обновляется. Атрибуты (имя, город, сегмент) живут в Satellite.
2.2 Link: связи как первоклассные граждане¶
Link - таблица, которая хранит бизнес-связи между Hub'ами. Если Hub - это «кто существует», то Link - это «кто с кем связан».
Структура Link:
| Колонка | Тип | Описание |
|---|---|---|
link_hash_key |
STRING (SHA-256) | Хэш комбинации всех бизнес-ключей |
hub_customer_hk |
STRING | FK → hub_customer |
hub_order_hk |
STRING | FK → hub_order |
load_ts |
TIMESTAMP | Когда связь впервые зафиксирована |
record_source |
STRING | Источник |
Важнейшее свойство Link: он тоже append-only и никогда не обновляется. Связь «клиент сделал заказ» - это факт, который не меняется. Если заказ отменён - это новое событие, которое попадёт в Satellite (атрибут заказа), а не удалит запись в Link.
Другая ключевая особенность: Link может связывать более двух Hub'ов. Например, lnk_order_product_store связывает три сущности: заказ, продукт и магазин. В Кимбалле это реализуется через дополнительные колонки в Fact Table. В Data Vault - через новый Link, не трогая существующие таблицы.
2.3 Satellite: атрибуты и история¶
Satellite - это место, где живут описательные атрибуты и история их изменений. Каждый Satellite привязан ровно к одному Hub или одному Link через их Hash Key.
Структура Satellite для Hub'а:
| Колонка | Тип | Описание |
|---|---|---|
hub_hash_key |
STRING | FK → родительский Hub |
load_ts |
TIMESTAMP | Когда эта версия атрибутов появилась |
load_end_ts |
TIMESTAMP | Когда эта версия перестала быть актуальной (NULL = текущая) |
hash_diff |
STRING | SHA-256 хэш всех атрибутов - признак изменения |
record_source |
STRING | Источник |
full_name |
STRING | ... бизнес-атрибуты ... |
city |
STRING | |
segment |
STRING |
Satellite - единственное место в Data Vault, где хранится история изменений. При получении новой версии атрибутов (клиент переехал) создаётся новая строка в Satellite, а старая остаётся без изменений. Это похоже на SCD Type 2, но без is_current флага - вместо него используется load_end_ts = NULL для текущей версии.
hash_diff - ключевое поле Satellite. Это SHA-256 от конкатенации всех атрибутов. Перед тем как вставить новую строку, система проверяет: равен ли hash_diff входящих данных тому, что уже есть в Satellite? Если да - данные не изменились, вставка не нужна. Это позволяет избежать дублей без дорогого колонка-за-колонкой сравнения.
3. Архитектура Data Vault 2.0 на практике¶
3.1 Полная схема для e-commerce¶
Каждый Hub, Link и Satellite - это отдельная физическая таблица Iceberg. Новый источник данных - просто новый Satellite или новый Link, не трогающий существующие таблицы. Именно в этом заключается гибкость Data Vault.
4. Магия SHA-256 хэш-ключей¶
4.1 От последовательных номеров к детерминированным хэшам¶
Data Vault 1.0 (Линстедт, 2000) использовал последовательные целочисленные Surrogate Keys - такие же, как в классическом DWH. Это требовало централизованного генератора (База данных или sequence), что создавало узкое горло при параллельной загрузке.
Data Vault 2.0 (2013) совершил революционный сдвиг: все ключи генерируются из бизнес-ключей через детерминированную хэш-функцию. sha2('CUST-99', 256) всегда вернёт одно и то же значение на любом Executor'е в любой момент времени. Никакой координации. Никакого центрального генератора.
4.2 Правила хэширования бизнес-ключей¶
Хэширование должно быть строго детерминированным: один и тот же бизнес-ключ всегда даёт один и тот же хэш. Для этого необходима стандартизация перед хэшированием:
from pyspark.sql import functions as F
# Правило 1: Нормализация кейса и пробелов
customer_hk = F.sha2(
F.trim(F.lower(F.col("customer_id"))),
256
)
# 'CUST-99', ' cust-99 ', 'Cust-99' - все дают одинаковый хэш
# Правило 2: Составной бизнес-ключ - конкатенация с фиксированным разделителем
# '||' как разделитель избегает коллизий: concat('AB','CD') vs concat('A','BCD')
order_product_lnk_hk = F.sha2(
F.concat_ws("||",
F.trim(F.lower(F.col("order_id"))),
F.trim(F.lower(F.col("product_id")))
),
256
)
# Правило 3: NULL в бизнес-ключе - НИКОГДА не хэшировать
# coalesce нельзя: sha2('') даст один хэш для всех NULL-клиентов → Data Skew!
# Правильно: фильтровать битые записи ДО хэширования
valid_df = bronze_df.filter(F.col("customer_id").isNotNull())
dead_letter_df = bronze_df.filter(F.col("customer_id").isNull())
Почему NULL опасен: если customer_id равен NULL для 10 000 строк, и мы делаем sha2(coalesce(customer_id, ''), 256), все 10 000 строк получают одинаковый хэш sha2('', 256). В Hub'е появится одна «фантомная» запись, к которой привяжутся тысячи Satellite-строк. При загрузке Spark произведёт Data Skew: весь этот трафик пойдёт в одну партицию. А аналитик получит «клиента без имени» с тысячами заказов.
4.3 Hash Diff: элегантное сравнение версий¶
hash_diff в Satellite - это SHA-256 от конкатенации всех бизнес-атрибутов. Он позволяет за O(1) определить, изменились ли данные:
# Считаем hash_diff как хэш от всех атрибутов сателлита
hash_diff = F.sha2(
F.concat_ws("||",
F.coalesce(F.col("full_name"), F.lit("NULL")),
F.coalesce(F.col("city"), F.lit("NULL")),
F.coalesce(F.col("segment"), F.lit("NULL")),
F.coalesce(F.col("email"), F.lit("NULL")),
),
256
)
Зачем coalesce(col, lit("NULL")): если использовать просто concat_ws, то NULL-поле просто пропускается. concat_ws('||', 'John', NULL, 'PREMIUM') = 'John||PREMIUM' - и concat_ws('||', 'John', 'Moscow', 'PREMIUM') даст другой результат. Но concat_ws('||', 'John', NULL, 'PREMIUM') равен concat_ws('||', 'John', 'Default', 'PREMIUM') при использовании coalesce с одинаковым значением. Поэтому явный плейсхолдер "NULL" делает сравнение предсказуемым.
При загрузке нового батча:
# Новые данные сравниваем с последней версией в Satellite
# Вставляем только те строки, где hash_diff изменился
new_records = incoming_df.join(
last_version_df.select("hub_customer_hk", "hash_diff"),
on="hub_customer_hk",
how="left_anti" # строки, которых нет в last_version
).union(
incoming_df.join(
last_version_df.select("hub_customer_hk", "hash_diff"),
on=["hub_customer_hk"],
how="inner"
).filter(F.col("incoming_hash_diff") != F.col("existing_hash_diff"))
)
5. Independent Load Patterns: параллельность как архитектурный принцип¶
5.1 Staging слой - общий для всех¶
Перед загрузкой в DV-таблицы создаётся Staging DataFrame: единый подготовленный набор данных, из которого все таблицы (Hub, Link, Satellite) читают независимо.
Staging DataFrame кэшируется (df.cache()) перед параллельным запуском нескольких Spark-action'ов. Без кэша каждая ветка заново читала бы Bronze - три чтения одних и тех же данных.
5.2 APPEND вместо MERGE: почему DV быстрее¶
Главное операционное преимущество Data Vault перед Кимбалловским пайплайном: Hub и Link используют APPEND (добавление), а не MERGE (обновление существующих строк).
Почему это значимо:
- MERGE в Iceberg (Copy-on-Write) перезаписывает Parquet-файлы, содержащие обновлённые строки. При частых обновлениях - постоянный Write Amplification.
- APPEND просто добавляет новый Parquet-файл. Операция O(1) с точки зрения записи.
- Дедупликация Hub'а происходит через
LEFT ANTI JOINсо Staging: вставляем только те бизнес-ключи, которых ещё нет в Hub.
# Эффективный APPEND с дедупликацией для Hub
existing_keys = spark.table("lakehouse.silver.hub_customer") \
.select("business_key")
new_hub_rows = staged_df \
.select("hk_customer", "customer_id", "load_ts", "record_source") \
.distinct() \
.join(existing_keys,
staged_df["customer_id"] == existing_keys["business_key"],
how="left_anti") # только те, которых ещё нет
new_hub_rows.writeTo("lakehouse.silver.hub_customer").append()
left_anti возвращает строки из левой таблицы, для которых нет совпадений в правой. Это эффективный способ вставить только новые бизнес-ключи, не трогая существующие.
6. Практика: полный Data Vault Pipeline¶
6.1 Полный код¶
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = SparkSession.builder \
.appName("DataVault-2.0-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()
RECORD_SOURCE = "ecom_orders_api"
# ─── 0. Создаём DV-таблицы ────────────────────────────────────────────────────
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.silver.hub_customer (
hub_hk STRING COMMENT 'SHA-256 от business_key',
business_key STRING COMMENT 'Оригинальный customer_id из источника',
load_ts TIMESTAMP,
record_source STRING
) USING iceberg
PARTITIONED BY (days(load_ts))
""")
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.silver.hub_order (
hub_hk STRING,
business_key STRING,
load_ts TIMESTAMP,
record_source STRING
) USING iceberg
PARTITIONED BY (days(load_ts))
""")
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.silver.hub_product (
hub_hk STRING,
business_key STRING,
load_ts TIMESTAMP,
record_source STRING
) USING iceberg
PARTITIONED BY (days(load_ts))
""")
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.silver.lnk_order_customer (
lnk_hk STRING COMMENT 'SHA-256 от order_id||customer_id',
hk_order STRING COMMENT 'FK → hub_order',
hk_customer STRING COMMENT 'FK → hub_customer',
load_ts TIMESTAMP,
record_source STRING
) USING iceberg
PARTITIONED BY (days(load_ts))
""")
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.silver.lnk_order_product (
lnk_hk STRING,
hk_order STRING,
hk_product STRING,
load_ts TIMESTAMP,
record_source STRING
) USING iceberg
PARTITIONED BY (days(load_ts))
""")
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.silver.sat_customer_details (
hk_customer STRING COMMENT 'FK → hub_customer',
load_ts TIMESTAMP COMMENT 'Начало действия этой версии',
load_end_ts TIMESTAMP COMMENT 'NULL = текущая версия',
hash_diff STRING COMMENT 'SHA-256 от всех атрибутов',
record_source STRING,
full_name STRING,
city STRING,
segment STRING,
email STRING
) USING iceberg
PARTITIONED BY (days(load_ts))
""")
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.silver.sat_order_context (
hk_order STRING,
load_ts TIMESTAMP,
load_end_ts TIMESTAMP,
hash_diff STRING,
record_source STRING,
product_name STRING,
amount DECIMAL(10,2),
status STRING
) USING iceberg
PARTITIONED BY (days(load_ts))
""")
# ─── 1. Сырые данные из Bronze ────────────────────────────────────────────────
bronze_data = [
# order_id, customer_id, product_id, customer_name, city, segment, email, product_name, amount, status
("ORD-100", "CUST-99", "PROD-001", "Ivan Ivanov", "Kazan", "STANDARD", "ivan@m.ru", "Laptop", 1200.0, "CREATED"),
("ORD-100", "CUST-99", "PROD-002", "Ivan Ivanov", "Kazan", "STANDARD", "ivan@m.ru", "Mouse", 50.0, "CREATED"),
("ORD-101", "CUST-42", "PROD-001", "Alice Smith", "London", "PREMIUM", "alice@m.com", "Laptop", 990.0, "PAID"),
# Битая запись: customer_id = NULL → должна попасть в DLQ
("ORD-102", None, "PROD-003", None, None, None, None, "Book", 15.0, "CREATED"),
]
cols = ["order_id", "customer_id", "product_id", "customer_name", "city",
"segment", "email", "product_name", "amount", "status"]
bronze_df = spark.createDataFrame(bronze_data, cols)
# ─── 2. Валидация: фильтрация битых записей в DLQ ─────────────────────────────
# Ключевое правило: НИКОГДА не хэшировать NULL бизнес-ключи
valid_df = bronze_df.filter(F.col("customer_id").isNotNull())
dead_letter_df = bronze_df.filter(F.col("customer_id").isNull())
if dead_letter_df.count() > 0:
print(f"⚠️ DLQ: {dead_letter_df.count()} записей с NULL business_key отфильтрованы")
# В production: сохранить dead_letter_df в отдельную таблицу для разбора
# ─── 3. Staging: вычисляем все хэши локально ─────────────────────────────────
load_ts = F.current_timestamp()
staged_df = valid_df \
.withColumn("hk_customer",
F.sha2(F.trim(F.lower(F.col("customer_id"))), 256)) \
.withColumn("hk_order",
F.sha2(F.trim(F.lower(F.col("order_id"))), 256)) \
.withColumn("hk_product",
F.sha2(F.trim(F.lower(F.col("product_id"))), 256)) \
.withColumn("lnk_order_customer_hk",
F.sha2(F.concat_ws("||",
F.trim(F.lower(F.col("order_id"))),
F.trim(F.lower(F.col("customer_id")))
), 256)) \
.withColumn("lnk_order_product_hk",
F.sha2(F.concat_ws("||",
F.trim(F.lower(F.col("order_id"))),
F.trim(F.lower(F.col("product_id")))
), 256)) \
.withColumn("hash_diff_customer",
F.sha2(F.concat_ws("||",
F.coalesce(F.col("customer_name"), F.lit("NULL")),
F.coalesce(F.col("city"), F.lit("NULL")),
F.coalesce(F.col("segment"), F.lit("NULL")),
F.coalesce(F.col("email"), F.lit("NULL")),
), 256)) \
.withColumn("hash_diff_order",
F.sha2(F.concat_ws("||",
F.coalesce(F.col("product_name"), F.lit("NULL")),
F.coalesce(F.col("amount").cast("string"), F.lit("NULL")),
F.coalesce(F.col("status"), F.lit("NULL")),
), 256)) \
.withColumn("load_ts", load_ts) \
.withColumn("record_source", F.lit(RECORD_SOURCE))
staged_df.cache() # кэшируем - от него пойдут параллельные ветки
# ─── 4. Load Pattern: Hub Customer ───────────────────────────────────────────
existing_cust_keys = spark.table("lakehouse.silver.hub_customer").select("business_key")
hub_customer_new = staged_df \
.select("hk_customer", "customer_id", "load_ts", "record_source") \
.withColumnRenamed("hk_customer", "hub_hk") \
.withColumnRenamed("customer_id", "business_key") \
.distinct() \
.join(existing_cust_keys, on="business_key", how="left_anti")
hub_customer_new.writeTo("lakehouse.silver.hub_customer").append()
print(f"hub_customer: +{hub_customer_new.count()} новых записей")
# ─── 5. Load Pattern: Hub Order ───────────────────────────────────────────────
existing_order_keys = spark.table("lakehouse.silver.hub_order").select("business_key")
hub_order_new = staged_df \
.select("hk_order", "order_id", "load_ts", "record_source") \
.withColumnRenamed("hk_order", "hub_hk") \
.withColumnRenamed("order_id", "business_key") \
.distinct() \
.join(existing_order_keys, on="business_key", how="left_anti")
hub_order_new.writeTo("lakehouse.silver.hub_order").append()
# ─── 6. Load Pattern: Hub Product ─────────────────────────────────────────────
existing_prod_keys = spark.table("lakehouse.silver.hub_product").select("business_key")
hub_product_new = staged_df \
.select("hk_product", "product_id", "load_ts", "record_source") \
.withColumnRenamed("hk_product", "hub_hk") \
.withColumnRenamed("product_id", "business_key") \
.distinct() \
.join(existing_prod_keys, on="business_key", how="left_anti")
hub_product_new.writeTo("lakehouse.silver.hub_product").append()
# ─── 7. Load Pattern: Link Order-Customer ────────────────────────────────────
existing_loc_links = spark.table("lakehouse.silver.lnk_order_customer").select("lnk_hk")
lnk_order_customer_new = staged_df \
.select("lnk_order_customer_hk", "hk_order", "hk_customer", "load_ts", "record_source") \
.withColumnRenamed("lnk_order_customer_hk", "lnk_hk") \
.distinct() \
.join(existing_loc_links, on="lnk_hk", how="left_anti")
lnk_order_customer_new.writeTo("lakehouse.silver.lnk_order_customer").append()
# ─── 8. Load Pattern: Link Order-Product ─────────────────────────────────────
existing_lop_links = spark.table("lakehouse.silver.lnk_order_product").select("lnk_hk")
lnk_order_product_new = staged_df \
.select("lnk_order_product_hk", "hk_order", "hk_product", "load_ts", "record_source") \
.withColumnRenamed("lnk_order_product_hk", "lnk_hk") \
.distinct() \
.join(existing_lop_links, on="lnk_hk", how="left_anti")
lnk_order_product_new.writeTo("lakehouse.silver.lnk_order_product").append()
# ─── 9. Load Pattern: Satellite Customer Details ──────────────────────────────
# Для Satellite: вставляем только строки с изменённым hash_diff
existing_sat_cust = spark.table("lakehouse.silver.sat_customer_details") \
.filter(F.col("load_end_ts").isNull()) \
.select("hk_customer", "hash_diff")
sat_customer_new = staged_df \
.select("hk_customer", "load_ts", "hash_diff_customer",
"record_source", "customer_name", "city", "segment", "email") \
.withColumn("load_end_ts", F.lit(None).cast("timestamp")) \
.withColumnRenamed("hash_diff_customer", "hash_diff") \
.withColumnRenamed("customer_name", "full_name") \
.distinct() \
.join(existing_sat_cust,
on=["hk_customer", "hash_diff"],
how="left_anti") # только строки с изменившимися атрибутами
sat_customer_new.writeTo("lakehouse.silver.sat_customer_details").append()
# ─── 10. Load Pattern: Satellite Order Context ────────────────────────────────
existing_sat_order = spark.table("lakehouse.silver.sat_order_context") \
.filter(F.col("load_end_ts").isNull()) \
.select("hk_order", "hash_diff")
sat_order_new = staged_df \
.select("hk_order", "load_ts", "hash_diff_order",
"record_source", "product_name", "amount", "status") \
.withColumn("load_end_ts", F.lit(None).cast("timestamp")) \
.withColumnRenamed("hash_diff_order", "hash_diff") \
.withColumn("amount", F.col("amount").cast("decimal(10,2)")) \
.distinct() \
.join(existing_sat_order,
on=["hk_order", "hash_diff"],
how="left_anti")
sat_order_new.writeTo("lakehouse.silver.sat_order_context").append()
staged_df.unpersist()
print("\n✅ Data Vault 2.0 загрузка завершена!")
# ─── 11. Пример реконструкции сущности ───────────────────────────────────────
print("\n=== Реконструкция: текущее состояние клиентов ===")
spark.sql("""
SELECT
h.business_key AS customer_id,
s.full_name,
s.city,
s.segment,
s.load_ts
FROM lakehouse.silver.hub_customer h
JOIN lakehouse.silver.sat_customer_details s
ON h.hub_hk = s.hk_customer
WHERE s.load_end_ts IS NULL
ORDER BY h.business_key
""").show(truncate=False)
7. От Data Vault к Gold: PIT-таблицы и Business Vault¶
7.1 Проблема BI-запросов в Data Vault¶
Data Vault - отличный интеграционный слой, но плохой serving слой. Чтобы получить «профиль клиента», аналитик должен написать запрос с JOIN'ами трёх-четырёх таблиц. BI-инструменты не справляются с такой сложностью нативно.
Решение - не менять Data Vault, а построить поверх него Gold-слой в виде классической Star Schema. Spark читает Hub + Link + Satellite, выполняет денормализацию и записывает результат в gold.dim_customer и gold.fact_sales.
7.2 PIT-таблица: Point-in-Time снапшот¶
При наличии нескольких Satellite'ов для одного Hub'а возникает проблема JOIN'а: Satellite'ы обновляются в разное время. Чтобы получить консистентный «снапшот» клиента на конкретную дату, нужно JOIN'ить каждый Satellite по условию load_ts <= snapshot_date AND load_end_ts > snapshot_date.
PIT-таблица (Point-in-Time) - предвычисленная вспомогательная таблица, которая для каждого Hub-ключа на каждую дату снапшота хранит точные load_ts из каждого Satellite:
-- PIT-таблица для hub_customer
SELECT
hc.hub_hk AS hk_customer,
snap.snap_date,
-- Для каждого satellite: ближайший load_ts не позже snap_date
MAX(sc.load_ts) AS sat_cust_details_load_ts,
MAX(se.load_ts) AS sat_cust_email_load_ts
FROM hub_customer hc
CROSS JOIN snapshot_dates snap
LEFT JOIN sat_customer_details sc
ON hc.hub_hk = sc.hk_customer
AND sc.load_ts <= snap.snap_date
AND (sc.load_end_ts > snap.snap_date OR sc.load_end_ts IS NULL)
LEFT JOIN sat_customer_email se
ON hc.hub_hk = se.hk_customer
AND se.load_ts <= snap.snap_date
AND (se.load_end_ts > snap.snap_date OR se.load_end_ts IS NULL)
GROUP BY hc.hub_hk, snap.snap_date
С PIT-таблицей JOIN в Gold-пайплайне становится простым и эффективным: используем заранее вычисленные load_ts вместо условия диапазона.
8. Когда Data Vault нужен, а когда избыточен¶
Data Vault нужен, когда:
- 10+ разнородных источников с разными схемами и идентификаторами
- Данные активно меняются, новые источники добавляются ежеквартально
- Требования финансового регулятора или GDPR требуют полного audit trail каждого поля
- Команда достаточно зрелая для поддержки архитектурной сложности
Data Vault избыточен, когда:
- Стартап с одним-двумя источниками и быстро меняющейся бизнес-логикой
- Команда из 2–3 инженеров без опыта DV
- Аналитика не требует point-in-time запросов на уровне отдельных атрибутов
- Схема источников стабильна и не меняется часто
Итог¶
Data Vault 2.0 - это не просто методология моделирования, это архитектурный принцип для высококонкурентных enterprise-платформ. Его ключевые свойства идеально соответствуют природе Apache Spark:
- SHA-256 хэш-ключи вместо централизованных последовательностей - каждый Executor вычисляет ключи независимо, без сетевой координации
- Independent Load Patterns - Hub, Link и Satellite загружаются параллельно из одного Staging DataFrame
- Append-only семантика для Hub и Link - минимальный Write Amplification в Iceberg
- hash_diff в Satellite - O(1) проверка изменений вместо колонка-за-колонкой сравнения
- Расширяемость без модификации - новый источник добавляет новый Satellite, не меняя существующих таблиц
Data Vault не заменяет Кимбалла - он дополняет его. Data Vault - это интеграционный Silver-слой. Star Schema - это serving Gold-слой. PIT-таблицы - мост между ними.