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; -- (факты можно не перезаливать — даты заказов не поменялись)