refactor(modeling): улучшены учебные материалы DWH по результатам ревью
- Зачем: - убрать путаницу, дублирование и неточности в демо-скриптах и домашке - Что: - 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 <noreply@anthropic.com>
This commit is contained in:
@@ -30,7 +30,7 @@
|
|||||||
Структура файла:
|
Структура файла:
|
||||||
|
|
||||||
```text
|
```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
|
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`);
|
- `status` — статус клиента в CRM (`new`, `active`, `vip`, `churned`);
|
||||||
- `event_ts` — момент, когда статус сменился в CRM;
|
- `event_ts` — момент, когда статус сменился в CRM;
|
||||||
- `_load_id` — идентификатор батча загрузки;
|
- `_load_id` — идентификатор батча загрузки;
|
||||||
- `load_ts` — момент, когда данные попали в DWH (в таблицах STG/ODS эта колонка будет называться `_load_ts`, но по смыслу это то же самое время загрузки).
|
- `_load_ts` — момент, когда данные попали в DWH.
|
||||||
|
|
||||||
Файл содержит несколько клиентов и несколько смен статуса по каждому — этого достаточно, чтобы отработать SCD2.
|
Файл содержит несколько клиентов и несколько смен статуса по каждому — этого достаточно, чтобы отработать 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');
|
('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
|
#### Вариант B: загрузить CSV
|
||||||
|
|
||||||
Можно загрузить файл `dwh-modeling/data/customer_status_events.csv` в таблицу `stg.customer_status_raw`:
|
Можно загрузить файл `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` на хост.
|
- **Через `psql` в контейнере (`./psql_sh`)**: без установки `psql` на хост.
|
||||||
|
|
||||||
Способ: передайте CSV в `psql` через STDIN и выполните `\copy ... FROM STDIN`:
|
Способ: передайте CSV в `psql` через STDIN и выполните `\copy ... FROM STDIN`:
|
||||||
@@ -126,7 +126,7 @@ SELECT * FROM stg.customer_status_raw LIMIT 10;
|
|||||||
- привести:
|
- привести:
|
||||||
- `customer_id` → `INT`,
|
- `customer_id` → `INT`,
|
||||||
- `status` → `VARCHAR(20)` (можно оставить как есть),
|
- `status` → `VARCHAR(20)` (можно оставить как есть),
|
||||||
- `event_ts` и `load_ts` → `TIMESTAMP` (в DWH-таблицах эта колонка будет лежать как `_load_ts`);
|
- `event_ts` и `_load_ts` → `TIMESTAMP`;
|
||||||
- аккуратно обработать возможные пустые значения (если бы они были);
|
- аккуратно обработать возможные пустые значения (если бы они были);
|
||||||
- заполнить `_load_id` и `_load_ts` в `ods.customer_status`.
|
- заполнить `_load_id` и `_load_ts` в `ods.customer_status`.
|
||||||
|
|
||||||
@@ -292,3 +292,11 @@ ORDER BY date_actual, status;
|
|||||||
- при желании — собрать простую витрину в `dm`.
|
- при желании — собрать простую витрину в `dm`.
|
||||||
|
|
||||||
Если что‑то не получается — можно разбирать решения по шагам вместе с ментором: от простого `SELECT` из STG до полноценного SCD2 в DDS.
|
Если что‑то не получается — можно разбирать решения по шагам вместе с ментором: от простого `SELECT` из STG до полноценного SCD2 в DDS.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 8. Эталонное решение
|
||||||
|
|
||||||
|
Когда выполните домашку и захотите сверить результат — готовое решение лежит в файле [`09_dml_hw_customer_status_solution.sql`](sql/09_dml_hw_customer_status_solution.sql).
|
||||||
|
|
||||||
|
Постарайтесь не подглядывать до того, как напишете свой вариант — основная ценность задания именно в самостоятельном разборе.
|
||||||
|
|||||||
@@ -720,7 +720,7 @@ SELECT 'OK' WHERE EXISTS (
|
|||||||
|
|
||||||
[`customers.csv`](data/customers.csv):
|
[`customers.csv`](data/customers.csv):
|
||||||
```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
|
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
|
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
|
101,b@ex.com,700,Москва,2024-05-16,batch_20240516_0800,2024-05-16 08:00
|
||||||
|
|||||||
@@ -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,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,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
|
101,vip,2024-05-10 11:00:00,batch_20240510_1200,2024-05-10 12:00:00
|
||||||
|
|||||||
|
@@ -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
|
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
|
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
|
103,active,2024-03-20 12:00:00,batch_20240320_1300,2024-03-20 13:00:00
|
||||||
|
|||||||
|
@@ -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
|
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
|
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
|
101,b@ex.com,700,Москва,2024-05-16,batch_20240516_0800,2024-05-16 08:00
|
||||||
|
|||||||
|
@@ -2,6 +2,13 @@
|
|||||||
-- DML-скрипт: загрузка и трансформация данных
|
-- 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)
|
-- 1. STG: имитация загрузки из источников (в реальности — COPY или INSERT из Kafka/NiFi)
|
||||||
-- ⚠️ В продакшене STG часто очищается перед загрузкой (TRUNCATE), либо используется партицирование по дате
|
-- ⚠️ В продакшене STG часто очищается перед загрузкой (TRUNCATE), либо используется партицирование по дате
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ CREATE TABLE dm.mart_customer_360 (
|
|||||||
customer_bk INT NOT NULL,
|
customer_bk INT NOT NULL,
|
||||||
first_order_date DATE,
|
first_order_date DATE,
|
||||||
last_order_date DATE,
|
last_order_date DATE,
|
||||||
total_orders INT NOT NULL,
|
total_line_items INT NOT NULL,
|
||||||
total_items INT NOT NULL,
|
total_items INT NOT NULL,
|
||||||
lifetime_value NUMERIC(18,2) NOT NULL,
|
lifetime_value NUMERIC(18,2) NOT NULL,
|
||||||
last_email VARCHAR(100),
|
last_email VARCHAR(100),
|
||||||
|
|||||||
@@ -33,14 +33,14 @@ GROUP BY d.date_actual, p.product_name,
|
|||||||
-- Считаем суммы по всей истории его покупок
|
-- Считаем суммы по всей истории его покупок
|
||||||
INSERT INTO dm.mart_customer_360 (
|
INSERT INTO dm.mart_customer_360 (
|
||||||
customer_bk, first_order_date, last_order_date,
|
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
|
last_email, last_city
|
||||||
)
|
)
|
||||||
SELECT
|
SELECT
|
||||||
c.customer_bk,
|
c.customer_bk,
|
||||||
MIN(d.date_actual) AS first_order_date,
|
MIN(d.date_actual) AS first_order_date,
|
||||||
MAX(d.date_actual) AS last_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.quantity) AS total_items,
|
||||||
SUM(f.amount) AS lifetime_value,
|
SUM(f.amount) AS lifetime_value,
|
||||||
-- Берём самый свежий email и город клиента
|
-- Берём самый свежий email и город клиента
|
||||||
|
|||||||
@@ -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;
|
||||||
Reference in New Issue
Block a user