From 2d78fb4459bd718c436fcdb19b32c68809829fe6 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 23 Nov 2025 18:31:43 +0300 Subject: [PATCH 1/5] =?UTF-8?q?=D0=94=D0=BE=D0=BC=D0=B0=D1=88=D0=BA=D0=B0?= =?UTF-8?q?=20=D0=BF=D0=BE=20=D0=BC=D0=BE=D0=B4=D0=B5=D0=BB=D0=B8=D1=80?= =?UTF-8?q?=D0=BE=D0=B2=D0=B0=D0=BD=D0=B8=D1=8E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- README.md | 1 + .../Homework_Customer_Status_DDS_DM.md | 241 ++++++++++++++++++ dwh-modeling/README.md | 7 +- dwh-modeling/data/customer_status_events.csv | 10 + .../sql/07_ddl_hw_customer_status.sql | 43 ++++ .../08_dml_hw_customer_status_template.sql | 48 ++++ 6 files changed, 349 insertions(+), 1 deletion(-) create mode 100644 dwh-modeling/Homework_Customer_Status_DDS_DM.md create mode 100644 dwh-modeling/data/customer_status_events.csv create mode 100644 dwh-modeling/sql/07_ddl_hw_customer_status.sql create mode 100644 dwh-modeling/sql/08_dml_hw_customer_status_template.sql diff --git a/README.md b/README.md index 2e2a34a..a240d9d 100644 --- a/README.md +++ b/README.md @@ -57,6 +57,7 @@ - [Базы данных. 1,2,3 нормальные формы. - Youtube](https://www.youtube.com/watch?v=zwQzL80U51c) - [Введение в структуру хранилища данных](dwh-modeling/README.md) - Теория про Slowly Changing Dimensions: [SCD](dwh-modeling/SCD.md) +- Практика по моделированию статусов клиента: [домашка STG → ODS → DDS → DM](dwh-modeling/Homework_Customer_Status_DDS_DM.md) - Еще про Data Vault: - [Введение в Data Vault - Хабр](https://habr.com/ru/articles/348188/) - [Основы Data Vault - Хабр](https://habr.com/ru/articles/502968/) diff --git a/dwh-modeling/Homework_Customer_Status_DDS_DM.md b/dwh-modeling/Homework_Customer_Status_DDS_DM.md new file mode 100644 index 0000000..dacb976 --- /dev/null +++ b/dwh-modeling/Homework_Customer_Status_DDS_DM.md @@ -0,0 +1,241 @@ +# Домашка: статусы клиента от STG до DDS (и немного DM) + +Небольшое практическое задание на 1–2 вечера: по данным о смене статусов клиента (CRM) построить цепочку слоёв `STG → ODS → DDS (SCD2)` и, по желанию, небольшую витрину в `dm`. + +Цель — потренировать **руками**: + +- работу со слоями DWH (stg / ods / dds / dm); +- проектирование и загрузку **измерения с историей (SCD Type 2)**; +- аккуратную работу со временем (`event_ts`, `valid_from`, `valid_to`). + +Исходим из того, что вы уже прошли основную статью `dwh-modeling/README.md` и познакомились с примером интернет‑магазина. + +--- + +## 1. Данные: события смены статуса клиента + +Представьте, что в CRM для каждого клиента хранится история статусов: + +- `new` — только что зарегистрировался; +- `active` — делал покупки недавно; +- `vip` — часто покупает и много тратит; +- `churned` — давно ничего не делал, считаем «отвалившимся». + +Эта информация приходит в DWH в виде **событий** (events): «у клиента X в момент времени Y статус стал Z». + +В репозитории в каталоге `dwh-modeling/data` лежит файл: + +- `customer_status_events.csv` + +Структура файла: + +```text +customer_id,status,event_ts,_load_id,load_ts +101,new,2024-01-01 09:00:00,batch_20250405_0800,2025-04-05 08:00:00 +... +``` + +Колонки: + +- `customer_id` — бизнес-ключ клиента (тот же, что и в основном примере — 101, 102, 103); +- `status` — статус клиента в CRM (`new`, `active`, `vip`, `churned`); +- `event_ts` — момент, когда статус сменился в CRM; +- `_load_id` — идентификатор батча загрузки; +- `load_ts` — момент, когда данные попали в DWH. + +Файл содержит несколько клиентов и несколько смен статуса по каждому — этого достаточно, чтобы отработать SCD2. + +--- + +## 2. Целевая схема: какие таблицы уже есть + +Чтобы не тратить время на DDL, структуры таблиц для домашки уже подготовлены в `dwh-modeling/sql`: + +- `07_ddl_hw_customer_status.sql` — создаёт дополнительные таблицы: + - `stg.customer_status_raw` — сырые события о статусе клиента; + - `ods.customer_status` — очищенные и типизированные события; + - `dds.dim_customer_status` — измерение статусов клиента в формате **SCD Type 2**. +- `08_dml_hw_customer_status_template.sql` — шаблон DML-скрипта с подсказками и заготовками блоков. + +Перед началом работы: + +1. Поднимите demo‑Postgres по инструкции из корневого `README.md`. +2. Выполните базовые скрипты DWH: + - `01_ddl_stg-dds.sql` + - `02_dml_stg-dds.sql` +3. Выполните DDL для домашки: + - `07_ddl_hw_customer_status.sql` + +После этого схемы `stg`, `ods`, `dds` уже существуют, а дополнительные таблицы для статусов созданы. + +--- + +## 3. Часть 1 — STG → ODS (обязательно) + +**Задача:** загрузить CSV в STG и переложить данные в ODS с приведением типов. + +### 3.1. STG: загрузка CSV + +1. Откройте `psql` в контейнере (см. инструкцию в `postgres-bookings/README.md`). +2. Загрузите файл `customer_status_events.csv` в таблицу `stg.customer_status_raw`: + - можно использовать `\copy` из `psql`; + - либо любой другой способ, к которому вы привыкли. +3. Убедитесь, что данные загрузились: + +```sql +SELECT * FROM stg.customer_status_raw LIMIT 10; +``` + +### 3.2. ODS: очистка и типизация + +В файле `08_dml_hw_customer_status_template.sql` найдите заготовку блока ODS и допишите SQL: + +- привести: + - `customer_id` → `INT`, + - `status` → `VARCHAR(20)` (можно оставить как есть), + - `event_ts` и `load_ts` → `TIMESTAMP`; +- аккуратно обработать возможные пустые значения (если бы они были); +- заполнить `_load_id` и `_load_ts` в `ods.customer_status`. + +Проверьте, что в `ods.customer_status` данные выглядят аккуратно: + +```sql +SELECT * +FROM ods.customer_status +ORDER BY customer_id, event_ts; +``` + +--- + +## 4. Часть 2 — ODS → DDS (SCD Type 2, обязательно) + +**Задача:** по событиям в `ods.customer_status` построить измерение `dds.dim_customer_status`, где каждая строка — период действия статуса. + +Целевая таблица уже создана (см. `07_ddl_hw_customer_status.sql`): + +- `customer_bk` — бизнес-ключ клиента (тот же, что `customer_id` в ODS); +- `status` — статус клиента; +- `hashdiff` — хэш от атрибутов (здесь достаточно самого `status`); +- `valid_from` / `valid_to` — период, когда статус был актуален; +- `is_current` — флаг актуальной строки; +- `created_at` / `updated_at` — технические поля. + +### 4.1. Начальная загрузка SCD2 + +В шаблоне `08_dml_hw_customer_status_template.sql` допишите блок начальной загрузки: + +1. Сформируйте промежуточный набор: + + - `customer_bk`, + - `status`, + - `event_ts` (как «время начала действия статуса»), + - `hashdiff` (например, `md5(status)` или `dds.customer_hash(status)` по аналогии с `dim_customer`). + +2. Для каждого клиента отсортируйте события по `event_ts` и с помощью `LEAD()` посчитайте: + + - `valid_from` — текущее `event_ts`, + - `valid_to` — следующее `event_ts - 1 second` (или `9999-12-31`, если следующего нет). + +3. Вставьте получившиеся строки в `dds.dim_customer_status`: + +```sql +INSERT INTO dds.dim_customer_status ( + customer_bk, status, hashdiff, + valid_from, valid_to, + is_current, created_at, updated_at +) +SELECT + ... +``` + +Проверьте результат: + +```sql +SELECT * +FROM dds.dim_customer_status +ORDER BY customer_bk, valid_from; +``` + +Ожидаемое поведение: + +- у клиента 101 несколько строк с разными статусами и непересекающимися периодами; +- `is_current = TRUE` только у самой свежей строки для каждого клиента. + +### 4.2. Проверка себя + +Примеры проверочных запросов (можно придумать свои): + +- «Какой статус был у клиента 101 на дату `2024-06-01`?» + → одна строка с нужным статусом. +- «Сколько клиентов были в статусе `active` на `2024-04-10`?» + → несколько строк, если статус *активен* для диапазона дат. + +--- + +## 5. Часть 3 — инкрементальная загрузка (по желанию) + +Если хочется потренироваться глубже: + +1. Добавьте в CSV ещё несколько событий смены статуса (например, переход части клиентов из `churned` обратно в `active`). +2. Загрузите новые строки только в `stg.customer_status_raw`. +3. Напишите логику инкрементального обновления `dds.dim_customer_status`: + +- ориентируйтесь на пример из `03_demo_increment.sql` для `dds.dim_customer`; +- важно: + - корректно «закрыть» старую актуальную строку (`valid_to`, `is_current = FALSE`); + - вставить новую строку с `is_current = TRUE`. + +Эта часть особенно полезна, если вы хотите почувствовать, как SCD2 живёт в реальном DWH. + +--- + +## 6. Часть 4 — витрина в DM (по желанию) + +Опциональное задание для закрепления: собрать небольшую витрину с количеством клиентов по статусам на каждую дату. + +Пример целевой таблицы: + +```sql +CREATE TABLE dm.mart_customer_status_daily ( + date_actual DATE NOT NULL, + status VARCHAR(20) NOT NULL, + customers_cnt INT NOT NULL +); +``` + +Идея: + +- использовать `dds.dim_date` как календарь; +- для каждой `date_actual` найти, какой статус был у клиента в этот день + (через `JOIN` на `dds.dim_customer_status` по диапазону `valid_from/valid_to`); +- агрегировать по `status`. + +Пример запроса к витрине: + +```sql +SELECT + date_actual, + status, + customers_cnt +FROM dm.mart_customer_status_daily +WHERE date_actual BETWEEN '2024-04-01' AND '2024-04-30' +ORDER BY date_actual, status; +``` + +--- + +## 7. Как вписать эту домашку в обучение + +Рекомендуемое место в дорожке: + +1. Пройти основную теорию по DWH и SCD: + - `dwh-modeling/README.md` + - `dwh-modeling/SCD.md` +2. Разобрать базовый пример интернет‑магазина (скрипты `01_`–`06_`). +3. Выполнить **эту домашку** как первую попытку «самостоятельного» моделирования и ETL: + - познакомиться с ещё одним измерением с историей (`dim_customer_status`); + - потренироваться аккуратно работать с датами и периодами; + - при желании — собрать простую витрину в `dm`. + +Если что‑то не получается — можно разбирать решения по шагам вместе с ментором: от простого `SELECT` из STG до полноценного SCD2 в DDS. + diff --git a/dwh-modeling/README.md b/dwh-modeling/README.md index 352b8dd..8fb6718 100644 --- a/dwh-modeling/README.md +++ b/dwh-modeling/README.md @@ -723,6 +723,12 @@ SELECT 'OK' WHERE EXISTS ( ## Приложения +### Дополнительные материалы и практика + +- [Домашка: статусы клиента от STG до DDS (и немного DM)](Homework_Customer_Status_DDS_DM.md) +- [SCD: как хранить историю изменений](SCD.md) +- [DataVault: как пережить бурную жизнь источников](DataVault.md) + ### 📚 Мини-глоссарий (RU / EN) | Термин | Пояснение | @@ -839,4 +845,3 @@ CREATE TABLE dds.fact_sales ( ``` > 💡 `date_key` — это `20240110`, а не `DATE`, чтобы не делать JOIN по диапазону в `fact → dim_date`. - diff --git a/dwh-modeling/data/customer_status_events.csv b/dwh-modeling/data/customer_status_events.csv new file mode 100644 index 0000000..493886a --- /dev/null +++ b/dwh-modeling/data/customer_status_events.csv @@ -0,0 +1,10 @@ +customer_id,status,event_ts,_load_id,load_ts +101,new,2024-01-01 09:00:00,batch_20250405_0800,2025-04-05 08:00:00 +101,active,2024-02-15 10:30:00,batch_20250405_0800,2025-04-05 08:00:00 +101,vip,2024-05-10 11:00:00,batch_20250405_1200,2025-04-05 12:00:00 +101,churned,2024-09-01 12:15:00,batch_20250405_1800,2025-04-05 18:00:00 +102,new,2024-03-05 14:00:00,batch_20250405_0800,2025-04-05 08:00:00 +102,active,2024-04-01 09:45:00,batch_20250405_1200,2025-04-05 12:00:00 +102,churned,2024-04-20 16:20:00,batch_20250405_1800,2025-04-05 18:00:00 +103,new,2024-03-10 10:10:00,batch_20250405_0800,2025-04-05 08:00:00 + diff --git a/dwh-modeling/sql/07_ddl_hw_customer_status.sql b/dwh-modeling/sql/07_ddl_hw_customer_status.sql new file mode 100644 index 0000000..c3c2946 --- /dev/null +++ b/dwh-modeling/sql/07_ddl_hw_customer_status.sql @@ -0,0 +1,43 @@ +-- =============================================== +-- DDL: дополнительные таблицы для домашки +-- Тема: статусы клиента (SCD2 поверх статуса) +-- Скрипт можно запускать после 01_ddl_stg-dds.sql +-- =============================================== + +-- 1. STG: сырые события о статусе клиента из CRM +DROP TABLE IF EXISTS stg.customer_status_raw; +CREATE TABLE stg.customer_status_raw ( + customer_id TEXT, + status TEXT, + event_ts TEXT, + _load_id TEXT, + _load_ts TIMESTAMP DEFAULT NOW() +); + +-- 2. ODS: очищенные и типизированные статусы +DROP TABLE IF EXISTS ods.customer_status; +CREATE TABLE ods.customer_status ( + customer_id INT NOT NULL, + status VARCHAR(20) NOT NULL, + event_ts TIMESTAMP NOT NULL, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL +); + +ALTER TABLE ods.customer_status + ADD PRIMARY KEY (customer_id, event_ts); + +-- 3. DDS: измерение статусов клиента с историей (SCD Type 2) +DROP TABLE IF EXISTS dds.dim_customer_status; +CREATE TABLE dds.dim_customer_status ( + customer_status_sk BIGSERIAL PRIMARY KEY, + customer_bk INT NOT NULL, + status VARCHAR(20) NOT NULL, + hashdiff TEXT NOT NULL, + valid_from TIMESTAMP NOT NULL, + valid_to TIMESTAMP NOT NULL, + is_current BOOLEAN NOT NULL DEFAULT TRUE, + created_at TIMESTAMP NOT NULL DEFAULT NOW(), + updated_at TIMESTAMP NOT NULL DEFAULT NOW() +); + diff --git a/dwh-modeling/sql/08_dml_hw_customer_status_template.sql b/dwh-modeling/sql/08_dml_hw_customer_status_template.sql new file mode 100644 index 0000000..32f5f2a --- /dev/null +++ b/dwh-modeling/sql/08_dml_hw_customer_status_template.sql @@ -0,0 +1,48 @@ +-- =============================================== +-- DML-шаблон для домашки +-- Тема: статусы клиента (STG → ODS → DDS SCD2) +-- Цель: по customer_status_events.csv построить историю статусов +-- =============================================== + +-- Подсказка: +-- 1) Загрузите CSV в stg.customer_status_raw (через COPY или \copy в psql). +-- См. пример структуры файла в dwh-modeling/data/customer_status_events.csv +-- 2) Переложите данные в ods.customer_status с приведением типов. +-- customer_id → INT, status → VARCHAR(20), event_ts / load_ts → TIMESTAMP. +-- 3) Постройте из ods.customer_status измерение dds.dim_customer_status в стиле SCD2: +-- - одна строка на период действия статуса (valid_from / valid_to); +-- - is_current = TRUE только у актуальной строки для клиента; +-- - hashdiff можно считать, например, от одного поля status. +-- 4) При желании добавьте инкрементальную логику (как в 03_demo_increment.sql). +-- 5) Опционально: соберите витрину dm.mart_customer_status_daily +-- с количеством клиентов по статусам на каждую дату. + +-- Ниже — ЗАГОТОВКИ блоков, которые можно дописать. +-- Они намеренно оставлены пустыми, чтобы вы написали SQL сами. + +-- 1. ODS: очистка и типизация +-- TRUNCATE ods.customer_status; +-- INSERT INTO ods.customer_status (...) +-- SELECT ... +-- FROM stg.customer_status_raw; + +-- 2. DDS: начальная загрузка SCD2 +-- Примерный план: +-- - рассчитать hashdiff по (status); +-- - по каждому клиенту отсортировать события по времени; +-- - построить для каждой строки valid_from и valid_to (LEAD() OVER ...); +-- - вставить в dds.dim_customer_status. + +-- 3. DDS: инкрементальная загрузка (по желанию) +-- Можно ориентироваться на примеры в 03_demo_increment.sql. + +-- 4. DM: витрина статусов клиентов по датам (по желанию) +-- Пример целевой структуры: +-- CREATE TABLE dm.mart_customer_status_daily ( +-- date_actual DATE NOT NULL, +-- status VARCHAR(20) NOT NULL, +-- customers_cnt INT NOT NULL +-- ); +-- Идея: на каждую дату взять актуальный статус клиента +-- через JOIN dds.dim_customer_status + dds.dim_date. + From 218c6b9fe781ea8073951ae397d80a831bfd8f44 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 29 Nov 2025 20:16:16 +0300 Subject: [PATCH 2/5] =?UTF-8?q?=D0=98=D1=81=D0=BF=D1=80=D0=B0=D0=B2=D0=BB?= =?UTF-8?q?=D0=B8=D0=B5=D0=BD=D0=B8=D0=B5=20=D0=BF=D1=80=D0=B8=D0=BC=D0=B5?= =?UTF-8?q?=D1=80=D0=B0=20=D1=81=20scd2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- dwh-modeling/README.md | 12 ++-- dwh-modeling/SCD.md | 7 +++ dwh-modeling/sql/02_dml_stg-dds.sql | 4 +- dwh-modeling/sql/03_demo_increment.sql | 80 +++++++++++++++++--------- 4 files changed, 68 insertions(+), 35 deletions(-) diff --git a/dwh-modeling/README.md b/dwh-modeling/README.md index 8fb6718..5d9b5b6 100644 --- a/dwh-modeling/README.md +++ b/dwh-modeling/README.md @@ -529,12 +529,12 @@ flowchart TD Все необходимые скрипты для построения хранилища находятся в папке [`sql/`](sql/): -- [`01_ddl_stg-dds.sql`](sql/01_ddl_stg-dds.sql) — создание схем и таблиц (STG, ODS, DDS) -- [`02_dml_stg-dds.sql`](sql/02_dml_stg-dds.sql) — загрузка данных и трансформация -- [`03_demo_increment.sql`](sql/03_demo_increment.sql) — инкрементальная загрузка и SCD2 -- [`04_validation.sql`](sql/04_validation.sql) — проверки качества данных -- [`05_ddl_dm.sql`](sql/05_ddl_dm.sql) — создание витрин (Data Marts) -- [`06_dml_dm.sql`](sql/06_dml_dm.sql) — наполнение витрин данными +- [`01_ddl_stg-dds.sql`](sql/01_ddl_stg-dds.sql) — создание схем и таблиц (STG, ODS, DDS); +- [`02_dml_stg-dds.sql`](sql/02_dml_stg-dds.sql) — первичная загрузка данных и демонстрация SCD2 через полный пересчёт (`full backfill`) из STG; +- [`03_demo_increment.sql`](sql/03_demo_increment.sql) — пример инкрементальной загрузки и SCD2 по последнему снимку в ODS; +- [`04_validation.sql`](sql/04_validation.sql) — проверки качества данных; +- [`05_ddl_dm.sql`](sql/05_ddl_dm.sql) — создание витрин (Data Marts); +- [`06_dml_dm.sql`](sql/06_dml_dm.sql) — наполнение витрин данными. ### Пример SQL-запроса для витрины diff --git a/dwh-modeling/SCD.md b/dwh-modeling/SCD.md index 8d674dc..2ab2bf8 100644 --- a/dwh-modeling/SCD.md +++ b/dwh-modeling/SCD.md @@ -230,6 +230,13 @@ WHERE customer_id = 1 AND is_current = true --- +В учебном проекте из папки `dwh-modeling/sql/` эти идеи можно увидеть «вживую»: + +- в [`02_dml_stg-dds.sql`](sql/02_dml_stg-dds.sql) собирается полная история клиентов (SCD2) из всех событий в `stg.customers_raw` — это пример **первичной загрузки** / `full backfill`; +- в [`03_demo_increment.sql`](sql/03_demo_increment.sql) реализован **инкрементальный SCD2**: в одной транзакции добавляются новые версии клиентов из снимка `ods.customers` и закрываются предыдущие актуальные строки в `dds.dim_customer`. + +--- + #### А что, если СУБД не позволяет UPDATE? (Trino, Hive, ClickHouse в режиме append-only) Некоторые аналитические системы (например, **Hive в формате ORC/Parquet**, **Trino**, **ClickHouse в режиме только вставки**) **не поддерживают UPDATE старых строк**. Как тогда реализовать SCD Type 2? diff --git a/dwh-modeling/sql/02_dml_stg-dds.sql b/dwh-modeling/sql/02_dml_stg-dds.sql index 3eeecb4..170cd18 100644 --- a/dwh-modeling/sql/02_dml_stg-dds.sql +++ b/dwh-modeling/sql/02_dml_stg-dds.sql @@ -128,6 +128,8 @@ SELECT product_id, name FROM ods.products; -- 5. DDS: dim_customer — первичная загрузка SCD2 (full backfill из STG) +-- Здесь мы пересчитываем всю историю клиента из событий в STG. +-- В реальном DWH такую полную перезагрузку делают редко; для инкремента см. 03_demo_increment.sql и SCD.md. TRUNCATE dds.dim_customer, dds.fact_sales; BEGIN; @@ -158,7 +160,7 @@ BEGIN; framed AS ( SELECT customer_bk, email, phone, city, hashdiff, - CASE WHEN rn = 1 THEN timestamp '1900-01-01' ELSE eff_ts END AS valid_from, + CASE WHEN rn = 1 THEN timestamp '1900-01-01' ELSE eff_ts END AS valid_from, -- первая версия: техническое "начало истории" lead(eff_ts) OVER (PARTITION BY customer_bk ORDER BY eff_ts) AS next_ts FROM changes ) diff --git a/dwh-modeling/sql/03_demo_increment.sql b/dwh-modeling/sql/03_demo_increment.sql index 5ae402d..e15d44b 100644 --- a/dwh-modeling/sql/03_demo_increment.sql +++ b/dwh-modeling/sql/03_demo_increment.sql @@ -43,7 +43,7 @@ SET email = EXCLUDED.email, WHERE COALESCE(EXCLUDED.event_ts, EXCLUDED._load_ts) > COALESCE(ods.customers.event_ts, ods.customers._load_ts); --- 2) Инкрементальное SCD2 из ODS +-- 2) Инкрементальное SCD2 из ODS (в одной транзакции, по последнему снимку в ODS) BEGIN; WITH delta AS ( SELECT @@ -53,34 +53,58 @@ BEGIN; dds.customer_hash(c.email, c.phone, c.city) AS hashdiff FROM ods.customers c ), - expired AS ( - UPDATE dds.dim_customer d - SET valid_to = LEAST(d.valid_to, delta.eff_ts - interval '1 second'), - is_current = FALSE, - updated_at = now() - FROM delta - WHERE d.customer_bk = delta.customer_bk - AND d.is_current = TRUE - AND d.hashdiff <> delta.hashdiff - AND delta.eff_ts >= d.valid_from - RETURNING d.customer_bk + current AS ( + -- текущие версии клиентов в измерении + SELECT d.* + FROM dds.dim_customer d + WHERE d.is_current = TRUE + ), + to_upsert AS ( + -- только новые BK или реально изменившиеся атрибуты + SELECT + d.customer_bk, + d.email, + d.phone, + d.city, + d.eff_ts, + d.hashdiff, + c.customer_sk AS current_sk + FROM delta d + LEFT JOIN current c + ON c.customer_bk = d.customer_bk + WHERE c.customer_sk IS NULL -- новый клиент + OR c.hashdiff <> d.hashdiff -- изменились атрибуты + ), + inserted AS ( + -- вставляем новые версии (одна строка на BK) + INSERT INTO dds.dim_customer ( + customer_bk, email, phone, city, hashdiff, + valid_from, valid_to, + is_current, created_at, updated_at + ) + SELECT + t.customer_bk, t.email, t.phone, t.city, t.hashdiff, + CASE WHEN t.current_sk IS NULL + THEN timestamp '1900-01-01' -- первая версия: техническое "начало истории" + ELSE t.eff_ts + END AS valid_from, + timestamp '9999-12-31' AS valid_to, + TRUE, now(), now() + FROM to_upsert t + ON CONFLICT (customer_bk, valid_from) DO NOTHING + RETURNING customer_bk, valid_from ) - INSERT INTO dds.dim_customer ( - customer_bk, email, phone, city, hashdiff, - valid_from, valid_to, - is_current, created_at, updated_at - ) - SELECT - s.customer_bk, s.email, s.phone, s.city, s.hashdiff, - CASE WHEN d.customer_bk IS NULL THEN timestamp '1900-01-01' ELSE s.eff_ts END AS valid_from, - timestamp '9999-12-31', - TRUE, now(), now() - FROM delta s - LEFT JOIN dds.dim_customer d - ON d.customer_bk = s.customer_bk AND d.is_current = TRUE - WHERE d.customer_bk IS NULL -- новый BK - OR d.hashdiff <> s.hashdiff -- изменившийся BK - ON CONFLICT (customer_bk, valid_from) DO NOTHING; + -- закрываем старые версии только для тех BK, по которым реально вставилась новая + UPDATE dds.dim_customer d + SET valid_to = LEAST(d.valid_to, t.eff_ts - interval '1 second'), + is_current = FALSE, + updated_at = now() + FROM to_upsert t + JOIN inserted i + ON i.customer_bk = t.customer_bk + WHERE d.customer_sk = t.current_sk + AND d.is_current = TRUE + AND t.eff_ts >= d.valid_from; COMMIT; -- (факты можно не перезаливать — даты заказов не поменялись) From e659fbb7315190d13c80e7aea7612a26305341a5 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 29 Nov 2025 20:46:03 +0300 Subject: [PATCH 3/5] =?UTF-8?q?=D0=98=D1=81=D0=BF=D1=80=D0=B0=D0=B2=D0=BB?= =?UTF-8?q?=D0=B5=D0=BD=D0=B8=D0=B5=20=D0=BD=D0=B5=D1=82=D0=BE=D1=87=D0=BD?= =?UTF-8?q?=D0=BE=D1=81=D1=82=D0=B5=D0=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../Homework_Customer_Status_DDS_DM.md | 9 +- dwh-modeling/sql/01_ddl_stg-dds.sql | 6 +- dwh-modeling/sql/02_dml_stg-dds.sql | 97 ++++++++++--------- dwh-modeling/sql/04_validation.sql | 3 +- dwh-modeling/sql/06_dml_dm.sql | 3 +- .../08_dml_hw_customer_status_template.sql | 4 +- 6 files changed, 61 insertions(+), 61 deletions(-) diff --git a/dwh-modeling/Homework_Customer_Status_DDS_DM.md b/dwh-modeling/Homework_Customer_Status_DDS_DM.md index dacb976..fb54cdd 100644 --- a/dwh-modeling/Homework_Customer_Status_DDS_DM.md +++ b/dwh-modeling/Homework_Customer_Status_DDS_DM.md @@ -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. +- `load_ts` — момент, когда данные попали в DWH (в таблицах STG/ODS эта колонка будет называться `_load_ts`, но по смыслу это то же самое время загрузки). Файл содержит несколько клиентов и несколько смен статуса по каждому — этого достаточно, чтобы отработать SCD2. @@ -93,7 +93,7 @@ SELECT * FROM stg.customer_status_raw LIMIT 10; - привести: - `customer_id` → `INT`, - `status` → `VARCHAR(20)` (можно оставить как есть), - - `event_ts` и `load_ts` → `TIMESTAMP`; + - `event_ts` и `load_ts` → `TIMESTAMP` (в DWH-таблицах эта колонка будет лежать как `_load_ts`); - аккуратно обработать возможные пустые значения (если бы они были); - заполнить `_load_id` и `_load_ts` в `ods.customer_status`. @@ -129,14 +129,14 @@ ORDER BY customer_id, event_ts; - `customer_bk`, - `status`, - `event_ts` (как «время начала действия статуса»), - - `hashdiff` (например, `md5(status)` или `dds.customer_hash(status)` по аналогии с `dim_customer`). + - `hashdiff` (например, `md5(status)`; можно вынести расчёт в отдельную функцию по аналогии с `dds.customer_hash` для клиентов). 2. Для каждого клиента отсортируйте события по `event_ts` и с помощью `LEAD()` посчитайте: - `valid_from` — текущее `event_ts`, - `valid_to` — следующее `event_ts - 1 second` (или `9999-12-31`, если следующего нет). -3. Вставьте получившиеся строки в `dds.dim_customer_status`: +3. Вставьте получившиеся строки в `dds.dim_customer_status`. У актуальной строки для каждого клиента задайте `is_current = TRUE` (например, там, где `valid_to = '9999-12-31'`), у остальных — `FALSE`: ```sql INSERT INTO dds.dim_customer_status ( @@ -238,4 +238,3 @@ ORDER BY date_actual, status; - при желании — собрать простую витрину в `dm`. Если что‑то не получается — можно разбирать решения по шагам вместе с ментором: от простого `SELECT` из STG до полноценного SCD2 в DDS. - diff --git a/dwh-modeling/sql/01_ddl_stg-dds.sql b/dwh-modeling/sql/01_ddl_stg-dds.sql index 1480378..d313a94 100644 --- a/dwh-modeling/sql/01_ddl_stg-dds.sql +++ b/dwh-modeling/sql/01_ddl_stg-dds.sql @@ -83,7 +83,7 @@ ALTER TABLE ods.products ADD PRIMARY KEY (product_id); -- 4. DDS: интегрированная модель --- dim_date: справочник дат (без первичного ключа — генерируется) +-- dim_date: справочник дат (ключ — суррогатный date_key) CREATE TABLE dds.dim_date ( date_key INT PRIMARY KEY, date_actual DATE NOT NULL, @@ -118,9 +118,9 @@ CREATE TABLE dds.dim_customer ( updated_at TIMESTAMP NOT NULL DEFAULT NOW() ); --- одна версия на момент времени +-- одна версия на момент времени (на одну пару BK+valid_from) ALTER TABLE dds.dim_customer - ADD CONSTRAINT uq_dim_customer_bk_from UNIQUE (customer_bk, valid_from); -- dim_product: измерение " + ADD CONSTRAINT uq_dim_customer_bk_from UNIQUE (customer_bk, valid_from); -- ускорители CREATE INDEX ix_dim_customer_bk_current ON dds.dim_customer (customer_bk) WHERE is_current; diff --git a/dwh-modeling/sql/02_dml_stg-dds.sql b/dwh-modeling/sql/02_dml_stg-dds.sql index 170cd18..8abdcae 100644 --- a/dwh-modeling/sql/02_dml_stg-dds.sql +++ b/dwh-modeling/sql/02_dml_stg-dds.sql @@ -128,56 +128,59 @@ SELECT product_id, name FROM ods.products; -- 5. DDS: dim_customer — первичная загрузка SCD2 (full backfill из STG) --- Здесь мы пересчитываем всю историю клиента из событий в STG. --- В реальном DWH такую полную перезагрузку делают редко; для инкремента см. 03_demo_increment.sql и SCD.md. +-- В ЭТОМ ДЕМО: dim_customer строится напрямую из stg.customers_raw, который играет роль +-- устойчивого event-лога (все события по клиенту в одном месте). +-- Это удобно для учебной первичной загрузки (full backfill), когда мы один раз +-- восстанавливаем всю историю клиента. +-- В РЕАЛЬНОМ DWH: так делают редко. Исторические измерения обычно строят +-- поверх очищенных и нормализованных слоёв (ODS / PSA / Data Vault). +-- Для примера инкрементальной заливки SCD2 по снимку из ODS см. 03_demo_increment.sql и SCD.md. TRUNCATE dds.dim_customer, dds.fact_sales; -BEGIN; - WITH src AS ( - SELECT - s.customer_id::INT AS customer_bk, - NULLIF(trim(s.email), '') AS email, - NULLIF(trim(s.phone), '') AS phone, - NULLIF(trim(s.city), '') AS city, - COALESCE(NULLIF(s.event_ts,'')::timestamp, s._load_ts, - to_timestamp(regexp_replace(s._load_id,'^batch_',''),'YYYYMMDD_HH24MI')) AS eff_ts, - dds.customer_hash(s.email, s.phone, s.city) AS hashdiff - FROM stg.customers_raw s - WHERE s.customer_id ~ '^\d+$' - ), - ordered AS ( - SELECT *, - lag(hashdiff) OVER (PARTITION BY customer_bk ORDER BY eff_ts) AS prev_hash, - row_number() OVER (PARTITION BY customer_bk ORDER BY eff_ts) AS rn - FROM src - ), - changes AS ( - -- только первые состояния и фактические изменения атрибутов - SELECT * - FROM ordered - WHERE prev_hash IS DISTINCT FROM hashdiff OR prev_hash IS NULL - ), - framed AS ( - SELECT - customer_bk, email, phone, city, hashdiff, - CASE WHEN rn = 1 THEN timestamp '1900-01-01' ELSE eff_ts END AS valid_from, -- первая версия: техническое "начало истории" - lead(eff_ts) OVER (PARTITION BY customer_bk ORDER BY eff_ts) AS next_ts - FROM changes - ) - INSERT INTO dds.dim_customer ( - customer_bk, email, phone, city, hashdiff, - valid_from, valid_to, is_current, - created_at, updated_at - ) +WITH src AS ( SELECT - customer_bk, email, phone, city, hashdiff, - valid_from, - COALESCE(next_ts - interval '1 second', timestamp '9999-12-31') AS valid_to, - (next_ts IS NULL) AS is_current, - now(), now() - FROM framed - ORDER BY customer_bk, valid_from; -COMMIT; + s.customer_id::INT AS customer_bk, + NULLIF(trim(s.email), '') AS email, + NULLIF(trim(s.phone), '') AS phone, + NULLIF(trim(s.city), '') AS city, + COALESCE(NULLIF(s.event_ts,'')::timestamp, s._load_ts, + to_timestamp(regexp_replace(s._load_id,'^batch_',''),'YYYYMMDD_HH24MI')) AS eff_ts, + dds.customer_hash(s.email, s.phone, s.city) AS hashdiff + FROM stg.customers_raw s + WHERE s.customer_id ~ '^\d+$' +), +ordered AS ( + SELECT *, + lag(hashdiff) OVER (PARTITION BY customer_bk ORDER BY eff_ts) AS prev_hash, + row_number() OVER (PARTITION BY customer_bk ORDER BY eff_ts) AS rn + FROM src +), +changes AS ( + -- только первые состояния и фактические изменения атрибутов + SELECT * + FROM ordered + WHERE prev_hash IS DISTINCT FROM hashdiff OR prev_hash IS NULL +), +framed AS ( + SELECT + customer_bk, email, phone, city, hashdiff, + CASE WHEN rn = 1 THEN timestamp '1900-01-01' ELSE eff_ts END AS valid_from, -- первая версия: техническое "начало истории" + lead(eff_ts) OVER (PARTITION BY customer_bk ORDER BY eff_ts) AS next_ts + FROM changes +) +INSERT INTO dds.dim_customer ( + customer_bk, email, phone, city, hashdiff, + valid_from, valid_to, is_current, + created_at, updated_at +) +SELECT + customer_bk, email, phone, city, hashdiff, + valid_from, + COALESCE(next_ts - interval '1 second', timestamp '9999-12-31') AS valid_to, + (next_ts IS NULL) AS is_current, + now(), now() +FROM framed +ORDER BY customer_bk, valid_from; -- 6. DDS: fact_sales — загрузка фактов с учётом SCD -- В продакшене — фильтруем по диапазону дат (инкрементально) diff --git a/dwh-modeling/sql/04_validation.sql b/dwh-modeling/sql/04_validation.sql index 3134451..a10cd12 100644 --- a/dwh-modeling/sql/04_validation.sql +++ b/dwh-modeling/sql/04_validation.sql @@ -1,6 +1,6 @@ -- =============================================== -- Проверки качества данных после загрузки STG→ODS→DDS --- Запускается после 02_dml.sql +-- Запускается после 02_dml_stg-dds.sql (и, при необходимости, 03_demo_increment.sql) -- =============================================== -- 1. Проверка: dim_customer не пуста @@ -40,4 +40,3 @@ BEGIN FORMAT('ОШИБКА: у клиента 101 только %s версия, ожидается ≥2 (должна быть история)', version_count); RAISE NOTICE '✅ SCD Type 2: клиент 101 имеет %s версий — история сохранена', version_count; END $$; - diff --git a/dwh-modeling/sql/06_dml_dm.sql b/dwh-modeling/sql/06_dml_dm.sql index bd8acfa..f7b5330 100644 --- a/dwh-modeling/sql/06_dml_dm.sql +++ b/dwh-modeling/sql/06_dml_dm.sql @@ -34,7 +34,6 @@ GROUP BY d.date_actual, p.product_name, -- Здесь — НЕ используем is_current! Нам нужна вся история для расчёта LTV -- Для кого: CRM-менеджер, retention-спец -- Пример использования: «Найти клиентов с LTV > 250 ₽ и email из Москвы для email-рассылки» -«Найти клиентов с LTV > 250 ₽ и email из Москвы для email-рассылки» INSERT INTO dm.mart_customer_360 ( customer_bk, first_order_date, last_order_date, total_orders, total_items, lifetime_value, @@ -44,7 +43,7 @@ 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_orders, -- считаем строки факта (продажи), не бизнес-заказы SUM(f.quantity) AS total_items, SUM(f.amount) AS lifetime_value, -- Берём email и город из самой свежей версии клиента diff --git a/dwh-modeling/sql/08_dml_hw_customer_status_template.sql b/dwh-modeling/sql/08_dml_hw_customer_status_template.sql index 32f5f2a..676527f 100644 --- a/dwh-modeling/sql/08_dml_hw_customer_status_template.sql +++ b/dwh-modeling/sql/08_dml_hw_customer_status_template.sql @@ -8,7 +8,8 @@ -- 1) Загрузите CSV в stg.customer_status_raw (через COPY или \copy в psql). -- См. пример структуры файла в dwh-modeling/data/customer_status_events.csv -- 2) Переложите данные в ods.customer_status с приведением типов. --- customer_id → INT, status → VARCHAR(20), event_ts / load_ts → TIMESTAMP. +-- customer_id → INT, status → VARCHAR(20), event_ts / load_ts → TIMESTAMP +-- (в STG/ODS эта колонка будет жить как _load_ts). -- 3) Постройте из ods.customer_status измерение dds.dim_customer_status в стиле SCD2: -- - одна строка на период действия статуса (valid_from / valid_to); -- - is_current = TRUE только у актуальной строки для клиента; @@ -45,4 +46,3 @@ -- ); -- Идея: на каждую дату взять актуальный статус клиента -- через JOIN dds.dim_customer_status + dds.dim_date. - From eb9215c4f16a7fa6028d45aa0444236d30148d6b Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 29 Nov 2025 21:24:08 +0300 Subject: [PATCH 4/5] =?UTF-8?q?=D0=A3=D1=83=D0=BB=D1=87=D1=88=D0=B5=D0=BD?= =?UTF-8?q?=D0=B8=D1=8F=20=D0=B2=D0=B8=D1=82=D1=80=D0=B8=D0=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- dwh-modeling/sql/06_dml_dm.sql | 25 +++++++++---------------- 1 file changed, 9 insertions(+), 16 deletions(-) diff --git a/dwh-modeling/sql/06_dml_dm.sql b/dwh-modeling/sql/06_dml_dm.sql index f7b5330..f7b02c2 100644 --- a/dwh-modeling/sql/06_dml_dm.sql +++ b/dwh-modeling/sql/06_dml_dm.sql @@ -6,17 +6,15 @@ -- 1. Очистка (full refresh — для простоты; в продакшене — incremental) TRUNCATE dm.mart_daily_sales, dm.mart_customer_360; --- 2. mart_daily_sales: агрегация по дням --- Используем актуальные версии измерений (is_current = true) --- Для кого: Маркетолог, продакт-аналитик --- Пример использования: «Как менялись продажи Phone в сегменте Premium по дням?» +-- 2. mart_daily_sales: продажи по датам и товарам +-- Одна строка = дата × товар × простой сегмент заказа INSERT INTO dm.mart_daily_sales ( date_actual, product_name, customer_segment, total_qty, total_revenue ) SELECT d.date_actual, p.product_name, - -- Простая сегментация по выручке за заказ + -- Делим заказы на Premium / Basic по сумме CASE WHEN f.amount >= 200 THEN 'Premium' ELSE 'Basic' @@ -26,14 +24,13 @@ SELECT FROM dds.fact_sales f JOIN dds.dim_date d ON f.date_key = d.date_key JOIN dds.dim_product p ON f.product_sk = p.product_sk -JOIN dds.dim_customer c ON f.customer_sk = c.customer_sk AND c.is_current +JOIN dds.dim_customer c ON f.customer_sk = c.customer_sk -- факт уже ссылается на нужную версию клиента GROUP BY d.date_actual, p.product_name, CASE WHEN f.amount >= 200 THEN 'Premium' ELSE 'Basic' END; --- 3. mart_customer_360: lifetime-портрет --- Здесь — НЕ используем is_current! Нам нужна вся история для расчёта LTV --- Для кого: CRM-менеджер, retention-спец --- Пример использования: «Найти клиентов с LTV > 250 ₽ и email из Москвы для email-рассылки» +-- 3. mart_customer_360: 360‑портрет клиента +-- Одна строка = один клиент +-- Считаем суммы по всей истории его покупок INSERT INTO dm.mart_customer_360 ( customer_bk, first_order_date, last_order_date, total_orders, total_items, lifetime_value, @@ -46,7 +43,7 @@ SELECT COUNT(DISTINCT f.sale_id) AS total_orders, -- считаем строки факта (продажи), не бизнес-заказы SUM(f.quantity) AS total_items, SUM(f.amount) AS lifetime_value, - -- Берём email и город из самой свежей версии клиента + -- Берём самый свежий email и город клиента (SELECT email FROM dds.dim_customer c2 WHERE c2.customer_bk = c.customer_bk ORDER BY c2.valid_from DESC @@ -58,10 +55,6 @@ SELECT FROM dds.fact_sales f JOIN dds.dim_date d ON f.date_key = d.date_key JOIN dds.dim_customer c ON f.customer_sk = c.customer_sk --- ⚠️ Важно: здесь НЕ фильтруем по is_current — нужна вся история! +-- Здесь НЕ фильтруем по is_current: нужна вся история фактов GROUP BY c.customer_bk; --- 4. Дополнительно: материализованное представление (альтернатива таблице) --- В PG 12+ можно использовать MATERIALIZED VIEW — обновляется по команде REFRESH --- CREATE MATERIALIZED VIEW dm.mv_monthly_sales AS --- SELECT ... (аналогично) From ba9cdff42167bb3d99eabee58f27af833770f85e Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 29 Nov 2025 21:37:03 +0300 Subject: [PATCH 5/5] =?UTF-8?q?=D0=9F=D1=80=D0=B8=D1=87=D0=B5=D1=81=D1=8B?= =?UTF-8?q?=D0=B2=D0=B0=D0=BD=D0=B8=D0=B5=20=D0=BA=D0=BE=D0=BC=D0=BC=D0=B5?= =?UTF-8?q?=D0=BD=D1=82=D0=B0=D1=80=D0=B8=D0=B5=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- dwh-modeling/sql/04_validation.sql | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/dwh-modeling/sql/04_validation.sql b/dwh-modeling/sql/04_validation.sql index a10cd12..5a044d6 100644 --- a/dwh-modeling/sql/04_validation.sql +++ b/dwh-modeling/sql/04_validation.sql @@ -12,7 +12,8 @@ BEGIN (SELECT COUNT(*) FROM dds.dim_customer); END $$; --- 3. Проверка: у каждого факта есть валидная дата (date_key существует) +-- 2. Проверка: каждая строка из ODS попала в DDS-факт +-- Сравниваем количество строк в ods.order_items и dds.fact_sales DO $$ DECLARE expected_count bigint; @@ -28,7 +29,7 @@ BEGIN RAISE NOTICE '✅ fact_sales: количество строк совпадает с ods.order_items (%)', actual_count; END $$; --- 4. Проверка SCD Type 2: у клиента 101 должно быть ≥2 версий (из-за смены email) +-- 3. Проверка SCD Type 2: у клиента 101 должно быть ≥2 версий (из-за смены email) DO $$ DECLARE version_count INT; BEGIN @@ -38,5 +39,5 @@ BEGIN -- ASSERT version_count >= 2, FORMAT('ОШИБКА: у клиента 101 только %s версия, ожидается ≥2 (должна быть история)', version_count); - RAISE NOTICE '✅ SCD Type 2: клиент 101 имеет %s версий — история сохранена', version_count; + RAISE NOTICE '✅ SCD Type 2: клиент 101 имеет % версий — история сохранена', version_count; END $$;