Исправлиение примера с scd2
This commit is contained in:
@@ -529,12 +529,12 @@ flowchart TD
|
|||||||
|
|
||||||
Все необходимые скрипты для построения хранилища находятся в папке [`sql/`](sql/):
|
Все необходимые скрипты для построения хранилища находятся в папке [`sql/`](sql/):
|
||||||
|
|
||||||
- [`01_ddl_stg-dds.sql`](sql/01_ddl_stg-dds.sql) — создание схем и таблиц (STG, ODS, DDS)
|
- [`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) — загрузка данных и трансформация
|
- [`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
|
- [`03_demo_increment.sql`](sql/03_demo_increment.sql) — пример инкрементальной загрузки и SCD2 по последнему снимку в ODS;
|
||||||
- [`04_validation.sql`](sql/04_validation.sql) — проверки качества данных
|
- [`04_validation.sql`](sql/04_validation.sql) — проверки качества данных;
|
||||||
- [`05_ddl_dm.sql`](sql/05_ddl_dm.sql) — создание витрин (Data Marts)
|
- [`05_ddl_dm.sql`](sql/05_ddl_dm.sql) — создание витрин (Data Marts);
|
||||||
- [`06_dml_dm.sql`](sql/06_dml_dm.sql) — наполнение витрин данными
|
- [`06_dml_dm.sql`](sql/06_dml_dm.sql) — наполнение витрин данными.
|
||||||
|
|
||||||
### Пример SQL-запроса для витрины
|
### Пример SQL-запроса для витрины
|
||||||
|
|
||||||
|
|||||||
@@ -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)
|
#### А что, если СУБД не позволяет UPDATE? (Trino, Hive, ClickHouse в режиме append-only)
|
||||||
|
|
||||||
Некоторые аналитические системы (например, **Hive в формате ORC/Parquet**, **Trino**, **ClickHouse в режиме только вставки**) **не поддерживают UPDATE старых строк**. Как тогда реализовать SCD Type 2?
|
Некоторые аналитические системы (например, **Hive в формате ORC/Parquet**, **Trino**, **ClickHouse в режиме только вставки**) **не поддерживают UPDATE старых строк**. Как тогда реализовать SCD Type 2?
|
||||||
|
|||||||
@@ -128,6 +128,8 @@ SELECT product_id, name
|
|||||||
FROM ods.products;
|
FROM ods.products;
|
||||||
|
|
||||||
-- 5. DDS: dim_customer — первичная загрузка SCD2 (full backfill из STG)
|
-- 5. DDS: dim_customer — первичная загрузка SCD2 (full backfill из STG)
|
||||||
|
-- Здесь мы пересчитываем всю историю клиента из событий в STG.
|
||||||
|
-- В реальном DWH такую полную перезагрузку делают редко; для инкремента см. 03_demo_increment.sql и SCD.md.
|
||||||
TRUNCATE dds.dim_customer, dds.fact_sales;
|
TRUNCATE dds.dim_customer, dds.fact_sales;
|
||||||
|
|
||||||
BEGIN;
|
BEGIN;
|
||||||
@@ -158,7 +160,7 @@ BEGIN;
|
|||||||
framed AS (
|
framed AS (
|
||||||
SELECT
|
SELECT
|
||||||
customer_bk, email, phone, city, hashdiff,
|
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
|
lead(eff_ts) OVER (PARTITION BY customer_bk ORDER BY eff_ts) AS next_ts
|
||||||
FROM changes
|
FROM changes
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -43,7 +43,7 @@ SET email = EXCLUDED.email,
|
|||||||
WHERE COALESCE(EXCLUDED.event_ts, EXCLUDED._load_ts) >
|
WHERE COALESCE(EXCLUDED.event_ts, EXCLUDED._load_ts) >
|
||||||
COALESCE(ods.customers.event_ts, ods.customers._load_ts);
|
COALESCE(ods.customers.event_ts, ods.customers._load_ts);
|
||||||
|
|
||||||
-- 2) Инкрементальное SCD2 из ODS
|
-- 2) Инкрементальное SCD2 из ODS (в одной транзакции, по последнему снимку в ODS)
|
||||||
BEGIN;
|
BEGIN;
|
||||||
WITH delta AS (
|
WITH delta AS (
|
||||||
SELECT
|
SELECT
|
||||||
@@ -53,34 +53,58 @@ BEGIN;
|
|||||||
dds.customer_hash(c.email, c.phone, c.city) AS hashdiff
|
dds.customer_hash(c.email, c.phone, c.city) AS hashdiff
|
||||||
FROM ods.customers c
|
FROM ods.customers c
|
||||||
),
|
),
|
||||||
expired AS (
|
current AS (
|
||||||
UPDATE dds.dim_customer d
|
-- текущие версии клиентов в измерении
|
||||||
SET valid_to = LEAST(d.valid_to, delta.eff_ts - interval '1 second'),
|
SELECT d.*
|
||||||
is_current = FALSE,
|
FROM dds.dim_customer d
|
||||||
updated_at = now()
|
WHERE d.is_current = TRUE
|
||||||
FROM delta
|
),
|
||||||
WHERE d.customer_bk = delta.customer_bk
|
to_upsert AS (
|
||||||
AND d.is_current = TRUE
|
-- только новые BK или реально изменившиеся атрибуты
|
||||||
AND d.hashdiff <> delta.hashdiff
|
SELECT
|
||||||
AND delta.eff_ts >= d.valid_from
|
d.customer_bk,
|
||||||
RETURNING 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 (
|
INSERT INTO dds.dim_customer (
|
||||||
customer_bk, email, phone, city, hashdiff,
|
customer_bk, email, phone, city, hashdiff,
|
||||||
valid_from, valid_to,
|
valid_from, valid_to,
|
||||||
is_current, created_at, updated_at
|
is_current, created_at, updated_at
|
||||||
)
|
)
|
||||||
SELECT
|
SELECT
|
||||||
s.customer_bk, s.email, s.phone, s.city, s.hashdiff,
|
t.customer_bk, t.email, t.phone, t.city, t.hashdiff,
|
||||||
CASE WHEN d.customer_bk IS NULL THEN timestamp '1900-01-01' ELSE s.eff_ts END AS valid_from,
|
CASE WHEN t.current_sk IS NULL
|
||||||
timestamp '9999-12-31',
|
THEN timestamp '1900-01-01' -- первая версия: техническое "начало истории"
|
||||||
|
ELSE t.eff_ts
|
||||||
|
END AS valid_from,
|
||||||
|
timestamp '9999-12-31' AS valid_to,
|
||||||
TRUE, now(), now()
|
TRUE, now(), now()
|
||||||
FROM delta s
|
FROM to_upsert t
|
||||||
LEFT JOIN dds.dim_customer d
|
ON CONFLICT (customer_bk, valid_from) DO NOTHING
|
||||||
ON d.customer_bk = s.customer_bk AND d.is_current = TRUE
|
RETURNING customer_bk, valid_from
|
||||||
WHERE d.customer_bk IS NULL -- новый BK
|
)
|
||||||
OR d.hashdiff <> s.hashdiff -- изменившийся BK
|
-- закрываем старые версии только для тех BK, по которым реально вставилась новая
|
||||||
ON CONFLICT (customer_bk, valid_from) DO NOTHING;
|
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;
|
COMMIT;
|
||||||
|
|
||||||
-- (факты можно не перезаливать — даты заказов не поменялись)
|
-- (факты можно не перезаливать — даты заказов не поменялись)
|
||||||
|
|||||||
Reference in New Issue
Block a user