diff --git a/dwh-modeling/data/customers.csv b/dwh-modeling/data/customers.csv index 0cd27b5..84a16cd 100644 --- a/dwh-modeling/data/customers.csv +++ b/dwh-modeling/data/customers.csv @@ -1,4 +1,5 @@ -customer_id,email,phone,city -101,a@ex.com,700,Москва -101,b@ex.com,700,Москва -102,c@ex.com,701,СПб \ No newline at end of file +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 diff --git a/dwh-modeling/sql/01_ddl.sql b/dwh-modeling/sql/01_ddl.sql index 4c88eb8..ca60891 100644 --- a/dwh-modeling/sql/01_ddl.sql +++ b/dwh-modeling/sql/01_ddl.sql @@ -15,11 +15,12 @@ CREATE SCHEMA dds; -- 2. STG: сырые данные (как пришли) CREATE TABLE stg.customers_raw ( - customer_id TEXT, + customer_id TEXT, -- может быть строкой или числом email TEXT, phone TEXT, city TEXT, - _load_ts TIMESTAMP DEFAULT NOW() + _load_id TEXT, -- идентификатор загрузки (обязательно!) + _load_ts TIMESTAMP DEFAULT NOW(), -- время получения в DWH ); CREATE TABLE stg.orders_raw ( @@ -43,10 +44,12 @@ CREATE TABLE stg.products_raw ( -- 3. ODS: очищенные данные CREATE TABLE ods.customers ( - customer_id INT, + customer_id INT NOT NULL, -- привели к INT email VARCHAR(100), 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 ( @@ -98,16 +101,27 @@ CREATE TABLE dds.dim_product ( -- dim_customer: измерение "Клиент" с SCD Type 2 CREATE TABLE dds.dim_customer ( - customer_sk BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY, - customer_bk INT NOT NULL, - email VARCHAR(100), - phone VARCHAR(20), - city VARCHAR(50), - valid_from DATE NOT NULL, - valid_to DATE NOT NULL DEFAULT '9999-12-31', - is_current BOOLEAN NOT NULL DEFAULT TRUE + customer_sk BIGSERIAL PRIMARY KEY, + customer_bk INT NOT NULL, -- бизнес-ключ + email TEXT, + phone TEXT, + city TEXT, + hashdiff TEXT NOT NULL, -- md5 по нормализованным атрибутам + valid_from TIMESTAMP NOT NULL, + 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: факт "Продажи" CREATE TABLE dds.fact_sales ( 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_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); + + +-- Через 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), '')) + ) + ); +$$; diff --git a/dwh-modeling/sql/02_dim.sql b/dwh-modeling/sql/02_dim.sql index 450b4b8..5e9f4b3 100644 --- a/dwh-modeling/sql/02_dim.sql +++ b/dwh-modeling/sql/02_dim.sql @@ -11,10 +11,13 @@ DELETE FROM stg.orders_raw; DELETE FROM stg.order_items_raw; DELETE FROM stg.products_raw; -INSERT INTO stg.customers_raw (customer_id, email, phone, city) VALUES -('101', 'a@ex.com', '700', 'Москва'), -('101', 'b@ex.com', '700', 'Москва'), -('102', 'c@ex.com', '701', 'СПб'); +-- STG (пример вставки с метками времени) +INSERT INTO stg.customers_raw (_load_id, _load_ts, event_ts, customer_id, email, phone, city) VALUES +('batch_20250405_0800', '2025-04-05 08:00', NULL, '101','a@ex.com','700','Москва'), +('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 ('5001', '2024-01-10', '101'), @@ -39,7 +42,9 @@ SELECT customer_id::INT, NULLIF(TRIM(email), ''), NULLIF(TRIM(phone), ''), - NULLIF(TRIM(city), '') + NULLIF(TRIM(city), ''), + _load_id, + _load_ts FROM stg.customers_raw WHERE customer_id ~ '^\d+$'; @@ -107,53 +112,71 @@ INSERT INTO dds.dim_product (product_bk, product_name) SELECT product_id, name FROM ods.products; --- 5. DDS: dim_customer — SCD Type 2 (упрощённая версия для обучения) --- В продакшене — используем алгоритм «детектирования изменений + UPSERT» --- Здесь: перестраиваем всю историю на основе STG (для детерминированности) +-- 5. DDS: dim_customer — SCD Type 2 +-- Загрузка SCD2 (идемпотентный апсерт) +-- Делаем внутри одной транзакции -DELETE FROM dds.dim_customer; - -WITH ranked AS ( +BEGIN; + -- Подготовим дельту: по BK одна запись на «текущую истину» из ODS/STG + WITH src AS ( SELECT - customer_id::INT AS customer_bk, - email, - phone, - city, + 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 customer_id - ORDER BY _load_ts, email + 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 - WHERE customer_id ~ '^\d+$' -), -changes AS ( + FROM stg.customers_raw s + WHERE s.customer_id ~ '^\d+$' + ), + 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 - customer_bk, - email, - phone, - city, - -- Имитируем хронологию: +1 день на каждое изменение - '2023-01-01'::DATE + (rn - 1) * INTERVAL '1 day' AS eff_from - FROM ranked -), -final AS ( - SELECT - customer_bk, - email, - phone, - city, - eff_from::DATE AS valid_from, - COALESCE( - LEAD(eff_from) OVER (PARTITION BY customer_bk ORDER BY eff_from) - INTERVAL '1 day', - '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; + 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; + -- Конец транзакции +COMMIT; -- 6. DDS: fact_sales — загрузка фактов с учётом SCD -- В продакшене — фильтруем по диапазону дат (инкрементально)