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. -