diff --git a/dwh-modeling/sql/09_dml_hw_customer_status_solution.sql b/dwh-modeling/sql/09_dml_hw_customer_status_solution.sql index 0fe0f99..37c3661 100644 --- a/dwh-modeling/sql/09_dml_hw_customer_status_solution.sql +++ b/dwh-modeling/sql/09_dml_hw_customer_status_solution.sql @@ -2,14 +2,19 @@ -- 09_dml_hw_customer_status_solution.sql -- Решение домашки: статусы клиента (STG -> ODS -> DDS SCD2 -> DM) -- +-- Что делает этот файл: +-- 1) Перекладывает события статусов в ODS (приводит типы, чистит пустое). +-- 2) Строит DDS-измерение со "встроенной историей" (SCD2): периоды valid_from/valid_to. +-- 3) Показывает пример обновления DDS маленькой порцией (инкремент): закрыть старое + вставить новое. +-- 4) Собирает простую витрину в DM: сколько клиентов в каком статусе по дням. +-- -- Как запускать: --- - можно запускать файл целиком; --- - блок 3 (инкремент) имеет смысл запускать повторно после того, --- как вы добавили новые события и обновили ODS. +-- - для первого знакомства можно запускать файл целиком; +-- - если хотите потренировать инкремент (п.3): добавьте новые события -> обновите ODS -> запустите блок 3 ещё раз. -- -- Важно: --- - в решении используются TRUNCATE (полная очистка) для ODS/DDS/DM, чтобы было проще повторять домашку. --- - в реальном DWH так делают не всегда: чаще грузят инкрементально и не трогают историю целиком. +-- - здесь часто используется TRUNCATE (полная очистка), чтобы было легко повторять домашку; +-- - в реальном DWH так делают не всегда, но для обучения это удобнее. -- -- Предусловия (DDL + данные в STG): -- 1) dwh-modeling/sql/01_ddl_stg-dds.sql @@ -22,8 +27,10 @@ -- 1) ODS: очистка и типизация (full refresh) -- ========================================================== --- Идея: STG хранит "как пришло" (TEXT), ODS хранит "аккуратно" (типы + базовая чистка). --- Для простоты делаем full refresh: пересобираем ods.customer_status целиком. +-- Идея: +-- - STG хранит "как пришло" (обычно TEXT); +-- - ODS хранит "аккуратно": правильные типы + простая чистка. +-- Для простоты пересобираем ODS с нуля. TRUNCATE ods.customer_status; @@ -46,11 +53,12 @@ WHERE s.customer_id ~ '^\d+$' -- ========================================================== -- Идея SCD2 простыми словами: --- - одна строка = один период, когда статус был одинаковым; --- - valid_from = с какого дня статус "действует"; --- - valid_to = с какого дня статус перестал действовать (NULL = текущая версия); --- - интервалы считаем как [valid_from, valid_to). --- - чтобы найти статус "на дату D": D >= valid_from AND (valid_to IS NULL OR D < valid_to) +-- - одна строка = один период, когда статус был одним и тем же; +-- - valid_from = с какого дня статус "начался"; +-- - valid_to = с какого дня статус "закончился" (NULL = текущий статус); +-- - интервалы считаем так: [valid_from, valid_to) (valid_to не включаем). +-- - чтобы найти статус "на дату D": +-- D >= valid_from AND (valid_to IS NULL OR D < valid_to) -- -- Упрощение для домашки: -- - считаем, что у клиента нет двух разных смен статуса в один день. @@ -58,8 +66,8 @@ WHERE s.customer_id ~ '^\d+$' TRUNCATE dds.dim_customer_status; WITH src AS ( - -- src: события из ODS + hashdiff. - -- В домашке hashdiff можно считать просто как md5(status): так удобно сравнить "изменился статус или нет". + -- src: события из ODS + "контрольная сумма" статуса. + -- Так проще проверять, поменялся статус или остался тем же. SELECT customer_id AS customer_bk, status, @@ -68,7 +76,7 @@ WITH src AS ( FROM ods.customer_status ), ordered AS ( - -- ordered: для каждого клиента смотрим предыдущий hashdiff (LAG) + -- ordered: для каждого клиента смотрим "какая версия была до этого" (LAG) SELECT *, lag(hashdiff) OVER ( @@ -84,7 +92,7 @@ changes AS ( WHERE prev_hash IS DISTINCT FROM hashdiff OR prev_hash IS NULL ), framed AS ( - -- framed: превращаем изменения в периоды (valid_to = следующий event_ts через LEAD) + -- framed: превращаем изменения в периоды (valid_to = дата следующего события через LEAD) SELECT customer_bk, status, @@ -112,15 +120,15 @@ ORDER BY customer_bk, valid_from; -- 3) DDS: инкрементальная загрузка SCD2 (по последним событиям) -- ========================================================== --- Этот блок нужен, чтобы показать "как живёт SCD2 в проде": --- после каждой новой порции событий мы: +-- Этот блок нужен, чтобы показать "как это обычно обновляют": +-- после новой порции событий мы: -- 1) берём по каждому клиенту самое позднее событие из ODS; -- 2) сравниваем его с текущей версией в DDS (valid_to IS NULL); -- 3) если статус изменился — закрываем старую версию и вставляем новую. -- --- Ограничение учебного варианта: --- - если пришло "задним числом" событие со старой датой, этот инкремент историю не пересоберёт; --- для такого кейса нужен другой алгоритм (это уже advanced). +-- Ограничение учебного варианта (в домашке можно игнорировать): +-- - если вы добавили событие "задним числом" со старой датой, этот блок не пересоберёт всю историю. +-- Для такого кейса обычно делают отдельную логику или full refresh. -- -- Примечание: -- - в этом файле блок 2 (full refresh) запускается раньше, поэтому сразу после него @@ -129,9 +137,8 @@ ORDER BY customer_bk, valid_from; BEGIN; -- 3.1) Закрываем предыдущую актуальную версию WITH ranked AS ( - -- ranked: пронумеровали события так, чтобы rn = 1 было "самое свежее" на клиента - -- event_ts::date превращает событие в "изменение в этот день" (daily-grain). - -- _load_ts используем как tie-breaker, если два события имеют одинаковый event_ts. + -- ranked: выбираем "самое свежее" событие на клиента. + -- Если event_ts одинаковый, берём то, что загрузилось позже (_load_ts). SELECT customer_id AS customer_bk, status, @@ -158,7 +165,8 @@ BEGIN; SET valid_to = x.eff_date, updated_at = now() FROM ( - -- x: кандидаты на "закрытие" текущей версии (клиент есть в DDS и статус изменился) + -- x: кого "закрываем": + -- клиент уже есть в DDS, и статус действительно изменился. SELECT t.customer_bk, t.eff_date, @@ -174,7 +182,7 @@ BEGIN; -- 3.2) Вставляем новую версию WITH ranked AS ( - -- ranked/delta/current повторяем отдельно, чтобы блок INSERT читался автономно + -- ranked/delta/current повторяем отдельно, чтобы блок INSERT читался отдельно от UPDATE SELECT customer_id AS customer_bk, status, @@ -196,9 +204,9 @@ BEGIN; WHERE d.valid_to IS NULL ), to_insert AS ( - -- to_insert: кандидаты на вставку + -- to_insert: кого "вставляем": -- 1) новый клиент (в current нет строки); - -- 2) изменившийся клиент (hashdiff поменялся). + -- 2) изменившийся клиент (статус поменялся). SELECT t.customer_bk, t.status, @@ -233,8 +241,9 @@ COMMIT; -- 4) DM: витрина статусов клиентов по датам (full refresh) -- ========================================================== --- Витрина "снимок на дату": сколько клиентов в каком статусе на каждый день. --- Используем календарь dds.dim_date и JOIN по диапазону [valid_from, valid_to). +-- Витрина "снимок на дату": +-- для каждого дня считаем, сколько клиентов было в каждом статусе. +-- Берём календарь dds.dim_date и подбираем статус по периоду valid_from/valid_to. CREATE TABLE IF NOT EXISTS dm.mart_customer_status_daily ( date_actual DATE NOT NULL,