docs(etl): расширена документация ETL DAG-ов ODS, DDS, DM
- Зачем:
- документация трёх DAG-ов была значительно беднее эталона (bookings_to_gp_stage.md):
отсутствовал пошаговый разбор задач, SQL-пути и ASCII-графы зависимостей.
- Что:
- добавлены секции "Как это работает внутри" с таблицами task_id → SQL-файл → паттерн.
- добавлены ASCII-графы зависимостей в ODS и DDS (в DM уже был).
- исправлено описание паттернов ODS: явно разделены TRUNCATE+INSERT для AO snapshot-справочников и SCD1 UPSERT для транзакционных таблиц.
- исправлено описание fact_flight_sales: убрана неточная отсылка к late-arriving dimensions, добавлено объяснение defensive LEFT JOIN и точных DQ-правил (0% / 1%).
- добавлена таблица grain и описание стратегий загрузки для каждой DM-витрины.
- план улучшений перенесён в docs/archive/.
- Проверка:
- визуально: открыть каждый docs/bookings_to_gp_*.md и убедиться, что секции присутствуют.
This commit is contained in:
@@ -0,0 +1,88 @@
|
|||||||
|
# План: выравнивание качества документации ETL DAG'ов
|
||||||
|
|
||||||
|
## Контекст
|
||||||
|
|
||||||
|
Документ `bookings_to_gp_stage.md` — эталон: пошаговый разбор задач, ASCII-граф,
|
||||||
|
ссылки на SQL-файлы, описание edge cases. Три остальных документа
|
||||||
|
(`_ods`, `_dds`, `_dm`) значительно беднее. Цель — подтянуть их до того же уровня.
|
||||||
|
|
||||||
|
## Единый шаблон секций (целевая структура)
|
||||||
|
|
||||||
|
Каждый документ должен содержать:
|
||||||
|
|
||||||
|
1. **Заголовок + вводный абзац** (что за DAG, какой слой, зачем)
|
||||||
|
2. **Что делает DAG** (буллеты, краткое описание)
|
||||||
|
3. **Что должно быть готово** (prerequisites)
|
||||||
|
4. **Как запустить** (UI + опциональные параметры)
|
||||||
|
5. **Граф зависимостей** (ASCII-диаграмма, не буллет-лист)
|
||||||
|
6. **Как это работает внутри (по шагам)** ← ГЛАВНОЕ ДОБАВЛЕНИЕ
|
||||||
|
- Пронумерованные шаги: task_id → SQL-файл → что делает → паттерн (SCD1/SCD2/rebuild/HWM)
|
||||||
|
- Учебные пояснения к нетривиальным паттернам
|
||||||
|
7. **Как проверить результат** (SQL-запросы)
|
||||||
|
8. **Типичные ошибки** (уже есть, оставляем)
|
||||||
|
|
||||||
|
## Что именно добавить/исправить в каждом документе
|
||||||
|
|
||||||
|
### A. `bookings_to_gp_ods.md`
|
||||||
|
|
||||||
|
**Текущее состояние:** 95 строк, нет пошагового разбора, нет ASCII-графа, нет SQL-путей.
|
||||||
|
|
||||||
|
| # | Что сделать | Детали |
|
||||||
|
|---|-------------|--------|
|
||||||
|
| A1 | ASCII-граф зависимостей | Заменить буллет-лист на диаграмму (как в stage). Показать параллельные ветки airports/airplanes, схождение на routes/seats, цепочку flights→segments→boarding_passes |
|
||||||
|
| A2 | Секция "Как это работает внутри" | 10 шагов: resolve_stg_batch_id (Python, INTERSECT-логика), затем 9 пар load→dq с указанием SQL-файлов |
|
||||||
|
| A3 | Пояснить паттерн SCD1 UPSERT | Кратко: TEMP TABLE → UPDATE (IS DISTINCT FROM) → INSERT. Одного абзаца достаточно, потом ссылка "паттерн одинаков для всех 9 таблиц" |
|
||||||
|
| A4 | Описать разницу snapshot vs HWM | Snapshot-справочники фильтруются по `stg_batch_id`; транзакционные таблицы — по HWM (`_load_ts`). Объяснить почему (чтобы не терять инкременты при повторных запусках STG) |
|
||||||
|
| A5 | Edge case: пустой батч | Для инкрементальных таблиц допустим; для snapshot — нет |
|
||||||
|
|
||||||
|
**Ожидаемый объём:** ~140–160 строк.
|
||||||
|
|
||||||
|
### B. `bookings_to_gp_dds.md`
|
||||||
|
|
||||||
|
**Текущее состояние:** 73 строки — самый бедный документ. Нет пошаговости, нет ASCII-графа, SCD2 не объяснён.
|
||||||
|
|
||||||
|
| # | Что сделать | Детали |
|
||||||
|
|---|-------------|--------|
|
||||||
|
| B1 | ASCII-граф зависимостей | calendar → параллельно 4 SCD1-измерения + dim_routes (после airports+airplanes) → fact (после всех dims) → summary |
|
||||||
|
| B2 | Секция "Как это работает внутри" | 8 шагов: calendar (rebuild), 4×SCD1-измерения, dim_routes (SCD2), fact_flight_sales, summary |
|
||||||
|
| B3 | Объяснить SCD2 для dim_routes | Учебный блок: hashdiff (MD5), закрытие старых версий, вставка новых, point-in-time valid_from. Это ключевой паттерн DDS — заслуживает 10–15 строк |
|
||||||
|
| B4 | Объяснить Phase 3 (денормализация) | Обновление SCD1-атрибутов (города, модель) во ВСЕХ версиях dim_routes. Зачем: чтобы не хранить устаревшие названия городов |
|
||||||
|
| B5 | Объяснить late-arriving dimensions | LEFT JOIN в fact_flight_sales: факт может прийти раньше справочника. 3–5 строк |
|
||||||
|
| B6 | Объяснить point-in-time join | Как факт привязывается к правильной версии SCD2-маршрута по дате рейса |
|
||||||
|
| B7 | Edge case: генерация SK | MAX() + ROW_NUMBER() безопасна только при concurrency=1 |
|
||||||
|
|
||||||
|
**Ожидаемый объём:** ~160–180 строк.
|
||||||
|
|
||||||
|
### C. `bookings_to_gp_dm.md`
|
||||||
|
|
||||||
|
**Текущее состояние:** 80 строк. ASCII-граф есть (хорошо!), описание стратегий загрузки есть (хорошо!), но нет пошагового разбора с SQL-файлами.
|
||||||
|
|
||||||
|
| # | Что сделать | Детали |
|
||||||
|
|---|-------------|--------|
|
||||||
|
| C1 | Секция "Как это работает внутри" | 5 шагов (по одному на витрину) + start_dm + finish_dm_summary. Для каждой витрины: task_id → SQL-файл → паттерн → зерно (grain) |
|
||||||
|
| C2 | Расширить описание sales_report | HWM-паттерн: какие даты пересчитываются, NULLIF-guard для boarding_rate, автоматическая "догонка" при первичной загрузке |
|
||||||
|
| C3 | Расширить описание route_performance | Почему Full Rebuild: AO Column (нет UPDATE/DELETE), таблица маленькая. 3-шаговый паттерн: TRUNCATE → агрегация по route_bk → JOIN с текущей версией SCD2 |
|
||||||
|
| C4 | Добавить grain для каждой витрины | sales_report: (flight_date, departure_airport_sk, arrival_airport_sk, tariff_sk); route_performance: route_bk; и т.д. |
|
||||||
|
| C5 | Добавить SQL-пути к задачам | Сейчас нигде не указаны пути к SQL-файлам |
|
||||||
|
|
||||||
|
**Ожидаемый объём:** ~140–160 строк.
|
||||||
|
|
||||||
|
### D. Мелкие правки в `bookings_to_gp_stage.md` (опционально)
|
||||||
|
|
||||||
|
| # | Что сделать | Детали |
|
||||||
|
|---|-------------|--------|
|
||||||
|
| D1 | Убедиться в консистентности шаблона | Если в ходе работы над ODS/DDS/DM выработается чуть лучшая структура — привести stage к тому же формату (только структурные правки, контент не менять) |
|
||||||
|
|
||||||
|
## Порядок работы
|
||||||
|
|
||||||
|
1. **ODS** (средняя сложность, знакомый паттерн SCD1)
|
||||||
|
2. **DDS** (наибольшая учебная ценность — SCD2, late-arriving dims)
|
||||||
|
3. **DM** (наименьший объём правок — ASCII-граф уже есть)
|
||||||
|
4. **Stage** — только если нужна косметика для консистентности
|
||||||
|
|
||||||
|
## Принципы
|
||||||
|
|
||||||
|
- Не раздувать: целевой объём каждого документа — 140–180 строк (stage = 157).
|
||||||
|
- Учебная ценность > полнота: объяснять «почему», а не перечислять все колонки.
|
||||||
|
- SQL-пути обязательны — это главный навигационный инструмент для студента.
|
||||||
|
- Один паттерн объясняем один раз подробно, дальше ссылаемся: "паттерн аналогичен X".
|
||||||
+128
-20
@@ -1,17 +1,19 @@
|
|||||||
# DAG `bookings_to_gp_dds`: `ods` -> `dds` в Greenplum
|
# DAG `bookings_to_gp_dds`: `ods` → `dds` в Greenplum
|
||||||
|
|
||||||
Этот DAG — учебный пример загрузки аналитического слоя **DDS** (Star Schema) из текущего состояния **ODS**.
|
Этот DAG — учебный пример загрузки аналитического слоя **DDS** (Star Schema) из текущего состояния **ODS**.
|
||||||
Логика: измерения + факт, проверки качества данных после каждой загрузки.
|
Здесь сосредоточены ключевые паттерны аналитического хранилища: SCD1, SCD2 с hashdiff,
|
||||||
|
point-in-time join, защитные LEFT JOIN для устойчивости к data quality аномалиям.
|
||||||
|
|
||||||
## Что делает DAG
|
## Что делает DAG
|
||||||
|
|
||||||
- Загружает измерения DDS:
|
- Загружает 6 измерений DDS:
|
||||||
- `dds.dim_calendar` (статическое измерение дат);
|
- `dds.dim_calendar` — статическое измерение дат (Full Rebuild);
|
||||||
- `dds.dim_airports`, `dds.dim_airplanes`, `dds.dim_tariffs`, `dds.dim_passengers` (SCD1 UPSERT);
|
- `dds.dim_airports`, `dds.dim_airplanes`, `dds.dim_tariffs`, `dds.dim_passengers` — SCD1 UPSERT;
|
||||||
- `dds.dim_routes` (SCD2 с `hashdiff`, `valid_from`, `valid_to`).
|
- `dds.dim_routes` — **SCD2** с `hashdiff`, `valid_from`, `valid_to` + денормализация.
|
||||||
- Загружает факт `dds.fact_flight_sales` инкрементальным UPSERT по зерну `(ticket_no, flight_id)`.
|
- Загружает факт `dds.fact_flight_sales` — инкрементальный UPSERT по зерну `(ticket_no, flight_id)`.
|
||||||
- Для каждой таблицы выполняет пару задач `load -> dq`.
|
- Для каждой таблицы выполняет пару задач `load → dq`.
|
||||||
- Использует `_load_id = {{ run_id }}` (DDS не требует `stg_batch_id`, потому что читает current state ODS).
|
- Использует `_load_id = {{ run_id }}`. DDS не требует `stg_batch_id`, потому что читает
|
||||||
|
текущее состояние ODS.
|
||||||
|
|
||||||
## Что должно быть готово перед запуском
|
## Что должно быть готово перед запуском
|
||||||
|
|
||||||
@@ -32,17 +34,116 @@ make up
|
|||||||
2) Если запускаете DDS впервые — выполните `bookings_dds_ddl`.
|
2) Если запускаете DDS впервые — выполните `bookings_dds_ddl`.
|
||||||
3) Запустите `bookings_to_gp_dds`.
|
3) Запустите `bookings_to_gp_dds`.
|
||||||
|
|
||||||
## Граф зависимостей (упрощённо)
|
## Граф зависимостей
|
||||||
|
|
||||||
- `load_dds_dim_calendar -> dq_dds_dim_calendar`
|
```
|
||||||
- После calendar параллельно:
|
load_dds_dim_calendar → dq_dds_dim_calendar
|
||||||
- `load_dds_dim_airports -> dq_dds_dim_airports`
|
├─ load_dds_dim_airports → dq_dds_dim_airports ─┐
|
||||||
- `load_dds_dim_airplanes -> dq_dds_dim_airplanes`
|
│ ├─ load_dds_dim_routes
|
||||||
- `load_dds_dim_tariffs -> dq_dds_dim_tariffs`
|
├─ load_dds_dim_airplanes → dq_dds_dim_airplanes ─┘ └─ dq_dds_dim_routes
|
||||||
- `load_dds_dim_passengers -> dq_dds_dim_passengers`
|
│ │
|
||||||
- `load_dds_dim_routes -> dq_dds_dim_routes`
|
├─ load_dds_dim_tariffs → dq_dds_dim_tariffs │
|
||||||
- Факт:
|
│ │
|
||||||
- `load_dds_fact_flight_sales -> dq_dds_fact_flight_sales -> finish_dds_summary`
|
└─ load_dds_dim_passengers → dq_dds_dim_passengers │
|
||||||
|
│
|
||||||
|
все 5 dq_dds_dim_* ─────────────────────────────────────────┘
|
||||||
|
└─ load_dds_fact_flight_sales
|
||||||
|
└─ dq_dds_fact_flight_sales
|
||||||
|
└─ finish_dds_summary
|
||||||
|
```
|
||||||
|
|
||||||
|
Ключевой момент: `dim_routes` зависит от `dim_airports` и `dim_airplanes` (денормализация),
|
||||||
|
а факт ждёт завершения **всех** пяти измерений.
|
||||||
|
|
||||||
|
## Как это работает внутри (по шагам)
|
||||||
|
|
||||||
|
### 1) `load_dds_dim_calendar` → `dq_dds_dim_calendar`
|
||||||
|
|
||||||
|
- **SQL:** `sql/dds/dim_calendar_load.sql`, `sql/dds/dim_calendar_dq.sql`
|
||||||
|
- **Паттерн:** Full Rebuild — каждый запуск пересоздаёт календарь целиком.
|
||||||
|
Измерение маленькое и детерминированное, дельту считать нет смысла.
|
||||||
|
|
||||||
|
### 2–5) SCD1-измерения (параллельно после calendar)
|
||||||
|
|
||||||
|
| # | Задача | SQL-файлы | Что загружает |
|
||||||
|
|---|--------|-----------|---------------|
|
||||||
|
| 2 | `load_dds_dim_airports` → `dq_dds_dim_airports` | `sql/dds/dim_airports_load.sql`, `sql/dds/dim_airports_dq.sql` | Аэропорты (код, город, координаты) |
|
||||||
|
| 3 | `load_dds_dim_airplanes` → `dq_dds_dim_airplanes` | `sql/dds/dim_airplanes_load.sql`, `sql/dds/dim_airplanes_dq.sql` | Самолёты (код, модель, кол-во мест) |
|
||||||
|
| 4 | `load_dds_dim_tariffs` → `dq_dds_dim_tariffs` | `sql/dds/dim_tariffs_load.sql`, `sql/dds/dim_tariffs_dq.sql` | Тарифы (класс обслуживания) |
|
||||||
|
| 5 | `load_dds_dim_passengers` → `dq_dds_dim_passengers` | `sql/dds/dim_passengers_load.sql`, `sql/dds/dim_passengers_dq.sql` | Пассажиры (ID, имя, контакты) |
|
||||||
|
|
||||||
|
Паттерн загрузки — SCD1 UPSERT: TEMP TABLE → UPDATE (IS DISTINCT FROM) → INSERT.
|
||||||
|
|
||||||
|
### 6) `load_dds_dim_routes` → `dq_dds_dim_routes` (SCD2)
|
||||||
|
|
||||||
|
- **SQL:** `sql/dds/dim_routes_load.sql`, `sql/dds/dim_routes_dq.sql`
|
||||||
|
- **Паттерн:** SCD2 — самый нетривиальный паттерн в проекте. Работает в 3 фазы:
|
||||||
|
|
||||||
|
**Фаза 1. Hashdiff и закрытие старых версий.**
|
||||||
|
Скрипт считает MD5-хеш от шести бизнес-атрибутов маршрута (`departure_airport`, `arrival_airport`,
|
||||||
|
`airplane_code`, `days_of_week`, `departure_time`, `duration`).
|
||||||
|
Если хеш текущей версии в DDS не совпадает с хешем из ODS — старая версия закрывается
|
||||||
|
(`valid_to = CURRENT_DATE`). Также закрываются маршруты, исчезнувшие из ODS.
|
||||||
|
|
||||||
|
**Фаза 2. Вставка новых версий.**
|
||||||
|
Для изменённых и совершенно новых маршрутов создаётся новая строка.
|
||||||
|
`valid_from` выставляется в `1900-01-01` для первой версии маршрута и `CURRENT_DATE` для версии 2+.
|
||||||
|
Суррогатный ключ (`route_sk`) генерируется через `MAX(route_sk) + ROW_NUMBER()`.
|
||||||
|
|
||||||
|
> **Важно:** такая генерация SK безопасна только при `max_active_runs=1` (Airflow гарантирует
|
||||||
|
> последовательный запуск). В боевых системах используют sequence.
|
||||||
|
|
||||||
|
**Фаза 3. Обновление денормализованных атрибутов.**
|
||||||
|
`dim_routes` хранит денормализованные SCD1-атрибуты из `dim_airports` (города)
|
||||||
|
и `dim_airplanes` (модель, кол-во мест). Если, например, город переименовали —
|
||||||
|
фаза 3 обновляет **все** версии маршрута (и текущие, и исторические),
|
||||||
|
при этом `_load_id` и `_load_ts` не перезаписываются (lineage версий сохраняется).
|
||||||
|
|
||||||
|
### 7) `load_dds_fact_flight_sales` → `dq_dds_fact_flight_sales`
|
||||||
|
|
||||||
|
- **SQL:** `sql/dds/fact_flight_sales_load.sql`, `sql/dds/fact_flight_sales_dq.sql`
|
||||||
|
- **Зерно:** `(ticket_no, flight_id)` — один билет на один рейс.
|
||||||
|
- **Паттерн:** инкрементальный UPSERT.
|
||||||
|
|
||||||
|
Три учебных приёма в этом скрипте:
|
||||||
|
|
||||||
|
**Защитные LEFT JOIN (defensive coding).**
|
||||||
|
Все JOIN-ы с измерениями — `LEFT JOIN`. DAG гарантирует, что все измерения загружены
|
||||||
|
и прошли DQ **до** старта факта (жёсткие зависимости в графе). Поэтому в штатном режиме
|
||||||
|
NULL SK не возникают. LEFT JOIN здесь — защита от data quality аномалий (например, если
|
||||||
|
в `ods.routes` появится маршрут с несуществующим аэропортом).
|
||||||
|
|
||||||
|
DQ-проверки факта отражают эту логику:
|
||||||
|
- `passenger_sk` и `tariff_sk` — **запрещены** NULL целиком (0 строк);
|
||||||
|
- route-related FK (`route_sk`, `airport_sk`, `airplane_sk`) и `calendar_sk` —
|
||||||
|
допускается до **1%** NULL (NOTICE-предупреждение), при превышении — EXCEPTION.
|
||||||
|
|
||||||
|
> Это **не** паттерн late-arriving dimensions (опаздывающих измерений) в классическом
|
||||||
|
> понимании: backfill NULL SK при повторном запуске не реализован.
|
||||||
|
> В боевых системах для этого используют «строку-заглушку» (unknown member, SK = 0)
|
||||||
|
> и отдельный процесс backfill.
|
||||||
|
|
||||||
|
**Point-in-time join для SCD2.**
|
||||||
|
Для `dim_routes` используется привязка по дате рейса:
|
||||||
|
|
||||||
|
```sql
|
||||||
|
LEFT JOIN dds.dim_routes AS rte
|
||||||
|
ON rte.route_bk = flt.route_no
|
||||||
|
AND flt.scheduled_departure::DATE >= rte.valid_from
|
||||||
|
AND (rte.valid_to IS NULL OR flt.scheduled_departure::DATE < rte.valid_to)
|
||||||
|
```
|
||||||
|
|
||||||
|
Это гарантирует, что факт привязывается к той версии маршрута, которая была актуальна
|
||||||
|
на дату рейса.
|
||||||
|
|
||||||
|
**UPDATE мутабельных полей.**
|
||||||
|
UPDATE обновляет только `seat_no`, `price`, `is_boarded` (данные, которые реально
|
||||||
|
могут измениться — посадка пассажира, корректировка цены). SK измерений не перезаписываются —
|
||||||
|
они зафиксированы на момент вставки.
|
||||||
|
|
||||||
|
### 8) `finish_dds_summary`
|
||||||
|
|
||||||
|
Ждёт завершения DQ факта и логирует сводку.
|
||||||
|
|
||||||
## Как проверить результат
|
## Как проверить результат
|
||||||
|
|
||||||
@@ -55,12 +156,19 @@ SELECT COUNT(*) FROM dds.dim_calendar;
|
|||||||
SELECT COUNT(*) FROM dds.dim_routes;
|
SELECT COUNT(*) FROM dds.dim_routes;
|
||||||
SELECT COUNT(*) FROM dds.fact_flight_sales;
|
SELECT COUNT(*) FROM dds.fact_flight_sales;
|
||||||
|
|
||||||
|
-- Проверка: кол-во строк факта ≈ кол-во строк ODS segments
|
||||||
SELECT
|
SELECT
|
||||||
(SELECT COUNT(*) FROM dds.fact_flight_sales) AS fact_rows,
|
(SELECT COUNT(*) FROM dds.fact_flight_sales) AS fact_rows,
|
||||||
(SELECT COUNT(*) FROM ods.segments) AS ods_rows;
|
(SELECT COUNT(*) FROM ods.segments) AS ods_rows;
|
||||||
|
|
||||||
|
-- Проверка SCD2: текущие версии маршрутов (valid_to IS NULL)
|
||||||
|
SELECT COUNT(*) AS current_versions,
|
||||||
|
(SELECT COUNT(*) FROM dds.dim_routes) AS total_versions
|
||||||
|
FROM dds.dim_routes
|
||||||
|
WHERE valid_to IS NULL;
|
||||||
```
|
```
|
||||||
|
|
||||||
Ожидаемо: `fact_rows = ods_rows`.
|
Ожидаемо: `fact_rows = ods_rows`, `current_versions ≤ total_versions`.
|
||||||
|
|
||||||
## Типичные ошибки
|
## Типичные ошибки
|
||||||
|
|
||||||
|
|||||||
+110
-18
@@ -1,17 +1,20 @@
|
|||||||
# DAG `bookings_to_gp_dm`: `dds` -> `dm` в Greenplum
|
# DAG `bookings_to_gp_dm`: `dds` → `dm` в Greenplum
|
||||||
|
|
||||||
Этот DAG — учебный пример загрузки слоя **DM** (Data Mart / витрины) из текущего состояния **DDS**.
|
Этот DAG — учебный пример загрузки слоя **DM** (Data Mart / витрины) из текущего состояния **DDS**.
|
||||||
Логика: все 5 витрин загружаются параллельно, для каждой — пара `load -> dq`.
|
Все 5 витрин загружаются **параллельно** и демонстрируют разные стратегии загрузки —
|
||||||
|
это ключевая учебная ценность данного DAG.
|
||||||
|
|
||||||
## Что делает DAG
|
## Что делает DAG
|
||||||
|
|
||||||
- Загружает витрины DM параллельно (паттерны загрузки разные — учебная демонстрация выбора стратегии):
|
Загружает 5 витрин параллельно, для каждой — пара `load → dq`:
|
||||||
- `dm.sales_report` — UPSERT по датам; DQ проверяет только строки текущего `run_id` (`_load_id`);
|
|
||||||
- `dm.route_performance` — Full Rebuild (TRUNCATE + INSERT): таблица маленькая, дельту считать дороже;
|
| Витрина | Зерно (grain) | Паттерн загрузки |
|
||||||
- `dm.passenger_loyalty` — инкрементальный UPSERT по «затронутым ключам» (HWM по `_load_ts`): пересчитываем агрегаты только для пассажиров с новыми фактами;
|
|---------|---------------|------------------|
|
||||||
- `dm.airport_traffic` — инкрементальный UPSERT по датам (HWM по `_load_ts`);
|
| `dm.sales_report` | (flight_date, departure_airport_sk, arrival_airport_sk, tariff_sk) | Инкрементальный UPSERT (HWM по датам) |
|
||||||
- `dm.monthly_overview` — инкрементальный UPSERT по месяцам (HWM по `_load_ts`).
|
| `dm.route_performance` | route_bk | Full Rebuild (TRUNCATE + INSERT) |
|
||||||
- Для каждой витрины выполняет пару задач `load -> dq`.
|
| `dm.passenger_loyalty` | passenger_sk | Инкрементальный UPSERT (HWM по затронутым ключам) |
|
||||||
|
| `dm.airport_traffic` | (traffic_date, airport_sk) | Инкрементальный UPSERT (HWM по датам) |
|
||||||
|
| `dm.monthly_overview` | (year_actual, month_actual, airplane_sk) | Инкрементальный UPSERT (HWM по месяцам) |
|
||||||
|
|
||||||
## Что должно быть готово перед запуском
|
## Что должно быть готово перед запуском
|
||||||
|
|
||||||
@@ -34,17 +37,106 @@ make up
|
|||||||
|
|
||||||
## Граф зависимостей
|
## Граф зависимостей
|
||||||
|
|
||||||
Все 5 веток запускаются параллельно от `start_dm`, затем сходятся в `finish_dm_summary`:
|
|
||||||
|
|
||||||
```
|
```
|
||||||
start_dm
|
start_dm
|
||||||
├── load_dm_sales_report -> dq_dm_sales_report -> finish_dm_summary
|
├── load_dm_sales_report → dq_dm_sales_report ─┐
|
||||||
├── load_dm_route_performance -> dq_dm_route_performance -> finish_dm_summary
|
├── load_dm_route_performance → dq_dm_route_performance ─┤
|
||||||
├── load_dm_passenger_loyalty -> dq_dm_passenger_loyalty -> finish_dm_summary
|
├── load_dm_passenger_loyalty → dq_dm_passenger_loyalty ─┼─ finish_dm_summary
|
||||||
├── load_dm_airport_traffic -> dq_dm_airport_traffic -> finish_dm_summary
|
├── load_dm_airport_traffic → dq_dm_airport_traffic ─┤
|
||||||
└── load_dm_monthly_overview -> dq_dm_monthly_overview -> finish_dm_summary
|
└── load_dm_monthly_overview → dq_dm_monthly_overview ─┘
|
||||||
```
|
```
|
||||||
|
|
||||||
|
Все 5 веток полностью независимы и работают параллельно.
|
||||||
|
|
||||||
|
## Как это работает внутри (по шагам)
|
||||||
|
|
||||||
|
### 1) `load_dm_sales_report` → `dq_dm_sales_report`
|
||||||
|
|
||||||
|
- **SQL:** `sql/dm/sales_report_load.sql`, `sql/dm/sales_report_dq.sql`
|
||||||
|
- **Паттерн:** инкрементальный UPSERT (HWM по датам).
|
||||||
|
|
||||||
|
Витрина агрегирует продажи билетов по дате, аэропортам вылета/прилёта и тарифу.
|
||||||
|
Инкрементальность работает через HWM: витрина сравнивает свой `MAX(_load_ts)` с `_load_ts`
|
||||||
|
фактов в DDS и пересчитывает агрегаты только для **затронутых дат**.
|
||||||
|
|
||||||
|
> Если витрина пуста — `1900-01-01` заберёт всю историю (первичная загрузка).
|
||||||
|
> Если DAG не запускался несколько дней — при следующем запуске витрина автоматически
|
||||||
|
> «догонит» всю накопленную дельту.
|
||||||
|
|
||||||
|
Учебные приёмы:
|
||||||
|
- **TEMP TABLE** для однократной агрегации (канон для MPP);
|
||||||
|
- **NULLIF** для защиты от деления на ноль (`boarding_rate = boarded / NULLIF(sold, 0)`);
|
||||||
|
- **Денормализация**: города и коды аэропортов тянутся в витрину из измерений.
|
||||||
|
|
||||||
|
### 2) `load_dm_route_performance` → `dq_dm_route_performance`
|
||||||
|
|
||||||
|
- **SQL:** `sql/dm/route_performance_load.sql`, `sql/dm/route_performance_dq.sql`
|
||||||
|
- **Паттерн:** Full Rebuild (TRUNCATE + INSERT).
|
||||||
|
|
||||||
|
Витрина агрегирует эффективность маршрутов за всю историю. Таблица маленькая (~1000 строк),
|
||||||
|
поэтому пересоздать её с нуля дешевле, чем вычислять дельту. Дополнительная причина:
|
||||||
|
таблица хранится в формате **AO Column Store**, который не поддерживает эффективный UPDATE/DELETE.
|
||||||
|
|
||||||
|
Трёхшаговый паттерн:
|
||||||
|
1. **TRUNCATE** — очистка (единственный эффективный способ для AO).
|
||||||
|
2. **Агрегация** по `route_bk` — факты суммируются через **все исторические версии** маршрута
|
||||||
|
(route_sk из SCD2), чтобы не терять данные при версионировании.
|
||||||
|
3. **JOIN** с текущей (актуальной, `valid_to IS NULL`) версией `dim_routes` для денормализации.
|
||||||
|
|
||||||
|
> Благодаря денормализации `dim_routes` — один JOIN вместо четырёх
|
||||||
|
> (аэропорты вылета/прилёта, самолёт уже хранятся в `dim_routes`).
|
||||||
|
|
||||||
|
Метрики: `avg_load_factor = total_boarded / (total_flights * total_seats)`,
|
||||||
|
`avg_ticket_price = total_revenue / total_tickets`.
|
||||||
|
|
||||||
|
### 3) `load_dm_passenger_loyalty` → `dq_dm_passenger_loyalty`
|
||||||
|
|
||||||
|
- **SQL:** `sql/dm/passenger_loyalty_load.sql`, `sql/dm/passenger_loyalty_dq.sql`
|
||||||
|
- **Паттерн:** инкрементальный UPSERT по «затронутым ключам».
|
||||||
|
|
||||||
|
В отличие от `sales_report` (где инкремент по датам), здесь HWM находит **конкретных пассажиров**
|
||||||
|
с новыми фактами, а затем пересчитывает для них всю историю. Это гарантирует точность
|
||||||
|
накопительных агрегатов (`total_spent`, `first/last_flight_date`).
|
||||||
|
|
||||||
|
Учебные приёмы:
|
||||||
|
- **DISTINCT ON** (PostgreSQL-специфика) для нахождения моды (самый частый тариф пассажира);
|
||||||
|
- **Агрегация SCD2 по BK**: при подсчёте уникальных маршрутов используем `route_bk`,
|
||||||
|
а не `route_sk`, т.к. один маршрут может иметь несколько версий;
|
||||||
|
- Фильтрация `passenger_sk IS NOT NULL` — защита от неконсистентных фактов (NULL SK
|
||||||
|
при data quality аномалиях в измерениях).
|
||||||
|
|
||||||
|
### 4) `load_dm_airport_traffic` → `dq_dm_airport_traffic`
|
||||||
|
|
||||||
|
- **SQL:** `sql/dm/airport_traffic_load.sql`, `sql/dm/airport_traffic_dq.sql`
|
||||||
|
- **Паттерн:** инкрементальный UPSERT по датам.
|
||||||
|
|
||||||
|
Витрина показывает пассажиропоток аэропортов по дням (вылеты + прилёты).
|
||||||
|
|
||||||
|
Учебный приём — **Dual-role dimension через UNION ALL**: один билет превращается
|
||||||
|
в два «события» (вылет из одного аэропорта и прилёт в другой).
|
||||||
|
Это позволяет собрать единую статистику аэропорта (departures + arrivals) в одном проходе.
|
||||||
|
|
||||||
|
### 5) `load_dm_monthly_overview` → `dq_dm_monthly_overview`
|
||||||
|
|
||||||
|
- **SQL:** `sql/dm/monthly_overview_load.sql`, `sql/dm/monthly_overview_dq.sql`
|
||||||
|
- **Паттерн:** инкрементальный UPSERT по месяцам.
|
||||||
|
|
||||||
|
Витрина показывает помесячную статистику по типам самолётов.
|
||||||
|
|
||||||
|
Учебный приём — **двухуровневая агрегация**: чтобы честно посчитать `avg_load_factor`,
|
||||||
|
сначала считаем load factor для каждого рейса (`boarded / total_seats`),
|
||||||
|
затем берём среднее по месяцу. Прямая агрегация `SUM(boarded) / SUM(seats)` дала бы
|
||||||
|
искажённый результат (взвешенный по числу билетов, а не рейсов).
|
||||||
|
|
||||||
|
> **Ограничение SCD1:** `total_seats` берётся из текущего состояния `dim_airplanes`.
|
||||||
|
> Если самолёт переоборудовали в прошлом, для точного исторического расчёта
|
||||||
|
> потребовалось бы SCD2-измерение.
|
||||||
|
|
||||||
|
### 6) `start_dm` / `finish_dm_summary`
|
||||||
|
|
||||||
|
- `start_dm` — стартовый sentinel, от которого расходятся все 5 параллельных веток.
|
||||||
|
- `finish_dm_summary` — ждёт завершения всех DQ-задач и логирует сводку.
|
||||||
|
|
||||||
## Как проверить результат
|
## Как проверить результат
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
@@ -58,10 +150,10 @@ SELECT COUNT(*) FROM dm.passenger_loyalty;
|
|||||||
SELECT COUNT(*) FROM dm.airport_traffic;
|
SELECT COUNT(*) FROM dm.airport_traffic;
|
||||||
SELECT COUNT(*) FROM dm.monthly_overview;
|
SELECT COUNT(*) FROM dm.monthly_overview;
|
||||||
|
|
||||||
-- Проверка инварианта sales_report: посаженных не больше, чем продано
|
-- Инвариант sales_report: посаженных не больше, чем продано
|
||||||
SELECT COUNT(*) FROM dm.sales_report WHERE tickets_sold < passengers_boarded;
|
SELECT COUNT(*) FROM dm.sales_report WHERE tickets_sold < passengers_boarded;
|
||||||
|
|
||||||
-- Проверка route_performance: нет дублей по бизнес-ключу
|
-- Нет дублей по бизнес-ключу route_performance
|
||||||
SELECT route_bk, COUNT(*) FROM dm.route_performance GROUP BY route_bk HAVING COUNT(*) > 1;
|
SELECT route_bk, COUNT(*) FROM dm.route_performance GROUP BY route_bk HAVING COUNT(*) > 1;
|
||||||
```
|
```
|
||||||
|
|
||||||
|
|||||||
+110
-28
@@ -1,23 +1,21 @@
|
|||||||
# DAG `bookings_to_gp_ods`: `stg` -> `ods` в Greenplum
|
# DAG `bookings_to_gp_ods`: `stg` → `ods` в Greenplum
|
||||||
|
|
||||||
Этот DAG — учебный пример загрузки типизированного слоя **ODS** из уже подготовленного слоя **STG**.
|
Этот DAG — учебный пример загрузки типизированного слоя **ODS** из уже подготовленного слоя **STG**.
|
||||||
Логика простая и каноничная: **SCD1 UPSERT** (обновляем изменившиеся записи, вставляем новые) + DQ-проверки.
|
Логика каноничная: **TRUNCATE + INSERT** для snapshot-справочников, **SCD1 UPSERT** для транзакционных
|
||||||
|
таблиц + DQ-проверки.
|
||||||
|
|
||||||
## Что делает DAG
|
## Что делает DAG
|
||||||
|
|
||||||
- Определяет `stg_batch_id`:
|
- Определяет `stg_batch_id` — последний согласованный батч, по которому все 4 snapshot-справочника
|
||||||
- берёт из `dag_run.conf["stg_batch_id"]`, если передан;
|
(`airports`, `airplanes`, `routes`, `seats`) уже приехали в STG.
|
||||||
- иначе берёт последний **согласованный** `_load_id`, который есть во всех snapshot-таблицах STG
|
- Загружает 9 таблиц ODS: `bookings`, `tickets`, `airports`, `airplanes`, `routes`, `seats`,
|
||||||
(`airports`, `airplanes`, `routes`, `seats`).
|
`flights`, `segments`, `boarding_passes`.
|
||||||
- Загружает 9 таблиц ODS (`airports`, `airplanes`, `routes`, `seats`, `bookings`, `tickets`,
|
- Для каждой таблицы выполняет пару задач `load → dq`.
|
||||||
`flights`, `segments`, `boarding_passes`).
|
- Snapshot-справочники фильтруются по `stg_batch_id`, транзакционные таблицы — по HWM (`_load_ts`).
|
||||||
- Для каждой таблицы выполняет пару задач `load -> dq`.
|
- Для snapshot-справочников дополнительно синхронизирует ключи (удаляет из ODS записи,
|
||||||
- На загрузке использует дедупликацию внутри батча + UPSERT (SCD1).
|
отсутствующие в выбранном STG-батче).
|
||||||
- Для snapshot-справочников (`airports`, `airplanes`, `routes`, `seats`) дополнительно
|
|
||||||
синхронизирует ключи (удаляет из ODS записи, отсутствующие в выбранном STG-батче).
|
|
||||||
- Для `flights` дополнительно добирает рейсы из истории `stg.flights`, если на них
|
- Для `flights` дополнительно добирает рейсы из истории `stg.flights`, если на них
|
||||||
ссылаются `stg.segments` выбранного батча (чтобы сохранить ссылочную целостность
|
ссылаются `stg.segments` (чтобы сохранить ссылочную целостность `segments.flight_id → flights.flight_id`).
|
||||||
`segments.flight_id -> flights.flight_id`).
|
|
||||||
|
|
||||||
## Что должно быть готово перед запуском
|
## Что должно быть готово перед запуском
|
||||||
|
|
||||||
@@ -35,7 +33,7 @@ make up
|
|||||||
3) ODS-таблицы созданы (один из вариантов):
|
3) ODS-таблицы созданы (один из вариантов):
|
||||||
|
|
||||||
- учебный: запустить DAG `bookings_ods_ddl`;
|
- учебный: запустить DAG `bookings_ods_ddl`;
|
||||||
- шорткат: `make ddl-gp` (в этом проекте он создаёт и STG, и ODS).
|
- шорткат: `make ddl-gp` (создаёт и STG, и ODS).
|
||||||
|
|
||||||
## Как запустить
|
## Как запустить
|
||||||
|
|
||||||
@@ -49,20 +47,103 @@ make up
|
|||||||
|
|
||||||
Если конфиг не передан, DAG автоматически возьмёт последний согласованный snapshot-батч.
|
Если конфиг не передан, DAG автоматически возьмёт последний согласованный snapshot-батч.
|
||||||
|
|
||||||
## Граф зависимостей (упрощённо)
|
## Граф зависимостей
|
||||||
|
|
||||||
- `resolve_stg_batch_id`
|
```
|
||||||
- Параллельно стартуют ветки:
|
resolve_stg_batch_id
|
||||||
- `bookings -> tickets`
|
├─ load_ods_bookings → dq_ods_bookings
|
||||||
- `airports`
|
│ └─ load_ods_tickets → dq_ods_tickets ──────────────────┐
|
||||||
- `airplanes`
|
│ │
|
||||||
- Далее:
|
├─ load_ods_airports → dq_ods_airports ─┐ │
|
||||||
- `routes` после `airports` и `airplanes`
|
│ ├─ load_ods_routes │
|
||||||
- `seats` после `airplanes`
|
├─ load_ods_airplanes → dq_ods_airplanes ─┤ └─ dq_ods_routes
|
||||||
- `flights` после `routes`
|
│ │ └─ load_ods_flights
|
||||||
- `segments` после `flights` и `tickets`
|
│ │ └─ dq_ods_flights ─┐
|
||||||
- `boarding_passes` после `segments`
|
│ │ │
|
||||||
- Финал: `finish_ods_summary` ждёт `dq_ods_boarding_passes` и `dq_ods_seats`.
|
│ │ dq_ods_flights + dq_ods_tickets
|
||||||
|
│ │ └─ load_ods_segments
|
||||||
|
│ │ └─ dq_ods_segments
|
||||||
|
│ │ └─ load_ods_boarding_passes
|
||||||
|
│ │ └─ dq_ods_boarding_passes ─┐
|
||||||
|
│ │ │
|
||||||
|
│ └─ load_ods_seats │
|
||||||
|
│ └─ dq_ods_seats ─────────────────────────────────┤
|
||||||
|
│ │
|
||||||
|
└────────────────────────────────────────────────────────────────── finish_ods_summary ◀──────────┘
|
||||||
|
```
|
||||||
|
|
||||||
|
Ветка `seats` работает параллельно с веткой `routes → flights → segments → boarding_passes`.
|
||||||
|
Обе ветки сходятся на `finish_ods_summary`.
|
||||||
|
|
||||||
|
## Как это работает внутри (по шагам)
|
||||||
|
|
||||||
|
### 1) `resolve_stg_batch_id` (Python)
|
||||||
|
|
||||||
|
Определяет, какой STG-батч использовать для snapshot-справочников.
|
||||||
|
Если `stg_batch_id` не передан через `dag_run.conf`, ищет последний **согласованный** батч —
|
||||||
|
`_load_id`, который есть одновременно во всех четырёх snapshot-таблицах
|
||||||
|
(`stg.airports`, `stg.airplanes`, `stg.routes`, `stg.seats`).
|
||||||
|
Для этого используется **INTERSECT** по `_load_id`.
|
||||||
|
|
||||||
|
> **Зачем согласованность?** Чтобы ODS загружал только те данные, для которых приехали
|
||||||
|
> ВСЕ связанные справочники. Иначе возможна потеря ссылочной целостности при сборке витрин.
|
||||||
|
|
||||||
|
### 2–10) Загрузка 9 таблиц: `load_ods_*` → `dq_ods_*`
|
||||||
|
|
||||||
|
Каждая пара задач работает одинаково:
|
||||||
|
|
||||||
|
| # | Задача | SQL-файл | Тип загрузки |
|
||||||
|
|---|--------|----------|--------------|
|
||||||
|
| 2 | `load_ods_bookings` → `dq_ods_bookings` | `sql/ods/bookings_load.sql`, `sql/ods/bookings_dq.sql` | HWM (инкремент) |
|
||||||
|
| 3 | `load_ods_tickets` → `dq_ods_tickets` | `sql/ods/tickets_load.sql`, `sql/ods/tickets_dq.sql` | HWM (инкремент) |
|
||||||
|
| 4 | `load_ods_airports` → `dq_ods_airports` | `sql/ods/airports_load.sql`, `sql/ods/airports_dq.sql` | snapshot по `stg_batch_id` |
|
||||||
|
| 5 | `load_ods_airplanes` → `dq_ods_airplanes` | `sql/ods/airplanes_load.sql`, `sql/ods/airplanes_dq.sql` | snapshot по `stg_batch_id` |
|
||||||
|
| 6 | `load_ods_routes` → `dq_ods_routes` | `sql/ods/routes_load.sql`, `sql/ods/routes_dq.sql` | snapshot по `stg_batch_id` |
|
||||||
|
| 7 | `load_ods_seats` → `dq_ods_seats` | `sql/ods/seats_load.sql`, `sql/ods/seats_dq.sql` | snapshot по `stg_batch_id` |
|
||||||
|
| 8 | `load_ods_flights` → `dq_ods_flights` | `sql/ods/flights_load.sql`, `sql/ods/flights_dq.sql` | HWM (инкремент) |
|
||||||
|
| 9 | `load_ods_segments` → `dq_ods_segments` | `sql/ods/segments_load.sql`, `sql/ods/segments_dq.sql` | HWM (инкремент) |
|
||||||
|
| 10 | `load_ods_boarding_passes` → `dq_ods_boarding_passes` | `sql/ods/boarding_passes_load.sql`, `sql/ods/boarding_passes_dq.sql` | HWM (инкремент) |
|
||||||
|
|
||||||
|
### Два паттерна загрузки
|
||||||
|
|
||||||
|
В ODS используются **два разных паттерна** — выбор зависит от типа данных и формата хранения:
|
||||||
|
|
||||||
|
**Snapshot-справочники** (`airports`, `airplanes`, `routes`, `seats`) — **TRUNCATE + INSERT**:
|
||||||
|
|
||||||
|
1. **TRUNCATE** — полная очистка таблицы.
|
||||||
|
2. **INSERT** — вставка всех строк из STG-батча (`_load_id = stg_batch_id`) с дедупликацией
|
||||||
|
через `ROW_NUMBER()`.
|
||||||
|
|
||||||
|
> Почему не UPSERT? Эти таблицы хранятся в формате **AO Row** (`appendonly=true`),
|
||||||
|
> который не поддерживает эффективный row-level UPDATE (вызывает bloat).
|
||||||
|
> Для маленьких справочников (~100–300 строк) полная перезагрузка быстрее и чище.
|
||||||
|
|
||||||
|
**Транзакционные таблицы** (`bookings`, `tickets`, `flights`, `segments`, `boarding_passes`) —
|
||||||
|
**SCD1 UPSERT**:
|
||||||
|
|
||||||
|
1. **TEMP TABLE** — собирает дельту (новые/изменённые строки) с дедупликацией внутри батча
|
||||||
|
через `ROW_NUMBER()`. Временная таблица автоматически удаляется (`ON COMMIT DROP`).
|
||||||
|
2. **UPDATE** — обновляет существующие строки. Использует `IS DISTINCT FROM` для корректного
|
||||||
|
сравнения `NULL`-значений (обычный `<>` не обнаружит изменение `NULL → значение`).
|
||||||
|
3. **INSERT** — добавляет новые строки (которых нет в ODS по бизнес-ключу).
|
||||||
|
|
||||||
|
Транзакционные таблицы фильтруются по **HWM** — `WHERE _load_ts > (SELECT MAX(_load_ts) FROM ods.table)`.
|
||||||
|
Это сделано, чтобы не потерять инкременты, если STG-DAG запускался несколько раз
|
||||||
|
до запуска ODS-DAG'а.
|
||||||
|
|
||||||
|
### DQ-проверки (одинаковый паттерн)
|
||||||
|
|
||||||
|
Каждый `*_dq.sql` — PL/pgSQL-блок (`DO $$...$$`), который проверяет:
|
||||||
|
- нет дублей по бизнес-ключу в ODS;
|
||||||
|
- все ключи из STG текущего батча присутствуют в ODS;
|
||||||
|
- обязательные поля не содержат NULL.
|
||||||
|
|
||||||
|
При ошибке — `RAISE EXCEPTION` с понятным текстом. Для инкрементальных таблиц пустой батч допустим.
|
||||||
|
|
||||||
|
### 11) `finish_ods_summary`
|
||||||
|
|
||||||
|
Ждёт завершения обеих параллельных веток (`dq_ods_boarding_passes` и `dq_ods_seats`)
|
||||||
|
и логирует краткую сводку.
|
||||||
|
|
||||||
## Как проверить результат
|
## Как проверить результат
|
||||||
|
|
||||||
@@ -75,6 +156,7 @@ SELECT COUNT(*) FROM ods.bookings;
|
|||||||
SELECT COUNT(*) FROM ods.tickets;
|
SELECT COUNT(*) FROM ods.tickets;
|
||||||
SELECT COUNT(*) FROM ods.flights;
|
SELECT COUNT(*) FROM ods.flights;
|
||||||
|
|
||||||
|
-- Проверка: в ODS не должно быть дублей по бизнес-ключу
|
||||||
SELECT book_ref, COUNT(*)
|
SELECT book_ref, COUNT(*)
|
||||||
FROM ods.bookings
|
FROM ods.bookings
|
||||||
GROUP BY 1
|
GROUP BY 1
|
||||||
|
|||||||
Reference in New Issue
Block a user