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