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
В предыдущем уроке мы спроектировали 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 требовала тяжёлого паттерна:
- Прочитать всю партицию (даже если меняется одна строка из миллиарда)
- Применить изменения в памяти
- Перезаписать всю партицию в новый файл
- Атомарно подменить старый файл новым (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:
- UPDATE существующей строки:
is_current = false,valid_to = сегодня - 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, не справится с этой ситуацией. Нужна отдельная логика для запоздавших событий:
- Найти строку, в диапазон
[valid_from, valid_to]которой попадает дата опоздавшего события - «Разрезать» этот интервал: закрыть найденную строку на дату события
- Вставить новую строку с атрибутами из события (
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-операций - обязательная часть эксплуатации, иначе накопленные мелкие файлы деградируют производительность чтения.