Dimensional Modeling (Kimball) на Spark: Star Schema, Snowflake, Grain, Surrogate Keys

Классическая многомерная модель Кимбалла в Data Lakehouse: проектирование таблиц фактов и измерений, определение зерна (Grain), генерация суррогатных ключей через SHA-256, Star vs Snowflake Schema, Broadcast Join и Dynamic Partition Pruning

lakehouse dimensional-modeling star-schema surrogate-keys

Многие инженеры, переходя из классических DWH (Oracle, Teradata, Greenplum) в мир Data Lakehouse (Apache Iceberg, Delta Lake, Apache Spark), совершают одну и ту же ошибку: выбрасывают всё, что знали о проектировании схем данных. «Теперь у нас S3 и Parquet, значит, сложим всё в одну гигантскую плоскую таблицу - и дело с концом».

Результат предсказуем: через полгода Gold Layer превращается в нагромождение несвязанных таблиц, BI-аналитики пишут запросы с пятью JOIN'ами к сырым таблицам, KPI в разных дашбордах не сходятся, а каждый новый отчёт требует недель разработки.

Классическая многомерная модель по Ральфу Кимбаллу (Dimensional Modeling) - не устаревший реликт эпохи Teradata. Это по-прежнему наиболее эффективный способ организовать аналитические данные так, чтобы BI-инструменты работали быстро, метрики считались правильно, а новые отчёты писались за часы, а не за недели. В Data Lakehouse эта модель получает новое дыхание: Spark умеет эффективно выполнять именно те паттерны Join'ов, которые Кимбалл описал ещё в 1990-х.


1. Зачем Кимбалл в Data Lakehouse

1.1 Миф о «плоской таблице»

One Big Table (OBT) - денормализованная таблица, где все атрибуты собраны в одну строку - кажется простым решением. Зачем JOIN'ить customer_dim к fact_sales, если можно хранить имя клиента прямо в таблице продаж?

У этого подхода есть реальные преимущества: нет JOIN-overhead, простые запросы, понятная структура. Но у него есть фундаментальные ограничения, которые становятся критичными на больших объёмах:

Проблема OBT Последствие
Дублирование атрибутов измерений customer_name повторяется в каждой строке продаж - 100 млн продаж × клиент с именем «Александр Николаевич Петров-Водкин»
Невозможность Broadcast Join Spark не может «скопировать» OBT в память Executor'ов - она слишком большая
Медленное изменение атрибутов Обновить customer_city для клиента - переписать 100 млн строк продаж
Смешение зерна данных Одна строка содержит и детали транзакции, и агрегаты по клиенту - это катастрофа
Dynamic Partition Pruning не работает DPP срабатывает только при Join с отдельной Dimension таблицей

Последний пункт особенно важен: механизм Dynamic Partition Pruning (DPP) в Spark SQL, который позволяет читать только нужные партиции Fact Table при фильтрации по измерению, работает только при наличии отдельной Dimension таблицы. При OBT Spark вынужден сканировать всю таблицу полностью.

1.2 Бизнес-ориентированность многомерной модели

Многомерная модель Кимбалла разделяет данные на два фундаментально разных типа:

  • Факты - «что произошло?» - событие, транзакция, измерение. Цифры, которые суммируются, усредняются, считаются. Выручка, количество кликов, продолжительность сессии.
  • Измерения - «в каком контексте это произошло?» - описательные атрибуты события. Кто купил, что купили, где, когда, через какой канал.

Аналитик, работающий с такой моделью, мыслит естественно: «покажи мне выручку (факт) по регионам (измерение) за последний квартал (измерение даты)». Структура таблиц отражает эту ментальную модель.


2. Анатомия модели: Grain, Факты и Измерения

2.1 Grain - самый важный шаг проектирования

Grain (зерно) - это точное определение того, что представляет одна строка в таблице фактов. Это не техническое решение - это бизнес-решение, которое принимается ещё до написания первой строки кода.

Зерно должно быть сформулировано на языке бизнеса, предельно конкретно:

«Одна строка = одна позиция товара в одном заказе в момент оформления»

Это не то же самое, что:

«Одна строка = один заказ» (другое зерно - агрегированнее) «Одна строка = одно сканирование штрих-кода на кассе» (другое зерно - детальнее)

Почему зерно так критично? Смешение строк разного зерна в одной таблице немедленно ломает любые агрегаты.

