Более понятная реализация scd2
This commit is contained in:
@@ -103,7 +103,7 @@ CREATE TABLE dds.dim_product (
|
||||
product_name VARCHAR(100) NOT NULL
|
||||
);
|
||||
|
||||
-- dim_customer: измерение "Клиент" с SCD Type 2
|
||||
-- dim_customer: измерение "Клиент" с историей (SCD Type 2)
|
||||
CREATE TABLE dds.dim_customer (
|
||||
customer_sk BIGSERIAL PRIMARY KEY,
|
||||
customer_bk INT NOT NULL, -- бизнес-ключ
|
||||
|
||||
@@ -33,7 +33,7 @@ INSERT INTO stg.products_raw (product_id, name) VALUES
|
||||
('9002', 'Case');
|
||||
|
||||
-- 2. ODS: очистка и типизация
|
||||
-- ⚠️ В продакшене используем UPSERT или incremental load, не TRUNCATE+INSERT
|
||||
-- ⚠️ В продакшене используем UPSERT (INSERT ... ON CONFLICT DO UPDATE) или incremental load, не TRUNCATE+INSERT
|
||||
|
||||
TRUNCATE ods.customers, ods.orders, ods.order_items, ods.products;
|
||||
|
||||
@@ -134,9 +134,20 @@ FROM ods.products;
|
||||
-- В РЕАЛЬНОМ DWH: так делают редко. Исторические измерения обычно строят
|
||||
-- поверх очищенных и нормализованных слоёв (ODS / PSA / Data Vault).
|
||||
-- Для примера инкрементальной заливки SCD2 по снимку из ODS см. 03_demo_increment.sql и SCD.md.
|
||||
--
|
||||
-- Идея SCD2 простыми словами:
|
||||
-- - одна строка = один период, когда атрибуты клиента (email/phone/city) были одинаковыми;
|
||||
-- - valid_from = дата, когда "стало так";
|
||||
-- - valid_to = дата следующего изменения (NULL = текущая версия).
|
||||
--
|
||||
-- Откуда берём дату изменения:
|
||||
-- - если в событии есть event_ts — считаем, что изменение произошло тогда;
|
||||
-- - если event_ts пустой — берём дату загрузки (_load_ts), чтобы не терять историю.
|
||||
--
|
||||
-- Важно для демо: считаем, что у клиента не бывает двух разных изменений в один и тот же день.
|
||||
TRUNCATE dds.dim_customer, dds.fact_sales;
|
||||
|
||||
WITH src AS (
|
||||
WITH src AS ( -- 1) Приводим типы, готовим дату изменения (eff_date) и считаем hashdiff атрибутов
|
||||
SELECT
|
||||
s.customer_id::INT AS customer_bk,
|
||||
NULLIF(trim(s.email), '') AS email,
|
||||
@@ -147,18 +158,17 @@ WITH src AS (
|
||||
FROM stg.customers_raw s
|
||||
WHERE s.customer_id ~ '^\d+$'
|
||||
),
|
||||
ordered AS (
|
||||
ordered AS ( -- 2) Сортируем по датам и смотрим "какой hashdiff был до этого" (LAG)
|
||||
SELECT *,
|
||||
lag(hashdiff) OVER (PARTITION BY customer_bk ORDER BY eff_date) AS prev_hash
|
||||
FROM src
|
||||
),
|
||||
changes AS (
|
||||
-- только первые состояния и фактические изменения атрибутов
|
||||
changes AS ( -- 3) Оставляем только первое состояние и реальные изменения (где hashdiff поменялся)
|
||||
SELECT *
|
||||
FROM ordered
|
||||
WHERE prev_hash IS DISTINCT FROM hashdiff OR prev_hash IS NULL
|
||||
),
|
||||
framed AS (
|
||||
framed AS ( -- 4) Превращаем изменения в периоды: valid_to = дата следующего изменения (LEAD)
|
||||
SELECT
|
||||
customer_bk, email, phone, city, hashdiff,
|
||||
eff_date AS valid_from,
|
||||
|
||||
@@ -8,7 +8,8 @@ INSERT INTO stg.customers_raw (_load_id, _load_ts, event_ts, customer_id, email,
|
||||
('batch_20241101_0800', '2024-11-01 08:00', '2024-11-01', '101','b@ex.com','700','Москва'), -- город вернулся
|
||||
('batch_20240310_0800', '2024-03-10 08:00', '2024-03-10', '103','d@ex.com','702','Казань'); -- новый клиент
|
||||
|
||||
-- 1) UPSERT в ODS последнего снимка по BK
|
||||
-- 1) UPSERT в ODS (вставка с обновлением по конфликту, INSERT ... ON CONFLICT DO UPDATE):
|
||||
-- сохраняем в ods.customers последнюю версию клиента по BK (бизнес-ключ = customer_id)
|
||||
WITH src AS (
|
||||
SELECT
|
||||
s.customer_id::INT AS customer_id,
|
||||
@@ -18,6 +19,8 @@ WITH src AS (
|
||||
NULLIF(s.event_ts,'')::timestamp AS event_ts,
|
||||
s._load_id,
|
||||
s._load_ts,
|
||||
-- eff_ts нужен, чтобы выбрать "самое свежее" событие по клиенту:
|
||||
-- если event_ts нет, используем время загрузки (_load_ts) как приближение.
|
||||
COALESCE(NULLIF(s.event_ts,'')::timestamp, s._load_ts) AS eff_ts
|
||||
FROM stg.customers_raw s
|
||||
WHERE s.customer_id ~ '^\d+$'
|
||||
@@ -25,6 +28,7 @@ WITH src AS (
|
||||
ranked AS (
|
||||
SELECT
|
||||
customer_id, email, phone, city, event_ts, _load_id, _load_ts,
|
||||
-- берём одну строку на клиента: с максимальным eff_ts (при равенстве — с максимальным _load_ts)
|
||||
row_number() OVER (PARTITION BY customer_id ORDER BY eff_ts DESC, _load_ts DESC) AS rn
|
||||
FROM src
|
||||
)
|
||||
@@ -44,10 +48,20 @@ WHERE COALESCE(EXCLUDED.event_ts, EXCLUDED._load_ts) >
|
||||
COALESCE(ods.customers.event_ts, ods.customers._load_ts);
|
||||
|
||||
-- 2) Инкрементальное SCD2 из ODS (в одной транзакции, по последнему снимку в ODS)
|
||||
-- Канон для курса: daily-grain, интервалы [valid_from, valid_to), текущая версия = valid_to IS NULL
|
||||
-- Канон для курса: считаем по дням (valid_from/valid_to — DATE), интервалы [valid_from, valid_to),
|
||||
-- текущая версия = valid_to IS NULL
|
||||
-- Предполагаем, что для клиента нет нескольких изменений в один день.
|
||||
--
|
||||
-- Идея (по шагам):
|
||||
-- 1) Берём текущий "снимок" клиента из ODS (одна строка на BK = customer_id).
|
||||
-- 2) Сравниваем его с текущей версией в dds.dim_customer (valid_to IS NULL) по hashdiff.
|
||||
-- 3) Если изменилось — закрываем текущую версию (ставим valid_to) и вставляем новую (valid_to = NULL).
|
||||
--
|
||||
-- Про даты:
|
||||
-- eff_date берём из event_ts, а если его нет — из _load_ts (как приближение).
|
||||
BEGIN;
|
||||
-- 2.1) Закрываем предыдущую актуальную версию (только если реально изменились атрибуты)
|
||||
-- delta/current считаем прямо в запросе (без временных таблиц) специально для читабельности.
|
||||
WITH delta AS (
|
||||
SELECT
|
||||
c.customer_id AS customer_bk,
|
||||
@@ -65,6 +79,8 @@ BEGIN;
|
||||
SET valid_to = x.eff_date,
|
||||
updated_at = now()
|
||||
FROM (
|
||||
-- x = кандидаты на "закрытие" текущей версии:
|
||||
-- клиент есть в DDS (current) и атрибуты изменились (hashdiff стал другим).
|
||||
SELECT
|
||||
t.customer_bk,
|
||||
t.eff_date,
|
||||
@@ -73,12 +89,13 @@ BEGIN;
|
||||
JOIN current c
|
||||
ON c.customer_bk = t.customer_bk
|
||||
WHERE c.hashdiff <> t.hashdiff
|
||||
AND t.eff_date > c.valid_from
|
||||
AND t.eff_date > c.valid_from -- не создаём период нулевой/отрицательной длины
|
||||
) x
|
||||
WHERE d.customer_sk = x.customer_sk
|
||||
AND d.valid_to IS NULL;
|
||||
|
||||
-- 2.2) Вставляем новую версию (только если новая или изменившаяся)
|
||||
-- delta/current повторяем ещё раз отдельно, чтобы блок вставки читался независимо от блока UPDATE.
|
||||
WITH delta AS (
|
||||
SELECT
|
||||
c.customer_id AS customer_bk,
|
||||
@@ -93,6 +110,10 @@ BEGIN;
|
||||
WHERE d.valid_to IS NULL
|
||||
),
|
||||
to_insert AS (
|
||||
-- to_insert = кандидаты на вставку:
|
||||
-- 1) новый клиент (в current нет строки);
|
||||
-- 2) изменившийся клиент (hashdiff поменялся).
|
||||
-- если атрибуты не менялись — клиент сюда не попадёт, и ничего делать не нужно.
|
||||
SELECT
|
||||
t.customer_bk,
|
||||
t.email,
|
||||
@@ -116,6 +137,7 @@ BEGIN;
|
||||
t.eff_date, NULL,
|
||||
now(), now()
|
||||
FROM to_insert t
|
||||
-- защита от повторного запуска: не вставляем одну и ту же версию (BK + valid_from) второй раз
|
||||
WHERE NOT EXISTS (
|
||||
SELECT 1
|
||||
FROM dds.dim_customer d
|
||||
|
||||
@@ -29,7 +29,7 @@ BEGIN
|
||||
RAISE NOTICE '✅ fact_sales: количество строк совпадает с ods.order_items (%)', actual_count;
|
||||
END $$;
|
||||
|
||||
-- 3. Проверка SCD Type 2: у клиента 101 должно быть ≥2 версий (из-за смены email)
|
||||
-- 3. Проверка SCD2 (Type 2): у клиента 101 должно быть ≥2 версий (из-за смены email)
|
||||
DO $$
|
||||
DECLARE version_count INT;
|
||||
BEGIN
|
||||
@@ -39,10 +39,10 @@ BEGIN
|
||||
--
|
||||
ASSERT version_count >= 2,
|
||||
FORMAT('ОШИБКА: у клиента 101 только %s версия, ожидается ≥2 (должна быть история)', version_count);
|
||||
RAISE NOTICE '✅ SCD Type 2: клиент 101 имеет % версий — история сохранена', version_count;
|
||||
RAISE NOTICE '✅ SCD2: клиент 101 имеет % версий — история сохранена', version_count;
|
||||
END $$;
|
||||
|
||||
-- 4. Проверка SCD Type 2: у каждого клиента ровно одна актуальная версия (valid_to IS NULL)
|
||||
-- 4. Проверка SCD2 (Type 2): у каждого клиента ровно одна актуальная версия (valid_to IS NULL)
|
||||
DO $$
|
||||
DECLARE
|
||||
customers_cnt BIGINT;
|
||||
@@ -56,10 +56,10 @@ BEGIN
|
||||
ASSERT current_cnt = customers_cnt,
|
||||
FORMAT('ОШИБКА: актуальных строк %s, а уникальных клиентов %s (ожидается 1 current на клиента)',
|
||||
current_cnt, customers_cnt);
|
||||
RAISE NOTICE '✅ SCD Type 2: current-строки = количеству клиентов (%)', current_cnt;
|
||||
RAISE NOTICE '✅ SCD2: current-строки = количеству клиентов (%)', current_cnt;
|
||||
END $$;
|
||||
|
||||
-- 5. Проверка SCD Type 2: периоды корректны (valid_to > valid_from или valid_to IS NULL)
|
||||
-- 5. Проверка SCD2 (Type 2): периоды корректны (valid_to > valid_from или valid_to IS NULL)
|
||||
DO $$
|
||||
BEGIN
|
||||
ASSERT NOT EXISTS (
|
||||
@@ -69,5 +69,5 @@ BEGIN
|
||||
AND valid_to <= valid_from
|
||||
),
|
||||
'ОШИБКА: найдены строки dim_customer с некорректным периодом (valid_to <= valid_from)';
|
||||
RAISE NOTICE '✅ SCD Type 2: периоды valid_from/valid_to корректны';
|
||||
RAISE NOTICE '✅ SCD2: периоды valid_from/valid_to корректны';
|
||||
END $$;
|
||||
|
||||
Reference in New Issue
Block a user