SCD2 для dim_customers
This commit is contained in:
@@ -1,4 +1,5 @@
|
|||||||
customer_id,email,phone,city
|
customer_id,email,phone,city,_load_id,load_ts
|
||||||
101,a@ex.com,700,Москва
|
101,a@ex.com,700,Москва,batch_20250405_0800,2025-04-05 08:00
|
||||||
101,b@ex.com,700,Москва
|
102,c@ex.com,701,СПб,batch_20250405_0800,2025-04-05 08:00
|
||||||
102,c@ex.com,701,СПб
|
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
|
||||||
|
+39
-12
@@ -15,11 +15,12 @@ CREATE SCHEMA dds;
|
|||||||
|
|
||||||
-- 2. STG: сырые данные (как пришли)
|
-- 2. STG: сырые данные (как пришли)
|
||||||
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,
|
||||||
_load_ts TIMESTAMP DEFAULT NOW()
|
_load_id TEXT, -- идентификатор загрузки (обязательно!)
|
||||||
|
_load_ts TIMESTAMP DEFAULT NOW(), -- время получения в DWH
|
||||||
);
|
);
|
||||||
|
|
||||||
CREATE TABLE stg.orders_raw (
|
CREATE TABLE stg.orders_raw (
|
||||||
@@ -43,10 +44,12 @@ CREATE TABLE stg.products_raw (
|
|||||||
|
|
||||||
-- 3. ODS: очищенные данные
|
-- 3. ODS: очищенные данные
|
||||||
CREATE TABLE ods.customers (
|
CREATE TABLE ods.customers (
|
||||||
customer_id 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),
|
||||||
|
_load_id TEXT NOT NULL, -- сохраняем для отладки и SCD
|
||||||
|
_load_ts TIMESTAMP NOT NULL -- время загрузки (копия из STG)
|
||||||
);
|
);
|
||||||
|
|
||||||
CREATE TABLE ods.orders (
|
CREATE TABLE ods.orders (
|
||||||
@@ -98,16 +101,27 @@ CREATE TABLE dds.dim_product (
|
|||||||
|
|
||||||
-- dim_customer: измерение "Клиент" с SCD Type 2
|
-- dim_customer: измерение "Клиент" с SCD Type 2
|
||||||
CREATE TABLE dds.dim_customer (
|
CREATE TABLE dds.dim_customer (
|
||||||
customer_sk BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
|
customer_sk BIGSERIAL PRIMARY KEY,
|
||||||
customer_bk INT NOT NULL,
|
customer_bk INT NOT NULL, -- бизнес-ключ
|
||||||
email VARCHAR(100),
|
email TEXT,
|
||||||
phone VARCHAR(20),
|
phone TEXT,
|
||||||
city VARCHAR(50),
|
city TEXT,
|
||||||
valid_from DATE NOT NULL,
|
hashdiff TEXT NOT NULL, -- md5 по нормализованным атрибутам
|
||||||
valid_to DATE NOT NULL DEFAULT '9999-12-31',
|
valid_from TIMESTAMP NOT NULL,
|
||||||
is_current BOOLEAN NOT NULL DEFAULT TRUE
|
valid_to TIMESTAMP NOT NULL,
|
||||||
|
is_current BOOLEAN NOT NULL DEFAULT TRUE,
|
||||||
|
created_at TIMESTAMP NOT NULL DEFAULT NOW(),
|
||||||
|
updated_at TIMESTAMP NOT NULL DEFAULT NOW()
|
||||||
);
|
);
|
||||||
|
|
||||||
|
-- одна версия на момент времени
|
||||||
|
ALTER TABLE dds.dim_customer
|
||||||
|
ADD CONSTRAINT uq_dim_customer_bk_from UNIQUE (customer_bk, valid_from); -- dim_product: измерение "
|
||||||
|
|
||||||
|
-- ускорители
|
||||||
|
CREATE INDEX ix_dim_customer_bk_current ON dds.dim_customer (customer_bk) WHERE is_current;
|
||||||
|
CREATE INDEX ix_dim_customer_bk_from_to ON dds.dim_customer (customer_bk, valid_from, valid_to);
|
||||||
|
|
||||||
-- fact_sales: факт "Продажи"
|
-- fact_sales: факт "Продажи"
|
||||||
CREATE TABLE dds.fact_sales (
|
CREATE TABLE dds.fact_sales (
|
||||||
sale_id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
|
sale_id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
|
||||||
@@ -123,3 +137,16 @@ ALTER TABLE dds.fact_sales
|
|||||||
ADD CONSTRAINT fk_fact_customer FOREIGN KEY (customer_sk) REFERENCES dds.dim_customer(customer_sk),
|
ADD CONSTRAINT fk_fact_customer FOREIGN KEY (customer_sk) REFERENCES dds.dim_customer(customer_sk),
|
||||||
ADD CONSTRAINT fk_fact_product FOREIGN KEY (product_sk) REFERENCES dds.dim_product(product_sk),
|
ADD CONSTRAINT fk_fact_product FOREIGN KEY (product_sk) REFERENCES dds.dim_product(product_sk),
|
||||||
ADD CONSTRAINT fk_fact_date FOREIGN KEY (date_key) REFERENCES dds.dim_date(date_key);
|
ADD CONSTRAINT fk_fact_date FOREIGN KEY (date_key) REFERENCES dds.dim_date(date_key);
|
||||||
|
|
||||||
|
|
||||||
|
-- Через md5 по нормализованным атрибутам
|
||||||
|
CREATE OR REPLACE FUNCTION dds.customer_hash(email TEXT, phone TEXT, city TEXT)
|
||||||
|
RETURNS TEXT LANGUAGE sql IMMUTABLE AS $$
|
||||||
|
SELECT md5(
|
||||||
|
concat_ws('||',
|
||||||
|
lower(coalesce(trim(email), '')),
|
||||||
|
lower(coalesce(trim(phone), '')),
|
||||||
|
lower(coalesce(trim(city), ''))
|
||||||
|
)
|
||||||
|
);
|
||||||
|
$$;
|
||||||
|
|||||||
+70
-47
@@ -11,10 +11,13 @@ DELETE FROM stg.orders_raw;
|
|||||||
DELETE FROM stg.order_items_raw;
|
DELETE FROM stg.order_items_raw;
|
||||||
DELETE FROM stg.products_raw;
|
DELETE FROM stg.products_raw;
|
||||||
|
|
||||||
INSERT INTO stg.customers_raw (customer_id, email, phone, city) VALUES
|
-- STG (пример вставки с метками времени)
|
||||||
('101', 'a@ex.com', '700', 'Москва'),
|
INSERT INTO stg.customers_raw (_load_id, _load_ts, event_ts, customer_id, email, phone, city) VALUES
|
||||||
('101', 'b@ex.com', '700', 'Москва'),
|
('batch_20250405_0800', '2025-04-05 08:00', NULL, '101','a@ex.com','700','Москва'),
|
||||||
('102', 'c@ex.com', '701', 'СПб');
|
('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','Санкт-Петербург');
|
||||||
|
|
||||||
|
|
||||||
INSERT INTO stg.orders_raw (order_id, order_date, customer_id) VALUES
|
INSERT INTO stg.orders_raw (order_id, order_date, customer_id) VALUES
|
||||||
('5001', '2024-01-10', '101'),
|
('5001', '2024-01-10', '101'),
|
||||||
@@ -39,7 +42,9 @@ 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), ''),
|
||||||
|
_load_id,
|
||||||
|
_load_ts
|
||||||
FROM stg.customers_raw
|
FROM stg.customers_raw
|
||||||
WHERE customer_id ~ '^\d+$';
|
WHERE customer_id ~ '^\d+$';
|
||||||
|
|
||||||
@@ -107,53 +112,71 @@ 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 — SCD Type 2
|
||||||
-- В продакшене — используем алгоритм «детектирования изменений + UPSERT»
|
-- Загрузка SCD2 (идемпотентный апсерт)
|
||||||
-- Здесь: перестраиваем всю историю на основе STG (для детерминированности)
|
-- Делаем внутри одной транзакции
|
||||||
|
|
||||||
DELETE FROM dds.dim_customer;
|
BEGIN;
|
||||||
|
-- Подготовим дельту: по BK одна запись на «текущую истину» из ODS/STG
|
||||||
WITH ranked AS (
|
WITH src AS (
|
||||||
SELECT
|
SELECT
|
||||||
customer_id::INT AS customer_bk,
|
s.customer_id::INT AS customer_bk,
|
||||||
email,
|
NULLIF(trim(s.email), '') AS email,
|
||||||
phone,
|
NULLIF(trim(s.phone), '') AS phone,
|
||||||
city,
|
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 (
|
ROW_NUMBER() OVER (
|
||||||
PARTITION BY customer_id
|
PARTITION BY s.customer_id
|
||||||
ORDER BY _load_ts, email
|
ORDER BY COALESCE(s.event_ts, s._load_ts) DESC, s._load_ts DESC
|
||||||
) AS rn
|
) AS rn
|
||||||
FROM stg.customers_raw
|
FROM stg.customers_raw s
|
||||||
WHERE customer_id ~ '^\d+$'
|
WHERE s.customer_id ~ '^\d+$'
|
||||||
),
|
),
|
||||||
changes AS (
|
delta AS (
|
||||||
|
-- берём по BK последнюю версию из поступивших данных
|
||||||
|
SELECT customer_bk, email, phone, city, hashdiff, eff_ts
|
||||||
|
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
|
||||||
|
)
|
||||||
SELECT
|
SELECT
|
||||||
customer_bk,
|
s.customer_bk, s.email, s.phone, s.city, s.hashdiff,
|
||||||
email,
|
COALESCE( -- у новых BK открываем «историю с вечности»
|
||||||
phone,
|
-- если нужна «вечность» именно как TIMESTAMP-величина:
|
||||||
city,
|
CASE WHEN d.customer_bk IS NULL THEN TIMESTAMP '1900-01-01' ELSE s.eff_ts END,
|
||||||
-- Имитируем хронологию: +1 день на каждое изменение
|
TIMESTAMP '1900-01-01'
|
||||||
'2023-01-01'::DATE + (rn - 1) * INTERVAL '1 day' AS eff_from
|
) AS valid_from,
|
||||||
FROM ranked
|
TIMESTAMP '9999-12-31' AS valid_to,
|
||||||
),
|
TRUE AS is_current,
|
||||||
final AS (
|
NOW(), NOW()
|
||||||
SELECT
|
FROM delta s
|
||||||
customer_bk,
|
LEFT JOIN dds.dim_customer d
|
||||||
email,
|
ON d.customer_bk = s.customer_bk AND d.is_current = TRUE
|
||||||
phone,
|
WHERE d.customer_bk IS NULL -- новый BK
|
||||||
city,
|
OR d.hashdiff <> s.hashdiff -- или изменившийся BK
|
||||||
eff_from::DATE AS valid_from,
|
ON CONFLICT (customer_bk, valid_from) DO NOTHING;
|
||||||
COALESCE(
|
-- Конец транзакции
|
||||||
LEAD(eff_from) OVER (PARTITION BY customer_bk ORDER BY eff_from) - INTERVAL '1 day',
|
COMMIT;
|
||||||
'9999-12-31'::DATE
|
|
||||||
) AS valid_to,
|
|
||||||
CASE WHEN LEAD(eff_from) OVER (PARTITION BY customer_b_k ORDER BY eff_from) IS NULL
|
|
||||||
THEN TRUE ELSE FALSE END AS is_current
|
|
||||||
FROM changes
|
|
||||||
)
|
|
||||||
INSERT INTO dds.dim_customer (customer_bk, email, phone, city, valid_from, valid_to, is_current)
|
|
||||||
SELECT customer_bk, email, phone, city, valid_from, valid_to, is_current
|
|
||||||
FROM final;
|
|
||||||
|
|
||||||
-- 6. DDS: fact_sales — загрузка фактов с учётом SCD
|
-- 6. DDS: fact_sales — загрузка фактов с учётом SCD
|
||||||
-- В продакшене — фильтруем по диапазону дат (инкрементально)
|
-- В продакшене — фильтруем по диапазону дат (инкрементально)
|
||||||
|
|||||||
Reference in New Issue
Block a user