From 45f8f743781c3956723bf096d2061c38c6b182ea Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 21 Feb 2026 14:31:10 +0300 Subject: [PATCH] =?UTF-8?q?refactor(modeling):=20=D1=83=D0=BB=D1=83=D1=87?= =?UTF-8?q?=D1=88=D0=B5=D0=BD=D1=8B=20=D1=83=D1=87=D0=B5=D0=B1=D0=BD=D1=8B?= =?UTF-8?q?=D0=B5=20=D0=BC=D0=B0=D1=82=D0=B5=D1=80=D0=B8=D0=B0=D0=BB=D1=8B?= =?UTF-8?q?=20DWH=20=D0=BF=D0=BE=20=D1=80=D0=B5=D0=B7=D1=83=D0=BB=D1=8C?= =?UTF-8?q?=D1=82=D0=B0=D1=82=D0=B0=D0=BC=20=D1=80=D0=B5=D0=B2=D1=8C=D1=8E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - убрать путаницу, дублирование и неточности в демо-скриптах и домашке - Что: - CSV: заголовок load_ts → _load_ts во всех файлах (совпадает с именем в таблицах) - домашка: убраны оговорки о расхождении load_ts/_load_ts, добавлена ссылка на эталонное решение - 09_dml_hw_customer_status_solution.sql: эталонное решение скопировано из ветки solution/hw_customer_status в основную - 02_dml: добавлена карта загрузки в шапку (что откуда строится) - 05_ddl_dm + 06_dml_dm: total_orders → total_line_items (название точнее отражает содержимое) - Проверка: - визуальная проверка diff, скрипты не запускались (демо-стенд не поднят) Co-Authored-By: Claude Opus 4.6 --- .../Homework_Customer_Status_DDS_DM.md | 18 +- dwh-modeling/README.md | 2 +- dwh-modeling/data/customer_status_events.csv | 2 +- .../data/customer_status_events_increment.csv | 2 +- dwh-modeling/data/customers.csv | 2 +- dwh-modeling/sql/02_dml_stg-dds.sql | 7 + dwh-modeling/sql/05_ddl_dm.sql | 2 +- dwh-modeling/sql/06_dml_dm.sql | 4 +- .../09_dml_hw_customer_status_solution.sql | 276 ++++++++++++++++++ 9 files changed, 303 insertions(+), 12 deletions(-) create mode 100644 dwh-modeling/sql/09_dml_hw_customer_status_solution.sql diff --git a/dwh-modeling/Homework_Customer_Status_DDS_DM.md b/dwh-modeling/Homework_Customer_Status_DDS_DM.md index 55b6c2f..41bc004 100644 --- a/dwh-modeling/Homework_Customer_Status_DDS_DM.md +++ b/dwh-modeling/Homework_Customer_Status_DDS_DM.md @@ -30,7 +30,7 @@ Структура файла: ```text -customer_id,status,event_ts,_load_id,load_ts +customer_id,status,event_ts,_load_id,_load_ts 101,new,2024-01-01 09:00:00,batch_20240101_1000,2024-01-01 10:00:00 ... ``` @@ -41,7 +41,7 @@ customer_id,status,event_ts,_load_id,load_ts - `status` — статус клиента в CRM (`new`, `active`, `vip`, `churned`); - `event_ts` — момент, когда статус сменился в CRM; - `_load_id` — идентификатор батча загрузки; -- `load_ts` — момент, когда данные попали в DWH (в таблицах STG/ODS эта колонка будет называться `_load_ts`, но по смыслу это то же самое время загрузки). +- `_load_ts` — момент, когда данные попали в DWH. Файл содержит несколько клиентов и несколько смен статуса по каждому — этого достаточно, чтобы отработать SCD2. @@ -94,13 +94,13 @@ INSERT INTO stg.customer_status_raw (customer_id, status, event_ts, _load_id, _l ('101','churned','2024-09-01 12:15:00','batch_20240901_1300','2024-09-01 13:00:00'); ``` -> 💡 Здесь `_load_ts` — это время загрузки (в CSV оно называется `load_ts`). +> 💡 Здесь `_load_ts` — это время загрузки. #### Вариант B: загрузить CSV Можно загрузить файл `dwh-modeling/data/customer_status_events.csv` в таблицу `stg.customer_status_raw`: -- **Через DBeaver**: Import Data → CSV → `stg.customer_status_raw` (колонку `load_ts` маппить в `_load_ts`). +- **Через DBeaver**: Import Data → CSV → `stg.customer_status_raw`. - **Через `psql` в контейнере (`./psql_sh`)**: без установки `psql` на хост. Способ: передайте CSV в `psql` через STDIN и выполните `\copy ... FROM STDIN`: @@ -126,7 +126,7 @@ SELECT * FROM stg.customer_status_raw LIMIT 10; - привести: - `customer_id` → `INT`, - `status` → `VARCHAR(20)` (можно оставить как есть), - - `event_ts` и `load_ts` → `TIMESTAMP` (в DWH-таблицах эта колонка будет лежать как `_load_ts`); + - `event_ts` и `_load_ts` → `TIMESTAMP`; - аккуратно обработать возможные пустые значения (если бы они были); - заполнить `_load_id` и `_load_ts` в `ods.customer_status`. @@ -292,3 +292,11 @@ ORDER BY date_actual, status; - при желании — собрать простую витрину в `dm`. Если что‑то не получается — можно разбирать решения по шагам вместе с ментором: от простого `SELECT` из STG до полноценного SCD2 в DDS. + +--- + +## 8. Эталонное решение + +Когда выполните домашку и захотите сверить результат — готовое решение лежит в файле [`09_dml_hw_customer_status_solution.sql`](sql/09_dml_hw_customer_status_solution.sql). + +Постарайтесь не подглядывать до того, как напишете свой вариант — основная ценность задания именно в самостоятельном разборе. diff --git a/dwh-modeling/README.md b/dwh-modeling/README.md index b3b9ef2..b269be5 100644 --- a/dwh-modeling/README.md +++ b/dwh-modeling/README.md @@ -720,7 +720,7 @@ SELECT 'OK' WHERE EXISTS ( [`customers.csv`](data/customers.csv): ```csv -customer_id,email,phone,city,event_ts,_load_id,load_ts +customer_id,email,phone,city,event_ts,_load_id,_load_ts 101,a@ex.com,700,Москва,2024-01-01,batch_20240101_0800,2024-01-01 08:00 102,c@ex.com,701,СПб,2024-01-01,batch_20240101_0800,2024-01-01 08:00 101,b@ex.com,700,Москва,2024-05-16,batch_20240516_0800,2024-05-16 08:00 diff --git a/dwh-modeling/data/customer_status_events.csv b/dwh-modeling/data/customer_status_events.csv index 40a586e..218656e 100644 --- a/dwh-modeling/data/customer_status_events.csv +++ b/dwh-modeling/data/customer_status_events.csv @@ -1,4 +1,4 @@ -customer_id,status,event_ts,_load_id,load_ts +customer_id,status,event_ts,_load_id,_load_ts 101,new,2024-01-01 09:00:00,batch_20240101_1000,2024-01-01 10:00:00 101,active,2024-02-15 10:30:00,batch_20240215_1100,2024-02-15 11:00:00 101,vip,2024-05-10 11:00:00,batch_20240510_1200,2024-05-10 12:00:00 diff --git a/dwh-modeling/data/customer_status_events_increment.csv b/dwh-modeling/data/customer_status_events_increment.csv index 83d32d6..5bae97e 100644 --- a/dwh-modeling/data/customer_status_events_increment.csv +++ b/dwh-modeling/data/customer_status_events_increment.csv @@ -1,4 +1,4 @@ -customer_id,status,event_ts,_load_id,load_ts +customer_id,status,event_ts,_load_id,_load_ts 101,active,2024-11-15 09:00:00,batch_20241115_1000,2024-11-15 10:00:00 102,active,2024-05-05 09:30:00,batch_20240505_1000,2024-05-05 10:00:00 103,active,2024-03-20 12:00:00,batch_20240320_1300,2024-03-20 13:00:00 diff --git a/dwh-modeling/data/customers.csv b/dwh-modeling/data/customers.csv index 229fb6f..e9f2d25 100644 --- a/dwh-modeling/data/customers.csv +++ b/dwh-modeling/data/customers.csv @@ -1,4 +1,4 @@ -customer_id,email,phone,city,event_ts,_load_id,load_ts +customer_id,email,phone,city,event_ts,_load_id,_load_ts 101,a@ex.com,700,Москва,2024-01-01,batch_20240101_0800,2024-01-01 08:00 102,c@ex.com,701,СПб,2024-01-01,batch_20240101_0800,2024-01-01 08:00 101,b@ex.com,700,Москва,2024-05-16,batch_20240516_0800,2024-05-16 08:00 diff --git a/dwh-modeling/sql/02_dml_stg-dds.sql b/dwh-modeling/sql/02_dml_stg-dds.sql index dd4ed1f..1cf9dbf 100644 --- a/dwh-modeling/sql/02_dml_stg-dds.sql +++ b/dwh-modeling/sql/02_dml_stg-dds.sql @@ -2,6 +2,13 @@ -- DML-скрипт: загрузка и трансформация данных -- Запускается ПОВТОРНО при каждой загрузке (идемпотентно!) -- =============================================== +-- +-- Карта загрузки в этом скрипте: +-- STG → ODS: orders, order_items, products, customers (снимок: одна строка на BK) +-- STG → DDS: dim_customer (SCD2, full backfill напрямую из STG — см. комментарий к п.5) +-- ODS → DDS: dim_product, fact_sales +-- отдельно: dim_date (генерация календаря) +-- -- 1. STG: имитация загрузки из источников (в реальности — COPY или INSERT из Kafka/NiFi) -- ⚠️ В продакшене STG часто очищается перед загрузкой (TRUNCATE), либо используется партицирование по дате diff --git a/dwh-modeling/sql/05_ddl_dm.sql b/dwh-modeling/sql/05_ddl_dm.sql index 642feba..214bb46 100644 --- a/dwh-modeling/sql/05_ddl_dm.sql +++ b/dwh-modeling/sql/05_ddl_dm.sql @@ -21,7 +21,7 @@ CREATE TABLE dm.mart_customer_360 ( customer_bk INT NOT NULL, first_order_date DATE, last_order_date DATE, - total_orders INT NOT NULL, + total_line_items INT NOT NULL, total_items INT NOT NULL, lifetime_value NUMERIC(18,2) NOT NULL, last_email VARCHAR(100), diff --git a/dwh-modeling/sql/06_dml_dm.sql b/dwh-modeling/sql/06_dml_dm.sql index 29d3b63..d4b81ab 100644 --- a/dwh-modeling/sql/06_dml_dm.sql +++ b/dwh-modeling/sql/06_dml_dm.sql @@ -33,14 +33,14 @@ GROUP BY d.date_actual, p.product_name, -- Считаем суммы по всей истории его покупок INSERT INTO dm.mart_customer_360 ( customer_bk, first_order_date, last_order_date, - total_orders, total_items, lifetime_value, + total_line_items, total_items, lifetime_value, last_email, last_city ) SELECT c.customer_bk, MIN(d.date_actual) AS first_order_date, MAX(d.date_actual) AS last_order_date, - COUNT(DISTINCT f.sale_id) AS total_orders, -- считаем строки факта (продажи), не бизнес-заказы + COUNT(DISTINCT f.sale_id) AS total_line_items, -- строки факта (позиции продаж), не бизнес-заказы SUM(f.quantity) AS total_items, SUM(f.amount) AS lifetime_value, -- Берём самый свежий email и город клиента diff --git a/dwh-modeling/sql/09_dml_hw_customer_status_solution.sql b/dwh-modeling/sql/09_dml_hw_customer_status_solution.sql new file mode 100644 index 0000000..37c3661 --- /dev/null +++ b/dwh-modeling/sql/09_dml_hw_customer_status_solution.sql @@ -0,0 +1,276 @@ +-- =============================================== +-- 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 ещё раз. +-- +-- Важно: +-- - здесь часто используется TRUNCATE (полная очистка), чтобы было легко повторять домашку; +-- - в реальном DWH так делают не всегда, но для обучения это удобнее. +-- +-- Предусловия (DDL + данные в STG): +-- 1) dwh-modeling/sql/01_ddl_stg-dds.sql +-- 2) dwh-modeling/sql/05_ddl_dm.sql +-- 3) dwh-modeling/sql/07_ddl_hw_customer_status.sql +-- 4) stg.customer_status_raw заполнена (см. dwh-modeling/Homework_Customer_Status_DDS_DM.md) +-- =============================================== + +-- ========================================================== +-- 1) ODS: очистка и типизация (full refresh) +-- ========================================================== + +-- Идея: +-- - STG хранит "как пришло" (обычно TEXT); +-- - ODS хранит "аккуратно": правильные типы + простая чистка. +-- Для простоты пересобираем ODS с нуля. + +TRUNCATE ods.customer_status; + +INSERT INTO ods.customer_status ( + customer_id, status, event_ts, _load_id, _load_ts +) +SELECT + s.customer_id::INT AS customer_id, + NULLIF(trim(s.status), '') AS status, + NULLIF(trim(s.event_ts), '')::TIMESTAMP AS event_ts, + s._load_id, + COALESCE(s._load_ts, now()) AS _load_ts +FROM stg.customer_status_raw s +WHERE s.customer_id ~ '^\d+$' + AND NULLIF(trim(s.event_ts), '') IS NOT NULL + AND NULLIF(trim(s.status), '') IS NOT NULL; + +-- ========================================================== +-- 2) DDS: начальная загрузка SCD2 (full refresh) +-- ========================================================== + +-- Идея SCD2 простыми словами: +-- - одна строка = один период, когда статус был одним и тем же; +-- - valid_from = с какого дня статус "начался"; +-- - valid_to = с какого дня статус "закончился" (NULL = текущий статус); +-- - интервалы считаем так: [valid_from, valid_to) (valid_to не включаем). +-- - чтобы найти статус "на дату D": +-- D >= valid_from AND (valid_to IS NULL OR D < valid_to) +-- +-- Упрощение для домашки: +-- - считаем, что у клиента нет двух разных смен статуса в один день. + +TRUNCATE dds.dim_customer_status; + +WITH src AS ( + -- src: события из ODS + "контрольная сумма" статуса. + -- Так проще проверять, поменялся статус или остался тем же. + SELECT + customer_id AS customer_bk, + status, + event_ts, + md5(lower(coalesce(status, ''))) AS hashdiff + FROM ods.customer_status +), +ordered AS ( + -- ordered: для каждого клиента смотрим "какая версия была до этого" (LAG) + SELECT + *, + lag(hashdiff) OVER ( + PARTITION BY customer_bk + ORDER BY event_ts + ) AS prev_hash + FROM src +), +changes AS ( + -- changes: оставляем только первое состояние и реальные изменения статуса + SELECT * + FROM ordered + WHERE prev_hash IS DISTINCT FROM hashdiff OR prev_hash IS NULL +), +framed AS ( + -- framed: превращаем изменения в периоды (valid_to = дата следующего события через LEAD) + SELECT + customer_bk, + status, + hashdiff, + event_ts::DATE AS valid_from, + lead(event_ts::DATE) OVER ( + PARTITION BY customer_bk + ORDER BY event_ts + ) AS valid_to + FROM changes +) +INSERT INTO dds.dim_customer_status ( + customer_bk, status, hashdiff, + valid_from, valid_to, + created_at, updated_at +) +SELECT + customer_bk, status, hashdiff, + valid_from, valid_to, + now(), now() +FROM framed +ORDER BY customer_bk, valid_from; + +-- ========================================================== +-- 3) DDS: инкрементальная загрузка SCD2 (по последним событиям) +-- ========================================================== + +-- Этот блок нужен, чтобы показать "как это обычно обновляют": +-- после новой порции событий мы: +-- 1) берём по каждому клиенту самое позднее событие из ODS; +-- 2) сравниваем его с текущей версией в DDS (valid_to IS NULL); +-- 3) если статус изменился — закрываем старую версию и вставляем новую. +-- +-- Ограничение учебного варианта (в домашке можно игнорировать): +-- - если вы добавили событие "задним числом" со старой датой, этот блок не пересоберёт всю историю. +-- Для такого кейса обычно делают отдельную логику или full refresh. +-- +-- Примечание: +-- - в этом файле блок 2 (full refresh) запускается раньше, поэтому сразу после него +-- блок 3, скорее всего, ничего не изменит. Зато его можно повторять после новых событий. + +BEGIN; + -- 3.1) Закрываем предыдущую актуальную версию + WITH ranked AS ( + -- ranked: выбираем "самое свежее" событие на клиента. + -- Если event_ts одинаковый, берём то, что загрузилось позже (_load_ts). + SELECT + customer_id AS customer_bk, + status, + event_ts::DATE AS eff_date, + md5(lower(coalesce(status, ''))) AS hashdiff, + row_number() OVER ( + PARTITION BY customer_id + ORDER BY event_ts DESC, _load_ts DESC + ) AS rn + FROM ods.customer_status + WHERE event_ts IS NOT NULL + ), + delta AS ( + -- delta: ровно одна строка на клиента (самое свежее событие) + SELECT * FROM ranked WHERE rn = 1 + ), + current AS ( + -- current: текущие версии в DDS (valid_to IS NULL) + SELECT d.* + FROM dds.dim_customer_status d + WHERE d.valid_to IS NULL + ) + UPDATE dds.dim_customer_status d + SET valid_to = x.eff_date, + updated_at = now() + FROM ( + -- x: кого "закрываем": + -- клиент уже есть в DDS, и статус действительно изменился. + SELECT + t.customer_bk, + t.eff_date, + c.customer_status_sk + FROM delta t + JOIN current c + ON c.customer_bk = t.customer_bk + WHERE c.hashdiff <> t.hashdiff + AND t.eff_date > c.valid_from -- не создаём период нулевой/отрицательной длины + ) x + WHERE d.customer_status_sk = x.customer_status_sk + AND d.valid_to IS NULL; + + -- 3.2) Вставляем новую версию + WITH ranked AS ( + -- ranked/delta/current повторяем отдельно, чтобы блок INSERT читался отдельно от UPDATE + SELECT + customer_id AS customer_bk, + status, + event_ts::DATE AS eff_date, + md5(lower(coalesce(status, ''))) AS hashdiff, + row_number() OVER ( + PARTITION BY customer_id + ORDER BY event_ts DESC, _load_ts DESC + ) AS rn + FROM ods.customer_status + WHERE event_ts IS NOT NULL + ), + delta AS ( + SELECT * FROM ranked WHERE rn = 1 + ), + current AS ( + SELECT d.* + FROM dds.dim_customer_status d + WHERE d.valid_to IS NULL + ), + to_insert AS ( + -- to_insert: кого "вставляем": + -- 1) новый клиент (в current нет строки); + -- 2) изменившийся клиент (статус поменялся). + SELECT + t.customer_bk, + t.status, + t.hashdiff, + t.eff_date + FROM delta t + LEFT JOIN current c + ON c.customer_bk = t.customer_bk + WHERE c.customer_status_sk IS NULL + OR (c.hashdiff <> t.hashdiff AND t.eff_date > c.valid_from) + ) + INSERT INTO dds.dim_customer_status ( + customer_bk, status, hashdiff, + valid_from, valid_to, + created_at, updated_at + ) + SELECT + t.customer_bk, t.status, t.hashdiff, + t.eff_date, NULL, + now(), now() + FROM to_insert t + -- защита от повторного запуска: не вставляем одну и ту же версию (BK + valid_from) второй раз + WHERE NOT EXISTS ( + SELECT 1 + FROM dds.dim_customer_status d + WHERE d.customer_bk = t.customer_bk + AND d.valid_from = t.eff_date + ); +COMMIT; + +-- ========================================================== +-- 4) DM: витрина статусов клиентов по датам (full refresh) +-- ========================================================== + +-- Витрина "снимок на дату": +-- для каждого дня считаем, сколько клиентов было в каждом статусе. +-- Берём календарь dds.dim_date и подбираем статус по периоду valid_from/valid_to. + +CREATE TABLE IF NOT EXISTS dm.mart_customer_status_daily ( + date_actual DATE NOT NULL, + status VARCHAR(20) NOT NULL, + customers_cnt INT NOT NULL +); + +TRUNCATE dm.mart_customer_status_daily; + +WITH bounds AS ( + SELECT + min(valid_from) AS date_from, + max(coalesce(valid_to, valid_from)) AS date_to + FROM dds.dim_customer_status +) +INSERT INTO dm.mart_customer_status_daily ( + date_actual, status, customers_cnt +) +SELECT + d.date_actual, + s.status, + COUNT(DISTINCT s.customer_bk) AS customers_cnt +FROM dds.dim_date d +JOIN bounds b + ON d.date_actual BETWEEN b.date_from AND b.date_to +JOIN dds.dim_customer_status s + ON d.date_actual >= s.valid_from + AND (s.valid_to IS NULL OR d.date_actual < s.valid_to) +GROUP BY d.date_actual, s.status +ORDER BY d.date_actual, s.status;