SCD в PySpark: Type 1, Type 2 и Type 3 через Iceberg MERGE INTO

Медленно меняющиеся измерения (Slowly Changing Dimensions): SCD Type 1 (перезапись), Type 2 (полная история строк с valid_from/valid_to), Type 3 (prev/current колонки), реализация через Iceberg MERGE INTO с паттерном Source Expansion, обработка Late-Arriving Events

lakehouse scd iceberg merge

В предыдущем уроке мы спроектировали Star Schema: таблица фактов с миллиардами строк продаж, окружённая компактными измерениями - клиентами, продуктами, датами. Модель красива и эффективна. Но есть проблема, которую мы обошли стороной: данные в измерениях меняются.

Клиент переезжает из Казани в Москву. Менеджер переходит в другой отдел. Продукт мигрирует из категории «Телефоны» в «Смартфоны» после ребрендинга. Тарифный план клиента повышается с «Базового» до «Премиум».

Возникает фундаментальный вопрос аналитики: если клиент купил товар в Казани, а потом переехал в Москву, и мы просто перезаписали его город - куда относить эту продажу в историческом отчёте? Отчёт за прошлый год покажет продажу в Москве, хотя фактически она была совершена в Казани. Для операционного отдела это ошибка. Для регионального маркетинга - катастрофа.

Этот класс проблем решает Slowly Changing Dimensions (SCD) - семейство паттернов для управления изменениями атрибутов измерений во времени. Ральф Кимбалл описал их ещё в 1990-х, но реализация в распределённых системах оставалась болезненной вплоть до появления Apache Iceberg с его настоящими ACID-транзакциями и оператором MERGE INTO.


1. Проблема изменчивости измерений

1.1 Что меняется и почему это важно

В любом реальном бизнесе данные измерений нестатичны. Вопрос не в том, изменятся ли они - вопрос в том, как система должна реагировать на изменение.

Это принципиальная развилка: нужна ли нам историческая точность (продажа в Казани, потому что тогда клиент жил там) или актуальная точность (ассоциировать все продажи клиента с его текущим городом)?

Ответ зависит от бизнес-требований, и именно он определяет выбор типа SCD.

1.2 Три кита SCD

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

Тип Стратегия История Сложность Когда использовать
SCD Type 1 Перезапись (Overwrite) Не хранится Низкая Исправление ошибок, нерелевантная история
SCD Type 2 Новая строка (Append Version) Полная Высокая Point-in-time аналитика, изменения с бизнес-значением
SCD Type 3 Колонки prev/current Один шаг назад Средняя «До/после» сравнения, ограниченная история

2. Под капотом Iceberg: как MERGE INTO делает SCD возможным

2.1 Боль старых Data Lake

До появления Apache Iceberg и Delta Lake реализация UPDATE и DELETE в Parquet/Hive требовала тяжёлого паттерна:

  1. Прочитать всю партицию (даже если меняется одна строка из миллиарда)
  2. Применить изменения в памяти
  3. Перезаписать всю партицию в новый файл
  4. Атомарно подменить старый файл новым (rename)

Это работало, но имело принципиальные ограничения: обновление одной строки стоило как полное сканирование терабайтной партиции. SCD Type 2 с ежедневными батч-обновлениями тысяч клиентов превращался в многочасовую операцию.

2.2 Copy-on-Write vs Merge-on-Read в Iceberg

Apache Iceberg поддерживает две стратегии выполнения MERGE INTO:

Copy-on-Write (CoW) - при обновлении Iceberg читает все затронутые Parquet-файлы, применяет изменения и записывает новые файлы. Старые файлы помечаются как удалённые в manifest. Результат: медленная запись, быстрое чтение - каждый Parquet-файл содержит только актуальные данные.

Merge-on-Read (MoR) - при обновлении Iceberg добавляет только небольшой «delete файл» с указанием, какие строки нужно скрыть, и «позиционные дельты» с новыми значениями. Чтение из таблицы объединяет базовые файлы с delete-файлами на лету. Результат: быстрая запись, чуть медленнее чтение - reader должен применять дельты.

Для SCD-таблиц, где запросов на чтение значительно больше, чем обновлений, Copy-on-Write предпочтительнее - читатели (BI-инструменты, Spark-запросы) не несут дополнительных затрат. После тяжёлых MoR-операций нужно регулярно запускать compaction (rewrite_data_files), чтобы схлопнуть накопленные delete-файлы.

