SCD2 теперь строится
This commit is contained in:
+67
-77
@@ -37,25 +37,30 @@ INSERT INTO stg.products_raw (product_id, name) VALUES
|
|||||||
|
|
||||||
TRUNCATE ods.customers, ods.orders, ods.order_items, ods.products;
|
TRUNCATE ods.customers, ods.orders, ods.order_items, ods.products;
|
||||||
|
|
||||||
|
-- берём по BK самую позднюю запись (event_ts > _load_ts > _load_id)
|
||||||
|
WITH src AS (
|
||||||
|
SELECT
|
||||||
|
s.customer_id::INT AS customer_id,
|
||||||
|
NULLIF(trim(s.email), '') AS email,
|
||||||
|
NULLIF(trim(s.phone), '') AS phone,
|
||||||
|
NULLIF(trim(s.city), '') AS city,
|
||||||
|
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
|
||||||
|
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
|
||||||
|
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)
|
INSERT INTO ods.customers (customer_id, email, phone, city, event_ts, _load_id, _load_ts)
|
||||||
SELECT
|
SELECT customer_id, email, phone, city, event_ts, _load_id, _load_ts
|
||||||
customer_id::INT,
|
FROM last_per_bk;
|
||||||
NULLIF(TRIM(email), ''),
|
|
||||||
NULLIF(TRIM(phone), ''),
|
|
||||||
NULLIF(TRIM(city), ''),
|
|
||||||
event_ts::timestamp,
|
|
||||||
_load_id,
|
|
||||||
_load_ts
|
|
||||||
FROM stg.customers_raw
|
|
||||||
WHERE customer_id ~ '^\d+$';
|
|
||||||
|
|
||||||
INSERT INTO ods.orders (order_id, order_date, customer_id)
|
|
||||||
SELECT
|
|
||||||
order_id::INT,
|
|
||||||
TO_DATE(order_date, 'YYYY-MM-DD'),
|
|
||||||
customer_id::INT
|
|
||||||
FROM stg.orders_raw
|
|
||||||
WHERE order_date IS NOT NULL AND customer_id ~ '^\d+$';
|
|
||||||
|
|
||||||
INSERT INTO ods.order_items (order_item_id, order_id, product_id, qty, price_at_sale)
|
INSERT INTO ods.order_items (order_item_id, order_id, product_id, qty, price_at_sale)
|
||||||
SELECT
|
SELECT
|
||||||
@@ -113,75 +118,60 @@ INSERT INTO dds.dim_product (product_bk, product_name)
|
|||||||
SELECT product_id, name
|
SELECT product_id, name
|
||||||
FROM ods.products;
|
FROM ods.products;
|
||||||
|
|
||||||
-- 5. DDS: dim_customer — SCD Type 2
|
-- 5. DDS: dim_customer — первичная загрузка SCD2 (full backfill из STG)
|
||||||
-- Загрузка SCD2 (идемпотентный апсерт)
|
TRUNCATE dds.dim_customer, dds.fact_sales;
|
||||||
-- Делаем внутри одной транзакции
|
|
||||||
|
|
||||||
BEGIN;
|
BEGIN;
|
||||||
-- Подготовим дельту: по BK одна запись на «текущую истину» из ODS/STG
|
WITH src AS (
|
||||||
WITH src AS (
|
|
||||||
SELECT
|
SELECT
|
||||||
s.customer_id::INT AS customer_bk,
|
s.customer_id::INT AS customer_bk,
|
||||||
NULLIF(trim(s.email), '') AS email,
|
NULLIF(trim(s.email), '') AS email,
|
||||||
NULLIF(trim(s.phone), '') AS phone,
|
NULLIF(trim(s.phone), '') AS phone,
|
||||||
NULLIF(trim(s.city), '') AS city,
|
NULLIF(trim(s.city), '') AS city,
|
||||||
COALESCE(s.event_ts, s._load_ts)::timestamp AS eff_ts,
|
COALESCE(NULLIF(s.event_ts,'')::timestamp, s._load_ts,
|
||||||
dds.customer_hash(s.email, s.phone, s.city) AS hashdiff,
|
to_timestamp(regexp_replace(s._load_id,'^batch_',''),'YYYYMMDD_HH24MI')) AS eff_ts,
|
||||||
ROW_NUMBER() OVER (
|
dds.customer_hash(s.email, s.phone, s.city) AS hashdiff
|
||||||
PARTITION BY s.customer_id
|
FROM stg.customers_raw s
|
||||||
ORDER BY COALESCE(s.event_ts, s._load_ts) DESC, s._load_ts DESC
|
WHERE s.customer_id ~ '^\d+$'
|
||||||
) AS rn
|
),
|
||||||
FROM ods.customers s
|
ordered AS (
|
||||||
),
|
SELECT *,
|
||||||
delta AS (
|
lag(hashdiff) OVER (PARTITION BY customer_bk ORDER BY eff_ts) AS prev_hash,
|
||||||
-- берём по BK последнюю версию из поступивших данных
|
row_number() OVER (PARTITION BY customer_bk ORDER BY eff_ts) AS rn
|
||||||
SELECT customer_bk, email, phone, city, hashdiff, eff_ts
|
|
||||||
FROM src
|
FROM src
|
||||||
WHERE rn = 1
|
),
|
||||||
),
|
changes AS (
|
||||||
-- Закрываем текущие версии там, где атрибуты изменились
|
-- только первые состояния и фактические изменения атрибутов
|
||||||
expired AS (
|
SELECT *
|
||||||
UPDATE dds.dim_customer d
|
FROM ordered
|
||||||
SET valid_to = LEAST(d.valid_to, delta.eff_ts - INTERVAL '1 second'),
|
WHERE prev_hash IS DISTINCT FROM hashdiff OR prev_hash IS NULL
|
||||||
is_current = FALSE,
|
),
|
||||||
updated_at = NOW()
|
framed AS (
|
||||||
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
|
|
||||||
)
|
|
||||||
-- Вставляем новые текущие версии:
|
|
||||||
-- а) для новых BK (раньше не было строки)
|
|
||||||
-- б) для изменившихся BK (после закрытия предыдущей версии)
|
|
||||||
INSERT INTO dds.dim_customer (
|
|
||||||
customer_bk, email, phone, city, hashdiff,
|
|
||||||
valid_from, valid_to, is_current, created_at, updated_at
|
|
||||||
)
|
|
||||||
SELECT
|
SELECT
|
||||||
s.customer_bk, s.email, s.phone, s.city, s.hashdiff,
|
customer_bk, email, phone, city, hashdiff,
|
||||||
COALESCE( -- у новых BK открываем «историю с вечности»
|
CASE WHEN rn = 1 THEN timestamp '1900-01-01' ELSE eff_ts END AS valid_from,
|
||||||
-- если нужна «вечность» именно как TIMESTAMP-величина:
|
lead(eff_ts) OVER (PARTITION BY customer_bk ORDER BY eff_ts) AS next_ts
|
||||||
CASE WHEN d.customer_bk IS NULL THEN TIMESTAMP '1900-01-01' ELSE s.eff_ts END,
|
FROM changes
|
||||||
TIMESTAMP '1900-01-01'
|
)
|
||||||
) AS valid_from,
|
INSERT INTO dds.dim_customer (
|
||||||
TIMESTAMP '9999-12-31' AS valid_to,
|
customer_bk, email, phone, city, hashdiff,
|
||||||
TRUE AS is_current,
|
valid_from, valid_to, is_current,
|
||||||
NOW(), NOW()
|
created_at, updated_at
|
||||||
FROM delta s
|
)
|
||||||
LEFT JOIN dds.dim_customer d
|
SELECT
|
||||||
ON d.customer_bk = s.customer_bk AND d.is_current = TRUE
|
customer_bk, email, phone, city, hashdiff,
|
||||||
WHERE d.customer_bk IS NULL -- новый BK
|
valid_from,
|
||||||
OR d.hashdiff <> s.hashdiff -- или изменившийся BK
|
COALESCE(next_ts - interval '1 second', timestamp '9999-12-31') AS valid_to,
|
||||||
ON CONFLICT (customer_bk, valid_from) DO NOTHING;
|
(next_ts IS NULL) AS is_current,
|
||||||
-- Конец транзакции
|
now(), now()
|
||||||
|
FROM framed
|
||||||
|
ORDER BY customer_bk, valid_from;
|
||||||
COMMIT;
|
COMMIT;
|
||||||
|
|
||||||
-- 6. DDS: fact_sales — загрузка фактов с учётом SCD
|
-- 6. DDS: fact_sales — загрузка фактов с учётом SCD
|
||||||
-- В продакшене — фильтруем по диапазону дат (инкрементально)
|
-- В продакшене — фильтруем по диапазону дат (инкрементально)
|
||||||
|
|
||||||
DELETE FROM dds.fact_sales;
|
--TRUNCATE dds.fact_sales;
|
||||||
|
|
||||||
INSERT INTO dds.fact_sales (customer_sk, product_sk, date_key, quantity, amount)
|
INSERT INTO dds.fact_sales (customer_sk, product_sk, date_key, quantity, amount)
|
||||||
SELECT
|
SELECT
|
||||||
|
|||||||
@@ -0,0 +1,86 @@
|
|||||||
|
-- ===============================================
|
||||||
|
-- 03_demo_increment.sql
|
||||||
|
-- Имитация новых событий + инкрементальный SCD2
|
||||||
|
-- ===============================================
|
||||||
|
|
||||||
|
-- 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','Казань'); -- новый клиент
|
||||||
|
|
||||||
|
-- 1) UPSERT в ODS последнего снимка по BK
|
||||||
|
WITH src AS (
|
||||||
|
SELECT
|
||||||
|
s.customer_id::INT AS customer_id,
|
||||||
|
NULLIF(trim(s.email), '') AS email,
|
||||||
|
NULLIF(trim(s.phone), '') AS phone,
|
||||||
|
NULLIF(trim(s.city), '') AS city,
|
||||||
|
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
|
||||||
|
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
|
||||||
|
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
|
||||||
|
ON CONFLICT (customer_id) DO UPDATE
|
||||||
|
SET email = EXCLUDED.email,
|
||||||
|
phone = EXCLUDED.phone,
|
||||||
|
city = EXCLUDED.city,
|
||||||
|
event_ts= EXCLUDED.event_ts,
|
||||||
|
_load_id= EXCLUDED._load_id,
|
||||||
|
_load_ts= EXCLUDED._load_ts
|
||||||
|
-- апдейтим только если пришло более «свежее» событие
|
||||||
|
WHERE COALESCE(EXCLUDED.event_ts, EXCLUDED._load_ts) >
|
||||||
|
COALESCE(ods.customers.event_ts, ods.customers._load_ts);
|
||||||
|
|
||||||
|
-- 2) Инкрементальное SCD2 из ODS
|
||||||
|
BEGIN;
|
||||||
|
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,
|
||||||
|
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
|
||||||
|
)
|
||||||
|
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;
|
||||||
|
COMMIT;
|
||||||
|
|
||||||
|
-- (факты можно не перезаливать — даты заказов не поменялись)
|
||||||
@@ -12,30 +12,23 @@ BEGIN
|
|||||||
(SELECT COUNT(*) FROM dds.dim_customer);
|
(SELECT COUNT(*) FROM dds.dim_customer);
|
||||||
END $$;
|
END $$;
|
||||||
|
|
||||||
-- 2. Проверка: fact_sales содержит все строки из order_items
|
|
||||||
DO $$
|
|
||||||
DECLARE
|
|
||||||
expected_count INT := (SELECT COUNT(*) FROM ods.order_items);
|
|
||||||
actual_count INT := (SELECT COUNT(*) FROM dds.fact_sales);
|
|
||||||
BEGIN
|
|
||||||
ASSERT actual_count = expected_count,
|
|
||||||
FORMAT('ОШИБКА: в fact_sales %s строк, а в ods.order_items — %s. Разница: %s',
|
|
||||||
actual_count, expected_count, expected_count - actual_count);
|
|
||||||
RAISE NOTICE '✅ fact_sales: количество строк совпадает с ods.order_items (%)', actual_count;
|
|
||||||
END $$;
|
|
||||||
|
|
||||||
-- 3. Проверка: у каждого факта есть валидная дата (date_key существует)
|
-- 3. Проверка: у каждого факта есть валидная дата (date_key существует)
|
||||||
DO $$
|
|
||||||
DECLARE missing_dates INT;
|
|
||||||
BEGIN
|
|
||||||
SELECT COUNT(*) INTO missing_dates
|
|
||||||
FROM dds.fact_sales f
|
|
||||||
LEFT JOIN dds.dim_date d ON f.date_key = d.date_key
|
|
||||||
WHERE d.date_key IS NULL;
|
|
||||||
|
|
||||||
ASSERT missing_dates = 0,
|
DO $$
|
||||||
FORMAT('ОШИБКА: %s фактов ссылаются на несуществующие даты (неверный date_key)', missing_dates);
|
DECLARE
|
||||||
RAISE NOTICE '✅ Все факты имеют валидные date_key';
|
expected_count bigint;
|
||||||
|
actual_count bigint;
|
||||||
|
BEGIN
|
||||||
|
SELECT COUNT(*) INTO expected_count FROM ods.order_items;
|
||||||
|
SELECT COUNT(*) INTO actual_count FROM dds.fact_sales;
|
||||||
|
--
|
||||||
|
ASSERT actual_count = expected_count,
|
||||||
|
format('ОШИБКА: в fact_sales %s строк, а в ods.order_items — %s. Разница: %s',
|
||||||
|
actual_count, expected_count, expected_count - actual_count);
|
||||||
|
--
|
||||||
|
RAISE NOTICE '✅ fact_sales: количество строк совпадает с ods.order_items (%)', actual_count;
|
||||||
END $$;
|
END $$;
|
||||||
|
|
||||||
-- 4. Проверка SCD Type 2: у клиента 101 должно быть ≥2 версий (из-за смены email)
|
-- 4. Проверка SCD Type 2: у клиента 101 должно быть ≥2 версий (из-за смены email)
|
||||||
@@ -45,7 +38,7 @@ BEGIN
|
|||||||
SELECT COUNT(*) INTO version_count
|
SELECT COUNT(*) INTO version_count
|
||||||
FROM dds.dim_customer
|
FROM dds.dim_customer
|
||||||
WHERE customer_bk = 101;
|
WHERE customer_bk = 101;
|
||||||
|
--
|
||||||
ASSERT version_count >= 2,
|
ASSERT version_count >= 2,
|
||||||
FORMAT('ОШИБКА: у клиента 101 только %s версия, ожидается ≥2 (должна быть история)', version_count);
|
FORMAT('ОШИБКА: у клиента 101 только %s версия, ожидается ≥2 (должна быть история)', version_count);
|
||||||
RAISE NOTICE '✅ SCD Type 2: клиент 101 имеет %s версий — история сохранена', version_count;
|
RAISE NOTICE '✅ SCD Type 2: клиент 101 имеет %s версий — история сохранена', version_count;
|
||||||
@@ -67,7 +60,7 @@ BEGIN
|
|||||||
WHERE ROUND(f.amount, 2) <> ROUND(oi.qty * oi.price_at_sale, 2)
|
WHERE ROUND(f.amount, 2) <> ROUND(oi.qty * oi.price_at_sale, 2)
|
||||||
LIMIT 10
|
LIMIT 10
|
||||||
) mismatches;
|
) mismatches;
|
||||||
|
--
|
||||||
ASSERT bad_rows = 0,
|
ASSERT bad_rows = 0,
|
||||||
'ОШИБКА: обнаружены расхождения между amount и qty * price_at_sale';
|
'ОШИБКА: обнаружены расхождения между amount и qty * price_at_sale';
|
||||||
RAISE NOTICE '✅ Все суммы рассчитаны верно (amount = qty × price_at_sale)';
|
RAISE NOTICE '✅ Все суммы рассчитаны верно (amount = qty × price_at_sale)';
|
||||||
@@ -76,11 +69,6 @@ END $$;
|
|||||||
-- ===============================================
|
-- ===============================================
|
||||||
-- Финальный отчёт для аналитика
|
-- Финальный отчёт для аналитика
|
||||||
-- ===============================================
|
-- ===============================================
|
||||||
|
|
||||||
RAISE NOTICE '──────────────────────────────';
|
|
||||||
RAISE NOTICE '📊 ОТЧЁТ: Выручка по клиентам (актуальные версии)';
|
|
||||||
RAISE NOTICE '──────────────────────────────';
|
|
||||||
|
|
||||||
SELECT
|
SELECT
|
||||||
dc.customer_bk AS "ID клиента",
|
dc.customer_bk AS "ID клиента",
|
||||||
dc.email AS "Email",
|
dc.email AS "Email",
|
||||||
@@ -90,9 +78,7 @@ SELECT
|
|||||||
FROM dds.fact_sales f
|
FROM dds.fact_sales f
|
||||||
JOIN dds.dim_customer dc
|
JOIN dds.dim_customer dc
|
||||||
ON f.customer_sk = dc.customer_sk
|
ON f.customer_sk = dc.customer_sk
|
||||||
AND dc.is_current -- только актуальная версия
|
--AND dc.is_current -- только актуальная версия
|
||||||
GROUP BY dc.customer_bk, dc.email, dc.city
|
GROUP BY dc.customer_bk, dc.email, dc.city
|
||||||
ORDER BY SUM(f.amount) DESC;
|
ORDER BY SUM(f.amount) DESC;
|
||||||
|
|
||||||
RAISE NOTICE '──────────────────────────────';
|
|
||||||
RAISE NOTICE '✅ Все проверки пройдены. DWH готов к построению витрин.';
|
|
||||||
Reference in New Issue
Block a user