From 57d1dd3186f39a02052c362bd7a9d39e0b48883d Mon Sep 17 00:00:00 2001 From: Dmitrii Date: Fri, 7 Nov 2025 15:46:52 +0300 Subject: [PATCH] =?UTF-8?q?=D0=A3=D1=82=D0=BE=D1=87=D0=BD=D0=B5=D0=BD?= =?UTF-8?q?=D0=B8=D1=8F=20=D0=BF=D0=BE=20SCD=202=20-=20=D0=B1=D0=B5=D1=85?= =?UTF-8?q?=20=D1=8E=D1=8D=D0=BA=D1=84=D0=B8=D0=BB=D0=BB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- dwh-modeling/data/customers.csv | 10 +++++----- dwh-modeling/sql/01_ddl.sql | 6 +++++- dwh-modeling/sql/02_dim.sql | 18 +++++++++--------- 3 files changed, 19 insertions(+), 15 deletions(-) diff --git a/dwh-modeling/data/customers.csv b/dwh-modeling/data/customers.csv index 84a16cd..843393c 100644 --- a/dwh-modeling/data/customers.csv +++ b/dwh-modeling/data/customers.csv @@ -1,5 +1,5 @@ -customer_id,email,phone,city,_load_id,load_ts -101,a@ex.com,700,Москва,batch_20250405_0800,2025-04-05 08:00 -102,c@ex.com,701,СПб,batch_20250405_0800,2025-04-05 08:00 -101,b@ex.com,700,Москва,batch_20250405_1200,2025-04-05 12:00 -101,b@ex.com,700,Санкт-Петербург,batch_20250405_1800,2025-04-05 18:00 \ No newline at end of file +customer_id,email,phone,city,event_ts,_load_id,load_ts +101,a@ex.com,700,Москва,NULL,batch_20250405_0800,2025-04-05 08:00 +102,c@ex.com,701,СПб,NULL,batch_20250405_0800,2025-04-05 08:00 +101,b@ex.com,700,Москва,NULL,batch_20250405_1200,2025-04-05 12:00 +101,b@ex.com,700,Санкт-Петербург,NULL,batch_20250405_1800,2025-04-05 18:00 \ No newline at end of file diff --git a/dwh-modeling/sql/01_ddl.sql b/dwh-modeling/sql/01_ddl.sql index ca60891..1480378 100644 --- a/dwh-modeling/sql/01_ddl.sql +++ b/dwh-modeling/sql/01_ddl.sql @@ -14,13 +14,15 @@ CREATE SCHEMA ods; CREATE SCHEMA dds; -- 2. STG: сырые данные (как пришли) +DROP TABLE IF EXISTS stg.customers_raw; CREATE TABLE stg.customers_raw ( customer_id TEXT, -- может быть строкой или числом email TEXT, phone TEXT, city TEXT, + event_ts TEXT, _load_id TEXT, -- идентификатор загрузки (обязательно!) - _load_ts TIMESTAMP DEFAULT NOW(), -- время получения в DWH + _load_ts TIMESTAMP DEFAULT NOW() -- время получения в DWH ); CREATE TABLE stg.orders_raw ( @@ -43,11 +45,13 @@ CREATE TABLE stg.products_raw ( ); -- 3. ODS: очищенные данные +DROP TABLE IF EXISTS ods.customers; CREATE TABLE ods.customers ( customer_id INT NOT NULL, -- привели к INT email VARCHAR(100), phone VARCHAR(20), city VARCHAR(50), + event_ts TIMESTAMP, _load_id TEXT NOT NULL, -- сохраняем для отладки и SCD _load_ts TIMESTAMP NOT NULL -- время загрузки (копия из STG) ); diff --git a/dwh-modeling/sql/02_dim.sql b/dwh-modeling/sql/02_dim.sql index 5e9f4b3..48dce75 100644 --- a/dwh-modeling/sql/02_dim.sql +++ b/dwh-modeling/sql/02_dim.sql @@ -37,12 +37,13 @@ INSERT INTO stg.products_raw (product_id, name) VALUES TRUNCATE ods.customers, ods.orders, ods.order_items, ods.products; -INSERT INTO ods.customers (customer_id, email, phone, city) +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 @@ -81,7 +82,7 @@ DELETE FROM dds.dim_date; WITH RECURSIVE dates AS ( SELECT DATE '2023-01-01' AS d UNION ALL - SELECT d + INTERVAL '1 day' + SELECT (d + INTERVAL '1 day')::DATE -- ← приведение к DATE FROM dates WHERE d + INTERVAL '1 day' <= DATE '2027-12-31' ) @@ -96,7 +97,7 @@ SELECT EXTRACT(QUARTER FROM d)::SMALLINT, EXTRACT(MONTH FROM d)::SMALLINT, EXTRACT(DAY FROM d)::SMALLINT, - TO_CHAR(d, 'Day'), + TRIM(TO_CHAR(d, 'Day')), -- ← TRIM — убрать trailing space EXTRACT(DOW FROM d)::SMALLINT, d BETWEEN DATE_TRUNC('month', d) AND DATE_TRUNC('month', d) + INTERVAL '1 month' - INTERVAL '1 day' @@ -130,8 +131,7 @@ BEGIN; PARTITION BY s.customer_id ORDER BY COALESCE(s.event_ts, s._load_ts) DESC, s._load_ts DESC ) AS rn - FROM stg.customers_raw s - WHERE s.customer_id ~ '^\d+$' + FROM ods.customers s ), delta AS ( -- берём по BK последнюю версию из поступивших данных @@ -192,11 +192,11 @@ SELECT oi.price_at_sale * oi.qty FROM ods.orders o JOIN ods.order_items oi ON o.order_id = oi.order_id -JOIN ods.products p ON oi.product_id = p.product_id +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 BETWEEN dc.valid_from AND dc.valid_to; +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; -- 7. Проверка — вывод итогов (не часть ETL, но полезно для отладки) -- В реальном пайплайне такие SELECT выносятся в отдельные скрипты или дашборды