Исправление неточностей

This commit is contained in:
2025-11-29 20:46:03 +03:00
parent 218c6b9fe7
commit e659fbb731
6 changed files with 61 additions and 61 deletions
@@ -41,7 +41,7 @@ customer_id,status,event_ts,_load_id,load_ts
- `status` — статус клиента в CRM (`new`, `active`, `vip`, `churned`); - `status` — статус клиента в CRM (`new`, `active`, `vip`, `churned`);
- `event_ts` — момент, когда статус сменился в CRM; - `event_ts` — момент, когда статус сменился в CRM;
- `_load_id` — идентификатор батча загрузки; - `_load_id` — идентификатор батча загрузки;
- `load_ts` — момент, когда данные попали в DWH. - `load_ts` — момент, когда данные попали в DWH (в таблицах STG/ODS эта колонка будет называться `_load_ts`, но по смыслу это то же самое время загрузки).
Файл содержит несколько клиентов и несколько смен статуса по каждому — этого достаточно, чтобы отработать SCD2. Файл содержит несколько клиентов и несколько смен статуса по каждому — этого достаточно, чтобы отработать SCD2.
@@ -93,7 +93,7 @@ SELECT * FROM stg.customer_status_raw LIMIT 10;
- привести: - привести:
- `customer_id``INT`, - `customer_id``INT`,
- `status``VARCHAR(20)` (можно оставить как есть), - `status``VARCHAR(20)` (можно оставить как есть),
- `event_ts` и `load_ts``TIMESTAMP`; - `event_ts` и `load_ts``TIMESTAMP` (в DWH-таблицах эта колонка будет лежать как `_load_ts`);
- аккуратно обработать возможные пустые значения (если бы они были); - аккуратно обработать возможные пустые значения (если бы они были);
- заполнить `_load_id` и `_load_ts` в `ods.customer_status`. - заполнить `_load_id` и `_load_ts` в `ods.customer_status`.
@@ -129,14 +129,14 @@ ORDER BY customer_id, event_ts;
- `customer_bk`, - `customer_bk`,
- `status`, - `status`,
- `event_ts` (как «время начала действия статуса»), - `event_ts` (как «время начала действия статуса»),
- `hashdiff` (например, `md5(status)` или `dds.customer_hash(status)` по аналогии с `dim_customer`). - `hashdiff` (например, `md5(status)`; можно вынести расчёт в отдельную функцию по аналогии с `dds.customer_hash` для клиентов).
2. Для каждого клиента отсортируйте события по `event_ts` и с помощью `LEAD()` посчитайте: 2. Для каждого клиента отсортируйте события по `event_ts` и с помощью `LEAD()` посчитайте:
- `valid_from` — текущее `event_ts`, - `valid_from` — текущее `event_ts`,
- `valid_to` — следующее `event_ts - 1 second` (или `9999-12-31`, если следующего нет). - `valid_to` — следующее `event_ts - 1 second` (или `9999-12-31`, если следующего нет).
3. Вставьте получившиеся строки в `dds.dim_customer_status`: 3. Вставьте получившиеся строки в `dds.dim_customer_status`. У актуальной строки для каждого клиента задайте `is_current = TRUE` (например, там, где `valid_to = '9999-12-31'`), у остальных — `FALSE`:
```sql ```sql
INSERT INTO dds.dim_customer_status ( INSERT INTO dds.dim_customer_status (
@@ -238,4 +238,3 @@ ORDER BY date_actual, status;
- при желании — собрать простую витрину в `dm`. - при желании — собрать простую витрину в `dm`.
Если что‑то не получается — можно разбирать решения по шагам вместе с ментором: от простого `SELECT` из STG до полноценного SCD2 в DDS. Если что‑то не получается — можно разбирать решения по шагам вместе с ментором: от простого `SELECT` из STG до полноценного SCD2 в DDS.
+3 -3
View File
@@ -83,7 +83,7 @@ ALTER TABLE ods.products ADD PRIMARY KEY (product_id);
-- 4. DDS: интегрированная модель -- 4. DDS: интегрированная модель
-- dim_date: справочник дат (без первичного ключагенерируется) -- dim_date: справочник дат (ключ — суррогатный date_key)
CREATE TABLE dds.dim_date ( CREATE TABLE dds.dim_date (
date_key INT PRIMARY KEY, date_key INT PRIMARY KEY,
date_actual DATE NOT NULL, date_actual DATE NOT NULL,
@@ -118,9 +118,9 @@ CREATE TABLE dds.dim_customer (
updated_at TIMESTAMP NOT NULL DEFAULT NOW() updated_at TIMESTAMP NOT NULL DEFAULT NOW()
); );
-- одна версия на момент времени -- одна версия на момент времени (на одну пару BK+valid_from)
ALTER TABLE dds.dim_customer ALTER TABLE dds.dim_customer
ADD CONSTRAINT uq_dim_customer_bk_from UNIQUE (customer_bk, valid_from); -- dim_product: измерение " ADD CONSTRAINT uq_dim_customer_bk_from UNIQUE (customer_bk, valid_from);
-- ускорители -- ускорители
CREATE INDEX ix_dim_customer_bk_current ON dds.dim_customer (customer_bk) WHERE is_current; CREATE INDEX ix_dim_customer_bk_current ON dds.dim_customer (customer_bk) WHERE is_current;
+50 -47
View File
@@ -128,56 +128,59 @@ SELECT product_id, name
FROM ods.products; FROM ods.products;
-- 5. DDS: dim_customer — первичная загрузка SCD2 (full backfill из STG) -- 5. DDS: dim_customer — первичная загрузка SCD2 (full backfill из STG)
-- Здесь мы пересчитываем всю историю клиента из событий в STG. -- В ЭТОМ ДЕМО: dim_customer строится напрямую из stg.customers_raw, который играет роль
-- В реальном DWH такую полную перезагрузку делают редко; для инкремента см. 03_demo_increment.sql и SCD.md. -- устойчивого event-лога (все события по клиенту в одном месте).
-- Это удобно для учебной первичной загрузки (full backfill), когда мы один раз
-- восстанавливаем всю историю клиента.
-- В РЕАЛЬНОМ DWH: так делают редко. Исторические измерения обычно строят
-- поверх очищенных и нормализованных слоёв (ODS / PSA / Data Vault).
-- Для примера инкрементальной заливки SCD2 по снимку из ODS см. 03_demo_increment.sql и SCD.md.
TRUNCATE dds.dim_customer, dds.fact_sales; TRUNCATE dds.dim_customer, dds.fact_sales;
BEGIN; 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(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
),
changes AS (
-- только первые состояния и фактические изменения атрибутов
SELECT *
FROM ordered
WHERE prev_hash IS DISTINCT FROM hashdiff OR prev_hash IS NULL
),
framed AS (
SELECT
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 SELECT
customer_bk, email, phone, city, hashdiff, s.customer_id::INT AS customer_bk,
valid_from, NULLIF(trim(s.email), '') AS email,
COALESCE(next_ts - interval '1 second', timestamp '9999-12-31') AS valid_to, NULLIF(trim(s.phone), '') AS phone,
(next_ts IS NULL) AS is_current, NULLIF(trim(s.city), '') AS city,
now(), now() COALESCE(NULLIF(s.event_ts,'')::timestamp, s._load_ts,
FROM framed to_timestamp(regexp_replace(s._load_id,'^batch_',''),'YYYYMMDD_HH24MI')) AS eff_ts,
ORDER BY customer_bk, valid_from; dds.customer_hash(s.email, s.phone, s.city) AS hashdiff
COMMIT; 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
),
changes AS (
-- только первые состояния и фактические изменения атрибутов
SELECT *
FROM ordered
WHERE prev_hash IS DISTINCT FROM hashdiff OR prev_hash IS NULL
),
framed AS (
SELECT
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;
-- 6. DDS: fact_sales — загрузка фактов с учётом SCD -- 6. DDS: fact_sales — загрузка фактов с учётом SCD
-- В продакшене — фильтруем по диапазону дат (инкрементально) -- В продакшене — фильтруем по диапазону дат (инкрементально)
+1 -2
View File
@@ -1,6 +1,6 @@
-- =============================================== -- ===============================================
-- Проверки качества данных после загрузки STG→ODS→DDS -- Проверки качества данных после загрузки STG→ODS→DDS
-- Запускается после 02_dml.sql -- Запускается после 02_dml_stg-dds.sql (и, при необходимости, 03_demo_increment.sql)
-- =============================================== -- ===============================================
-- 1. Проверка: dim_customer не пуста -- 1. Проверка: dim_customer не пуста
@@ -40,4 +40,3 @@ BEGIN
FORMAT('ОШИБКА: у клиента 101 только %s версия, ожидается ≥2 (должна быть история)', version_count); FORMAT('ОШИБКА: у клиента 101 только %s версия, ожидается ≥2 (должна быть история)', version_count);
RAISE NOTICE '✅ SCD Type 2: клиент 101 имеет %s версий — история сохранена', version_count; RAISE NOTICE '✅ SCD Type 2: клиент 101 имеет %s версий — история сохранена', version_count;
END $$; END $$;
+1 -2
View File
@@ -34,7 +34,6 @@ GROUP BY d.date_actual, p.product_name,
-- Здесь — НЕ используем is_current! Нам нужна вся история для расчёта LTV -- Здесь — НЕ используем is_current! Нам нужна вся история для расчёта LTV
-- Для кого: CRM-менеджер, retention-спец -- Для кого: CRM-менеджер, retention-спец
-- Пример использования: «Найти клиентов с LTV > 250 ₽ и email из Москвы для email-рассылки» -- Пример использования: «Найти клиентов с LTV > 250 ₽ и email из Москвы для email-рассылки»
«Найти клиентов с LTV > 250 и email из Москвы для email-рассылки»
INSERT INTO dm.mart_customer_360 ( INSERT INTO dm.mart_customer_360 (
customer_bk, first_order_date, last_order_date, customer_bk, first_order_date, last_order_date,
total_orders, total_items, lifetime_value, total_orders, total_items, lifetime_value,
@@ -44,7 +43,7 @@ SELECT
c.customer_bk, c.customer_bk,
MIN(d.date_actual) AS first_order_date, MIN(d.date_actual) AS first_order_date,
MAX(d.date_actual) AS last_order_date, MAX(d.date_actual) AS last_order_date,
COUNT(DISTINCT f.sale_id) AS total_orders, COUNT(DISTINCT f.sale_id) AS total_orders, -- считаем строки факта (продажи), не бизнес-заказы
SUM(f.quantity) AS total_items, SUM(f.quantity) AS total_items,
SUM(f.amount) AS lifetime_value, SUM(f.amount) AS lifetime_value,
-- Берём email и город из самой свежей версии клиента -- Берём email и город из самой свежей версии клиента
@@ -8,7 +8,8 @@
-- 1) Загрузите CSV в stg.customer_status_raw (через COPY или \copy в psql). -- 1) Загрузите CSV в stg.customer_status_raw (через COPY или \copy в psql).
-- См. пример структуры файла в dwh-modeling/data/customer_status_events.csv -- См. пример структуры файла в dwh-modeling/data/customer_status_events.csv
-- 2) Переложите данные в ods.customer_status с приведением типов. -- 2) Переложите данные в ods.customer_status с приведением типов.
-- customer_id → INT, status → VARCHAR(20), event_ts / load_ts → TIMESTAMP. -- customer_id → INT, status → VARCHAR(20), event_ts / load_ts → TIMESTAMP
-- (в STG/ODS эта колонка будет жить как _load_ts).
-- 3) Постройте из ods.customer_status измерение dds.dim_customer_status в стиле SCD2: -- 3) Постройте из ods.customer_status измерение dds.dim_customer_status в стиле SCD2:
-- - одна строка на период действия статуса (valid_from / valid_to); -- - одна строка на период действия статуса (valid_from / valid_to);
-- - is_current = TRUE только у актуальной строки для клиента; -- - is_current = TRUE только у актуальной строки для клиента;
@@ -45,4 +46,3 @@
-- ); -- );
-- Идея: на каждую дату взять актуальный статус клиента -- Идея: на каждую дату взять актуальный статус клиента
-- через JOIN dds.dim_customer_status + dds.dim_date. -- через JOIN dds.dim_customer_status + dds.dim_date.