2.3 Анатомия MERGE INTO в Spark SQL

Под капотом оператор MERGE INTO в Spark транслируется в следующий физический план:

Весь процесс атомарен: либо commit записывается полностью, либо таблица остаётся в предыдущем состоянии. Частичных обновлений не бывает - это то самое «A» в ACID.


3. SCD Type 1: перезапись без истории

3.1 Механика и применение

Type 1 - простейшая стратегия: при изменении атрибута старое значение просто заменяется новым. Исторические данные при этом теряются безвозвратно.

Когда это правильный выбор:

  • Исправление опечатки: имя «Иванов Иван Иваноыч» → «Иванов Иван Иванович». Хранить историю опечатки не нужно.
  • Технические атрибуты: внутренний идентификатор счёта в банковской системе изменился при миграции. История не важна.
  • Атрибуты без аналитической ценности: email клиента изменился - хранить старый email незачем (и небезопасно с точки зрения GDPR).
  • Некорректные данные: источник передал неверный country_code, потом прислал правильный. Исправляем, не сохраняя «неверную» версию.

Когда это неправильный выбор:

  • Когда исторические факты должны быть привязаны к состоянию измерения на момент события
  • Когда бизнес требует аудит-трейла изменений
  • Когда аналитика строится в разрезе «как было раньше vs как стало»

3.2 Реализация Type 1

Type 1 - это стандартный UPSERT: нашли совпадение по бизнес-ключу - обновляем, не нашли - вставляем новую строку.

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

spark = SparkSession.builder \
    .appName("SCD-Type1") \
    .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()

# Создаём Dimension таблицу
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.gold.dim_customer_scd1 (
        customer_id   STRING     COMMENT 'Бизнес-ключ',
        full_name     STRING,
        email         STRING,
        city          STRING,
        country       STRING,
        segment       STRING,
        _updated_at   TIMESTAMP
    ) USING iceberg
""")

# Начальная загрузка
spark.sql("""
    INSERT INTO lakehouse.gold.dim_customer_scd1 VALUES
    ('cust-001', 'Иван Иванов',  'ivan@mail.ru',  'Казань',  'RU', 'STANDARD', TIMESTAMP '2026-01-01 00:00:00'),
    ('cust-002', 'Alice Smith',  'alice@mail.com', 'London',  'GB', 'PREMIUM',  TIMESTAMP '2026-01-01 00:00:00')
""")

# Батч обновлений из Silver
updates_data = [
    # Иван исправил email (опечатка) + переехал в Москву
    ("cust-001", "Иван Иванов", "ivan.ivanov@mail.ru", "Москва", "RU", "PREMIUM"),
    # Новый клиент
    ("cust-003", "Bob Johnson",  "bob@example.com",     "Berlin", "DE", "STANDARD"),
]
updates_df = spark.createDataFrame(
    updates_data,
    ["customer_id", "full_name", "email", "city", "country", "segment"]
).withColumn("_updated_at", F.current_timestamp())

updates_df.createOrReplaceTempView("scd1_updates")

# SCD Type 1: простой UPSERT
spark.sql("""
    MERGE INTO lakehouse.gold.dim_customer_scd1 AS target
    USING scd1_updates AS source
        ON target.customer_id = source.customer_id
    WHEN MATCHED THEN
        UPDATE SET
            target.full_name   = source.full_name,
            target.email       = source.email,
            target.city        = source.city,
            target.country     = source.country,
            target.segment     = source.segment,
            target._updated_at = source._updated_at
    WHEN NOT MATCHED THEN
        INSERT *
""")

print("SCD Type 1 результат:")
spark.table("lakehouse.gold.dim_customer_scd1").show(truncate=False)
# Иван теперь в Москве. История "Казань" потеряна навсегда.

3.3 Оптимизация Type 1: обновлять только изменившиеся строки

Важная деталь: в WHEN MATCHED лучше добавить проверку на реальное изменение данных. Без неё MERGE будет перезаписывать строки даже если данные не изменились, что создаёт лишние file rewrites в Iceberg.

# Добавляем в WHEN MATCHED условие, что данные реально изменились
spark.sql("""
    MERGE INTO lakehouse.gold.dim_customer_scd1 AS target
    USING scd1_updates AS source
        ON target.customer_id = source.customer_id
    WHEN MATCHED AND (
        target.full_name != source.full_name OR
        target.email     != source.email     OR
        target.city      != source.city      OR
        target.segment   != source.segment
    ) THEN
        UPDATE SET ...
    WHEN NOT MATCHED THEN
        INSERT *
