feat(ods): реализован ODS слой и DAG загрузки из STG

- Зачем:
  - подготовлена учебная реализация ODS слоя с типизацией, UPSERT и DQ, чтобы продолжить работу от STG к DDS/DM.
- Что:
  - добавлены SQL-скрипты `sql/ods/*_ddl.sql`, `sql/ods/*_load.sql`, `sql/ods/*_dq.sql` для 9 сущностей bookings.
  - добавлены DAG `bookings_ods_ddl` и `bookings_to_gp_ods`, а также smoke-тесты для новых графов.
  - ODS DDL интегрирован в `sql/ddl_gp.sql`; документация и план обновлены под единый запуск через `make ddl-gp`.
- Проверка:
  - `make test`.
  - `make ddl-gp`.
This commit is contained in:
2026-02-28 22:13:11 +03:00
parent 2d629b6e80
commit 85fea19818
37 changed files with 2704 additions and 27 deletions
+2
View File
@@ -8,6 +8,7 @@
- [Учебные задания](../educational-tasks.md)
- [План тестирования и проверки](../TESTING.md)
- [Главный учебный DAG: bookings → stg](bookings_to_gp_stage.md)
- [Учебный DAG: stg -> ods](bookings_to_gp_ods.md)
## Технические детали (опционально)
@@ -15,4 +16,5 @@
- [Единые конвенции нейминга DWH (служебные поля и SCD)](internal/naming_conventions.md)
- [PXF в этом проекте (проектная реализация)](internal/pxf_bookings.md)
- [Дизайн stg для bookings (черновик)](internal/bookings_stg_design.md)
- [Дизайн ods для bookings (черновик)](internal/bookings_ods_design.md)
- [Про время/UTC в bookings (черновик)](internal/bookings_tz.md)
+91
View File
@@ -0,0 +1,91 @@
# DAG `bookings_to_gp_ods`: `stg` -> `ods` в Greenplum
Этот DAG — учебный пример загрузки типизированного слоя **ODS** из уже подготовленного слоя **STG**.
Логика простая и каноничная: **SCD1 UPSERT** (обновляем изменившиеся записи, вставляем новые) + DQ-проверки.
## Что делает DAG
- Определяет `stg_batch_id`:
- берёт из `dag_run.conf["stg_batch_id"]`, если передан;
- иначе берёт последний **согласованный** `batch_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-батче).
## Что должно быть готово перед запуском
1) Стек поднят:
```bash
make up
```
2) STG-слой создан и заполнен:
- запущен `bookings_stg_ddl` (или `make ddl-gp`);
- хотя бы один раз выполнен DAG `bookings_to_gp_stage`.
3) ODS-таблицы созданы (один из вариантов):
- учебный: запустить DAG `bookings_ods_ddl`;
- шорткат: `make ddl-gp` (в этом проекте он создаёт и STG, и ODS).
## Как запустить
1) Откройте Airflow UI: http://localhost:8080.
2) Запустите DAG `bookings_to_gp_ods`.
3) (Опционально) передайте `stg_batch_id` в конфиге запуска:
```json
{"stg_batch_id": "manual__2026-02-22T12:00:00+00:00"}
```
Если конфиг не передан, 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`.
## Как проверить результат
```bash
make gp-psql
```
```sql
SELECT COUNT(*) FROM ods.bookings;
SELECT COUNT(*) FROM ods.tickets;
SELECT COUNT(*) FROM ods.flights;
SELECT book_ref, COUNT(*)
FROM ods.bookings
GROUP BY 1
HAVING COUNT(*) > 1;
```
Ожидаемо: в последнем запросе `0` строк.
## Типичные ошибки
- `stg_batch_id не найден`:
- передайте `stg_batch_id` в `dag_run.conf`, или
- сначала загрузите STG через `bookings_to_gp_stage`.
- Ошибки `relation "ods...." does not exist`:
- не применён ODS DDL (`bookings_ods_ddl` / `make ddl-gp`).
- Ошибки DQ по ссылочной целостности:
- проверьте, что ODS DAG выполнялся с корректным `stg_batch_id` и без пропуска upstream задач.
+57 -15
View File
@@ -26,6 +26,8 @@ ODS в учебном проекте — это:
- приведение типов (`TEXT -> TIMESTAMPTZ/NUMERIC/INT/BOOLEAN/...`);
- дедупликацию внутри батча;
- `UPSERT` (SCD Type 1): обновляем текущую запись при изменении, вставляем новые.
- для snapshot-справочников (`airports`, `airplanes`, `routes`, `seats`) синхронизацию ключей:
удаляем из ODS записи, которых нет в выбранном `stg_batch_id`.
Не делаем в ODS (в базовом эталоне):
- SCD Type 2 с периодами действия;
@@ -264,23 +266,47 @@ log = logging.getLogger(__name__)
GREENPLUM_CONN_ID = "greenplum_conn"
def _resolve_stg_batch_id(**context):
"""Определяем stg_batch_id: из dag_run.conf или последний загруженный в STG."""
"""Определяем stg_batch_id: из dag_run.conf или последний согласованный snapshot-батч."""
conf = context["dag_run"].conf or {}
stg_batch_id = conf.get("stg_batch_id")
if not stg_batch_id:
# Берём batch_id с самым свежим load_dttm (TIMESTAMP, монотонно растёт).
# MAX(batch_id) ненадёжен: run_id — строка вида "manual__2024-...",
# лексикографическая сортировка не гарантирует хронологический порядок.
# Берём batch_id, который присутствует во всех snapshot-таблицах STG:
# airports, airplanes, routes, seats. Это защищает от частично успешных запусков.
hook = PostgresHook(postgres_conn_id=GREENPLUM_CONN_ID)
result = hook.get_first(
"SELECT batch_id FROM stg.bookings ORDER BY load_dttm DESC LIMIT 1"
'''
WITH candidate_batches AS (
SELECT batch_id FROM stg.airports WHERE batch_id IS NOT NULL GROUP BY batch_id
INTERSECT
SELECT batch_id FROM stg.airplanes WHERE batch_id IS NOT NULL GROUP BY batch_id
INTERSECT
SELECT batch_id FROM stg.routes WHERE batch_id IS NOT NULL GROUP BY batch_id
INTERSECT
SELECT batch_id FROM stg.seats WHERE batch_id IS NOT NULL GROUP BY batch_id
),
batch_ready AS (
SELECT
c.batch_id,
GREATEST(
(SELECT MAX(load_dttm) FROM stg.airports a WHERE a.batch_id = c.batch_id),
(SELECT MAX(load_dttm) FROM stg.airplanes a WHERE a.batch_id = c.batch_id),
(SELECT MAX(load_dttm) FROM stg.routes r WHERE r.batch_id = c.batch_id),
(SELECT MAX(load_dttm) FROM stg.seats s WHERE s.batch_id = c.batch_id)
) AS ready_dttm
FROM candidate_batches c
)
SELECT batch_id
FROM batch_ready
ORDER BY ready_dttm DESC
LIMIT 1
'''
)
stg_batch_id = result[0] if result and result[0] else None
if not stg_batch_id:
raise ValueError(
"stg_batch_id не найден: передайте в conf или сначала загрузите STG"
"stg_batch_id не найден: передайте в conf или сначала выполните bookings_to_gp_stage"
)
log.info("Используем stg_batch_id = %s", stg_batch_id)
@@ -382,6 +408,19 @@ WHERE NOT EXISTS (
WHERE o.airport_code = s.airport_code
);
-- Statement 3: DELETE ключей, которых нет в snapshot текущего батча
WITH src_keys AS (
SELECT DISTINCT airport_code
FROM stg.airports
WHERE batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text
)
DELETE FROM ods.airports o
WHERE NOT EXISTS (
SELECT 1
FROM src_keys s
WHERE s.airport_code = o.airport_code
);
-- Обновляем статистику для оптимизатора запросов Greenplum
ANALYZE ods.airports;
```
@@ -454,7 +493,10 @@ ANALYZE ods.bookings;
### 6.3. Идемпотентность паттерна
Паттерн UPDATE + INSERT WHERE NOT EXISTS — **натурально идемпотентен**: повторный запуск с тем же `stg_batch_id` не создаст дублей и не потеряет данные. UPDATE обновит только если атрибуты изменились, INSERT вставит только если бизнес-ключа нет. Это одно из преимуществ подхода.
- Для инкрементальных таблиц паттерн `UPDATE + INSERT WHERE NOT EXISTS`**натурально идемпотентен**:
повторный запуск с тем же `stg_batch_id` не создаст дублей и не потеряет данные.
- Для snapshot-справочников идемпотентность сохраняется паттерном
`UPDATE + INSERT + DELETE not in snapshot`: повторный запуск приводит ODS к тому же состоянию.
### 6.4. Поведение при пустом батче
@@ -541,14 +583,14 @@ sql/ods/
├── boarding_passes_load.sql
└── boarding_passes_dq.sql
sql/ddl_gp_ods.sql
sql/ddl_gp.sql (+ подключение sql/ods/*_ddl.sql)
airflow/dags/
├── bookings_ods_ddl.py
└── bookings_to_gp_ods.py
docs/bookings_to_gp_ods.md
Makefile (+ ddl-gp-ods)
Makefile (ddl-gp включает ODS DDL)
tests/test_dags_smoke.py (+ smoke для 2 новых DAG)
```
@@ -591,8 +633,8 @@ resolve_stg_batch_id
## 10) Порядок реализации
1. Подготовить DDL в `sql/ods/*_ddl.sql`.
2. Сделать мастер-скрипт `sql/ddl_gp_ods.sql`.
3. Добавить `Makefile`-таргет `ddl-gp-ods`.
2. Подключить `sql/ods/*_ddl.sql` в общий `sql/ddl_gp.sql`.
3. Использовать существующий `Makefile`-таргет `ddl-gp` для STG+ODS.
4. Создать DAG `bookings_ods_ddl.py`.
5. Реализовать `sql/ods/*_load.sql` (SCD1 UPSERT).
6. Реализовать `sql/ods/*_dq.sql`.
@@ -607,7 +649,7 @@ resolve_stg_batch_id
Готово, если:
1. Оба новых DAG парсятся и проходят smoke-тесты (`make test`).
2. `make ddl-gp-ods` создаёт объекты без ошибок.
2. `make ddl-gp` создаёт объекты STG+ODS без ошибок.
3. Для тестового `stg_batch_id` ODS-загрузка завершается успешно.
4. Все DQ-задачи зелёные и реально валят DAG при искусственной ошибке.
5. В ODS нет дублей по бизнес-ключам.
@@ -619,10 +661,10 @@ resolve_stg_batch_id
```bash
make up
make ddl-gp # создать STG-объекты
make ddl-gp-ods # создать ODS-объекты
make ddl-gp # создать STG+ODS-объекты
# Trigger bookings_to_gp_stage (загрузить STG)
# Получить batch_id: SELECT batch_id FROM stg.bookings ORDER BY load_dttm DESC LIMIT 1;
# Передать stg_batch_id в conf (рекомендуется) или дать ODS DAG выбрать
# последний согласованный batch автоматически.
# Trigger bookings_to_gp_ods с conf: {"stg_batch_id": "<значение>"}
make gp-psql
```