В левом варианте total_revenue - это атрибут всего заказа (зерно «заказ»), но он помещён в таблицу с зерном «позиция». При суммировании по позициям мы дважды учитываем выручку заказа. Это классическая ловушка, в которую попадают новички.

2.2 Таблицы фактов: числа и ключи

Таблица фактов содержит только два типа колонок:

  1. Внешние ключи (FK) к таблицам измерений - суррогатные ключи, указывающие на конкретный контекст события
  2. Числовые метрики - то, что мы будем агрегировать, усреднять, считать

Никаких текстовых описаний в таблице фактов! Если видите колонку customer_name в Fact Table - это нарушение модели.

Числовые метрики классифицируются по аддитивности:

Тип метрики Описание Пример Суммируемость
Аддитивная Суммируется по любому измерению revenue, quantity_sold По любому срезу
Полуаддитивная Суммируется только по некоторым измерениям account_balance, inventory_level По транзакциям, но не по времени
Неаддитивная Нельзя суммировать price_per_unit, conversion_rate Требует пересчёта через AVG/RATIO

Неаддитивные метрики особенно коварны. conversion_rate = 5% по каналу A и 8% по каналу B не суммируются в 13% по обоим каналам - нужно брать взвешенное среднее.

2.3 Таблицы измерений: широкие и маленькие

Таблица измерений - полная противоположность таблице фактов по своим характеристикам:

  • Широкая - десятки и сотни колонок с описательными атрибутами. Чем больше контекста, тем богаче возможности BI-срезов.
  • Небольшая по строкам - клиентов у компании может быть 10 млн, продуктов - 100 тысяч. Это несравнимо с миллиардами строк в Fact Table.
  • Денормализованная - все атрибуты одной концепции собраны в одну таблицу, даже если они образуют иерархию.

Пример богатого dim_customer:

CREATE TABLE lakehouse.gold.dim_customer (
    customer_sk      STRING     COMMENT 'Суррогатный ключ (SHA-256)',
    src_customer_id  STRING     COMMENT 'Оригинальный ID из CRM',
    full_name        STRING,
    email            STRING,
    phone            STRING,
    city             STRING,
    region           STRING,
    country          STRING,
    country_iso      STRING,
    postal_code      STRING,
    segment          STRING     COMMENT 'PREMIUM | STANDARD | TRIAL',
    acquisition_date DATE,
    acquisition_channel STRING,
    age_group        STRING     COMMENT '18-24 | 25-34 | 35-44 | 45+',
    gender           STRING,
    is_active        BOOLEAN,
    lifetime_orders  INT        COMMENT 'Денормализованный агрегат из Silver'
) USING iceberg

Обратите внимание: lifetime_orders - это денормализованный агрегат. В строгой нормализации он не должен быть в измерении. Но в Кимбалловской модели это допустимо и даже желательно - аналитик получает контекст прямо в JOIN'е, без дополнительного агрегирующего подзапроса.


3. Star Schema vs Snowflake Schema

3.1 Схема «Звезда» (Star Schema)

Star Schema - это когда таблица фактов окружена полностью денормализованными таблицами измерений. Иерархия атрибутов (город → регион → страна) сплющена в одну плоскую Dimension таблицу.

Почему Star Schema эффективна в Spark:

  • Broadcast Join - Dimension таблицы небольшие (мегабайты), Spark копирует их целиком в память каждого Executor'а. JOIN с Fact Table происходит полностью локально, без Shuffle по сети.
  • Dynamic Partition Pruning - при фильтре WHERE customer.city = 'Москва' Spark сначала находит все customer_sk московских клиентов из dim_customer, затем читает из fact_sales только те партиции, где встречаются эти ключи.
  • Columnar scan - Fact Table читается колонками: если запрос обращается только к revenue и customer_sk, остальные колонки (quantity, discount_amount) не читаются с диска.

3.2 Схема «Снежинка» (Snowflake Schema)

В Snowflake Schema измерения дополнительно нормализуются: иерархические атрибуты выносятся в отдельные таблицы. Например, вместо dim_product с колонками category и subcategory создаётся отдельная dim_category, и dim_product содержит только внешний ключ category_sk.

Почему Snowflake - анти-паттерн для Spark:

