merge(dwh-modeling): влита ветка dwh-modeling-basics с материалами по DWH
- Зачем: - объединить улучшенные учебные материалы по DWH-моделированию в основную ветку. - Что: - добавлен COMMIT_RULES.md с правилами оформления коммитов. - добавлено решение домашнего задания (09_dml_hw_customer_status_solution.sql). - улучшена документация: расширены разделы 3NF и Звезда, добавлены пояснения по ODS. - исправлены опечатки и SQL-скрипты по результатам ревью. - обновлены CSV-данные для корректной работы примеров. - Проверка: - git log --oneline -5. - просмотр изменённых файлов в dwh-modeling/.
This commit is contained in:
@@ -2,7 +2,7 @@
|
|||||||
|
|
||||||
## Project Structure & Module Organization
|
## Project Structure & Module Organization
|
||||||
- Root `README.md` describes the learning roadmap (RU).
|
- Root `README.md` describes the learning roadmap (RU).
|
||||||
- `dwh-modeling/` contains the article and demo DWH model; SQL lives in `dwh-modeling/sql` as ordered scripts `01_...sql`–`06_...sql`.
|
- `dwh-modeling/` contains the article and demo DWH model; SQL lives in `dwh-modeling/sql` as ordered scripts `01_...sql`–`09_...sql` (07–09 are homework DDL, template and solution).
|
||||||
- `postgres-bookings/` is a Dockerized PostgreSQL + demo “bookings” DB; start it first, then apply DWH scripts against the `demo` database.
|
- `postgres-bookings/` is a Dockerized PostgreSQL + demo “bookings” DB; start it first, then apply DWH scripts against the `demo` database.
|
||||||
|
|
||||||
## Build, Test, and Development Commands
|
## Build, Test, and Development Commands
|
||||||
@@ -26,8 +26,10 @@
|
|||||||
- For `postgres-bookings`, after modifications run `docker compose up -d && ./psql_sh` and verify simple queries such as `SELECT COUNT(*) FROM bookings.flights;`.
|
- For `postgres-bookings`, after modifications run `docker compose up -d && ./psql_sh` and verify simple queries such as `SELECT COUNT(*) FROM bookings.flights;`.
|
||||||
|
|
||||||
## Commit & Pull Request Guidelines
|
## Commit & Pull Request Guidelines
|
||||||
- Commit messages are short, imperative or descriptive phrases (often in Russian), e.g. `Добавлено оглавление`, `Переработка структуры`; group related edits into a single commit.
|
|
||||||
- Pull requests should focus on one topic, include a brief context, list of changes, and manual steps to reproduce or validate (commands you ran, expected results).
|
**Required:** Read [COMMIT_RULES.md](COMMIT_RULES.md) before making commits.
|
||||||
|
|
||||||
|
Pull requests should focus on one topic, include a brief context, list of changes, and manual steps to reproduce or validate (commands you ran, expected results).
|
||||||
|
|
||||||
## Security & Configuration Tips
|
## Security & Configuration Tips
|
||||||
- Do not commit personal `.env` files or credentials; use local overrides only.
|
- Do not commit personal `.env` files or credentials; use local overrides only.
|
||||||
|
|||||||
+220
@@ -0,0 +1,220 @@
|
|||||||
|
# Commit Rules
|
||||||
|
|
||||||
|
Unified commit style for all project contributors. Follows [Conventional Commits](https://www.conventionalcommits.org/) specification.
|
||||||
|
|
||||||
|
## Language
|
||||||
|
|
||||||
|
- **Primary language**: Russian
|
||||||
|
- If language is not specified, use Russian
|
||||||
|
- For AI-generated commits, Russian is mandatory unless task explicitly sets `lang:en`
|
||||||
|
- English is allowed only by explicit instruction (`lang:en`) or external collaboration requirements
|
||||||
|
- Do not mix languages in free-text parts of one commit message (subject + body + footer)
|
||||||
|
- Conventional Commit `type(scope)` stays in English
|
||||||
|
- Technical terms (PostgreSQL, SQL, DWH, DDL, DML) keep as-is
|
||||||
|
|
||||||
|
## Header Format
|
||||||
|
|
||||||
|
```
|
||||||
|
<type>(<scope>): <short description>
|
||||||
|
```
|
||||||
|
|
||||||
|
- Maximum header length: 72 characters
|
||||||
|
- For Russian subject, use result form (e.g. "добавлено", "исправлено", "обновлено")
|
||||||
|
- For English subject, use imperative present form (e.g. "add", "fix", "update")
|
||||||
|
- For English subject, do not use past forms (e.g. "added", "fixed", "updated")
|
||||||
|
- No trailing period
|
||||||
|
- Keep subject specific; avoid vague messages like "update", "fix bug", "changes"
|
||||||
|
|
||||||
|
### Allowed `type`
|
||||||
|
|
||||||
|
| Type | Description |
|
||||||
|
|------|-------------|
|
||||||
|
| `feat` | New feature |
|
||||||
|
| `fix` | Bug fix |
|
||||||
|
| `refactor` | Code restructuring without behavior change |
|
||||||
|
| `docs` | Documentation only |
|
||||||
|
| `test` | Tests, checks, validations |
|
||||||
|
| `chore` | Maintenance (configs, scripts, hooks) |
|
||||||
|
| `ci` | CI/CD changes |
|
||||||
|
| `perf` | Performance optimization |
|
||||||
|
| `revert` | Revert previous commit |
|
||||||
|
|
||||||
|
### Recommended `scope` for this repo
|
||||||
|
|
||||||
|
| Scope | Used for |
|
||||||
|
|-------|----------|
|
||||||
|
| `sql` | SQL scripts in `dwh-modeling/sql/` |
|
||||||
|
| `modeling` | DWH modeling docs, articles, schemas in `dwh-modeling/` |
|
||||||
|
| `bookings` | Docker PostgreSQL demo in `postgres-bookings/` |
|
||||||
|
| `docs` | Documentation, README, guides |
|
||||||
|
| `data` | Data files (CSV, fixtures) |
|
||||||
|
|
||||||
|
## Body Structure
|
||||||
|
|
||||||
|
For non-trivial changes, body is required. Use bullet points for readability.
|
||||||
|
|
||||||
|
Body is considered required when at least one condition is true:
|
||||||
|
- behavior or API/contract changed
|
||||||
|
- migration, rollback risk, or compatibility impact exists
|
||||||
|
- more than one meaningful file/module changed
|
||||||
|
- fix is non-obvious from header alone
|
||||||
|
|
||||||
|
### Multiline body in CLI (important)
|
||||||
|
|
||||||
|
- Do not pass body as one quoted string with `\n` (it will be stored literally).
|
||||||
|
- Use multiple `-m` flags, or `-F` with heredoc.
|
||||||
|
|
||||||
|
Correct:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
git commit \
|
||||||
|
-m "feat(sql): добавлена валидация данных для DWH" \
|
||||||
|
-m "- Зачем:
|
||||||
|
- нужна проверка целостности перед загрузкой
|
||||||
|
- Что:
|
||||||
|
- добавлен скрипт 04_validation.sql
|
||||||
|
- добавлены проверки на NULL и уникальность
|
||||||
|
- Проверка:
|
||||||
|
- psql -f dwh-modeling/sql/04_validation.sql"
|
||||||
|
```
|
||||||
|
|
||||||
|
Also correct:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
git commit -F- <<'MSG'
|
||||||
|
feat(sql): добавлена валидация данных для DWH
|
||||||
|
|
||||||
|
- Зачем:
|
||||||
|
- нужна проверка целостности перед загрузкой
|
||||||
|
- Что:
|
||||||
|
- добавлен скрипт 04_validation.sql
|
||||||
|
- добавлены проверки на NULL и уникальность
|
||||||
|
- Проверка:
|
||||||
|
- psql -f dwh-modeling/sql/04_validation.sql
|
||||||
|
MSG
|
||||||
|
```
|
||||||
|
|
||||||
|
### Template (Russian - default)
|
||||||
|
|
||||||
|
```
|
||||||
|
<type>(<scope>): <краткое описание результата>
|
||||||
|
|
||||||
|
- Зачем:
|
||||||
|
- причина изменения
|
||||||
|
- Что:
|
||||||
|
- ключевое изменение 1
|
||||||
|
- ключевое изменение 2
|
||||||
|
- Проверка:
|
||||||
|
- как проверено
|
||||||
|
```
|
||||||
|
|
||||||
|
### Template (English - only with `lang:en`)
|
||||||
|
|
||||||
|
```
|
||||||
|
<type>(<scope>): <short action description>
|
||||||
|
|
||||||
|
- Why:
|
||||||
|
- reason for change
|
||||||
|
- What:
|
||||||
|
- key change 1
|
||||||
|
- key change 2
|
||||||
|
- Check:
|
||||||
|
- how verified (command/test/smoke-check)
|
||||||
|
```
|
||||||
|
|
||||||
|
## Commit Scope Rules
|
||||||
|
|
||||||
|
- One commit = one logical task
|
||||||
|
- Don't mix feature changes with large refactoring
|
||||||
|
- Update docs in the same commit where behavior changes
|
||||||
|
|
||||||
|
## Breaking Changes
|
||||||
|
|
||||||
|
Use `!` in header for breaking changes:
|
||||||
|
```
|
||||||
|
feat(sql)!: rename stg_orders column contract
|
||||||
|
```
|
||||||
|
|
||||||
|
Add footer:
|
||||||
|
```
|
||||||
|
BREAKING CHANGE: column order_date renamed to created_at
|
||||||
|
```
|
||||||
|
|
||||||
|
## Examples
|
||||||
|
|
||||||
|
### Good examples
|
||||||
|
|
||||||
|
```
|
||||||
|
feat(sql): добавлен скрипт загрузки DM-слоя
|
||||||
|
|
||||||
|
- Зачем:
|
||||||
|
- нужны витрины для аналитики
|
||||||
|
- Что:
|
||||||
|
- добавлен 06_dml_dm.sql с загрузкой фактов и измерений
|
||||||
|
- добавлены индексы для оптимизации запросов
|
||||||
|
- Проверка:
|
||||||
|
- psql -f dwh-modeling/sql/06_dml_dm.sql
|
||||||
|
- SELECT COUNT(*) FROM dm.fact_orders;
|
||||||
|
```
|
||||||
|
|
||||||
|
```
|
||||||
|
fix(bookings): исправлен порт в docker-compose.yml
|
||||||
|
|
||||||
|
- Зачем:
|
||||||
|
- конфликт с локальным PostgreSQL на 5432
|
||||||
|
- Что:
|
||||||
|
- порт хоста изменен на 5433
|
||||||
|
- Проверка:
|
||||||
|
- docker compose up -d
|
||||||
|
- psql -h localhost -p 5433 -U postgres
|
||||||
|
```
|
||||||
|
|
||||||
|
```
|
||||||
|
docs(modeling): обновлена схема Data Vault после ревью
|
||||||
|
```
|
||||||
|
|
||||||
|
```
|
||||||
|
chore(docs): синхронизировано оглавление README
|
||||||
|
```
|
||||||
|
|
||||||
|
### Bad examples (don't do this)
|
||||||
|
|
||||||
|
```
|
||||||
|
❌ added sql script # no type, past tense
|
||||||
|
❌ feat: добавлен скрипт # no scope
|
||||||
|
❌ fix: исправлен баг # no scope, vague and non-actionable
|
||||||
|
❌ feat(sql): added new table # past tense in English subject
|
||||||
|
❌ feat(sql): add script and fix validation and update docs # multiple concerns
|
||||||
|
❌ feat(docs): add README и почини SQL # mixed languages in one message
|
||||||
|
```
|
||||||
|
|
||||||
|
## Quick Reference
|
||||||
|
|
||||||
|
```bash
|
||||||
|
# Feature
|
||||||
|
feat(scope): добавлена новая возможность
|
||||||
|
|
||||||
|
# Bug fix
|
||||||
|
fix(scope): исправлена проблема
|
||||||
|
|
||||||
|
# Documentation
|
||||||
|
docs(scope): обновлена документация
|
||||||
|
|
||||||
|
# Refactoring
|
||||||
|
refactor(scope): упрощена структура без изменения поведения
|
||||||
|
|
||||||
|
# Performance
|
||||||
|
perf(scope): ускорено выполнение
|
||||||
|
|
||||||
|
# Maintenance
|
||||||
|
chore(scope): обновлены служебные настройки
|
||||||
|
|
||||||
|
# Feature (lang:en)
|
||||||
|
feat(scope): add new capability
|
||||||
|
|
||||||
|
# Bug fix (lang:en)
|
||||||
|
fix(scope): correct response parsing
|
||||||
|
|
||||||
|
# Documentation (lang:en)
|
||||||
|
docs(scope): update setup guide
|
||||||
|
```
|
||||||
@@ -30,7 +30,7 @@
|
|||||||
Структура файла:
|
Структура файла:
|
||||||
|
|
||||||
```text
|
```text
|
||||||
customer_id,status,event_ts,_load_id,load_ts
|
customer_id,status,event_ts,_load_id,_load_ts
|
||||||
101,new,2024-01-01 09:00:00,batch_20240101_1000,2024-01-01 10:00:00
|
101,new,2024-01-01 09:00:00,batch_20240101_1000,2024-01-01 10:00:00
|
||||||
...
|
...
|
||||||
```
|
```
|
||||||
@@ -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 (в таблицах STG/ODS эта колонка будет называться `_load_ts`, но по смыслу это то же самое время загрузки).
|
- `_load_ts` — момент, когда данные попали в DWH.
|
||||||
|
|
||||||
Файл содержит несколько клиентов и несколько смен статуса по каждому — этого достаточно, чтобы отработать SCD2.
|
Файл содержит несколько клиентов и несколько смен статуса по каждому — этого достаточно, чтобы отработать SCD2.
|
||||||
|
|
||||||
@@ -62,11 +62,12 @@ customer_id,status,event_ts,_load_id,load_ts
|
|||||||
1. Поднимите demo‑Postgres по инструкции из корневого `README.md`.
|
1. Поднимите demo‑Postgres по инструкции из корневого `README.md`.
|
||||||
2. Выполните базовые скрипты DWH:
|
2. Выполните базовые скрипты DWH:
|
||||||
- `01_ddl_stg-dds.sql`
|
- `01_ddl_stg-dds.sql`
|
||||||
- `02_dml_stg-dds.sql`
|
- `02_dml_stg-dds.sql` (нужен как минимум для `dds.dim_date`)
|
||||||
|
- `05_ddl_dm.sql` (создаёт схему `dm` для витрин)
|
||||||
3. Выполните DDL для домашки:
|
3. Выполните DDL для домашки:
|
||||||
- `07_ddl_hw_customer_status.sql`
|
- `07_ddl_hw_customer_status.sql`
|
||||||
|
|
||||||
После этого схемы `stg`, `ods`, `dds` уже существуют, а дополнительные таблицы для статусов созданы.
|
После этого схемы `stg`, `ods`, `dds`, `dm` уже существуют, а дополнительные таблицы для статусов созданы.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -94,13 +95,13 @@ INSERT INTO stg.customer_status_raw (customer_id, status, event_ts, _load_id, _l
|
|||||||
('101','churned','2024-09-01 12:15:00','batch_20240901_1300','2024-09-01 13:00:00');
|
('101','churned','2024-09-01 12:15:00','batch_20240901_1300','2024-09-01 13:00:00');
|
||||||
```
|
```
|
||||||
|
|
||||||
> 💡 Здесь `_load_ts` — это время загрузки (в CSV оно называется `load_ts`).
|
> 💡 Здесь `_load_ts` — это время загрузки.
|
||||||
|
|
||||||
#### Вариант B: загрузить CSV
|
#### Вариант B: загрузить CSV
|
||||||
|
|
||||||
Можно загрузить файл `dwh-modeling/data/customer_status_events.csv` в таблицу `stg.customer_status_raw`:
|
Можно загрузить файл `dwh-modeling/data/customer_status_events.csv` в таблицу `stg.customer_status_raw`:
|
||||||
|
|
||||||
- **Через DBeaver**: Import Data → CSV → `stg.customer_status_raw` (колонку `load_ts` маппить в `_load_ts`).
|
- **Через DBeaver**: Import Data → CSV → `stg.customer_status_raw`.
|
||||||
- **Через `psql` в контейнере (`./psql_sh`)**: без установки `psql` на хост.
|
- **Через `psql` в контейнере (`./psql_sh`)**: без установки `psql` на хост.
|
||||||
|
|
||||||
Способ: передайте CSV в `psql` через STDIN и выполните `\copy ... FROM STDIN`:
|
Способ: передайте CSV в `psql` через STDIN и выполните `\copy ... FROM STDIN`:
|
||||||
@@ -121,12 +122,14 @@ SELECT * FROM stg.customer_status_raw LIMIT 10;
|
|||||||
|
|
||||||
### 3.2. ODS: очистка и типизация
|
### 3.2. ODS: очистка и типизация
|
||||||
|
|
||||||
|
> 💡 Обратите внимание: в основном примере `ods.customers` хранит **снимок** (одна строка на клиента, PK = `customer_id`), а здесь `ods.customer_status` хранит **все события** (PK = `customer_id + event_ts`). Это не ошибка, а сознательный выбор: источник данных о статусах - поток событий, и ODS сохраняет эту природу. Подробнее - в комментариях к решению.
|
||||||
|
|
||||||
В файле `08_dml_hw_customer_status_template.sql` найдите заготовку блока ODS и допишите SQL:
|
В файле `08_dml_hw_customer_status_template.sql` найдите заготовку блока ODS и допишите SQL:
|
||||||
|
|
||||||
- привести:
|
- привести:
|
||||||
- `customer_id` → `INT`,
|
- `customer_id` → `INT`,
|
||||||
- `status` → `VARCHAR(20)` (можно оставить как есть),
|
- `status` → `VARCHAR(20)` (можно оставить как есть),
|
||||||
- `event_ts` и `load_ts` → `TIMESTAMP` (в DWH-таблицах эта колонка будет лежать как `_load_ts`);
|
- `event_ts` и `_load_ts` → `TIMESTAMP`;
|
||||||
- аккуратно обработать возможные пустые значения (если бы они были);
|
- аккуратно обработать возможные пустые значения (если бы они были);
|
||||||
- заполнить `_load_id` и `_load_ts` в `ods.customer_status`.
|
- заполнить `_load_id` и `_load_ts` в `ods.customer_status`.
|
||||||
|
|
||||||
@@ -242,20 +245,7 @@ cat dwh-modeling/data/customer_status_events_increment.csv | ./postgres-bookings
|
|||||||
|
|
||||||
Опциональное задание для закрепления: собрать небольшую витрину с количеством клиентов по статусам на каждую дату.
|
Опциональное задание для закрепления: собрать небольшую витрину с количеством клиентов по статусам на каждую дату.
|
||||||
|
|
||||||
Перед началом убедитесь, что слой DM создан (схема `dm` и таблицы):
|
DDL витрины уже создан в `07_ddl_hw_customer_status.sql` (таблица `dm.mart_customer_status_daily`).
|
||||||
|
|
||||||
- выполните `dwh-modeling/sql/05_ddl_dm.sql` (один раз);
|
|
||||||
- затем можно собирать витрину.
|
|
||||||
|
|
||||||
Пример целевой таблицы:
|
|
||||||
|
|
||||||
```sql
|
|
||||||
CREATE TABLE dm.mart_customer_status_daily (
|
|
||||||
date_actual DATE NOT NULL,
|
|
||||||
status VARCHAR(20) NOT NULL,
|
|
||||||
customers_cnt INT NOT NULL
|
|
||||||
);
|
|
||||||
```
|
|
||||||
|
|
||||||
Идея:
|
Идея:
|
||||||
|
|
||||||
@@ -292,3 +282,11 @@ ORDER BY date_actual, status;
|
|||||||
- при желании — собрать простую витрину в `dm`.
|
- при желании — собрать простую витрину в `dm`.
|
||||||
|
|
||||||
Если что‑то не получается — можно разбирать решения по шагам вместе с ментором: от простого `SELECT` из STG до полноценного SCD2 в DDS.
|
Если что‑то не получается — можно разбирать решения по шагам вместе с ментором: от простого `SELECT` из STG до полноценного SCD2 в DDS.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 8. Эталонное решение
|
||||||
|
|
||||||
|
Когда выполните домашку и захотите сверить результат — готовое решение лежит в файле [`09_dml_hw_customer_status_solution.sql`](sql/09_dml_hw_customer_status_solution.sql).
|
||||||
|
|
||||||
|
Постарайтесь не подглядывать до того, как напишете свой вариант — основная ценность задания именно в самостоятельном разборе.
|
||||||
|
|||||||
+88
-93
@@ -180,7 +180,7 @@ flowchart TD
|
|||||||
|---------|------------|
|
|---------|------------|
|
||||||
| `dds.dim_customer` | Измерение «Клиент» с историей (SCD Type 2) |
|
| `dds.dim_customer` | Измерение «Клиент» с историей (SCD Type 2) |
|
||||||
| `dds.dim_product` | Измерение «Товар» |
|
| `dds.dim_product` | Измерение «Товар» |
|
||||||
| `dds.dim_date` | Готовый календарь на 10 лет вперёд (день/неделя/месяц/квартал) |
|
| `dds.dim_date` | Готовый календарь на 5 лет вперёд (день/неделя/месяц/квартал) |
|
||||||
| `dds.fact_sales` | Факт «Продажа» — строка заказа с суммой и количеством |
|
| `dds.fact_sales` | Факт «Продажа» — строка заказа с суммой и количеством |
|
||||||
|
|
||||||
💡 **Суррогатный ключ (Surrogate Key, SK)** — это `BIGINT`, который мы генерируем сами (например, `customer_sk = 1001`).
|
💡 **Суррогатный ключ (Surrogate Key, SK)** — это `BIGINT`, который мы генерируем сами (например, `customer_sk = 1001`).
|
||||||
@@ -281,7 +281,7 @@ AND (dim_customer.valid_to IS NULL OR fact_sales.order_date < dim_customer.valid
|
|||||||
```
|
```
|
||||||
и получаем актуальный на тот день email и город.
|
и получаем актуальный на тот день email и город.
|
||||||
|
|
||||||
> 🔍 Подробнее про SCD — в отдельной статье [Slow Changing Dimensions](SCD.md) (сравнение Type 1/2/3, паттерны обновления).
|
> 🔍 Подробнее про SCD — в отдельной статье [Slowly Changing Dimensions](SCD.md) (сравнение Type 1/2/3, паттерны обновления).
|
||||||
|
|
||||||
Теперь, когда мы разобрались, что такое факты, измерения и SCD, давайте посмотрим, как именно можно устроить слой DDS внутри — есть несколько вариантов.
|
Теперь, когда мы разобрались, что такое факты, измерения и SCD, давайте посмотрим, как именно можно устроить слой DDS внутри — есть несколько вариантов.
|
||||||
|
|
||||||
@@ -291,7 +291,7 @@ AND (dim_customer.valid_to IS NULL OR fact_sales.order_date < dim_customer.valid
|
|||||||
|
|
||||||
В DDS мы можем хранить данные по-разному. Это не «правильно/неправильно», а **выбор под задачу**.
|
В DDS мы можем хранить данные по-разному. Это не «правильно/неправильно», а **выбор под задачу**.
|
||||||
|
|
||||||
### 1. 3NF (третья нормальная форма)
|
### 1. 3NF (третья нормальная форма)
|
||||||
*Источник: Билл Инмон (Bill Inmon)*
|
*Источник: Билл Инмон (Bill Inmon)*
|
||||||
|
|
||||||
Если упростить, 3NF - это когда данные о разных бизнес-сущностях хранятся в отдельных таблицах и связываются ключами: клиент, заказ, город, регион, страна и т.д. Вместо одной большой таблицы с большим числом дублирующихся данных мы получаем цепочку таблиц, связанных ключами: `Заказ → Клиент → Город → Регион → Страна`. JOIN-ов становится больше, зато одно и то же свойство (например, название города) хранится в одном месте, а не дублируется в каждой строке заказа.
|
Если упростить, 3NF - это когда данные о разных бизнес-сущностях хранятся в отдельных таблицах и связываются ключами: клиент, заказ, город, регион, страна и т.д. Вместо одной большой таблицы с большим числом дублирующихся данных мы получаем цепочку таблиц, связанных ключами: `Заказ → Клиент → Город → Регион → Страна`. JOIN-ов становится больше, зато одно и то же свойство (например, название города) хранится в одном месте, а не дублируется в каждой строке заказа.
|
||||||
@@ -303,6 +303,53 @@ AND (dim_customer.valid_to IS NULL OR fact_sales.order_date < dim_customer.valid
|
|||||||
* Данные о сущностях разнесены по отдельным таблицам: клиент, заказ, продукт.
|
* Данные о сущностях разнесены по отдельным таблицам: клиент, заказ, продукт.
|
||||||
* Минимум дублирования: общие атрибуты хранятся в одном месте, таблицы связаны ключами.
|
* Минимум дублирования: общие атрибуты хранятся в одном месте, таблицы связаны ключами.
|
||||||
|
|
||||||
|
#### Как выглядел бы наш магазин в 3NF
|
||||||
|
|
||||||
|
В нашем примере `city` лежит прямо в `dim_customer`. В 3NF город стал бы отдельной таблицей, чтобы название хранилось в одном месте:
|
||||||
|
|
||||||
|
```mermaid
|
||||||
|
erDiagram
|
||||||
|
dim_city ||--o{ dim_customer : "город"
|
||||||
|
dim_customer ||--o{ fact_sales : "клиент"
|
||||||
|
|
||||||
|
dim_city {
|
||||||
|
int city_id PK
|
||||||
|
varchar city_name "Москва, СПб, ..."
|
||||||
|
}
|
||||||
|
|
||||||
|
dim_customer {
|
||||||
|
bigint customer_sk PK
|
||||||
|
int customer_bk
|
||||||
|
varchar email
|
||||||
|
int city_id FK "ссылка на dim_city"
|
||||||
|
date valid_from
|
||||||
|
date valid_to
|
||||||
|
}
|
||||||
|
|
||||||
|
fact_sales {
|
||||||
|
bigint sale_id PK
|
||||||
|
bigint customer_sk FK
|
||||||
|
int date_key FK
|
||||||
|
int quantity
|
||||||
|
decimal amount
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
Теперь, чтобы узнать **«сколько потратил клиент из Москвы за январь 2024?»**, нужно пройти по цепочке:
|
||||||
|
|
||||||
|
```sql
|
||||||
|
-- 3NF: три JOIN, чтобы добраться до города
|
||||||
|
SELECT SUM(f.amount)
|
||||||
|
FROM dds.fact_sales f
|
||||||
|
JOIN dds.dim_customer c ON f.customer_sk = c.customer_sk
|
||||||
|
JOIN dds.dim_city ct ON c.city_id = ct.city_id
|
||||||
|
JOIN dds.dim_date d ON f.date_key = d.date_key
|
||||||
|
WHERE ct.city_name = 'Москва'
|
||||||
|
AND d.year = 2024 AND d.month = 1;
|
||||||
|
```
|
||||||
|
|
||||||
|
Запрос читаемый, но JOIN-ов уже три - и это для простого вопроса. В реальном ядре цепочка может быть длиннее: `Клиент → Город → Регион → Страна`.
|
||||||
|
|
||||||
**Плюсы:**
|
**Плюсы:**
|
||||||
|
|
||||||
* Удобно поддерживать **единую «карту бизнеса»**: где живут «клиент», «заказ», «договор» и как они связаны;
|
* Удобно поддерживать **единую «карту бизнеса»**: где живут «клиент», «заказ», «договор» и как они связаны;
|
||||||
@@ -319,22 +366,36 @@ AND (dim_customer.valid_to IS NULL OR fact_sales.order_date < dim_customer.valid
|
|||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
### 2. Звезда (Star Schema)
|
### 2. Звезда (Star Schema)
|
||||||
*Источник: Ральф Кимболл (Ralph Kimball)*
|
*Источник: Ральф Кимболл (Ralph Kimball)*
|
||||||
|
|
||||||
Самый распространённый способ построения таблиц для слоя витрин.
|
Самый распространённый способ построения таблиц для слоя витрин. Именно эту модель мы использовали в [разделе 5](#5-базовые-понятия-факты-измерения-scd): схема `fact_sales` + `dim_date` / `dim_customer` / `dim_product` - это и есть Звезда.
|
||||||
|
|
||||||
Структура:
|
Структура:
|
||||||
|
|
||||||
* в центре - таблица фактов (события и метрики);
|
* в центре - таблица фактов (события и метрики);
|
||||||
* вокруг - измерения, обычно денормализованные («плоские») - широкие таблицы со всеми атрибутами сущности. Мы сознательно избегаем цепочек справочников ради простоты запросов.
|
* вокруг - измерения, обычно денормализованные («плоские») - широкие таблицы со всеми атрибутами сущности. Мы сознательно избегаем цепочек справочников ради простоты запросов.
|
||||||
|
|
||||||
Пример: `fact_sales` + `dim_date`, `dim_customer`, `dim_product`.
|
|
||||||
|
|
||||||

|

|
||||||
|
|
||||||
Методология Ральфа Кимбалла (Ralph Kimball) как раз делает упор на такие звёздные схемы: витрины, которые максимально просты для чтения и понятны аналитикам и BI-инструментам.
|
Методология Ральфа Кимбалла (Ralph Kimball) как раз делает упор на такие звёздные схемы: витрины, которые максимально просты для чтения и понятны аналитикам и BI-инструментам.
|
||||||
|
|
||||||
|
#### Тот же вопрос - в Звезде
|
||||||
|
|
||||||
|
В Звезде `city` лежит прямо в `dim_customer` (денормализовано). Тот же отчёт выглядит проще:
|
||||||
|
|
||||||
|
```sql
|
||||||
|
-- Звезда: два JOIN, город - прямо в измерении
|
||||||
|
SELECT SUM(f.amount)
|
||||||
|
FROM dds.fact_sales f
|
||||||
|
JOIN dds.dim_customer c ON f.customer_sk = c.customer_sk
|
||||||
|
JOIN dds.dim_date d ON f.date_key = d.date_key
|
||||||
|
WHERE c.city = 'Москва'
|
||||||
|
AND d.year = 2024 AND d.month = 1;
|
||||||
|
```
|
||||||
|
|
||||||
|
На один JOIN меньше, и не нужно знать, где именно хранится город: он лежит прямо в карточке клиента. Для аналитика или BI-инструмента это большая разница.
|
||||||
|
|
||||||
**Плюсы:**
|
**Плюсы:**
|
||||||
|
|
||||||
* проста для понимания: аналитикам и BI-инструментам удобно работать с такой моделью;
|
* проста для понимания: аналитикам и BI-инструментам удобно работать с такой моделью;
|
||||||
@@ -343,7 +404,7 @@ AND (dim_customer.valid_to IS NULL OR fact_sales.order_date < dim_customer.valid
|
|||||||
|
|
||||||
**Минусы:**
|
**Минусы:**
|
||||||
|
|
||||||
* измерения денормализованы, поэтому атрибуты дублируются (например, регион повторяется у всех клиентов региона);
|
* измерения денормализованы, поэтому атрибуты дублируются (например, название города повторяется у всех клиентов из этого города);
|
||||||
* изменения атрибутов могут требовать обновлять много строк в измерении.
|
* изменения атрибутов могут требовать обновлять много строк в измерении.
|
||||||
|
|
||||||
|
|
||||||
@@ -374,82 +435,9 @@ AND (dim_customer.valid_to IS NULL OR fact_sales.order_date < dim_customer.valid
|
|||||||
- новые источники проще прикручивать;
|
- новые источники проще прикручивать;
|
||||||
- меньше шансов «сломать» старые отчёты.
|
- меньше шансов «сломать» старые отчёты.
|
||||||
|
|
||||||
---
|
Data Vault хорошо подходит там, где много разнородных источников, нужна полная история изменений и прозрачный аудит. За гибкость приходится платить сложностью модели и количеством таблиц — поэтому для небольших проектов (2–5 источников, маленькая команда) DV почти наверняка избыточен.
|
||||||
|
|
||||||
#### Чем DV отличается от 3NF и Звезды
|
> 🔍 Подробнее про Data Vault — сравнение с 3NF/Звездой, Raw и Business Vault, когда внедрять — в отдельной статье [DataVault: как пережить бурную жизнь источников](DataVault.md).
|
||||||
|
|
||||||
Если сильно упростить:
|
|
||||||
|
|
||||||
- В **3NF/Звезде** мы часто смешиваем:
|
|
||||||
- бизнес-ключ,
|
|
||||||
- текущие атрибуты,
|
|
||||||
- историю (SCD2)
|
|
||||||
— всё это в одной таблице измерения.
|
|
||||||
|
|
||||||
- В **Data Vault** это *разнесено*:
|
|
||||||
- Hub — только бизнес-ключ;
|
|
||||||
- Satellite — только атрибуты + история;
|
|
||||||
- Link — только связи между сущностями.
|
|
||||||
|
|
||||||
За это приходится платить сложностью модели и количеством таблиц. Зато DV хорошо выдерживает:
|
|
||||||
- много разнородных источников;
|
|
||||||
- «грязные» данные;
|
|
||||||
- жёсткие требования по аудиту и трассировке.
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
#### Raw Vault и Business Vault — два слоя
|
|
||||||
|
|
||||||
Часто говорят «Raw Vault» и «Business Vault». Грубо:
|
|
||||||
|
|
||||||
- **Raw Vault** — «как прилетело из источников».
|
|
||||||
Хабы, линкы и сателлиты, максимально близкие к исходным данным.
|
|
||||||
Задача: надёжно собрать и сохранить **полную историю**.
|
|
||||||
|
|
||||||
- **Business Vault** — «как удобно считать дальше».
|
|
||||||
На основе Raw Vault появляются:
|
|
||||||
- служебные таблицы (PIT, Bridge и т.п.),
|
|
||||||
- подготовленные представления под витрины и отчёты,
|
|
||||||
- бизнес-правила (например, что считать «активным клиентом»).
|
|
||||||
|
|
||||||
Дальше поверх этого уже строятся **обычные витрины в формате Звезды**, с которыми работают аналитики.
|
|
||||||
|
|
||||||
Если примерить это к классическим слоям `stg → ods → dds → dm`, то **очень грубо** можно думать так:
|
|
||||||
|
|
||||||
- `stg` всё равно остаётся как «приземление» (landing) из источников;
|
|
||||||
- **Raw Vault** по духу ближе к **ODS**: мало бизнес-логики, зато полная история и интеграция из разных систем;
|
|
||||||
- **Business Vault** ближе к **DDS**: здесь уже живут бизнес-правила и подготовка данных к витринам;
|
|
||||||
- `dm` по-прежнему остаётся витринами в формате Звезды, с которыми работают аналитики и BI.
|
|
||||||
|
|
||||||
Важно: это именно *аналогия для понимания*, а не жёсткое правило проектирования.
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
#### Когда DV вам, скорее всего, рано
|
|
||||||
|
|
||||||
Если у вас:
|
|
||||||
|
|
||||||
- 2–5 источников,
|
|
||||||
- небольшая команда (1–2 инженера + аналитик),
|
|
||||||
- задачи уровня «сделать первые отчёты»,
|
|
||||||
|
|
||||||
то **Data Vault почти наверняка избыточен**.
|
|
||||||
Чаще всего хватает связки:
|
|
||||||
|
|
||||||
> `stg → ods → dds (3NF или простая Звезда с SCD2) → dm (Звезда)`
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
#### Что важно запомнить из этой статьи
|
|
||||||
|
|
||||||
Для этой статьи достаточно:
|
|
||||||
|
|
||||||
- знать, что **Data Vault** — это способ строить хранилище как **конструктор из Hub/Link/Satellite**,
|
|
||||||
- понимать, что он нужен в первую очередь там, где:
|
|
||||||
- много систем-источников,
|
|
||||||
- нужна *полная* история и прозрачный аудит.
|
|
||||||
|
|
||||||
Детали (Raw vs Business Vault, PIT/Bridge, DV 1.0 vs 2.0 и т.п.) — это уже тема для отдельной, взрослой статьи.
|
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -556,11 +544,17 @@ flowchart TD
|
|||||||
|
|
||||||
```sql
|
```sql
|
||||||
-- mart_daily_sales: ежедневные продажи с сегментацией
|
-- mart_daily_sales: ежедневные продажи с сегментацией
|
||||||
CREATE MATERIALIZED VIEW dm.mart_daily_sales AS
|
-- Полная пересборка (full refresh) - для простоты; в продакшене бывает incremental.
|
||||||
|
TRUNCATE dm.mart_daily_sales;
|
||||||
|
|
||||||
|
INSERT INTO dm.mart_daily_sales (
|
||||||
|
date_actual, product_name, customer_segment, total_qty, total_revenue
|
||||||
|
)
|
||||||
SELECT
|
SELECT
|
||||||
d.date_actual AS order_date,
|
d.date_actual,
|
||||||
p.product_name,
|
p.product_name,
|
||||||
c.customer_segment, -- например: 'Premium', 'Basic'
|
-- Сегмент определяем по сумме строки (в реальности может быть атрибутом клиента)
|
||||||
|
CASE WHEN f.amount >= 200 THEN 'Premium' ELSE 'Basic' END AS customer_segment,
|
||||||
SUM(f.quantity) AS total_qty,
|
SUM(f.quantity) AS total_qty,
|
||||||
SUM(f.amount) AS total_revenue
|
SUM(f.amount) AS total_revenue
|
||||||
FROM dds.fact_sales f
|
FROM dds.fact_sales f
|
||||||
@@ -569,13 +563,14 @@ JOIN dds.dim_date d
|
|||||||
JOIN dds.dim_product p
|
JOIN dds.dim_product p
|
||||||
ON f.product_sk = p.product_sk
|
ON f.product_sk = p.product_sk
|
||||||
JOIN dds.dim_customer c
|
JOIN dds.dim_customer c
|
||||||
ON f.customer_sk = c.customer_sk
|
ON f.customer_sk = c.customer_sk -- факт ссылается на нужную версию SK
|
||||||
AND f.order_date >= c.valid_from
|
GROUP BY d.date_actual, p.product_name,
|
||||||
AND (c.valid_to IS NULL OR f.order_date < c.valid_to) -- SCD!
|
CASE WHEN f.amount >= 200 THEN 'Premium' ELSE 'Basic' END;
|
||||||
GROUP BY d.date_actual, p.product_name, c.customer_segment;
|
|
||||||
```
|
```
|
||||||
|
|
||||||
> 💡 **Материализованное представление (MATERIALIZED VIEW)** — это «кэш» результата. Обновляется по расписанию (например, ночью).
|
> 💡 В продакшене витрину иногда оформляют как **MATERIALIZED VIEW** - «кэш» результата запроса, который обновляется по расписанию. В нашем примере используем обычную таблицу с `TRUNCATE` + `INSERT` - для учебных целей это нагляднее.
|
||||||
|
|
||||||
|
✏️ **Попробуйте сами:** [Домашка: статусы клиента от STG до DDS (и немного DM)](Homework_Customer_Status_DDS_DM.md) — пройдёте тот же путь, но самостоятельно.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
@@ -628,7 +623,7 @@ GROUP BY d.date_actual, p.product_name, c.customer_segment;
|
|||||||
Почему: ему нужны готовые метрики без сложных JOIN’ов. Звезда даёт понятные таблицы: «продажи по дням и товарам» — без углубления в атомарные сущности.
|
Почему: ему нужны готовые метрики без сложных JOIN’ов. Звезда даёт понятные таблицы: «продажи по дням и товарам» — без углубления в атомарные сущности.
|
||||||
|
|
||||||
— **BI-разработчик в Power BI / Tableau** → **Звезда**
|
— **BI-разработчик в Power BI / Tableau** → **Звезда**
|
||||||
Почему: все инструменты визуализации оптимизированы под star schema. Один факт + несколько измерений = быстрые отчёты и простою модель.
|
Почему: все инструменты визуализации оптимизированы под star schema. Один факт + несколько измерений = быстрые отчёты и простую модель.
|
||||||
|
|
||||||
— **Инженер ML (Data Scientist / ML-инженер)** → **3NF или сырые ODS-таблицы**
|
— **Инженер ML (Data Scientist / ML-инженер)** → **3NF или сырые ODS-таблицы**
|
||||||
Почему: для фичей нужны атомарные события и детальные атрибуты. Машинное обучение ценит полноту и детализацию данных больше, чем удобство отчётов.
|
Почему: для фичей нужны атомарные события и детальные атрибуты. Машинное обучение ценит полноту и детализацию данных больше, чем удобство отчётов.
|
||||||
@@ -791,7 +786,7 @@ SELECT 'OK' WHERE EXISTS (
|
|||||||
|
|
||||||
[`customers.csv`](data/customers.csv):
|
[`customers.csv`](data/customers.csv):
|
||||||
```csv
|
```csv
|
||||||
customer_id,email,phone,city,event_ts,_load_id,load_ts
|
customer_id,email,phone,city,event_ts,_load_id,_load_ts
|
||||||
101,a@ex.com,700,Москва,2024-01-01,batch_20240101_0800,2024-01-01 08:00
|
101,a@ex.com,700,Москва,2024-01-01,batch_20240101_0800,2024-01-01 08:00
|
||||||
102,c@ex.com,701,СПб,2024-01-01,batch_20240101_0800,2024-01-01 08:00
|
102,c@ex.com,701,СПб,2024-01-01,batch_20240101_0800,2024-01-01 08:00
|
||||||
101,b@ex.com,700,Москва,2024-05-16,batch_20240516_0800,2024-05-16 08:00
|
101,b@ex.com,700,Москва,2024-05-16,batch_20240516_0800,2024-05-16 08:00
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
customer_id,status,event_ts,_load_id,load_ts
|
customer_id,status,event_ts,_load_id,_load_ts
|
||||||
101,new,2024-01-01 09:00:00,batch_20240101_1000,2024-01-01 10:00:00
|
101,new,2024-01-01 09:00:00,batch_20240101_1000,2024-01-01 10:00:00
|
||||||
101,active,2024-02-15 10:30:00,batch_20240215_1100,2024-02-15 11:00:00
|
101,active,2024-02-15 10:30:00,batch_20240215_1100,2024-02-15 11:00:00
|
||||||
101,vip,2024-05-10 11:00:00,batch_20240510_1200,2024-05-10 12:00:00
|
101,vip,2024-05-10 11:00:00,batch_20240510_1200,2024-05-10 12:00:00
|
||||||
|
|||||||
|
@@ -1,4 +1,4 @@
|
|||||||
customer_id,status,event_ts,_load_id,load_ts
|
customer_id,status,event_ts,_load_id,_load_ts
|
||||||
101,active,2024-11-15 09:00:00,batch_20241115_1000,2024-11-15 10:00:00
|
101,active,2024-11-15 09:00:00,batch_20241115_1000,2024-11-15 10:00:00
|
||||||
102,active,2024-05-05 09:30:00,batch_20240505_1000,2024-05-05 10:00:00
|
102,active,2024-05-05 09:30:00,batch_20240505_1000,2024-05-05 10:00:00
|
||||||
103,active,2024-03-20 12:00:00,batch_20240320_1300,2024-03-20 13:00:00
|
103,active,2024-03-20 12:00:00,batch_20240320_1300,2024-03-20 13:00:00
|
||||||
|
|||||||
|
@@ -1,4 +1,4 @@
|
|||||||
customer_id,email,phone,city,event_ts,_load_id,load_ts
|
customer_id,email,phone,city,event_ts,_load_id,_load_ts
|
||||||
101,a@ex.com,700,Москва,2024-01-01,batch_20240101_0800,2024-01-01 08:00
|
101,a@ex.com,700,Москва,2024-01-01,batch_20240101_0800,2024-01-01 08:00
|
||||||
102,c@ex.com,701,СПб,2024-01-01,batch_20240101_0800,2024-01-01 08:00
|
102,c@ex.com,701,СПб,2024-01-01,batch_20240101_0800,2024-01-01 08:00
|
||||||
101,b@ex.com,700,Москва,2024-05-16,batch_20240516_0800,2024-05-16 08:00
|
101,b@ex.com,700,Москва,2024-05-16,batch_20240516_0800,2024-05-16 08:00
|
||||||
|
|||||||
|
@@ -2,6 +2,13 @@
|
|||||||
-- DML-скрипт: загрузка и трансформация данных
|
-- DML-скрипт: загрузка и трансформация данных
|
||||||
-- Запускается ПОВТОРНО при каждой загрузке (идемпотентно!)
|
-- Запускается ПОВТОРНО при каждой загрузке (идемпотентно!)
|
||||||
-- ===============================================
|
-- ===============================================
|
||||||
|
--
|
||||||
|
-- Карта загрузки в этом скрипте:
|
||||||
|
-- STG → ODS: orders, order_items, products, customers (снимок: одна строка на BK)
|
||||||
|
-- STG → DDS: dim_customer (SCD2, full backfill напрямую из STG — см. комментарий к п.5)
|
||||||
|
-- ODS → DDS: dim_product, fact_sales
|
||||||
|
-- отдельно: dim_date (генерация календаря)
|
||||||
|
--
|
||||||
|
|
||||||
-- 1. STG: имитация загрузки из источников (в реальности — COPY или INSERT из Kafka/NiFi)
|
-- 1. STG: имитация загрузки из источников (в реальности — COPY или INSERT из Kafka/NiFi)
|
||||||
-- ⚠️ В продакшене STG часто очищается перед загрузкой (TRUNCATE), либо используется партицирование по дате
|
-- ⚠️ В продакшене STG часто очищается перед загрузкой (TRUNCATE), либо используется партицирование по дате
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ CREATE TABLE dm.mart_customer_360 (
|
|||||||
customer_bk INT NOT NULL,
|
customer_bk INT NOT NULL,
|
||||||
first_order_date DATE,
|
first_order_date DATE,
|
||||||
last_order_date DATE,
|
last_order_date DATE,
|
||||||
total_orders INT NOT NULL,
|
total_line_items INT NOT NULL,
|
||||||
total_items INT NOT NULL,
|
total_items INT NOT NULL,
|
||||||
lifetime_value NUMERIC(18,2) NOT NULL,
|
lifetime_value NUMERIC(18,2) NOT NULL,
|
||||||
last_email VARCHAR(100),
|
last_email VARCHAR(100),
|
||||||
|
|||||||
@@ -33,14 +33,14 @@ GROUP BY d.date_actual, p.product_name,
|
|||||||
-- Считаем суммы по всей истории его покупок
|
-- Считаем суммы по всей истории его покупок
|
||||||
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_line_items, total_items, lifetime_value,
|
||||||
last_email, last_city
|
last_email, last_city
|
||||||
)
|
)
|
||||||
SELECT
|
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_line_items, -- строки факта (позиции продаж), не бизнес-заказы
|
||||||
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 и город клиента
|
||||||
|
|||||||
@@ -46,3 +46,11 @@ ALTER TABLE dds.dim_customer_status
|
|||||||
CREATE INDEX ix_dim_customer_status_bk_current
|
CREATE INDEX ix_dim_customer_status_bk_current
|
||||||
ON dds.dim_customer_status (customer_bk)
|
ON dds.dim_customer_status (customer_bk)
|
||||||
WHERE valid_to IS NULL;
|
WHERE valid_to IS NULL;
|
||||||
|
|
||||||
|
-- 4. DM: витрина статусов клиентов по датам (опциональная часть домашки)
|
||||||
|
DROP TABLE IF EXISTS dm.mart_customer_status_daily;
|
||||||
|
CREATE TABLE dm.mart_customer_status_daily (
|
||||||
|
date_actual DATE NOT NULL,
|
||||||
|
status VARCHAR(20) NOT NULL,
|
||||||
|
customers_cnt INT NOT NULL
|
||||||
|
);
|
||||||
|
|||||||
@@ -38,11 +38,6 @@
|
|||||||
-- Можно ориентироваться на примеры в 03_demo_increment.sql.
|
-- Можно ориентироваться на примеры в 03_demo_increment.sql.
|
||||||
|
|
||||||
-- 4. DM: витрина статусов клиентов по датам (по желанию)
|
-- 4. DM: витрина статусов клиентов по датам (по желанию)
|
||||||
-- Пример целевой структуры:
|
-- DDL витрины уже создан в 07_ddl_hw_customer_status.sql (dm.mart_customer_status_daily).
|
||||||
-- CREATE TABLE dm.mart_customer_status_daily (
|
|
||||||
-- date_actual DATE NOT NULL,
|
|
||||||
-- status VARCHAR(20) NOT NULL,
|
|
||||||
-- customers_cnt INT NOT NULL
|
|
||||||
-- );
|
|
||||||
-- Идея: на каждую дату взять актуальный статус клиента
|
-- Идея: на каждую дату взять актуальный статус клиента
|
||||||
-- через JOIN dds.dim_customer_status + dds.dim_date.
|
-- через JOIN dds.dim_customer_status + dds.dim_date.
|
||||||
|
|||||||
@@ -0,0 +1,340 @@
|
|||||||
|
-- ===============================================
|
||||||
|
-- 09_dml_hw_customer_status_solution.sql
|
||||||
|
-- Решение домашки: статусы клиента (STG -> ODS -> DDS SCD2 -> DM)
|
||||||
|
--
|
||||||
|
-- Это ЭТАЛОННОЕ РЕШЕНИЕ. Если вы ещё не пробовали решить домашку сами -
|
||||||
|
-- вернитесь к заданию (Homework_Customer_Status_DDS_DM.md) и шаблону (08_dml_hw_customer_status_template.sql).
|
||||||
|
-- Основная ценность задания - в самостоятельном разборе.
|
||||||
|
--
|
||||||
|
-- Что делает этот файл:
|
||||||
|
-- 1) Перекладывает события статусов в ODS (приводит типы, чистит пустое).
|
||||||
|
-- 2) Строит DDS-измерение со "встроенной историей" (SCD2): периоды valid_from/valid_to.
|
||||||
|
-- 3) Загружает инкрементальную порцию событий и обновляет ODS + DDS.
|
||||||
|
-- 4) Собирает простую витрину в DM: сколько клиентов в каком статусе по дням.
|
||||||
|
--
|
||||||
|
-- Как запускать:
|
||||||
|
-- - для первого знакомства можно запускать файл целиком;
|
||||||
|
-- - если хотите потренировать инкремент (п.3): добавьте свои события -> запустите блок 3 ещё раз.
|
||||||
|
--
|
||||||
|
-- Важно:
|
||||||
|
-- - здесь часто используется TRUNCATE (полная очистка), чтобы было легко повторять домашку;
|
||||||
|
-- - в реальном DWH так делают не всегда, но для обучения это удобнее.
|
||||||
|
--
|
||||||
|
-- Предусловия (DDL + данные в STG):
|
||||||
|
-- 1) dwh-modeling/sql/01_ddl_stg-dds.sql
|
||||||
|
-- 2) dwh-modeling/sql/02_dml_stg-dds.sql (нужен dim_date)
|
||||||
|
-- 3) dwh-modeling/sql/05_ddl_dm.sql
|
||||||
|
-- 4) dwh-modeling/sql/07_ddl_hw_customer_status.sql
|
||||||
|
-- 5) stg.customer_status_raw заполнена (см. dwh-modeling/Homework_Customer_Status_DDS_DM.md)
|
||||||
|
-- ===============================================
|
||||||
|
|
||||||
|
-- ==========================================================
|
||||||
|
-- 1) ODS: очистка и типизация (full refresh)
|
||||||
|
-- ==========================================================
|
||||||
|
|
||||||
|
-- Идея:
|
||||||
|
-- - STG хранит "как пришло" (обычно TEXT);
|
||||||
|
-- - ODS хранит "аккуратно": правильные типы + простая чистка.
|
||||||
|
-- Для простоты пересобираем ODS с нуля.
|
||||||
|
--
|
||||||
|
-- Обратите внимание: ods.customer_status хранит ВСЕ события (PK = customer_id + event_ts),
|
||||||
|
-- а не только последнее состояние, как ods.customers (PK = customer_id).
|
||||||
|
-- Причина: источник данных здесь - поток событий ("статус стал X в момент Y"),
|
||||||
|
-- а не снимок ("вот текущие данные клиента"). ODS сохраняет природу источника:
|
||||||
|
-- снимок остаётся снимком, события остаются событиями.
|
||||||
|
-- Благодаря этому full backfill SCD2 (блок 2) строится прямо из ODS, а не из STG.
|
||||||
|
|
||||||
|
TRUNCATE ods.customer_status;
|
||||||
|
|
||||||
|
INSERT INTO ods.customer_status (
|
||||||
|
customer_id, status, event_ts, _load_id, _load_ts
|
||||||
|
)
|
||||||
|
SELECT
|
||||||
|
s.customer_id::INT AS customer_id,
|
||||||
|
NULLIF(trim(s.status), '') AS status,
|
||||||
|
NULLIF(trim(s.event_ts), '')::TIMESTAMP AS event_ts,
|
||||||
|
s._load_id,
|
||||||
|
COALESCE(s._load_ts, now()) AS _load_ts
|
||||||
|
FROM stg.customer_status_raw s
|
||||||
|
WHERE s.customer_id ~ '^\d+$'
|
||||||
|
AND NULLIF(trim(s.event_ts), '') IS NOT NULL
|
||||||
|
AND NULLIF(trim(s.status), '') IS NOT NULL;
|
||||||
|
|
||||||
|
-- Проверка: что получилось в ODS
|
||||||
|
SELECT 'ods.customer_status count = ' || COUNT(*) FROM ods.customer_status;
|
||||||
|
SELECT * FROM ods.customer_status ORDER BY customer_id, event_ts;
|
||||||
|
|
||||||
|
-- ==========================================================
|
||||||
|
-- 2) DDS: начальная загрузка SCD2 (full refresh)
|
||||||
|
-- ==========================================================
|
||||||
|
|
||||||
|
-- Идея SCD2 простыми словами:
|
||||||
|
-- - одна строка = один период, когда статус был одним и тем же;
|
||||||
|
-- - valid_from = с какого дня статус "начался";
|
||||||
|
-- - valid_to = с какого дня статус "закончился" (NULL = текущий статус);
|
||||||
|
-- - интервалы считаем так: [valid_from, valid_to) (valid_to не включаем).
|
||||||
|
-- - чтобы найти статус "на дату D":
|
||||||
|
-- D >= valid_from AND (valid_to IS NULL OR D < valid_to)
|
||||||
|
--
|
||||||
|
-- Упрощение для домашки:
|
||||||
|
-- - считаем, что у клиента нет двух разных смен статуса в один день.
|
||||||
|
|
||||||
|
TRUNCATE dds.dim_customer_status;
|
||||||
|
|
||||||
|
WITH src AS (
|
||||||
|
-- src: события из ODS + "контрольная сумма" статуса.
|
||||||
|
-- Так проще проверять, поменялся статус или остался тем же.
|
||||||
|
SELECT
|
||||||
|
customer_id AS customer_bk,
|
||||||
|
status,
|
||||||
|
event_ts,
|
||||||
|
md5(lower(coalesce(status, ''))) AS hashdiff
|
||||||
|
FROM ods.customer_status
|
||||||
|
),
|
||||||
|
ordered AS (
|
||||||
|
-- ordered: для каждого клиента смотрим "какая версия была до этого" (LAG)
|
||||||
|
SELECT
|
||||||
|
*,
|
||||||
|
lag(hashdiff) OVER (
|
||||||
|
PARTITION BY customer_bk
|
||||||
|
ORDER BY event_ts
|
||||||
|
) AS prev_hash
|
||||||
|
FROM src
|
||||||
|
),
|
||||||
|
changes AS (
|
||||||
|
-- changes: оставляем только первое состояние и реальные изменения статуса
|
||||||
|
SELECT *
|
||||||
|
FROM ordered
|
||||||
|
WHERE prev_hash IS DISTINCT FROM hashdiff OR prev_hash IS NULL
|
||||||
|
),
|
||||||
|
framed AS (
|
||||||
|
-- framed: превращаем изменения в периоды (valid_to = дата следующего события через LEAD)
|
||||||
|
SELECT
|
||||||
|
customer_bk,
|
||||||
|
status,
|
||||||
|
hashdiff,
|
||||||
|
event_ts::DATE AS valid_from,
|
||||||
|
lead(event_ts::DATE) OVER (
|
||||||
|
PARTITION BY customer_bk
|
||||||
|
ORDER BY event_ts
|
||||||
|
) AS valid_to
|
||||||
|
FROM changes
|
||||||
|
)
|
||||||
|
INSERT INTO dds.dim_customer_status (
|
||||||
|
customer_bk, status, hashdiff,
|
||||||
|
valid_from, valid_to,
|
||||||
|
created_at, updated_at
|
||||||
|
)
|
||||||
|
SELECT
|
||||||
|
customer_bk, status, hashdiff,
|
||||||
|
valid_from, valid_to,
|
||||||
|
now(), now()
|
||||||
|
FROM framed
|
||||||
|
ORDER BY customer_bk, valid_from;
|
||||||
|
|
||||||
|
-- Проверка: периоды в DDS (у клиента 101 должно быть 4 строки: new -> active -> vip -> churned)
|
||||||
|
SELECT 'dim_customer_status count = ' || COUNT(*) FROM dds.dim_customer_status;
|
||||||
|
SELECT * FROM dds.dim_customer_status ORDER BY customer_bk, valid_from;
|
||||||
|
|
||||||
|
-- ==========================================================
|
||||||
|
-- 3) Инкрементальная загрузка: STG -> ODS -> DDS
|
||||||
|
-- ==========================================================
|
||||||
|
|
||||||
|
-- Имитируем приход новой порции событий (customer_status_events_increment.csv):
|
||||||
|
-- - клиент 101: churned -> active (вернулся)
|
||||||
|
-- - клиент 102: churned -> active
|
||||||
|
-- - клиент 103: new -> active
|
||||||
|
-- - клиент 104: новый клиент, статус new
|
||||||
|
|
||||||
|
-- 3.0) Новые события в STG
|
||||||
|
-- При повторном запуске эти строки добавятся в STG ещё раз (дубли).
|
||||||
|
-- Для демо это не страшно: ODS-вставка ниже использует ON CONFLICT DO NOTHING,
|
||||||
|
-- а SCD2-блок защищён от повторных вставок через NOT EXISTS.
|
||||||
|
-- В продакшене STG обычно очищается перед каждой загрузкой (TRUNCATE / партиция по дате).
|
||||||
|
INSERT INTO stg.customer_status_raw (customer_id, status, event_ts, _load_id, _load_ts) VALUES
|
||||||
|
('101','active','2024-11-15 09:00:00','batch_20241115_1000','2024-11-15 10:00:00'),
|
||||||
|
('102','active','2024-05-05 09:30:00','batch_20240505_1000','2024-05-05 10:00:00'),
|
||||||
|
('103','active','2024-03-20 12:00:00','batch_20240320_1300','2024-03-20 13:00:00'),
|
||||||
|
('104','new', '2024-06-01 08:00:00','batch_20240601_0900','2024-06-01 09:00:00');
|
||||||
|
|
||||||
|
-- 3.1) UPSERT в ODS: добавляем новые события (не трогаем старые)
|
||||||
|
-- PK в ods.customer_status = (customer_id, event_ts), поэтому каждое уникальное
|
||||||
|
-- событие встаёт отдельной строкой. Дубли (одинаковый customer_id + event_ts) игнорируем.
|
||||||
|
INSERT INTO ods.customer_status (
|
||||||
|
customer_id, status, event_ts, _load_id, _load_ts
|
||||||
|
)
|
||||||
|
SELECT
|
||||||
|
s.customer_id::INT,
|
||||||
|
NULLIF(trim(s.status), ''),
|
||||||
|
NULLIF(trim(s.event_ts), '')::TIMESTAMP,
|
||||||
|
s._load_id,
|
||||||
|
COALESCE(s._load_ts, now())
|
||||||
|
FROM stg.customer_status_raw s
|
||||||
|
WHERE s.customer_id ~ '^\d+$'
|
||||||
|
AND NULLIF(trim(s.event_ts), '') IS NOT NULL
|
||||||
|
AND NULLIF(trim(s.status), '') IS NOT NULL
|
||||||
|
ON CONFLICT (customer_id, event_ts) DO NOTHING;
|
||||||
|
|
||||||
|
-- Проверка: в ODS должны появиться новые строки
|
||||||
|
SELECT 'ods.customer_status after increment = ' || COUNT(*) FROM ods.customer_status;
|
||||||
|
|
||||||
|
-- 3.2) Инкрементальное обновление DDS (SCD2)
|
||||||
|
-- Идея:
|
||||||
|
-- 1) берём по каждому клиенту самое позднее событие из ODS;
|
||||||
|
-- 2) сравниваем с текущей версией в DDS (valid_to IS NULL);
|
||||||
|
-- 3) если статус изменился - закрываем старую версию и вставляем новую.
|
||||||
|
--
|
||||||
|
-- Ограничение учебного варианта:
|
||||||
|
-- - если добавили событие "задним числом" со старой датой, этот блок не пересоберёт всю историю.
|
||||||
|
-- Для такого кейса обычно делают full refresh (блок 2).
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
-- 3.2a) Закрываем предыдущую актуальную версию
|
||||||
|
WITH ranked AS (
|
||||||
|
-- ranked: выбираем "самое свежее" событие на клиента
|
||||||
|
SELECT
|
||||||
|
customer_id AS customer_bk,
|
||||||
|
status,
|
||||||
|
event_ts::DATE AS eff_date,
|
||||||
|
md5(lower(coalesce(status, ''))) AS hashdiff,
|
||||||
|
row_number() OVER (
|
||||||
|
PARTITION BY customer_id
|
||||||
|
ORDER BY event_ts DESC, _load_ts DESC
|
||||||
|
) AS rn
|
||||||
|
FROM ods.customer_status
|
||||||
|
WHERE event_ts IS NOT NULL
|
||||||
|
),
|
||||||
|
delta AS (
|
||||||
|
-- delta: ровно одна строка на клиента (самое свежее событие)
|
||||||
|
SELECT * FROM ranked WHERE rn = 1
|
||||||
|
),
|
||||||
|
current_ver AS (
|
||||||
|
-- current_ver: текущие версии в DDS (valid_to IS NULL)
|
||||||
|
SELECT d.*
|
||||||
|
FROM dds.dim_customer_status d
|
||||||
|
WHERE d.valid_to IS NULL
|
||||||
|
)
|
||||||
|
UPDATE dds.dim_customer_status d
|
||||||
|
SET valid_to = x.eff_date,
|
||||||
|
updated_at = now()
|
||||||
|
FROM (
|
||||||
|
-- x: кого "закрываем":
|
||||||
|
-- клиент уже есть в DDS, и статус действительно изменился.
|
||||||
|
SELECT
|
||||||
|
t.customer_bk,
|
||||||
|
t.eff_date,
|
||||||
|
c.customer_status_sk
|
||||||
|
FROM delta t
|
||||||
|
JOIN current_ver c
|
||||||
|
ON c.customer_bk = t.customer_bk
|
||||||
|
WHERE c.hashdiff <> t.hashdiff
|
||||||
|
AND t.eff_date > c.valid_from -- не создаём период нулевой/отрицательной длины
|
||||||
|
) x
|
||||||
|
WHERE d.customer_status_sk = x.customer_status_sk
|
||||||
|
AND d.valid_to IS NULL;
|
||||||
|
|
||||||
|
-- 3.2b) Вставляем новую версию
|
||||||
|
WITH ranked AS (
|
||||||
|
-- ranked/delta/current_ver повторяем отдельно, чтобы блок INSERT читался отдельно от UPDATE
|
||||||
|
SELECT
|
||||||
|
customer_id AS customer_bk,
|
||||||
|
status,
|
||||||
|
event_ts::DATE AS eff_date,
|
||||||
|
md5(lower(coalesce(status, ''))) AS hashdiff,
|
||||||
|
row_number() OVER (
|
||||||
|
PARTITION BY customer_id
|
||||||
|
ORDER BY event_ts DESC, _load_ts DESC
|
||||||
|
) AS rn
|
||||||
|
FROM ods.customer_status
|
||||||
|
WHERE event_ts IS NOT NULL
|
||||||
|
),
|
||||||
|
delta AS (
|
||||||
|
SELECT * FROM ranked WHERE rn = 1
|
||||||
|
),
|
||||||
|
current_ver AS (
|
||||||
|
SELECT d.*
|
||||||
|
FROM dds.dim_customer_status d
|
||||||
|
WHERE d.valid_to IS NULL
|
||||||
|
),
|
||||||
|
to_insert AS (
|
||||||
|
-- to_insert: кого "вставляем":
|
||||||
|
-- 1) новый клиент (в current_ver нет строки);
|
||||||
|
-- 2) изменившийся клиент (статус поменялся).
|
||||||
|
SELECT
|
||||||
|
t.customer_bk,
|
||||||
|
t.status,
|
||||||
|
t.hashdiff,
|
||||||
|
t.eff_date
|
||||||
|
FROM delta t
|
||||||
|
LEFT JOIN current_ver c
|
||||||
|
ON c.customer_bk = t.customer_bk
|
||||||
|
WHERE c.customer_status_sk IS NULL
|
||||||
|
OR (c.hashdiff <> t.hashdiff AND t.eff_date > c.valid_from)
|
||||||
|
)
|
||||||
|
INSERT INTO dds.dim_customer_status (
|
||||||
|
customer_bk, status, hashdiff,
|
||||||
|
valid_from, valid_to,
|
||||||
|
created_at, updated_at
|
||||||
|
)
|
||||||
|
SELECT
|
||||||
|
t.customer_bk, t.status, t.hashdiff,
|
||||||
|
t.eff_date, NULL,
|
||||||
|
now(), now()
|
||||||
|
FROM to_insert t
|
||||||
|
-- защита от повторного запуска: не вставляем одну и ту же версию (BK + valid_from) второй раз
|
||||||
|
WHERE NOT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM dds.dim_customer_status d
|
||||||
|
WHERE d.customer_bk = t.customer_bk
|
||||||
|
AND d.valid_from = t.eff_date
|
||||||
|
);
|
||||||
|
COMMIT;
|
||||||
|
|
||||||
|
-- Проверка: у клиента 101 должна появиться 5-я строка (active с 2024-11-15),
|
||||||
|
-- у 104 - первая строка (new с 2024-06-01)
|
||||||
|
SELECT 'dim_customer_status after increment = ' || COUNT(*) FROM dds.dim_customer_status;
|
||||||
|
SELECT * FROM dds.dim_customer_status ORDER BY customer_bk, valid_from;
|
||||||
|
|
||||||
|
-- ==========================================================
|
||||||
|
-- 4) DM: витрина статусов клиентов по датам (full refresh)
|
||||||
|
-- ==========================================================
|
||||||
|
|
||||||
|
-- Витрина "снимок на дату":
|
||||||
|
-- для каждого дня считаем, сколько клиентов было в каждом статусе.
|
||||||
|
-- Берём календарь dds.dim_date и подбираем статус по периоду valid_from/valid_to.
|
||||||
|
-- DDL витрины - в 07_ddl_hw_customer_status.sql.
|
||||||
|
|
||||||
|
TRUNCATE dm.mart_customer_status_daily;
|
||||||
|
|
||||||
|
WITH bounds AS (
|
||||||
|
SELECT
|
||||||
|
min(valid_from) AS date_from,
|
||||||
|
-- CURRENT_DATE для открытых интервалов (valid_to IS NULL = текущий статус),
|
||||||
|
-- иначе витрина не покроет даты после последней смены статуса.
|
||||||
|
-- Нюанс: количество строк в витрине зависит от даты запуска (каждый
|
||||||
|
-- день добавляется ещё один день). Для учебных целей это приемлемо.
|
||||||
|
max(coalesce(valid_to, CURRENT_DATE)) AS date_to
|
||||||
|
FROM dds.dim_customer_status
|
||||||
|
)
|
||||||
|
INSERT INTO dm.mart_customer_status_daily (
|
||||||
|
date_actual, status, customers_cnt
|
||||||
|
)
|
||||||
|
SELECT
|
||||||
|
d.date_actual,
|
||||||
|
s.status,
|
||||||
|
COUNT(DISTINCT s.customer_bk) AS customers_cnt
|
||||||
|
FROM dds.dim_date d
|
||||||
|
JOIN bounds b
|
||||||
|
ON d.date_actual BETWEEN b.date_from AND b.date_to
|
||||||
|
JOIN dds.dim_customer_status s
|
||||||
|
ON d.date_actual >= s.valid_from
|
||||||
|
AND (s.valid_to IS NULL OR d.date_actual < s.valid_to)
|
||||||
|
GROUP BY d.date_actual, s.status
|
||||||
|
ORDER BY d.date_actual, s.status;
|
||||||
|
|
||||||
|
-- Проверка: выборочные даты из витрины
|
||||||
|
SELECT 'mart_customer_status_daily count = ' || COUNT(*) FROM dm.mart_customer_status_daily;
|
||||||
|
SELECT *
|
||||||
|
FROM dm.mart_customer_status_daily
|
||||||
|
WHERE date_actual IN ('2024-01-15', '2024-04-10', '2024-09-15', '2024-12-01')
|
||||||
|
ORDER BY date_actual, status;
|
||||||
Reference in New Issue
Block a user