переписана реализация scd2

This commit is contained in:
2025-12-21 18:50:50 +03:00
parent 87301921a9
commit ae00475df5
15 changed files with 205 additions and 183 deletions
+3 -4
View File
@@ -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: факт "Продажи"
+21 -24
View File
@@ -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;
+73 -56
View File
@@ -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;
-- (факты можно не перезаливать — даты заказов не поменялись)
+30
View File
@@ -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 $$;
+1 -2
View File
@@ -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;
@@ -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;
@@ -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: инкрементальная загрузка (по желанию)