Запрос «продажи по отделам» требует цепочки JOIN'ов: fact_sales → dim_product → dim_category → dim_department. Каждый JOIN в Spark - это потенциальный SortMergeJoin с Shuffle. Три промежуточных JOIN'а вместо одного - это три Shuffle'а вместо нуля (при Broadcast Join).

Кроме того, таблицы dim_category и dim_department ещё меньше, чем dim_product, но их Broadcast Join с промежуточными результатами работает менее эффективно, чем один Broadcast Join целой денормализованной dim_product.

Вывод: в Spark всегда выбирайте Star Schema. Денормализуйте измерения в плоскую таблицу. Экономия места от нормализации на современных объёмах хранения несущественна, а потеря производительности - очень даже существенна.


4. Суррогатные ключи в распределённой системе

4.1 Зачем нужны суррогатные ключи

Суррогатный ключ (Surrogate Key, SK) - это технический идентификатор строки в Dimension таблице, созданный системой DWH независимо от источника данных. Он не несёт бизнес-смысла и генерируется инженером.

Зачем он нужен, если у нас уже есть бизнес-ключ из источника (customer_id из CRM)?

  1. Изоляция от источника - CRM может переформатировать customer_idC-12345 на 12345). Все Fact таблицы по-прежнему ссылаются на customer_sk, а не на источниковый ключ.
  2. Поддержка SCD Type 2 - когда клиент меняет город, создаётся новая строка в dim_customer с новым customer_sk. Исторические факты остаются привязаны к старому customer_sk (Москва), новые - к новому (Питер).
  3. Унификация нескольких источников - клиент в CRM имеет cust_001, в ERP - ERP-2024-001. Оба могут быть сопоставлены с единым customer_sk.
  4. Уменьшение размера Join - числовой SK занимает 4–8 байт, строковый бизнес-ключ - десятки байт. При миллиардах Join'ов разница существенна.

4.2 Почему нет AUTO_INCREMENT в Spark

Разработчики, пришедшие из реляционных СУБД, первым делом ищут аналог SERIAL или AUTO_INCREMENT. Его нет, и это осознанное решение.

AUTO_INCREMENT в PostgreSQL работает через атомарный счётчик в памяти сервера: каждый INSERT обращается к нему, получает следующее число и инкрементирует счётчик. Это работает, пока есть один сервер.

В Spark сотни Executor'ов параллельно генерируют строки на разных машинах. Чтобы обеспечить глобально уникальные последовательные числа, потребовалась бы централизованная координация через Driver - каждый Executor отправлял бы запрос Driver'у, получал следующее число, ждал ответа. Это превращает параллельные вычисления в последовательные и создаёт узкое горло.

4.3 monotonically_increasing_id(): быстро, но опасно

Spark предоставляет функцию monotonically_increasing_id(), которая возвращает монотонно возрастающие целые числа без Shuffle. Она присваивает номера локально на каждом Executor'е с префиксом из номера партиции.

from pyspark.sql import functions as F

df = spark.createDataFrame([("a",), ("b",), ("c",)], ["val"])
df.withColumn("id", F.monotonically_increasing_id()).show()
# Результат примерно такой:
# | val | id              |
# |-----|-----------------|
# |  a  | 0               |
# |  b  | 1               |
# |  c  | 8589934592      |  ← прыжок - другая партиция

Проблемы для Dimension таблиц:

  • Числа непоследовательны (8589934592 вместо 2) - нельзя предсказать следующее значение
  • При перезапуске с другим количеством партиций те же строки получат другие ID - это катастрофа для Foreign Keys в Fact таблицах
  • Нельзя deterministically присвоить ID новой строке без чтения всей таблицы

monotonically_increasing_id() подходит только для временных технических идентификаторов внутри одного Spark-задания, но не для персистентных Surrogate Keys в DWH.

4.4 SHA-256 хэш как идеальный Surrogate Key

Хэш-функция SHA-256 от бизнес-ключа - лучшее решение для Surrogate Keys в Spark. Каждый Executor вычисляет хэш локально, без каких-либо сетевых обращений.

from pyspark.sql import functions as F

# Бизнес-ключ: src_customer_id
# Суррогатный ключ: sha2(normalized_business_key, 256)

dim_customers = silver_df \
    .select("src_customer_id", "customer_name", "city", "region") \
    .distinct() \
    .withColumn(
        "customer_sk",
        F.sha2(
            F.trim(F.lower(F.col("src_customer_id"))),  # нормализация перед хэшированием
            256
        )
    )