""")

Это WHEN MATCHED AND (condition) называется conditional MATCHED clause - Iceberg выполнит UPDATE только при реальном изменении, что снижает Write Amplification.


4. SCD Type 3: колонки «предыдущее / текущее»

4.1 Механика и применение

Type 3 - компромисс между простотой Type 1 и полнотой Type 2. Вместо хранения полной истории строк, мы добавляем дополнительные колонки для хранения одного предыдущего значения:

current_segment    STRING   -- текущий тарифный план: PREMIUM
previous_segment   STRING   -- предыдущий: STANDARD
segment_changed_at DATE     -- когда произошло изменение

Когда это правильный выбор:

  • Нужно сравнение «до/после» для анализа конверсии (было STANDARD → стало PREMIUM)
  • Бизнес интересует только последнее изменение, не вся история
  • Нет ресурсов на полноценный SCD Type 2, но простого Type 1 недостаточно
  • Аналитика типа «показать клиентов, которые апгрейдились с BASIC на PREMIUM в последнем квартале»

Ограничение: тип 3 хранит только один шаг истории. Если клиент сменил тариф трижды (BASIC → STANDARD → PREMIUM), в таблице будет только «предыдущий = STANDARD, текущий = PREMIUM». История первого изменения (BASIC) безвозвратно потеряна.

4.2 Реализация Type 3

# Создаём таблицу со структурой Type 3
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.gold.dim_customer_scd3 (
        customer_id        STRING,
        full_name          STRING,
        current_segment    STRING,
        previous_segment   STRING   COMMENT 'NULL если изменений не было',
        segment_changed_at DATE     COMMENT 'Дата последнего изменения сегмента',
        current_city       STRING,
        previous_city      STRING,
        city_changed_at    DATE
    ) USING iceberg
""")

# Начальная загрузка
spark.sql("""
    INSERT INTO lakehouse.gold.dim_customer_scd3 VALUES
    ('cust-001', 'Иван Иванов', 'STANDARD', NULL, NULL, 'Казань', NULL, NULL),
    ('cust-002', 'Alice Smith',  'PREMIUM',  NULL, NULL, 'London', NULL, NULL)
""")

# Батч: Иван апгрейдился до PREMIUM, переехал в Москву
scd3_updates = [
    ("cust-001", "Иван Иванов", "PREMIUM", "Москва"),
]
scd3_df = spark.createDataFrame(
    scd3_updates,
    ["customer_id", "full_name", "new_segment", "new_city"]
)
scd3_df.createOrReplaceTempView("scd3_updates")

# SCD Type 3: при MATCHED - сдвигаем current → previous, записываем new → current
spark.sql("""
    MERGE INTO lakehouse.gold.dim_customer_scd3 AS target
    USING scd3_updates AS source
        ON target.customer_id = source.customer_id
    WHEN MATCHED AND target.current_segment != source.new_segment THEN
        UPDATE SET
            target.previous_segment   = target.current_segment,
            target.current_segment    = source.new_segment,
            target.segment_changed_at = CURRENT_DATE(),
            target.previous_city      = target.current_city,
            target.current_city       = source.new_city,
            target.city_changed_at    = CURRENT_DATE()
    WHEN NOT MATCHED THEN
        INSERT (customer_id, full_name, current_segment, previous_segment,
                segment_changed_at, current_city, previous_city, city_changed_at)
        VALUES (source.customer_id, source.full_name, source.new_segment, NULL,
                CURRENT_DATE(), source.new_city, NULL, NULL)
""")

print("SCD Type 3 результат:")
spark.table("lakehouse.gold.dim_customer_scd3").show(truncate=False)

Ключевая строка - target.previous_segment = target.current_segment. Мы читаем текущее значение цели и записываем его в колонку previous. Это «сдвиг» истории на один шаг назад - изящно и без сложного SQL.


5. SCD Type 2: полная история строк

