Data Vault 2.0 на Spark: Hubs, Links, Satellites и параллельный инжект

Философия Data Vault 2.0 Дэна Линстедта, анатомия Hub/Link/Satellite, детерминированные SHA-256 хэш-ключи без центрального генератора, Independent Load Patterns и параллельная запись в Iceberg, PIT-таблицы и мост к Gold-витринам

lakehouse data-vault dv20 hubs satellites

Представьте 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.

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-таблицы - мост между ними.