Почему нормализация бизнес-ключа обязательна: sha2("cust_99", 256) ≠ sha2("CUST_99", 256) ≠ sha2(" cust_99", 256). Один и тот же клиент из разных источников может иметь разный регистр или пробелы. trim() + lower() гарантируют, что один клиент всегда получит один и тот же хэш.

Составной бизнес-ключ: когда сущность идентифицируется несколькими полями, конкатенируйте их с разделителем:

# Составной ключ: country_code + region_id
.withColumn(
    "region_sk",
    F.sha2(
        F.concat_ws("|", F.col("country_code"), F.col("region_id")),
        256
    )
)

Разделитель | нужен, чтобы избежать коллизий: concat("AB", "CD") = concat("A", "BCD") = "ABCD", но concat_ws("|", "AB", "CD") = "AB|CD" ≠ "A|BCD".

Отвечаем на вопрос «почему SHA-256 выгоднее INT»:

SHA-256 строка занимает 64 байта - в 8 раз больше, чем INT64. Казалось бы, это расточительство. Но при сравнении с альтернативами:

  • Стоимость S3 хранения колонки с customer_sk в Parquet - копейки (dictionary encoding сжимает повторяющиеся значения до 1–2 байт на строку)
  • Стоимость координации через Driver для генерации последовательных INT - дорого (сетевые round-trip'ы, блокировка Executors)
  • Стоимость ошибки в SK при перезапуске - катастрофически дорого (нарушены все Foreign Keys, нужен полный пересчёт всех Fact таблиц)

Диск дёшев, сеть - нет. SHA-256 - правильный выбор.


5. Join фактов и измерений в Spark

5.1 Broadcast Join: ключевая оптимизация Star Schema

Когда Spark видит JOIN между большой таблицей (терабайты) и маленькой (мегабайты), он может использовать Broadcast Hash Join: целиком скопировать маленькую таблицу в память каждого Executor'а, а затем для каждой строки большой таблицы сделать локальный lookup. Никакого Shuffle.

Spark автоматически использует Broadcast Join, если таблица меньше порогового значения:

# Порог по умолчанию - 10 МБ. Для типичных Dimension таблиц его нужно увеличить:
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 100 * 1024 * 1024)  # 100 МБ

# Явное указание через hint - когда нужно гарантировать Broadcast:
spark.sql("""
    SELECT f.revenue, c.city, p.category
    FROM fact_sales f
    JOIN /*+ BROADCAST(c) */ dim_customer c ON f.customer_sk = c.customer_sk
    JOIN /*+ BROADCAST(p) */ dim_product p  ON f.product_sk  = p.product_sk
""")

5.2 Dynamic Partition Pruning (DPP)

Dynamic Partition Pruning - механизм Spark 3.x, который позволяет читать из Fact Table только те партиции, которые содержат данные для отфильтрованных Dimension строк.

Пример: запрос WHERE customer.city = 'Москва'. Без DPP Spark сканирует всю fact_sales. С DPP:

  1. Spark выполняет subquery к dim_customer и получает список customer_sk для Москвы
  2. Spark проверяет метаданные партиций fact_sales - какие содержат эти customer_sk
  3. Spark читает только нужные партиции

DPP работает автоматически при следующих условиях:

  • Включён spark.sql.optimizer.dynamicPartitionPruning.enabled (по умолчанию true)
  • fact_sales партиционирована
  • JOIN происходит через Broadcast Hash Join
# Убедиться, что DPP включён
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true")

6. Практика: Star Schema в Data Lakehouse

6.1 Полный код: от сырых данных к Star Schema

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

spark = SparkSession.builder \
    .appName("Kimball-Star-Schema") \
    .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") \
    .config("spark.sql.autoBroadcastJoinThreshold", str(100 * 1024 * 1024))
    .getOrCreate()