5.1 Механика и философия

SCD Type 2 - самый мощный и самый сложный тип. Вместо обновления строки при изменении атрибута мы создаём новую строку с новыми значениями, а старую помечаем как «закрытую».

В результате в Dimension таблице будет несколько строк с одинаковым customer_id - каждая представляет определённый период жизни клиента.

Point-in-time запрос - это ключевая возможность SCD Type 2. Можно взять любую дату в прошлом и получить точное состояние измерения на этот момент. Именно это нужно для корректного исторического анализа: продажа в январе 2026 ассоциируется с клиентом из Казани, а не из Москвы.

5.2 Аудит-колонки: valid_from, valid_to, is_current

valid_from (дата начала действия версии) - дата, с которой эта строка стала актуальной. Для первой загрузки обычно берётся «дата начала бизнеса» или 1970-01-01.

valid_to (дата окончания действия версии) - дата, с которой версия перестала быть актуальной. Для текущей версии устанавливается «магическая дата» 9999-12-31 (иногда 2099-12-31). Использование NULL в этой колонке - анти-паттерн: NULL усложняет range-запросы и требует специальной обработки (COALESCE(valid_to, '9999-12-31')).

is_current (булев флаг текущей версии) - избыточная, но удобная колонка. Технически её значение вычисляется из valid_to = '9999-12-31', но наличие флага делает запросы «текущего состояния» значительно быстрее: WHERE is_current = true вместо WHERE valid_to = '9999-12-31'. Iceberg может эффективно пропускать файлы, в которых нет TRUE значений.

5.3 Проблема: почему стандартный MERGE не работает для SCD Type 2

Это ключевое техническое препятствие. Рассмотрим наивную попытку реализовать SCD Type 2:

-- Этот код НЕ работает для SCD Type 2!
MERGE INTO dim_customer AS target
USING updates AS source ON target.customer_id = source.customer_id AND target.is_current = true
WHEN MATCHED AND data_changed THEN
    -- Хотим: закрыть старую строку И вставить новую
    UPDATE SET target.is_current = false, target.valid_to = today
    -- INSERT новой строки здесь невозможен!
WHEN NOT MATCHED THEN
    INSERT new row

Проблема фундаментальна: в стандартном SQL MERGE оператор WHEN MATCHED может только обновить или удалить найденную строку, но не вставить дополнительную новую строку. Для SCD Type 2 нам нужно при изменении данных сделать два действия с одним customer_id:

  1. UPDATE существующей строки: is_current = false, valid_to = сегодня
  2. INSERT новой строки: новые атрибуты, is_current = true, valid_from = сегодня

Стандартный MERGE не позволяет этого сделать за один проход.

5.4 Решение: паттерн Source Expansion (обогащение источника)

Элегантное решение - перед MERGE удвоить строки в источнике обновлений: для каждого изменения создать две записи - одну для UPDATE старой строки, другую для INSERT новой.

Ключевой трюк: первая запись имеет реальный customer_id (совпадёт с WHEN MATCHED), вторая имеет NULL в поле merge_key (не совпадёт ни с кем, попадёт в WHEN NOT MATCHED).


6. Практика: полный SCD Type 2 Pipeline

6.1 Полный код

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

