Уточнения по SCD 2 - бех юэкфилл
This commit is contained in:
@@ -1,5 +1,5 @@
|
|||||||
customer_id,email,phone,city,_load_id,load_ts
|
customer_id,email,phone,city,event_ts,_load_id,load_ts
|
||||||
101,a@ex.com,700,Москва,batch_20250405_0800,2025-04-05 08:00
|
101,a@ex.com,700,Москва,NULL,batch_20250405_0800,2025-04-05 08:00
|
||||||
102,c@ex.com,701,СПб,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,Москва,batch_20250405_1200,2025-04-05 12:00
|
101,b@ex.com,700,Москва,NULL,batch_20250405_1200,2025-04-05 12:00
|
||||||
101,b@ex.com,700,Санкт-Петербург,batch_20250405_1800,2025-04-05 18:00
|
101,b@ex.com,700,Санкт-Петербург,NULL,batch_20250405_1800,2025-04-05 18:00
|
||||||
|
@@ -14,13 +14,15 @@ CREATE SCHEMA ods;
|
|||||||
CREATE SCHEMA dds;
|
CREATE SCHEMA dds;
|
||||||
|
|
||||||
-- 2. STG: сырые данные (как пришли)
|
-- 2. STG: сырые данные (как пришли)
|
||||||
|
DROP TABLE IF EXISTS stg.customers_raw;
|
||||||
CREATE TABLE stg.customers_raw (
|
CREATE TABLE stg.customers_raw (
|
||||||
customer_id TEXT, -- может быть строкой или числом
|
customer_id TEXT, -- может быть строкой или числом
|
||||||
email TEXT,
|
email TEXT,
|
||||||
phone TEXT,
|
phone TEXT,
|
||||||
city TEXT,
|
city TEXT,
|
||||||
|
event_ts TEXT,
|
||||||
_load_id TEXT, -- идентификатор загрузки (обязательно!)
|
_load_id TEXT, -- идентификатор загрузки (обязательно!)
|
||||||
_load_ts TIMESTAMP DEFAULT NOW(), -- время получения в DWH
|
_load_ts TIMESTAMP DEFAULT NOW() -- время получения в DWH
|
||||||
);
|
);
|
||||||
|
|
||||||
CREATE TABLE stg.orders_raw (
|
CREATE TABLE stg.orders_raw (
|
||||||
@@ -43,11 +45,13 @@ CREATE TABLE stg.products_raw (
|
|||||||
);
|
);
|
||||||
|
|
||||||
-- 3. ODS: очищенные данные
|
-- 3. ODS: очищенные данные
|
||||||
|
DROP TABLE IF EXISTS ods.customers;
|
||||||
CREATE TABLE ods.customers (
|
CREATE TABLE ods.customers (
|
||||||
customer_id INT NOT NULL, -- привели к INT
|
customer_id INT NOT NULL, -- привели к INT
|
||||||
email VARCHAR(100),
|
email VARCHAR(100),
|
||||||
phone VARCHAR(20),
|
phone VARCHAR(20),
|
||||||
city VARCHAR(50),
|
city VARCHAR(50),
|
||||||
|
event_ts TIMESTAMP,
|
||||||
_load_id TEXT NOT NULL, -- сохраняем для отладки и SCD
|
_load_id TEXT NOT NULL, -- сохраняем для отладки и SCD
|
||||||
_load_ts TIMESTAMP NOT NULL -- время загрузки (копия из STG)
|
_load_ts TIMESTAMP NOT NULL -- время загрузки (копия из STG)
|
||||||
);
|
);
|
||||||
|
|||||||
@@ -37,12 +37,13 @@ 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;
|
||||||
|
|
||||||
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
|
SELECT
|
||||||
customer_id::INT,
|
customer_id::INT,
|
||||||
NULLIF(TRIM(email), ''),
|
NULLIF(TRIM(email), ''),
|
||||||
NULLIF(TRIM(phone), ''),
|
NULLIF(TRIM(phone), ''),
|
||||||
NULLIF(TRIM(city), ''),
|
NULLIF(TRIM(city), ''),
|
||||||
|
event_ts::timestamp,
|
||||||
_load_id,
|
_load_id,
|
||||||
_load_ts
|
_load_ts
|
||||||
FROM stg.customers_raw
|
FROM stg.customers_raw
|
||||||
@@ -81,7 +82,7 @@ DELETE FROM dds.dim_date;
|
|||||||
WITH RECURSIVE dates AS (
|
WITH RECURSIVE dates AS (
|
||||||
SELECT DATE '2023-01-01' AS d
|
SELECT DATE '2023-01-01' AS d
|
||||||
UNION ALL
|
UNION ALL
|
||||||
SELECT d + INTERVAL '1 day'
|
SELECT (d + INTERVAL '1 day')::DATE -- ← приведение к DATE
|
||||||
FROM dates
|
FROM dates
|
||||||
WHERE d + INTERVAL '1 day' <= DATE '2027-12-31'
|
WHERE d + INTERVAL '1 day' <= DATE '2027-12-31'
|
||||||
)
|
)
|
||||||
@@ -96,7 +97,7 @@ SELECT
|
|||||||
EXTRACT(QUARTER FROM d)::SMALLINT,
|
EXTRACT(QUARTER FROM d)::SMALLINT,
|
||||||
EXTRACT(MONTH FROM d)::SMALLINT,
|
EXTRACT(MONTH FROM d)::SMALLINT,
|
||||||
EXTRACT(DAY 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,
|
EXTRACT(DOW FROM d)::SMALLINT,
|
||||||
d BETWEEN DATE_TRUNC('month', d)
|
d BETWEEN DATE_TRUNC('month', d)
|
||||||
AND DATE_TRUNC('month', d) + INTERVAL '1 month' - INTERVAL '1 day'
|
AND DATE_TRUNC('month', d) + INTERVAL '1 month' - INTERVAL '1 day'
|
||||||
@@ -130,8 +131,7 @@ BEGIN;
|
|||||||
PARTITION BY s.customer_id
|
PARTITION BY s.customer_id
|
||||||
ORDER BY COALESCE(s.event_ts, s._load_ts) DESC, s._load_ts DESC
|
ORDER BY COALESCE(s.event_ts, s._load_ts) DESC, s._load_ts DESC
|
||||||
) AS rn
|
) AS rn
|
||||||
FROM stg.customers_raw s
|
FROM ods.customers s
|
||||||
WHERE s.customer_id ~ '^\d+$'
|
|
||||||
),
|
),
|
||||||
delta AS (
|
delta AS (
|
||||||
-- берём по BK последнюю версию из поступивших данных
|
-- берём по BK последнюю версию из поступивших данных
|
||||||
@@ -192,11 +192,11 @@ SELECT
|
|||||||
oi.price_at_sale * oi.qty
|
oi.price_at_sale * oi.qty
|
||||||
FROM ods.orders o
|
FROM ods.orders o
|
||||||
JOIN ods.order_items oi ON o.order_id = oi.order_id
|
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_product dp ON p.product_id = dp.product_bk
|
||||||
JOIN dds.dim_customer dc
|
JOIN dds.dim_customer dc
|
||||||
ON o.customer_id = dc.customer_bk
|
ON o.customer_id = dc.customer_bk
|
||||||
AND o.order_date BETWEEN dc.valid_from AND dc.valid_to;
|
AND o.order_date::timestamp BETWEEN dc.valid_from AND dc.valid_to;
|
||||||
|
|
||||||
-- 7. Проверка — вывод итогов (не часть ETL, но полезно для отладки)
|
-- 7. Проверка — вывод итогов (не часть ETL, но полезно для отладки)
|
||||||
-- В реальном пайплайне такие SELECT выносятся в отдельные скрипты или дашборды
|
-- В реальном пайплайне такие SELECT выносятся в отдельные скрипты или дашборды
|
||||||
|
|||||||
Reference in New Issue
Block a user