# ─── 0. Исходные данные (Silver уровень) ─────────────────────────────────────
# Зерно Silver: одна строка = одна позиция в заказе
raw_data = [
    # order_id, date,       cust_id,   name,          city,      product,   price, qty, category
    ("ord-100", "2026-05-22", "cust-99", "John Doe",    "Moscow",  "Laptop",  1200.0, 1, "Electronics"),
    ("ord-100", "2026-05-22", "cust-99", "John Doe",    "Moscow",  "Mouse",     50.0, 2, "Electronics"),
    ("ord-101", "2026-05-23", "cust-42", "Alice Smith",  "Kazan",   "Book",      15.0, 5, "Books"),
    ("ord-102", "2026-05-23", "cust-99", "John Doe",    "Moscow",  "Keyboard",  75.0, 1, "Electronics"),
    ("ord-103", "2026-05-24", "cust-77", "Bob Johnson", "Novosibirsk", "Book",  20.0, 3, "Books"),
]
cols = ["order_id", "order_date", "src_cust_id", "customer_name",
        "city", "product_name", "price", "quantity", "category"]

silver_df = spark.createDataFrame(raw_data, cols)

# ─── 1. Строим dim_customer ───────────────────────────────────────────────────
# Grain Dimension: один клиент = одна строка (текущее состояние)
dim_customer = silver_df \
    .select("src_cust_id", "customer_name", "city") \
    .distinct() \
    .withColumn(
        "customer_sk",
        F.sha2(F.trim(F.lower(F.col("src_cust_id"))), 256)
    ) \
    .select(
        F.col("customer_sk"),
        F.col("src_cust_id").alias("src_customer_id"),
        F.col("customer_name"),
        F.col("city"),
    )

# Сохраняем Dimension - НЕ партиционируем (маленькая таблица)
dim_customer.writeTo("lakehouse.gold.dim_customer") \
    .using("iceberg") \
    .createOrReplace()

print("dim_customer:")
dim_customer.show(truncate=False)

# ─── 2. Строим dim_product ────────────────────────────────────────────────────
# В реальном проекте - из Silver product таблицы. Здесь имитируем из заказов.
dim_product = silver_df \
    .select("product_name", "category") \
    .distinct() \
    .withColumn(
        "product_sk",
        F.sha2(F.trim(F.lower(F.col("product_name"))), 256)
    ) \
    .select(
        F.col("product_sk"),
        F.col("product_name"),
        F.col("category"),
    )

dim_product.writeTo("lakehouse.gold.dim_product") \
    .using("iceberg") \
    .createOrReplace()

# ─── 3. Строим dim_date ───────────────────────────────────────────────────────
# Date Dimension генерируется один раз на 10-30 лет вперёд и не меняется
date_range = spark.sql("SELECT sequence(DATE '2026-01-01', DATE '2030-12-31', INTERVAL 1 DAY) AS dates") \
    .select(F.explode(F.col("dates")).alias("full_date"))

dim_date = date_range.select(
    F.sha2(F.col("full_date").cast("string"), 256).alias("date_sk"),
    F.col("full_date"),
    F.dayofmonth(F.col("full_date")).alias("day_of_month"),
    F.dayofweek(F.col("full_date")).alias("day_of_week"),
    F.date_format(F.col("full_date"), "EEEE").alias("day_name"),
    F.weekofyear(F.col("full_date")).alias("week_of_year"),
    F.month(F.col("full_date")).alias("month_number"),
    F.date_format(F.col("full_date"), "MMMM").alias("month_name"),
    F.quarter(F.col("full_date")).alias("quarter"),
    F.year(F.col("full_date")).alias("year"),
    F.when(F.dayofweek(F.col("full_date")).isin(1, 7), True).otherwise(False).alias("is_weekend"),
)

dim_date.writeTo("lakehouse.gold.dim_date") \
    .using("iceberg") \
    .createOrReplace()

# ─── 4. Строим fact_sales ─────────────────────────────────────────────────────
# Зерно: одна строка = одна позиция товара в заказе
fact_sales = silver_df \
    .withColumn(
        "customer_sk",
        F.sha2(F.trim(F.lower(F.col("src_cust_id"))), 256)
    ) \
    .withColumn(
        "product_sk",
        F.sha2(F.trim(F.lower(F.col("product_name"))), 256)
    ) \
    .withColumn(
        "date_sk",
        F.sha2(F.col("order_date"), 256)
    ) \
    .withColumn(
        "fact_sales_sk",
        F.sha2(
            F.concat_ws("|", F.col("order_id"), F.col("product_name")),
            256
        )
    ) \
    .select(
        F.col("fact_sales_sk"),       # PK факта
        F.col("order_id"),
        F.col("customer_sk"),         # FK → dim_customer
        F.col("product_sk"),          # FK → dim_product
        F.col("date_sk"),             # FK → dim_date
        F.col("order_date").cast("date").alias("order_date"),
        # ─── Числовые метрики ────────────────────────────────────────────────
        F.col("price").cast("decimal(10,2)").alias("unit_price"),
        F.col("quantity").cast("int"),
        (F.col("price") * F.col("quantity")).cast("decimal(10,2)").alias("revenue"),
    )