spark = SparkSession.builder \
    .appName("SCD-Type2-Full-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()

# ─── 1. Создаём SCD Type 2 таблицу ───────────────────────────────────────────
spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.gold.dim_users_scd2 (
        user_id       STRING     COMMENT 'Бизнес-ключ',
        full_name     STRING,
        city          STRING,
        segment       STRING,
        valid_from    DATE       COMMENT 'Начало действия этой версии',
        valid_to      DATE       COMMENT '9999-12-31 для текущей версии',
        is_current    BOOLEAN    COMMENT 'true только для одной строки по user_id'
    ) USING iceberg
""")

# ─── 2. Начальная загрузка ────────────────────────────────────────────────────
initial_date = "2026-01-01"
spark.sql(f"""
    INSERT INTO lakehouse.gold.dim_users_scd2 VALUES
    ('usr-001', 'Иван Иванов',  'Казань', 'STANDARD', DATE '{initial_date}', DATE '9999-12-31', true),
    ('usr-002', 'Alice Smith',  'London', 'PREMIUM',  DATE '{initial_date}', DATE '9999-12-31', true),
    ('usr-003', 'Bob Johnson',  'Berlin', 'STANDARD', DATE '{initial_date}', DATE '9999-12-31', true)
""")

print("=== Состояние ПОСЛЕ начальной загрузки ===")
spark.table("lakehouse.gold.dim_users_scd2").orderBy("user_id", "valid_from").show(truncate=False)


# ─── 3. Функция SCD Type 2 обновления ─────────────────────────────────────────
def apply_scd2_update(updates_df, effective_date: str):
    """
    Применяет SCD Type 2 обновление к dim_users_scd2.

    Паттерн Source Expansion:
    - Первая копия обновлений (с реальным user_id) → WHEN MATCHED → UPDATE (закрыть старую строку)
    - Вторая копия (с NULL merge_key) → WHEN NOT MATCHED → INSERT (открыть новую версию)

    Условие совпадения ONLY для is_current=true гарантирует, что мы работаем
    только с актуальной версией, не трогая историческиe записи.
    """
    updates_df.createOrReplaceTempView("raw_scd2_updates")

    # Source Expansion через UNION ALL
    # Первая часть: строки для UPDATE (merge_key = реальный user_id, совпадёт в WHEN MATCHED)
    # Вторая часть: строки для INSERT (merge_key = NULL, не совпадёт, попадёт в WHEN NOT MATCHED)
    spark.sql(f"""
        CREATE OR REPLACE TEMP VIEW staged_scd2_updates AS

        -- Строки для обновления существующих записей (WHEN MATCHED → UPDATE)
        SELECT
            u.user_id    AS merge_key,
            u.user_id,
            u.full_name,
            u.city,
            u.segment,
            DATE '{effective_date}' AS effective_date
        FROM raw_scd2_updates u

        UNION ALL

        -- Строки для вставки новых версий (WHEN NOT MATCHED → INSERT).
        -- Выбираем только записи, где данные РЕАЛЬНО изменились.
        -- NULL в merge_key гарантирует, что MERGE не найдёт совпадений
        -- и отправит эти строки в ветку WHEN NOT MATCHED.
        SELECT
            NULL          AS merge_key,
            u.user_id,
            u.full_name,
            u.city,
            u.segment,
            DATE '{effective_date}' AS effective_date
        FROM raw_scd2_updates u
        JOIN lakehouse.gold.dim_users_scd2 t
            ON u.user_id = t.user_id AND t.is_current = true
        WHERE
            u.city     <> t.city     OR
            u.segment  <> t.segment  OR
            u.full_name <> t.full_name
    """)

    # Основной MERGE
    spark.sql(f"""
        MERGE INTO lakehouse.gold.dim_users_scd2 AS target
        USING staged_scd2_updates AS source
            ON target.user_id = source.merge_key AND target.is_current = true
        WHEN MATCHED AND (
            target.city      <> source.city      OR
            target.segment   <> source.segment   OR
            target.full_name <> source.full_name
        ) THEN
            -- Закрываем текущую версию строки
            UPDATE SET
                target.valid_to   = source.effective_date,
                target.is_current = false
        WHEN NOT MATCHED THEN
            -- Вставляем новую актуальную версию
            INSERT (user_id, full_name, city, segment, valid_from, valid_to, is_current)
            VALUES (
                source.user_id,
                source.full_name,
                source.city,
                source.segment,
                source.effective_date,
                DATE '9999-12-31',
                true
            )
    """)


# ─── 4. Батч 1: Иван переезжает в Москву ─────────────────────────────────────
print("\n=== Батч 1: Иван переезжает в Москву (2026-03-01) ===")
batch1_data = [
    ("usr-001", "Иван Иванов", "Москва",  "STANDARD"),  # переезд
    ("usr-002", "Alice Smith", "London",  "PREMIUM"),   # не изменился
]
batch1_df = spark.createDataFrame(batch1_data, ["user_id", "full_name", "city", "segment"])
apply_scd2_update(batch1_df, "2026-03-01")

spark.table("lakehouse.gold.dim_users_scd2").orderBy("user_id", "valid_from").show(truncate=False)


# ─── 5. Батч 2: Иван апгрейдится до PREMIUM ──────────────────────────────────
print("\n=== Батч 2: Иван апгрейдился до PREMIUM (2026-05-15) ===")
batch2_data = [
    ("usr-001", "Иван Иванов", "Москва", "PREMIUM"),   # апгрейд
    ("usr-004", "New User",    "Paris",  "TRIAL"),     # новый клиент
]
batch2_df = spark.createDataFrame(batch2_data, ["user_id", "full_name", "city", "segment"])
apply_scd2_update(batch2_df, "2026-05-15")

print("Финальное состояние таблицы:")
spark.table("lakehouse.gold.dim_users_scd2").orderBy("user_id", "valid_from").show(truncate=False)


# ─── 6. Point-in-time запросы ─────────────────────────────────────────────────
print("\n=== Point-in-time запрос: состояние клиентов на 15 февраля 2026 ===")
spark.sql("""
    SELECT user_id, full_name, city, segment, valid_from, valid_to
    FROM lakehouse.gold.dim_users_scd2
    WHERE DATE '2026-02-15' BETWEEN valid_from AND valid_to
    ORDER BY user_id
""").show(truncate=False)
# Иван должен быть в Казани, Alice в London - это состояние до его переезда

print("\n=== Текущее состояние всех клиентов ===")
spark.sql("""
    SELECT user_id, full_name, city, segment
    FROM lakehouse.gold.dim_users_scd2
    WHERE is_current = true
    ORDER BY user_id
""").show(truncate=False)

6.2 Ожидаемый результат

После двух батчей таблица должна выглядеть так:

user_id city segment valid_from valid_to is_current
usr-001 Казань STANDARD 2026-01-01 2026-03-01 FALSE
usr-001 Москва STANDARD 2026-03-01 2026-05-15 FALSE
usr-001 Москва PREMIUM 2026-05-15 9999-12-31 TRUE
usr-002 London PREMIUM 2026-01-01 9999-12-31 TRUE
usr-003 Berlin STANDARD 2026-01-01 9999-12-31 TRUE
usr-004 Paris TRIAL 2026-05-15 9999-12-31 TRUE

Запрос Point-in-time на 2026-02-15 вернёт Ивана с city='Казань' - именно там он жил в тот момент. Это и есть историческая точность SCD Type 2.


7. Промышленный кейс: Late-Arriving Events

7.1 Проблема запоздавших данных

SCD Type 2 строится на предположении, что события приходят в хронологическом порядке. Реальность сложнее: из-за сетевых сбоев, перезапусков CDC-коннекторов или задержек в очередях события могут прийти позже, чем произошли.

Сценарий: переезд Ивана в Москву произошёл 15 мая. Из-за сбоя API событие пришло только 23 мая. К тому моменту мы уже обработали батч от 20 мая, в котором Иван переехал в Питер.

7.2 Решение: вставка в середину временно́го ряда

Стандартный SCD Type 2 пайплайн, завязанный на is_current = true, не справится с этой ситуацией. Нужна отдельная логика для запоздавших событий:

  1. Найти строку, в диапазон [valid_from, valid_to] которой попадает дата опоздавшего события
  2. «Разрезать» этот интервал: закрыть найденную строку на дату события
  3. Вставить новую строку с атрибутами из события (valid_from = дата события, valid_to = старый valid_to найденной строки)
def apply_late_arriving_scd2(late_event_user_id: str, late_event_city: str,
                              late_event_date: str):
    """
    Обрабатывает запоздавшее событие для SCD Type 2.
    Вставляет новую версию строго в нужный временно́й слот, не нарушая хронологию.
    """
    # Находим строку, чей интервал содержит дату опоздавшего события
    target_row = spark.sql(f"""
        SELECT user_id, full_name, city, segment, valid_from, valid_to, is_current
        FROM lakehouse.gold.dim_users_scd2
        WHERE user_id = '{late_event_user_id}'
          AND DATE '{late_event_date}' >= valid_from
          AND DATE '{late_event_date}' < valid_to
    """).collect()

    if not target_row:
        print(f"Нет строки для user_id={late_event_user_id} на дату {late_event_date}")
        return

    row = target_row[0]
    original_valid_to = row["valid_to"]
    original_is_current = row["is_current"]

    # Шаг 1: Закрываем найденную строку на дату опоздавшего события
    spark.sql(f"""
        UPDATE lakehouse.gold.dim_users_scd2
        SET valid_to = DATE '{late_event_date}', is_current = false
        WHERE user_id = '{late_event_user_id}'
          AND valid_from = DATE '{row["valid_from"]}'
          AND valid_to = DATE '{original_valid_to}'
    """)

    # Шаг 2: Вставляем новую строку для опоздавшего события
    # valid_to берём из СТАРОЙ строки (может быть как 9999-12-31, так и промежуточной датой)
    spark.sql(f"""
        INSERT INTO lakehouse.gold.dim_users_scd2 VALUES (
            '{late_event_user_id}',
            '{row["full_name"]}',
            '{late_event_city}',
            '{row["segment"]}',
            DATE '{late_event_date}',
            DATE '{original_valid_to}',
            {str(original_is_current).lower()}
        )
    """)
    print(f"Late-arriving event обработан: {late_event_user_id} {late_event_city} {late_event_date}")


# Применяем: московский переезд (15 мая) пришёл после питерского (20 мая)
apply_late_arriving_scd2("usr-001", "Москва", "2026-05-15")

print("\n=== Временна́я шкала после обработки запоздавшего события ===")
spark.sql("""
    SELECT user_id, city, segment, valid_from, valid_to, is_current
    FROM lakehouse.gold.dim_users_scd2
    WHERE user_id = 'usr-001'
    ORDER BY valid_from
""").show(truncate=False)
# Ожидаем: Казань → Москва → Питер (хронологически правильно)

8. Как SCD влияет на производительность Iceberg

8.1 Write Amplification при MERGE

Каждый MERGE INTO в Iceberg (Copy-on-Write режим) перезаписывает все Parquet-файлы, содержащие затронутые строки. Если в одном файле 1 000 000 строк и изменилось 10 из них - весь файл перезаписывается заново.

Для SCD Type 2 это означает: частые батч-обновления (например, ежечасные) создают множество мелких новых файлов. Через несколько дней таблица может содержать тысячи маленьких Parquet-файлов, что резко замедляет все последующие чтения (проблема small files).

Решение: регулярный compaction после тяжёлых MERGE-операций:

# Запускается после нескольких SCD-обновлений, например раз в сутки
spark.sql("""
    CALL lakehouse.system.rewrite_data_files(
        table => 'lakehouse.gold.dim_users_scd2',
        strategy => 'binpack',
        options => map('target-file-size-bytes', '134217728')
    )
""")

# Очищаем старые snapshot'ы (метаданные прошлых версий)
spark.sql("""
    CALL lakehouse.system.expire_snapshots(
        table => 'lakehouse.gold.dim_users_scd2',
        older_than => TIMESTAMP '2026-05-01 00:00:00',
        retain_last => 5
    )
""")

8.2 Партиционирование SCD Type 2 таблиц

Не партиционируйте SCD Dimension таблицы по valid_from или user_id - это создаст огромное количество маленьких партиций.

Для большинства Dimension таблиц вообще не нужно партиционирование - они достаточно малы для полного Broadcast в Spark. Если таблица действительно огромная (десятки миллионов клиентов с длинной историей), рассмотрите партиционирование по is_current - это даёт быстрый доступ к актуальному срезу через WHERE is_current = true.


9. Выбор типа SCD: checklist


Итог

Slowly Changing Dimensions - это не опциональная «приятная фича», а фундаментальное требование к любой аналитической платформе, претендующей на историческую точность:

  • SCD Type 1 - простой и эффективный для атрибутов без аналитической ценности истории. Стандартный UPSERT с условием на реальное изменение данных.
  • SCD Type 3 - компромисс для сравнений «до/после», когда хранить полную историю избыточно, но иметь одно предыдущее значение необходимо.
  • SCD Type 2 - полная историческая точность через valid_from/valid_to/is_current. Требует паттерна Source Expansion для корректного MERGE. Является основой Point-in-time аналитики.
  • Late-Arriving Events - самый сложный сценарий, требующий отдельной логики «вставки в середину» временно́го ряда.
  • Iceberg MERGE INTO - делает все эти операции атомарными и масштабируемыми. Без него каждый SCD Type 2 батч требовал перезаписи целых партиций.

Регулярный compaction после MERGE-тяжёлых SCD-операций - обязательная часть эксплуатации, иначе накопленные мелкие файлы деградируют производительность чтения.