diff --git a/.gitignore b/.gitignore index e58e5a2..4b89c7d 100644 --- a/.gitignore +++ b/.gitignore @@ -1,2 +1,3 @@ .env .idea +.internal/ diff --git a/dwh-modeling/DataVault.md b/dwh-modeling/DataVault.md index b9c0abc..f40ceb4 100644 --- a/dwh-modeling/DataVault.md +++ b/dwh-modeling/DataVault.md @@ -293,7 +293,7 @@ PIT (Point-in-Time) решает очень конкретную боль: > «Покажи, как объект выглядел **на дату X**, но так, чтобы запрос был простым». Если у нас есть несколько сателлитов с историей (например, `sat_customer_info`, `sat_customer_segment`, `sat_customer_risk`), то без PIT любой запрос превращается в пачку условий -`BETWEEN valid_from AND valid_to` (в effectivity-сателлитах или любых таблицах с периодами действия) или оконных функций. +`as_of_date >= valid_from AND (valid_to IS NULL OR as_of_date < valid_to)` (в effectivity-сателлитах или любых таблицах с периодами действия) или оконных функций. **Идея PIT:** diff --git a/dwh-modeling/Homework_Customer_Status_DDS_DM.md b/dwh-modeling/Homework_Customer_Status_DDS_DM.md index fb54cdd..15926e3 100644 --- a/dwh-modeling/Homework_Customer_Status_DDS_DM.md +++ b/dwh-modeling/Homework_Customer_Status_DDS_DM.md @@ -31,7 +31,7 @@ ```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 +101,new,2024-01-01 09:00:00,batch_20240101_1000,2024-01-01 10:00:00 ... ``` @@ -117,7 +117,6 @@ ORDER BY customer_id, event_ts; - `status` — статус клиента; - `hashdiff` — хэш от атрибутов (здесь достаточно самого `status`); - `valid_from` / `valid_to` — период, когда статус был актуален; -- `is_current` — флаг актуальной строки; - `created_at` / `updated_at` — технические поля. ### 4.1. Начальная загрузка SCD2 @@ -131,18 +130,18 @@ ORDER BY customer_id, event_ts; - `event_ts` (как «время начала действия статуса»), - `hashdiff` (например, `md5(status)`; можно вынести расчёт в отдельную функцию по аналогии с `dds.customer_hash` для клиентов). -2. Для каждого клиента отсортируйте события по `event_ts` и с помощью `LEAD()` посчитайте: +2. Для каждого клиента отсортируйте события по `event_ts` и с помощью `LEAD()` посчитайте (в учебном варианте считаем, что обновление DWH идёт раз в день, поэтому используем `DATE`): - - `valid_from` — текущее `event_ts`, - - `valid_to` — следующее `event_ts - 1 second` (или `9999-12-31`, если следующего нет). + - `valid_from` — `event_ts::date`, + - `valid_to` — следующий `event_ts::date` (а у последней версии `valid_to = NULL`). -3. Вставьте получившиеся строки в `dds.dim_customer_status`. У актуальной строки для каждого клиента задайте `is_current = TRUE` (например, там, где `valid_to = '9999-12-31'`), у остальных — `FALSE`: +3. Вставьте получившиеся строки в `dds.dim_customer_status`. Актуальная строка для клиента — та, где `valid_to IS NULL`. ```sql INSERT INTO dds.dim_customer_status ( customer_bk, status, hashdiff, valid_from, valid_to, - is_current, created_at, updated_at + created_at, updated_at ) SELECT ... @@ -159,7 +158,7 @@ ORDER BY customer_bk, valid_from; Ожидаемое поведение: - у клиента 101 несколько строк с разными статусами и непересекающимися периодами; -- `is_current = TRUE` только у самой свежей строки для каждого клиента. +- `valid_to IS NULL` только у самой свежей строки для каждого клиента. ### 4.2. Проверка себя @@ -182,8 +181,8 @@ ORDER BY customer_bk, valid_from; - ориентируйтесь на пример из `03_demo_increment.sql` для `dds.dim_customer`; - важно: - - корректно «закрыть» старую актуальную строку (`valid_to`, `is_current = FALSE`); - - вставить новую строку с `is_current = TRUE`. + - корректно «закрыть» старую актуальную строку (заполнить `valid_to` датой начала новой версии); + - вставить новую строку с `valid_to = NULL`. Эта часть особенно полезна, если вы хотите почувствовать, как SCD2 живёт в реальном DWH. diff --git a/dwh-modeling/README.md b/dwh-modeling/README.md index 5d9b5b6..d29d747 100644 --- a/dwh-modeling/README.md +++ b/dwh-modeling/README.md @@ -234,13 +234,12 @@ erDiagram dim_customer { bigint customer_sk PK "суррогатный ключ" - varchar customer_bk "бизнес-ключ, напр. '101'" + int customer_bk "бизнес-ключ, напр. 101" varchar customer_name varchar email varchar city date valid_from "SCD2: с какой даты запись актуальна" - date valid_to "SCD2: по какую дату актуальна" - boolean is_current "SCD2: текущая версия?" + date valid_to "SCD2: по какую дату актуальна (NULL = сейчас)" } dim_product { @@ -269,15 +268,16 @@ erDiagram В `dim_customer` это будет **три строки**: -| customer_sk | customer_bk | email | city | valid_from | valid_to | is_current | -|-------------|-------------|-------|------|------------|----------|------------| -| 1001 | 101 | a@ex.com | Москва | 2023-01-01 | 2023-05-15 | false | -| 1002 | 101 | b@ex.com | Москва | 2023-05-16 | 2023-09-30 | false | -| 1003 | 101 | b@ex.com | СПб | 2023-10-01 | 9999-12-31 | true | +| customer_sk | customer_bk | email | city | valid_from | valid_to | +|-------------|-------------|-------|------|------------|----------| +| 1001 | 101 | a@ex.com | Москва | 2023-01-01 | 2023-05-16 | +| 1002 | 101 | b@ex.com | Москва | 2023-05-16 | 2023-10-01 | +| 1003 | 101 | b@ex.com | СПб | 2023-10-01 | NULL | Когда мы считаем продажи за **12 января** — джойним `fact_sales` к той строке `dim_customer`, где: ```sql -fact_sales.order_date BETWEEN dim_customer.valid_from AND dim_customer.valid_to +fact_sales.order_date >= dim_customer.valid_from +AND (dim_customer.valid_to IS NULL OR fact_sales.order_date < dim_customer.valid_to) ``` и получаем актуальный на тот день email и город. @@ -554,8 +554,8 @@ 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 f.order_date BETWEEN c.valid_from AND c.valid_to -- SCD! -WHERE c.is_current = true -- или не фильтровать — тогда будет история + AND f.order_date >= c.valid_from + AND (c.valid_to IS NULL OR f.order_date < c.valid_to) -- SCD! GROUP BY d.date_actual, p.product_name, c.customer_segment; ``` @@ -775,10 +775,11 @@ SELECT 'OK' WHERE EXISTS ( [`customers.csv`](data/customers.csv): ```csv -customer_id,email,phone,city -101,a@ex.com,700,Москва -101,b@ex.com,700,Москва -102,c@ex.com,701,СПб +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 +101,b@ex.com,700,Санкт-Петербург,2024-10-01,batch_20241001_0800,2024-10-01 08:00 ``` [`orders.csv`](data/orders.csv): @@ -807,7 +808,8 @@ product_id,name ```csv product_id,valid_from,valid_to,price 9001,2023-12-01,2024-01-31,100 -9001,2024-02-01,2999-12-31,110 +9001,2024-02-01,,110 +9002,2023-01-01,,50 ``` > 📂 Все SQL-скрипты для построения хранилища находятся в папке [`sql/`](sql/). @@ -824,13 +826,12 @@ product_id,valid_from,valid_to,price -- DDS: измерение клиента (SCD Type 2) CREATE TABLE dds.dim_customer ( customer_sk BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY, - customer_bk VARCHAR(50) NOT NULL, -- напр. '101' + customer_bk INT NOT NULL, -- напр. 101 email VARCHAR(100), phone VARCHAR(20), city VARCHAR(50), valid_from DATE NOT NULL, - valid_to DATE DEFAULT '9999-12-31', - is_current BOOLEAN NOT NULL DEFAULT TRUE + valid_to DATE ); -- DDS: факт продаж (гранулярность: строка заказа) diff --git a/dwh-modeling/SCD.md b/dwh-modeling/SCD.md index 2ab2bf8..ad8356d 100644 --- a/dwh-modeling/SCD.md +++ b/dwh-modeling/SCD.md @@ -107,7 +107,6 @@ WHERE customer_id = 1; - **`customer_id`** — бизнес-идентификатор (например, из CRM). Он **не меняется** и связывает все версии одного клиента. - **`valid_from`** — дата, **с которой** эта версия стала актуальной. - **`valid_to`** — дата, **по которую** эта версия была актуальной. Если `NULL` — значит, версия **актуальна сейчас**. -- **`is_current`** — флаг (`true`/`false`), который быстро показывает, какая строка — последняя. Удобен для запросов. Пример структуры: @@ -118,33 +117,31 @@ CREATE TABLE customers_scd2 ( name TEXT NOT NULL, category TEXT NOT NULL, valid_from DATE NOT NULL, -- с какой даты действует - valid_to DATE, -- по какую дату действовала (NULL = сейчас) - is_current BOOLEAN NOT NULL -- true = актуальная версия + valid_to DATE -- по какую дату действовала (NULL = сейчас) ); ``` **Шаг 1.** Добавляем начальную запись (клиент зарегистрировался 1 января 2024): ```sql -INSERT INTO customers_scd2 (customer_id, name, category, valid_from, valid_to, is_current) -VALUES (1, 'Иван Петров', 'Regular', '2024-01-01', NULL, true); +INSERT INTO customers_scd2 (customer_id, name, category, valid_from, valid_to) +VALUES (1, 'Иван Петров', 'Regular', DATE '2024-01-01', NULL); ``` **Шаг 2.** 15 апреля 2025 клиент становится VIP. Мы делаем **два действия**: -1. **Закрываем старую запись**: указываем, что она была актуальна **до 14 апреля**. +1. **Закрываем старую запись**: указываем дату начала новой версии (интервал `[valid_from, valid_to)`). 2. **Добавляем новую запись**: она начинает действовать **с 15 апреля** и пока актуальна. ```sql -- 1. Завершаем предыдущую версию UPDATE customers_scd2 -SET valid_to = '2025-04-14', - is_current = false -WHERE customer_id = 1 AND is_current = true; +SET valid_to = DATE '2025-04-15' +WHERE customer_id = 1 AND valid_to IS NULL; -- 2. Вставляем новую версию -INSERT INTO customers_scd2 (customer_id, name, category, valid_from, valid_to, is_current) -VALUES (1, 'Иван Петров', 'VIP', '2025-04-15', NULL, true); +INSERT INTO customers_scd2 (customer_id, name, category, valid_from, valid_to) +VALUES (1, 'Иван Петров', 'VIP', DATE '2025-04-15', NULL); ``` Теперь в таблице две строки для одного клиента. И мы можем спросить: @@ -155,36 +152,14 @@ VALUES (1, 'Иван Петров', 'VIP', '2025-04-15', NULL, true); SELECT category FROM customers_scd2 WHERE customer_id = 1 - AND '2025-04-10' >= valid_from - AND '2025-04-10' <= COALESCE(valid_to, '9999-12-31'); + AND DATE '2025-04-10' >= valid_from + AND (valid_to IS NULL OR DATE '2025-04-10' < valid_to); ``` Результат: `'Regular'` — правильно! -> 💡 Почему `COALESCE(valid_to, '9999-12-31')`? -> Потому что у актуальной записи `valid_to IS NULL`. Чтобы условие работало, мы временно заменяем `NULL` на очень далёкую дату. - ---- - -### Type 2: можно ли обойтись без `is_current`? - -Поле `is_current` **не обязательно**. Оно удобно для быстрого поиска актуальной версии (например, `WHERE is_current = true`), но **всю ту же информацию несёт пара `valid_from` / `valid_to`**. - -Многие реализации SCD Type 2 обходятся без этого флага — особенно в системах, где важна нормализация или минимизация избыточности. В таких случаях актуальность определяют по условию: - -```sql -WHERE valid_to IS NULL -``` - -Или, если используется «закрытый» интервал (например, `valid_to = '2025-04-14'` для прошлой версии, а у новой — `valid_from = '2025-04-15'`, `valid_to = '9999-12-31'`), то: - -```sql -WHERE CURRENT_DATE BETWEEN valid_from AND valid_to -``` - -Так что `is_current` — это **опциональное упрощение**, а не часть стандарта. - ---- +> 💡 Почему не `BETWEEN valid_from AND valid_to`? +> Потому что у актуальной записи `valid_to IS NULL`. В SQL сравнения с `NULL` не дают `TRUE`, поэтому для “текущей” версии обычно пишут `valid_to IS NULL OR ...`. ### Type 2 через логику, похожую на UPSERT @@ -206,21 +181,21 @@ WHERE CURRENT_DATE BETWEEN valid_from AND valid_to -- Шаг 1: вставляем новую версию, только если есть изменения WITH last_version AS ( SELECT * FROM customers_scd2 - WHERE customer_id = 1 AND is_current = true + WHERE customer_id = 1 AND valid_to IS NULL ) -INSERT INTO customers_scd2 (customer_id, name, category, valid_from, valid_to, is_current) -SELECT 1, 'Иван Петров', 'VIP', '2025-04-15', NULL, true +INSERT INTO customers_scd2 (customer_id, name, category, valid_from, valid_to) +SELECT 1, 'Иван Петров', 'VIP', DATE '2025-04-15', NULL WHERE EXISTS ( SELECT 1 FROM last_version WHERE category != 'VIP' ); -- Шаг 2: если вставка произошла — закрываем старую запись UPDATE customers_scd2 -SET valid_to = '2025-04-14', is_current = false -WHERE customer_id = 1 AND is_current = true +SET valid_to = DATE '2025-04-15' +WHERE customer_id = 1 AND valid_to IS NULL AND EXISTS ( SELECT 1 FROM customers_scd2 - WHERE customer_id = 1 AND category = 'VIP' AND valid_from = '2025-04-15' + WHERE customer_id = 1 AND category = 'VIP' AND valid_from = DATE '2025-04-15' ); ``` @@ -408,8 +383,8 @@ WHERE rn = 1; Это и есть «путешествие во времени» (time travel) в системах без встроенной поддержки этой функции. -> **Почему не использовать `is_current`?** -> В append-only системах флаг `is_current` становится устаревшим сразу после новой вставки. Его очень сложно поддерживать без `UPDATE`, поэтому в таких архитектурах его чаще **не используют**, полагаясь полностью на даты и оконные функции. Это делает модель данных более чистой и идемпотентной. +> **Почему не использовать флаг «текущая версия»?** +> В append-only системах любой флаг актуальности становится устаревшим сразу после новой вставки. Без `UPDATE` его сложно поддерживать, поэтому в таких архитектурах чаще полагаются на даты и оконные функции — так модель остаётся идемпотентной. Таким образом, SCD Type 2 не только совместим с append-only системами, но и является для них **естественным выбором**, так как его логика основана исключительно на добавлении данных, а не на их изменении. @@ -431,8 +406,8 @@ WHERE rn = 1; ## 6. Подводные камни и советы - **Не используйте `customer_id` как первичный ключ в Type 2**. Он повторяется! Вместо этого — `customer_key` (surrogate key). -- Всегда задавайте `valid_to` как `NULL` для актуальной записи, если это допустимо в вашей СУБД — это упрощает запросы. -- Используйте `COALESCE(valid_to, '9999-12-31')` в условиях, чтобы избежать `NULL`-проблем. +- Всегда задавайте `valid_to` как `NULL` для актуальной записи, если это допустимо в вашей СУБД — это упрощает модель. +- Для условий “актуально на дату” используйте паттерн: `d >= valid_from AND (valid_to IS NULL OR d < valid_to)`. - Type 2 увеличивает объём данных — но для аналитики это нормально. - В связке с фактами: в таблице фактов храните **`customer_key`**, а не `customer_id` — иначе не получится соединить с нужной версией. diff --git a/dwh-modeling/data/customer_status_events.csv b/dwh-modeling/data/customer_status_events.csv index 493886a..40a586e 100644 --- a/dwh-modeling/data/customer_status_events.csv +++ b/dwh-modeling/data/customer_status_events.csv @@ -1,10 +1,9 @@ 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 - +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 +101,churned,2024-09-01 12:15:00,batch_20240901_1300,2024-09-01 13:00:00 +102,new,2024-03-05 14:00:00,batch_20240305_1500,2024-03-05 15:00:00 +102,active,2024-04-01 09:45:00,batch_20240401_1000,2024-04-01 10:00:00 +102,churned,2024-04-20 16:20:00,batch_20240420_1700,2024-04-20 17:00:00 +103,new,2024-03-10 10:10:00,batch_20240310_1100,2024-03-10 11:00:00 diff --git a/dwh-modeling/data/customers.csv b/dwh-modeling/data/customers.csv index 843393c..229fb6f 100644 --- a/dwh-modeling/data/customers.csv +++ b/dwh-modeling/data/customers.csv @@ -1,5 +1,5 @@ customer_id,email,phone,city,event_ts,_load_id,load_ts -101,a@ex.com,700,Москва,NULL,batch_20250405_0800,2025-04-05 08:00 -102,c@ex.com,701,СПб,NULL,batch_20250405_0800,2025-04-05 08:00 -101,b@ex.com,700,Москва,NULL,batch_20250405_1200,2025-04-05 12:00 -101,b@ex.com,700,Санкт-Петербург,NULL,batch_20250405_1800,2025-04-05 18:00 \ No newline at end of file +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 +101,b@ex.com,700,Санкт-Петербург,2024-10-01,batch_20241001_0800,2024-10-01 08:00 diff --git a/dwh-modeling/data/prices.csv b/dwh-modeling/data/prices.csv index 6f1c6ac..c51561e 100644 --- a/dwh-modeling/data/prices.csv +++ b/dwh-modeling/data/prices.csv @@ -1,4 +1,4 @@ product_id,valid_from,valid_to,price 9001,2023-12-01,2024-01-31,100.00 -9001,2024-02-01,2999-12-31,110.00 -9002,2023-01-01,2999-12-31,50.00 \ No newline at end of file +9001,2024-02-01,,110.00 +9002,2023-01-01,,50.00 diff --git a/dwh-modeling/sql/01_ddl_stg-dds.sql b/dwh-modeling/sql/01_ddl_stg-dds.sql index d313a94..a73efd5 100644 --- a/dwh-modeling/sql/01_ddl_stg-dds.sql +++ b/dwh-modeling/sql/01_ddl_stg-dds.sql @@ -111,9 +111,8 @@ CREATE TABLE dds.dim_customer ( phone TEXT, city TEXT, hashdiff TEXT NOT NULL, -- md5 по нормализованным атрибутам - valid_from TIMESTAMP NOT NULL, - valid_to TIMESTAMP NOT NULL, - is_current BOOLEAN NOT NULL DEFAULT TRUE, + valid_from DATE NOT NULL, + valid_to DATE, created_at TIMESTAMP NOT NULL DEFAULT NOW(), updated_at TIMESTAMP NOT NULL DEFAULT NOW() ); @@ -123,7 +122,7 @@ ALTER TABLE dds.dim_customer 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; +CREATE INDEX ix_dim_customer_bk_current ON dds.dim_customer (customer_bk) WHERE valid_to IS NULL; CREATE INDEX ix_dim_customer_bk_from_to ON dds.dim_customer (customer_bk, valid_from, valid_to); -- fact_sales: факт "Продажи" diff --git a/dwh-modeling/sql/02_dml_stg-dds.sql b/dwh-modeling/sql/02_dml_stg-dds.sql index 8abdcae..5e26699 100644 --- a/dwh-modeling/sql/02_dml_stg-dds.sql +++ b/dwh-modeling/sql/02_dml_stg-dds.sql @@ -13,10 +13,10 @@ DELETE FROM stg.products_raw; -- STG (пример вставки с метками времени) INSERT INTO stg.customers_raw (_load_id, _load_ts, event_ts, customer_id, email, phone, city) VALUES -('batch_20250405_0800', '2025-04-05 08:00', NULL, '101','a@ex.com','700','Москва'), -('batch_20250405_0800', '2025-04-05 08:00', NULL, '102','c@ex.com','701','СПб'), -('batch_20250405_1200', '2025-04-05 12:00', NULL, '101','b@ex.com','700','Москва'), -('batch_20250405_1800', '2025-04-05 18:00', NULL, '101','b@ex.com','700','Санкт-Петербург'); +('batch_20240101_0800', '2024-01-01 08:00', '2024-01-01', '101','a@ex.com','700','Москва'), +('batch_20240101_0800', '2024-01-01 08:00', '2024-01-01', '102','c@ex.com','701','СПб'), +('batch_20240516_0800', '2024-05-16 08:00', '2024-05-16', '101','b@ex.com','700','Москва'), +('batch_20241001_0800', '2024-10-01 08:00', '2024-10-01', '101','b@ex.com','700','Санкт-Петербург'); INSERT INTO stg.orders_raw (order_id, order_date, customer_id) VALUES @@ -45,7 +45,7 @@ SELECT FROM stg.orders_raw WHERE order_date IS NOT NULL AND customer_id ~ '^\d+$'; --- берём по BK самую позднюю запись (event_ts > _load_ts > _load_id) +-- берём по BK самую позднюю запись (по дате события, иначе по дате загрузки) WITH src AS ( SELECT s.customer_id::INT AS customer_id, @@ -55,21 +55,20 @@ WITH src AS ( NULLIF(s.event_ts, '')::timestamp AS event_ts, s._load_id, s._load_ts, - COALESCE(NULLIF(s.event_ts, '')::timestamp, s._load_ts, - to_timestamp(regexp_replace(s._load_id,'^batch_',''),'YYYYMMDD_HH24MI')) - AS eff_ts + COALESCE(NULLIF(s.event_ts, '')::date, s._load_ts::date) AS eff_date FROM stg.customers_raw s WHERE s.customer_id ~ '^\d+$' ), -last_per_bk AS ( - SELECT DISTINCT ON (customer_id) - customer_id, email, phone, city, event_ts, _load_id, _load_ts +ranked AS ( + SELECT + customer_id, email, phone, city, event_ts, _load_id, _load_ts, + row_number() OVER (PARTITION BY customer_id ORDER BY eff_date DESC, _load_ts DESC) AS rn FROM src - ORDER BY customer_id, eff_ts DESC, _load_ts DESC ) INSERT INTO ods.customers (customer_id, email, phone, city, event_ts, _load_id, _load_ts) SELECT customer_id, email, phone, city, event_ts, _load_id, _load_ts -FROM last_per_bk; +FROM ranked +WHERE rn = 1; INSERT INTO ods.order_items (order_item_id, order_id, product_id, qty, price_at_sale) SELECT @@ -143,16 +142,14 @@ WITH src AS ( 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, + COALESCE(NULLIF(s.event_ts, '')::date, s._load_ts::date) AS eff_date, 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 + lag(hashdiff) OVER (PARTITION BY customer_bk ORDER BY eff_date) AS prev_hash FROM src ), changes AS ( @@ -164,20 +161,19 @@ changes AS ( 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 + eff_date AS valid_from, + lead(eff_date) OVER (PARTITION BY customer_bk ORDER BY eff_date) AS valid_to FROM changes ) INSERT INTO dds.dim_customer ( customer_bk, email, phone, city, hashdiff, - valid_from, valid_to, is_current, + valid_from, valid_to, 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, + valid_to, now(), now() FROM framed ORDER BY customer_bk, valid_from; @@ -200,10 +196,11 @@ JOIN ods.products p ON oi.product_id = p.product_id JOIN dds.dim_product dp ON p.product_id = dp.product_bk JOIN dds.dim_customer dc ON o.customer_id = dc.customer_bk - AND o.order_date::timestamp BETWEEN dc.valid_from AND dc.valid_to; + AND o.order_date >= dc.valid_from + AND (dc.valid_to IS NULL OR o.order_date < dc.valid_to); -- 7. Проверка — вывод итогов (не часть ETL, но полезно для отладки) -- В реальном пайплайне такие SELECT выносятся в отдельные скрипты или дашборды - SELECT 'dim_customer count = ' || COUNT(*) FROM dds.dim_customer WHERE is_current ; + SELECT 'dim_customer current = ' || COUNT(*) FROM dds.dim_customer WHERE valid_to IS NULL; SELECT 'fact_sales count = ' || COUNT(*) FROM dds.fact_sales; diff --git a/dwh-modeling/sql/03_demo_increment.sql b/dwh-modeling/sql/03_demo_increment.sql index e15d44b..d54558c 100644 --- a/dwh-modeling/sql/03_demo_increment.sql +++ b/dwh-modeling/sql/03_demo_increment.sql @@ -5,8 +5,8 @@ -- 0. Новые события в STG (пример) INSERT INTO stg.customers_raw (_load_id, _load_ts, event_ts, customer_id, email, phone, city) VALUES -('batch_20250406_0900', '2025-04-06 09:00', NULL, '101','b@ex.com','700','Москва'), -- город вернулся -('batch_20250406_1200', '2025-04-06 12:00', NULL, '103','d@ex.com','702','Казань'); -- новый клиент +('batch_20241101_0800', '2024-11-01 08:00', '2024-11-01', '101','b@ex.com','700','Москва'), -- город вернулся +('batch_20240310_0800', '2024-03-10 08:00', '2024-03-10', '103','d@ex.com','702','Казань'); -- новый клиент -- 1) UPSERT в ODS последнего снимка по BK WITH src AS ( @@ -18,20 +18,20 @@ WITH src AS ( NULLIF(s.event_ts,'')::timestamp AS event_ts, s._load_id, s._load_ts, - COALESCE(NULLIF(s.event_ts,'')::timestamp, s._load_ts, - to_timestamp(regexp_replace(s._load_id,'^batch_',''),'YYYYMMDD_HH24MI')) AS eff_ts + COALESCE(NULLIF(s.event_ts,'')::timestamp, s._load_ts) AS eff_ts FROM stg.customers_raw s WHERE s.customer_id ~ '^\d+$' ), -last_per_bk AS ( - SELECT DISTINCT ON (customer_id) - customer_id, email, phone, city, event_ts, _load_id, _load_ts, eff_ts +ranked AS ( + SELECT + customer_id, email, phone, city, event_ts, _load_id, _load_ts, + row_number() OVER (PARTITION BY customer_id ORDER BY eff_ts DESC, _load_ts DESC) AS rn FROM src - ORDER BY customer_id, eff_ts DESC, _load_ts DESC ) INSERT INTO ods.customers (customer_id, email, phone, city, event_ts, _load_id, _load_ts) SELECT customer_id, email, phone, city, event_ts, _load_id, _load_ts -FROM last_per_bk +FROM ranked +WHERE rn = 1 ON CONFLICT (customer_id) DO UPDATE SET email = EXCLUDED.email, phone = EXCLUDED.phone, @@ -44,67 +44,84 @@ WHERE COALESCE(EXCLUDED.event_ts, EXCLUDED._load_ts) > COALESCE(ods.customers.event_ts, ods.customers._load_ts); -- 2) Инкрементальное SCD2 из ODS (в одной транзакции, по последнему снимку в ODS) +-- Канон для курса: daily-grain, интервалы [valid_from, valid_to), текущая версия = valid_to IS NULL +-- Предполагаем, что для клиента нет нескольких изменений в один день. BEGIN; + -- 2.1) Закрываем предыдущую актуальную версию (только если реально изменились атрибуты) WITH delta AS ( SELECT c.customer_id AS customer_bk, c.email, c.phone, c.city, - COALESCE(c.event_ts, c._load_ts) AS eff_ts, + COALESCE(c.event_ts::date, c._load_ts::date) AS eff_date, dds.customer_hash(c.email, c.phone, c.city) AS hashdiff FROM ods.customers c ), 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 + WHERE d.valid_to IS NULL ) - -- закрываем старые версии только для тех BK, по которым реально вставилась новая UPDATE dds.dim_customer d - SET valid_to = LEAST(d.valid_to, t.eff_ts - interval '1 second'), - is_current = FALSE, + SET valid_to = x.eff_date, 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; + FROM ( + SELECT + t.customer_bk, + t.eff_date, + c.customer_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_sk = x.customer_sk + AND d.valid_to IS NULL; + + -- 2.2) Вставляем новую версию (только если новая или изменившаяся) + WITH delta AS ( + SELECT + c.customer_id AS customer_bk, + c.email, c.phone, c.city, + COALESCE(c.event_ts::date, c._load_ts::date) AS eff_date, + dds.customer_hash(c.email, c.phone, c.city) AS hashdiff + FROM ods.customers c + ), + current AS ( + SELECT d.* + FROM dds.dim_customer d + WHERE d.valid_to IS NULL + ), + to_insert AS ( + SELECT + t.customer_bk, + t.email, + t.phone, + t.city, + t.hashdiff, + t.eff_date + FROM delta t + LEFT JOIN current c + ON c.customer_bk = t.customer_bk + WHERE c.customer_sk IS NULL + OR (c.hashdiff <> t.hashdiff AND t.eff_date > c.valid_from) + ) + INSERT INTO dds.dim_customer ( + customer_bk, email, phone, city, hashdiff, + valid_from, valid_to, + created_at, updated_at + ) + SELECT + t.customer_bk, t.email, t.phone, t.city, t.hashdiff, + t.eff_date, NULL, + now(), now() + FROM to_insert t + WHERE NOT EXISTS ( + SELECT 1 + FROM dds.dim_customer d + WHERE d.customer_bk = t.customer_bk + AND d.valid_from = t.eff_date + ); COMMIT; -- (факты можно не перезаливать — даты заказов не поменялись) diff --git a/dwh-modeling/sql/04_validation.sql b/dwh-modeling/sql/04_validation.sql index 5a044d6..1810410 100644 --- a/dwh-modeling/sql/04_validation.sql +++ b/dwh-modeling/sql/04_validation.sql @@ -41,3 +41,33 @@ BEGIN FORMAT('ОШИБКА: у клиента 101 только %s версия, ожидается ≥2 (должна быть история)', version_count); RAISE NOTICE '✅ SCD Type 2: клиент 101 имеет % версий — история сохранена', version_count; END $$; + +-- 4. Проверка SCD Type 2: у каждого клиента ровно одна актуальная версия (valid_to IS NULL) +DO $$ +DECLARE + customers_cnt BIGINT; + current_cnt BIGINT; +BEGIN + SELECT COUNT(DISTINCT customer_bk) INTO customers_cnt FROM dds.dim_customer; + SELECT COUNT(*) INTO current_cnt + FROM dds.dim_customer + WHERE valid_to IS NULL; + -- + ASSERT current_cnt = customers_cnt, + FORMAT('ОШИБКА: актуальных строк %s, а уникальных клиентов %s (ожидается 1 current на клиента)', + current_cnt, customers_cnt); + RAISE NOTICE '✅ SCD Type 2: current-строки = количеству клиентов (%)', current_cnt; +END $$; + +-- 5. Проверка SCD Type 2: периоды корректны (valid_to > valid_from или valid_to IS NULL) +DO $$ +BEGIN + ASSERT NOT EXISTS ( + SELECT 1 + FROM dds.dim_customer + WHERE valid_to IS NOT NULL + AND valid_to <= valid_from + ), + 'ОШИБКА: найдены строки dim_customer с некорректным периодом (valid_to <= valid_from)'; + RAISE NOTICE '✅ SCD Type 2: периоды valid_from/valid_to корректны'; +END $$; diff --git a/dwh-modeling/sql/06_dml_dm.sql b/dwh-modeling/sql/06_dml_dm.sql index f7b02c2..29d3b63 100644 --- a/dwh-modeling/sql/06_dml_dm.sql +++ b/dwh-modeling/sql/06_dml_dm.sql @@ -55,6 +55,5 @@ 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: нужна вся история фактов +-- Здесь не фильтруем по valid_to: нужна вся история фактов GROUP BY c.customer_bk; - diff --git a/dwh-modeling/sql/07_ddl_hw_customer_status.sql b/dwh-modeling/sql/07_ddl_hw_customer_status.sql index c3c2946..e38b444 100644 --- a/dwh-modeling/sql/07_ddl_hw_customer_status.sql +++ b/dwh-modeling/sql/07_ddl_hw_customer_status.sql @@ -34,10 +34,15 @@ CREATE TABLE dds.dim_customer_status ( 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, + valid_from DATE NOT NULL, + valid_to DATE, created_at TIMESTAMP NOT NULL DEFAULT NOW(), updated_at TIMESTAMP NOT NULL DEFAULT NOW() ); +ALTER TABLE dds.dim_customer_status + ADD CONSTRAINT uq_dim_customer_status_bk_from UNIQUE (customer_bk, valid_from); + +CREATE INDEX ix_dim_customer_status_bk_current + ON dds.dim_customer_status (customer_bk) + WHERE valid_to IS NULL; 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 676527f..8dc7b01 100644 --- a/dwh-modeling/sql/08_dml_hw_customer_status_template.sql +++ b/dwh-modeling/sql/08_dml_hw_customer_status_template.sql @@ -12,7 +12,7 @@ -- (в STG/ODS эта колонка будет жить как _load_ts). -- 3) Постройте из ods.customer_status измерение dds.dim_customer_status в стиле SCD2: -- - одна строка на период действия статуса (valid_from / valid_to); --- - is_current = TRUE только у актуальной строки для клиента; +-- - актуальная строка для клиента — та, где valid_to IS NULL; -- - hashdiff можно считать, например, от одного поля status. -- 4) При желании добавьте инкрементальную логику (как в 03_demo_increment.sql). -- 5) Опционально: соберите витрину dm.mart_customer_status_daily @@ -31,7 +31,7 @@ -- Примерный план: -- - рассчитать hashdiff по (status); -- - по каждому клиенту отсортировать события по времени; --- - построить для каждой строки valid_from и valid_to (LEAD() OVER ...); +-- - построить для каждой строки valid_from и valid_to (LEAD() OVER ...), последняя valid_to = NULL; -- - вставить в dds.dim_customer_status. -- 3. DDS: инкрементальная загрузка (по желанию)