diff --git a/dwh-modeling/sql/01_ddl_stg-dds.sql b/dwh-modeling/sql/01_ddl_stg-dds.sql index a73efd5..b8588e9 100644 --- a/dwh-modeling/sql/01_ddl_stg-dds.sql +++ b/dwh-modeling/sql/01_ddl_stg-dds.sql @@ -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, -- бизнес-ключ diff --git a/dwh-modeling/sql/02_dml_stg-dds.sql b/dwh-modeling/sql/02_dml_stg-dds.sql index 5e26699..dd4ed1f 100644 --- a/dwh-modeling/sql/02_dml_stg-dds.sql +++ b/dwh-modeling/sql/02_dml_stg-dds.sql @@ -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, diff --git a/dwh-modeling/sql/03_demo_increment.sql b/dwh-modeling/sql/03_demo_increment.sql index d54558c..f0bbb8f 100644 --- a/dwh-modeling/sql/03_demo_increment.sql +++ b/dwh-modeling/sql/03_demo_increment.sql @@ -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 diff --git a/dwh-modeling/sql/04_validation.sql b/dwh-modeling/sql/04_validation.sql index 1810410..864a8dd 100644 --- a/dwh-modeling/sql/04_validation.sql +++ b/dwh-modeling/sql/04_validation.sql @@ -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 $$;