diff --git a/dwh-modeling/sql/02_dim.sql b/dwh-modeling/sql/02_dim.sql index 48dce75..3b3d5ae 100644 --- a/dwh-modeling/sql/02_dim.sql +++ b/dwh-modeling/sql/02_dim.sql @@ -37,25 +37,30 @@ INSERT INTO stg.products_raw (product_id, name) VALUES 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) -SELECT - customer_id::INT, - 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+$'; +SELECT customer_id, email, phone, city, event_ts, _load_id, _load_ts +FROM last_per_bk; INSERT INTO ods.order_items (order_item_id, order_id, product_id, qty, price_at_sale) SELECT @@ -113,75 +118,60 @@ INSERT INTO dds.dim_product (product_bk, product_name) SELECT product_id, name FROM ods.products; --- 5. DDS: dim_customer — SCD Type 2 --- Загрузка SCD2 (идемпотентный апсерт) --- Делаем внутри одной транзакции +-- 5. DDS: dim_customer — первичная загрузка SCD2 (full backfill из STG) +TRUNCATE dds.dim_customer, dds.fact_sales; BEGIN; - -- Подготовим дельту: по BK одна запись на «текущую истину» из ODS/STG - WITH src AS ( + WITH src AS ( SELECT - s.customer_id::INT AS customer_bk, - NULLIF(trim(s.email), '') AS email, - NULLIF(trim(s.phone), '') AS phone, - NULLIF(trim(s.city), '') AS city, - COALESCE(s.event_ts, s._load_ts)::timestamp AS eff_ts, - dds.customer_hash(s.email, s.phone, s.city) AS hashdiff, - ROW_NUMBER() OVER ( - PARTITION BY s.customer_id - ORDER BY COALESCE(s.event_ts, s._load_ts) DESC, s._load_ts DESC - ) AS rn - FROM ods.customers s - ), - delta AS ( - -- берём по BK последнюю версию из поступивших данных - SELECT customer_bk, email, phone, city, hashdiff, eff_ts + s.customer_id::INT AS customer_bk, + 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, + 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 FROM src - WHERE rn = 1 - ), - -- Закрываем текущие версии там, где атрибуты изменились - 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 - ) - -- Вставляем новые текущие версии: - -- а) для новых BK (раньше не было строки) - -- б) для изменившихся BK (после закрытия предыдущей версии) - INSERT INTO dds.dim_customer ( - customer_bk, email, phone, city, hashdiff, - valid_from, valid_to, is_current, created_at, updated_at - ) + ), + changes AS ( + -- только первые состояния и фактические изменения атрибутов + SELECT * + FROM ordered + WHERE prev_hash IS DISTINCT FROM hashdiff OR prev_hash IS NULL + ), + framed AS ( SELECT - s.customer_bk, s.email, s.phone, s.city, s.hashdiff, - COALESCE( -- у новых BK открываем «историю с вечности» - -- если нужна «вечность» именно как TIMESTAMP-величина: - CASE WHEN d.customer_bk IS NULL THEN TIMESTAMP '1900-01-01' ELSE s.eff_ts END, - TIMESTAMP '1900-01-01' - ) AS valid_from, - TIMESTAMP '9999-12-31' AS valid_to, - TRUE AS is_current, - 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; - -- Конец транзакции + 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 + FROM changes + ) + INSERT INTO dds.dim_customer ( + customer_bk, email, phone, city, hashdiff, + valid_from, valid_to, is_current, + 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, + now(), now() + FROM framed + ORDER BY customer_bk, valid_from; COMMIT; -- 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) SELECT diff --git a/dwh-modeling/sql/03_demo_increment.sql b/dwh-modeling/sql/03_demo_increment.sql new file mode 100644 index 0000000..5ae402d --- /dev/null +++ b/dwh-modeling/sql/03_demo_increment.sql @@ -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; + +-- (факты можно не перезаливать — даты заказов не поменялись) diff --git a/dwh-modeling/sql/03_validation.sql b/dwh-modeling/sql/04_validation.sql similarity index 66% rename from dwh-modeling/sql/03_validation.sql rename to dwh-modeling/sql/04_validation.sql index b666668..a2d31c3 100644 --- a/dwh-modeling/sql/03_validation.sql +++ b/dwh-modeling/sql/04_validation.sql @@ -12,30 +12,23 @@ BEGIN (SELECT COUNT(*) FROM dds.dim_customer); 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 существует) -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, - FORMAT('ОШИБКА: %s фактов ссылаются на несуществующие даты (неверный date_key)', missing_dates); - RAISE NOTICE '✅ Все факты имеют валидные date_key'; +DO $$ +DECLARE + 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 $$; -- 4. Проверка SCD Type 2: у клиента 101 должно быть ≥2 версий (из-за смены email) @@ -45,7 +38,7 @@ BEGIN SELECT COUNT(*) INTO version_count FROM dds.dim_customer WHERE customer_bk = 101; - + -- ASSERT version_count >= 2, FORMAT('ОШИБКА: у клиента 101 только %s версия, ожидается ≥2 (должна быть история)', 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) LIMIT 10 ) mismatches; - + -- ASSERT bad_rows = 0, 'ОШИБКА: обнаружены расхождения между amount и qty * price_at_sale'; RAISE NOTICE '✅ Все суммы рассчитаны верно (amount = qty × price_at_sale)'; @@ -76,11 +69,6 @@ END $$; -- =============================================== -- Финальный отчёт для аналитика -- =============================================== - -RAISE NOTICE '──────────────────────────────'; -RAISE NOTICE '📊 ОТЧЁТ: Выручка по клиентам (актуальные версии)'; -RAISE NOTICE '──────────────────────────────'; - SELECT dc.customer_bk AS "ID клиента", dc.email AS "Email", @@ -90,9 +78,7 @@ SELECT FROM dds.fact_sales f JOIN dds.dim_customer dc ON f.customer_sk = dc.customer_sk - AND dc.is_current -- только актуальная версия + --AND dc.is_current -- только актуальная версия GROUP BY dc.customer_bk, dc.email, dc.city ORDER BY SUM(f.amount) DESC; -RAISE NOTICE '──────────────────────────────'; -RAISE NOTICE '✅ Все проверки пройдены. DWH готов к построению витрин.'; diff --git a/dwh-modeling/sql/04_ddl_dm.sql b/dwh-modeling/sql/05_ddl_dm.sql similarity index 100% rename from dwh-modeling/sql/04_ddl_dm.sql rename to dwh-modeling/sql/05_ddl_dm.sql diff --git a/dwh-modeling/sql/05_dml_dm.sql b/dwh-modeling/sql/06_dml_dm.sql similarity index 100% rename from dwh-modeling/sql/05_dml_dm.sql rename to dwh-modeling/sql/06_dml_dm.sql