# Партиционируем по дате - это обеспечит DPP и эффективный File Skipping
fact_sales.writeTo("lakehouse.gold.fact_sales") \
    .using("iceberg") \
    .partitionedBy(F.col("order_date")) \
    .createOrReplace()

print("\nfact_sales:")
fact_sales.show(truncate=False)

# ─── 5. Аналитический запрос: Star Schema JOIN с Broadcast ───────────────────
result = spark.sql("""
    SELECT
        c.city,
        p.category,
        d.month_name,
        SUM(f.revenue)  AS total_revenue,
        SUM(f.quantity) AS total_units,
        COUNT(DISTINCT f.order_id) AS unique_orders
    FROM lakehouse.gold.fact_sales      f
    JOIN /*+ BROADCAST(c) */ lakehouse.gold.dim_customer c ON f.customer_sk = c.customer_sk
    JOIN /*+ BROADCAST(p) */ lakehouse.gold.dim_product  p ON f.product_sk  = p.product_sk
    JOIN /*+ BROADCAST(d) */ lakehouse.gold.dim_date     d ON f.date_sk     = d.date_sk
    GROUP BY c.city, p.category, d.month_name
    ORDER BY total_revenue DESC
""")

print("\nАналитический отчёт (Star Schema):")
result.show(truncate=False)

6.2 Проверка плана выполнения

# Убеждаемся, что Spark выбрал Broadcast Hash Join для всех Dimension JOIN'ов
spark.sql("""
    SELECT f.revenue, c.city, p.category
    FROM lakehouse.gold.fact_sales      f
    JOIN lakehouse.gold.dim_customer c ON f.customer_sk = c.customer_sk
    JOIN lakehouse.gold.dim_product  p ON f.product_sk  = p.product_sk
""").explain(mode="formatted")

В плане ищите строки BroadcastHashJoin и BroadcastExchange. Если видите SortMergeJoin - Dimension таблица слишком большая для Broadcast (превышает autoBroadcastJoinThreshold). Решение: либо увеличить порог, либо использовать явный /*+ BROADCAST */ hint.


7. Размер измерений и граница Broadcast

Типичные размеры Dimension таблиц в реальных проектах:

Измерение Строк Размер в Parquet Broadcast?
dim_date (10 лет) 3 650 ~500 KB ✅ Всегда
dim_store 500 ~200 KB ✅ Всегда
dim_product 100 000 ~5 MB ✅ Легко
dim_customer 5 000 000 ~200 MB ⚠️ Нужно настраивать порог
dim_customer 50 000 000 ~2 GB ❌ SortMergeJoin

Когда dim_customer не умещается в Broadcast, рассмотрите:

  • Предварительную фильтрацию Dimension (если запрос обращается к конкретному сегменту клиентов)
  • Партиционирование Fact Table по customer_sk hash bucket и Bucket Join (без Shuffle, но требует одинаковой разбивки)
  • Денормализацию наиболее часто используемых атрибутов прямо в Fact Table (частичная денормализация)

8. Checklist проектирования Star Schema


9. Анти-паттерны Dimensional Modeling в Spark


Итог

Dimensional Modeling по Кимбаллу - это не музейный экспонат из эпохи Teradata, а инженерный фундамент, который делает Data Lakehouse платформой, которой доверяет бизнес:

  • Grain - определяется в первую очередь, на языке бизнеса, до любого кода. Неправильное зерно ломает все агрегаты.
  • Star Schema - денормализованные измерения вокруг центральной Fact Table. В Spark это означает Broadcast Hash Join без Shuffle вместо SortMergeJoin.
  • SHA-256 Surrogate Keys - детерминированная генерация без координации Executors. Дороже по байтам, дешевле по вычислениям и надёжнее при перезапусках.
  • Broadcast Hint + DPP - два механизма, которые вместе позволяют Star Schema запросам читать минимальный объём данных из гигантской Fact Table.

В следующем уроке разберём Slowly Changing Dimensions (SCD) - как корректно хранить исторические изменения атрибутов измерений, сохраняя точность исторических фактов.