diff --git a/docs/archive/docs-etl-quality-plan.md b/docs/archive/docs-etl-quality-plan.md new file mode 100644 index 0000000..51545da --- /dev/null +++ b/docs/archive/docs-etl-quality-plan.md @@ -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". diff --git a/docs/bookings_to_gp_dds.md b/docs/bookings_to_gp_dds.md index 9cb0c75..d2f28b7 100644 --- a/docs/bookings_to_gp_dds.md +++ b/docs/bookings_to_gp_dds.md @@ -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**. -Логика: измерения + факт, проверки качества данных после каждой загрузки. +Здесь сосредоточены ключевые паттерны аналитического хранилища: SCD1, SCD2 с hashdiff, +point-in-time join, защитные LEFT JOIN для устойчивости к data quality аномалиям. ## Что делает DAG -- Загружает измерения DDS: - - `dds.dim_calendar` (статическое измерение дат); - - `dds.dim_airports`, `dds.dim_airplanes`, `dds.dim_tariffs`, `dds.dim_passengers` (SCD1 UPSERT); - - `dds.dim_routes` (SCD2 с `hashdiff`, `valid_from`, `valid_to`). -- Загружает факт `dds.fact_flight_sales` инкрементальным UPSERT по зерну `(ticket_no, flight_id)`. -- Для каждой таблицы выполняет пару задач `load -> dq`. -- Использует `_load_id = {{ run_id }}` (DDS не требует `stg_batch_id`, потому что читает current state ODS). +- Загружает 6 измерений DDS: + - `dds.dim_calendar` — статическое измерение дат (Full Rebuild); + - `dds.dim_airports`, `dds.dim_airplanes`, `dds.dim_tariffs`, `dds.dim_passengers` — SCD1 UPSERT; + - `dds.dim_routes` — **SCD2** с `hashdiff`, `valid_from`, `valid_to` + денормализация. +- Загружает факт `dds.fact_flight_sales` — инкрементальный UPSERT по зерну `(ticket_no, flight_id)`. +- Для каждой таблицы выполняет пару задач `load → dq`. +- Использует `_load_id = {{ run_id }}`. DDS не требует `stg_batch_id`, потому что читает + текущее состояние ODS. ## Что должно быть готово перед запуском @@ -32,17 +34,116 @@ make up 2) Если запускаете DDS впервые — выполните `bookings_dds_ddl`. 3) Запустите `bookings_to_gp_dds`. -## Граф зависимостей (упрощённо) +## Граф зависимостей -- `load_dds_dim_calendar -> dq_dds_dim_calendar` -- После calendar параллельно: - - `load_dds_dim_airports -> dq_dds_dim_airports` - - `load_dds_dim_airplanes -> dq_dds_dim_airplanes` - - `load_dds_dim_tariffs -> dq_dds_dim_tariffs` - - `load_dds_dim_passengers -> dq_dds_dim_passengers` - - `load_dds_dim_routes -> dq_dds_dim_routes` -- Факт: - - `load_dds_fact_flight_sales -> dq_dds_fact_flight_sales -> finish_dds_summary` +``` +load_dds_dim_calendar → dq_dds_dim_calendar + ├─ load_dds_dim_airports → dq_dds_dim_airports ─┐ + │ ├─ load_dds_dim_routes + ├─ load_dds_dim_airplanes → dq_dds_dim_airplanes ─┘ └─ dq_dds_dim_routes + │ │ + ├─ load_dds_dim_tariffs → dq_dds_dim_tariffs │ + │ │ + └─ 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.fact_flight_sales; +-- Проверка: кол-во строк факта ≈ кол-во строк ODS segments SELECT (SELECT COUNT(*) FROM dds.fact_flight_sales) AS fact_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`. ## Типичные ошибки diff --git a/docs/bookings_to_gp_dm.md b/docs/bookings_to_gp_dm.md index 07c2c66..d635709 100644 --- a/docs/bookings_to_gp_dm.md +++ b/docs/bookings_to_gp_dm.md @@ -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**. -Логика: все 5 витрин загружаются параллельно, для каждой — пара `load -> dq`. +Все 5 витрин загружаются **параллельно** и демонстрируют разные стратегии загрузки — +это ключевая учебная ценность данного DAG. ## Что делает DAG -- Загружает витрины DM параллельно (паттерны загрузки разные — учебная демонстрация выбора стратегии): - - `dm.sales_report` — UPSERT по датам; DQ проверяет только строки текущего `run_id` (`_load_id`); - - `dm.route_performance` — Full Rebuild (TRUNCATE + INSERT): таблица маленькая, дельту считать дороже; - - `dm.passenger_loyalty` — инкрементальный UPSERT по «затронутым ключам» (HWM по `_load_ts`): пересчитываем агрегаты только для пассажиров с новыми фактами; - - `dm.airport_traffic` — инкрементальный UPSERT по датам (HWM по `_load_ts`); - - `dm.monthly_overview` — инкрементальный UPSERT по месяцам (HWM по `_load_ts`). -- Для каждой витрины выполняет пару задач `load -> dq`. +Загружает 5 витрин параллельно, для каждой — пара `load → dq`: + +| Витрина | Зерно (grain) | Паттерн загрузки | +|---------|---------------|------------------| +| `dm.sales_report` | (flight_date, departure_airport_sk, arrival_airport_sk, tariff_sk) | Инкрементальный UPSERT (HWM по датам) | +| `dm.route_performance` | route_bk | Full Rebuild (TRUNCATE + INSERT) | +| `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 -├── load_dm_sales_report -> dq_dm_sales_report -> finish_dm_summary -├── load_dm_route_performance -> dq_dm_route_performance -> 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_monthly_overview -> dq_dm_monthly_overview -> finish_dm_summary +├── load_dm_sales_report → dq_dm_sales_report ─┐ +├── load_dm_route_performance → dq_dm_route_performance ─┤ +├── load_dm_passenger_loyalty → dq_dm_passenger_loyalty ─┼─ finish_dm_summary +├── load_dm_airport_traffic → dq_dm_airport_traffic ─┤ +└── 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 @@ -58,10 +150,10 @@ SELECT COUNT(*) FROM dm.passenger_loyalty; SELECT COUNT(*) FROM dm.airport_traffic; SELECT COUNT(*) FROM dm.monthly_overview; --- Проверка инварианта sales_report: посаженных не больше, чем продано +-- Инвариант sales_report: посаженных не больше, чем продано 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; ``` diff --git a/docs/bookings_to_gp_ods.md b/docs/bookings_to_gp_ods.md index ecb47d3..3ecc718 100644 --- a/docs/bookings_to_gp_ods.md +++ b/docs/bookings_to_gp_ods.md @@ -1,23 +1,21 @@ -# DAG `bookings_to_gp_ods`: `stg` -> `ods` в Greenplum +# DAG `bookings_to_gp_ods`: `stg` → `ods` в Greenplum Этот DAG — учебный пример загрузки типизированного слоя **ODS** из уже подготовленного слоя **STG**. -Логика простая и каноничная: **SCD1 UPSERT** (обновляем изменившиеся записи, вставляем новые) + DQ-проверки. +Логика каноничная: **TRUNCATE + INSERT** для snapshot-справочников, **SCD1 UPSERT** для транзакционных +таблиц + DQ-проверки. ## Что делает DAG -- Определяет `stg_batch_id`: - - берёт из `dag_run.conf["stg_batch_id"]`, если передан; - - иначе берёт последний **согласованный** `_load_id`, который есть во всех snapshot-таблицах STG - (`airports`, `airplanes`, `routes`, `seats`). -- Загружает 9 таблиц ODS (`airports`, `airplanes`, `routes`, `seats`, `bookings`, `tickets`, - `flights`, `segments`, `boarding_passes`). -- Для каждой таблицы выполняет пару задач `load -> dq`. -- На загрузке использует дедупликацию внутри батча + UPSERT (SCD1). -- Для snapshot-справочников (`airports`, `airplanes`, `routes`, `seats`) дополнительно - синхронизирует ключи (удаляет из ODS записи, отсутствующие в выбранном STG-батче). +- Определяет `stg_batch_id` — последний согласованный батч, по которому все 4 snapshot-справочника + (`airports`, `airplanes`, `routes`, `seats`) уже приехали в STG. +- Загружает 9 таблиц ODS: `bookings`, `tickets`, `airports`, `airplanes`, `routes`, `seats`, + `flights`, `segments`, `boarding_passes`. +- Для каждой таблицы выполняет пару задач `load → dq`. +- Snapshot-справочники фильтруются по `stg_batch_id`, транзакционные таблицы — по HWM (`_load_ts`). +- Для snapshot-справочников дополнительно синхронизирует ключи (удаляет из ODS записи, + отсутствующие в выбранном STG-батче). - Для `flights` дополнительно добирает рейсы из истории `stg.flights`, если на них - ссылаются `stg.segments` выбранного батча (чтобы сохранить ссылочную целостность - `segments.flight_id -> flights.flight_id`). + ссылаются `stg.segments` (чтобы сохранить ссылочную целостность `segments.flight_id → flights.flight_id`). ## Что должно быть готово перед запуском @@ -35,7 +33,7 @@ make up 3) ODS-таблицы созданы (один из вариантов): - учебный: запустить DAG `bookings_ods_ddl`; -- шорткат: `make ddl-gp` (в этом проекте он создаёт и STG, и ODS). +- шорткат: `make ddl-gp` (создаёт и STG, и ODS). ## Как запустить @@ -49,20 +47,103 @@ make up Если конфиг не передан, DAG автоматически возьмёт последний согласованный snapshot-батч. -## Граф зависимостей (упрощённо) +## Граф зависимостей -- `resolve_stg_batch_id` -- Параллельно стартуют ветки: - - `bookings -> tickets` - - `airports` - - `airplanes` -- Далее: - - `routes` после `airports` и `airplanes` - - `seats` после `airplanes` - - `flights` после `routes` - - `segments` после `flights` и `tickets` - - `boarding_passes` после `segments` -- Финал: `finish_ods_summary` ждёт `dq_ods_boarding_passes` и `dq_ods_seats`. +``` +resolve_stg_batch_id + ├─ load_ods_bookings → dq_ods_bookings + │ └─ load_ods_tickets → dq_ods_tickets ──────────────────┐ + │ │ + ├─ load_ods_airports → dq_ods_airports ─┐ │ + │ ├─ load_ods_routes │ + ├─ load_ods_airplanes → dq_ods_airplanes ─┤ └─ dq_ods_routes + │ │ └─ load_ods_flights + │ │ └─ dq_ods_flights ─┐ + │ │ │ + │ │ 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.flights; +-- Проверка: в ODS не должно быть дублей по бизнес-ключу SELECT book_ref, COUNT(*) FROM ods.bookings GROUP BY 1