diff --git a/README.md b/README.md index 2b16b5b..2cdb372 100644 --- a/README.md +++ b/README.md @@ -12,9 +12,9 @@ - построения ETL/ELT; - работы с Airflow и Greenplum. -В курсовой у нас один источник данных — демо‑БД **bookings**. В стенде уже есть готовый учебный пример -загрузки **bookings → stg в Greenplum**, чтобы вы могли сфокусироваться на DWH‑части (ODS/DDS/DM) и -не тратить время на инфраструктуру. +В курсовой у нас один источник данных — демо‑БД **bookings**. В стенде уже есть готовые учебные примеры +загрузки **bookings → stg** и **stg -> ods** в Greenplum, чтобы вы могли сфокусироваться на DWH‑части +(ODS/DDS/DM) и не тратить время на инфраструктуру. ## Что внутри @@ -44,7 +44,7 @@ Про PXF и технические детали стенда: [docs/stack.md](docs/stack.md). -## Быстрый старт (основной сценарий: bookings → stg) +## Быстрый старт (основной сценарий: bookings -> stg -> ods) 1) Скопируйте настройки: @@ -65,16 +65,16 @@ make up make bookings-init ``` -Важно: генератор `bookings` в этом стенде поддерживается только в режиме `BOOKINGS_JOBS=1`. +4) Подготовьте STG/ODS-объекты в Greenplum (выберите один вариант): -4) Подготовьте STG‑объекты в Greenplum (выберите один вариант): - -- Учебный вариант: в Airflow UI запустите DAG `bookings_stg_ddl`; -- Технический шорткат: `make ddl-gp` (применяет все DDL разом вручную). +- Учебный вариант: в Airflow UI запустите DAG `bookings_stg_ddl`, затем `bookings_ods_ddl`; +- Технический шорткат: `make ddl-gp` (применяет DDL для STG и ODS разом вручную). 5) Запустите основной DAG `bookings_to_gp_stage`. -6) Проверьте результат в Greenplum: +6) Запустите DAG `bookings_to_gp_ods`. + +7) Проверьте результат в Greenplum: ```bash make gp-psql @@ -82,6 +82,8 @@ make gp-psql SELECT COUNT(*) FROM stg.bookings; SELECT COUNT(*) FROM stg.tickets; SELECT * FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10; +SELECT COUNT(*) FROM ods.bookings; +SELECT COUNT(*) FROM ods.tickets; ``` Подробнее про логику DAG и проверки — `docs/bookings_to_gp_stage.md`. @@ -90,10 +92,11 @@ SELECT * FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10; Основные (для потока bookings → DWH): -- `bookings_stg_ddl` — создаёт `stg.bookings_ext`/`stg.bookings` и `stg.tickets_ext`/`stg.tickets` в Greenplum; - `bookings_stg_ddl` — создаёт/обновляет весь STG слой для bookings (9 таблиц: bookings, tickets, airports, airplanes, routes, seats, flights, segments, boarding_passes; включая внешние `*_ext` через PXF); - `bookings_to_gp_stage` — генерирует учебный день в `bookings-db`, затем загружает данные в STG и выполняет DQ‑проверки. +- `bookings_ods_ddl` — создаёт/обновляет ODS-таблицы по домену bookings. +- `bookings_to_gp_ods` — загружает данные из STG в ODS (SCD1 UPSERT) и выполняет DQ‑проверки. Вспомогательные (побочный трек с CSV): @@ -108,7 +111,7 @@ make up # поднять стек make logs # логи airflow-webserver и airflow-scheduler make gp-psql # psql в Greenplum make bookings-psql # psql в демо-БД bookings (Postgres) -make ddl-gp # применить DDL к Greenplum вручную (вместо DDL-DAG) +make ddl-gp # применить DDL STG+ODS к Greenplum вручную (вместо DDL-DAG) make down # остановить и удалить контейнеры/сети (volumes сохраняются) make clean # полный reset: удалить контейнеры/сети и volumes (данные будут потеряны) ``` @@ -138,6 +141,7 @@ make clean # полный reset: удалить контейнер - Учебные задания: `educational-tasks.md`). - План тестирования/проверок и негативные кейсы: `TESTING.md`. - Дополнительные заметки и технические детали: `docs/README.md`. +- Детали по ODS DAG: `docs/bookings_to_gp_ods.md`. ## Типичные проблемы и решения diff --git a/airflow/dags/bookings_ods_ddl.py b/airflow/dags/bookings_ods_ddl.py new file mode 100644 index 0000000..80dde5b --- /dev/null +++ b/airflow/dags/bookings_ods_ddl.py @@ -0,0 +1,98 @@ +from __future__ import annotations + +""" +Учебный DAG: создаёт/обновляет слой ods в Greenplum для домена bookings. + +Запускается вручную перед DAG загрузки `bookings_to_gp_ods` или после изменения ODS DDL. +Создаёт 9 ODS-таблиц: airports, airplanes, routes, seats, bookings, tickets, +flights, segments, boarding_passes. +""" + +from datetime import timedelta + +import pendulum +from airflow.providers.postgres.operators.postgres import PostgresOperator + +from airflow import DAG + +GREENPLUM_CONN_ID = "greenplum_conn" + +default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)} + +with DAG( + dag_id="bookings_ods_ddl", + start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), + schedule=None, + catchup=False, + template_searchpath="/sql", + default_args=default_args, + tags=["demo", "greenplum", "ddl", "bookings", "ods"], + description="Учебный DDL DAG: создаёт/обновляет ods.* для bookings", +) as dag: + apply_ods_airports_ddl = PostgresOperator( + task_id="apply_ods_airports_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/airports_ddl.sql", + ) + + apply_ods_airplanes_ddl = PostgresOperator( + task_id="apply_ods_airplanes_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/airplanes_ddl.sql", + ) + + apply_ods_routes_ddl = PostgresOperator( + task_id="apply_ods_routes_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/routes_ddl.sql", + ) + + apply_ods_seats_ddl = PostgresOperator( + task_id="apply_ods_seats_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/seats_ddl.sql", + ) + + apply_ods_bookings_ddl = PostgresOperator( + task_id="apply_ods_bookings_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/bookings_ddl.sql", + ) + + apply_ods_tickets_ddl = PostgresOperator( + task_id="apply_ods_tickets_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/tickets_ddl.sql", + ) + + apply_ods_flights_ddl = PostgresOperator( + task_id="apply_ods_flights_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/flights_ddl.sql", + ) + + apply_ods_segments_ddl = PostgresOperator( + task_id="apply_ods_segments_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/segments_ddl.sql", + ) + + apply_ods_boarding_passes_ddl = PostgresOperator( + task_id="apply_ods_boarding_passes_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/boarding_passes_ddl.sql", + ) + + # DDL применяем последовательно, чтобы порядок был понятным для новичков, + # а ошибки — воспроизводимыми (в логах сразу видно, на каком объекте упали). + ( + apply_ods_airports_ddl + >> apply_ods_airplanes_ddl + >> apply_ods_routes_ddl + >> apply_ods_seats_ddl + >> apply_ods_bookings_ddl + >> apply_ods_tickets_ddl + >> apply_ods_flights_ddl + >> apply_ods_segments_ddl + >> apply_ods_boarding_passes_ddl + ) diff --git a/airflow/dags/bookings_to_gp_ods.py b/airflow/dags/bookings_to_gp_ods.py new file mode 100644 index 0000000..63d2204 --- /dev/null +++ b/airflow/dags/bookings_to_gp_ods.py @@ -0,0 +1,261 @@ +from __future__ import annotations + +""" +Учебный DAG: загрузка из STG в ODS (Greenplum) по домену bookings. + +Ключевая идея: +- весь запуск ODS работает с одним stg_batch_id; +- для каждой сущности выполняем пару задач load -> dq; +- загрузка реализована как SCD1 UPSERT (UPDATE изменившихся + INSERT новых). +""" + +from datetime import timedelta +from logging import getLogger + +import pendulum +from airflow.operators.python import PythonOperator +from airflow.providers.postgres.hooks.postgres import PostgresHook +from airflow.providers.postgres.operators.postgres import PostgresOperator + +from airflow import DAG + +GREENPLUM_CONN_ID = "greenplum_conn" + +log = getLogger(__name__) + +default_args = { + "owner": "airflow", + "retries": 1, + "retry_delay": timedelta(seconds=30), +} + + +def _resolve_stg_batch_id(**context) -> str: + """ + Возвращает stg_batch_id из dag_run.conf или вычисляет последний согласованный батч. + + Согласованным считаем batch_id, который есть во всех snapshot-справочниках STG: + airports, airplanes, routes, seats. + """ + conf = context["dag_run"].conf or {} + stg_batch_id = conf.get("stg_batch_id") + + if not stg_batch_id: + hook = PostgresHook(postgres_conn_id=GREENPLUM_CONN_ID) + result = hook.get_first( + """ + WITH candidate_batches AS ( + SELECT batch_id + FROM stg.airports + WHERE batch_id IS NOT NULL AND batch_id <> '' + GROUP BY batch_id + INTERSECT + SELECT batch_id + FROM stg.airplanes + WHERE batch_id IS NOT NULL AND batch_id <> '' + GROUP BY batch_id + INTERSECT + SELECT batch_id + FROM stg.routes + WHERE batch_id IS NOT NULL AND batch_id <> '' + GROUP BY batch_id + INTERSECT + SELECT batch_id + FROM stg.seats + WHERE batch_id IS NOT NULL AND batch_id <> '' + GROUP BY batch_id + ), + batch_ready AS ( + SELECT + c.batch_id, + GREATEST( + COALESCE( + (SELECT MAX(load_dttm) FROM stg.airports a WHERE a.batch_id = c.batch_id), + TIMESTAMP '1900-01-01 00:00:00' + ), + COALESCE( + (SELECT MAX(load_dttm) FROM stg.airplanes a WHERE a.batch_id = c.batch_id), + TIMESTAMP '1900-01-01 00:00:00' + ), + COALESCE( + (SELECT MAX(load_dttm) FROM stg.routes r WHERE r.batch_id = c.batch_id), + TIMESTAMP '1900-01-01 00:00:00' + ), + COALESCE( + (SELECT MAX(load_dttm) FROM stg.seats s WHERE s.batch_id = c.batch_id), + TIMESTAMP '1900-01-01 00:00:00' + ) + ) 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 не найден: передайте stg_batch_id в conf или " + "сначала выполните bookings_to_gp_stage для snapshot-справочников" + ) + + log.info("Используем stg_batch_id=%s", stg_batch_id) + return stg_batch_id + + +def _finish_summary() -> None: + """Логирует краткий итог выполнения ODS-ветки.""" + log.info("DAG bookings_to_gp_ods завершён. Подробности смотрите в логах задач.") + + +with DAG( + dag_id="bookings_to_gp_ods", + start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), + schedule=None, + catchup=False, + max_active_runs=1, + template_searchpath="/sql", + default_args=default_args, + tags=["demo", "bookings", "greenplum", "ods"], + description="Учебный DAG: загрузка STG -> ODS (SCD1 UPSERT) + DQ проверки", +) as dag: + resolve_stg_batch_id = PythonOperator( + task_id="resolve_stg_batch_id", + python_callable=_resolve_stg_batch_id, + ) + + load_ods_bookings = PostgresOperator( + task_id="load_ods_bookings", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/bookings_load.sql", + ) + + dq_ods_bookings = PostgresOperator( + task_id="dq_ods_bookings", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/bookings_dq.sql", + ) + + load_ods_tickets = PostgresOperator( + task_id="load_ods_tickets", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/tickets_load.sql", + ) + + dq_ods_tickets = PostgresOperator( + task_id="dq_ods_tickets", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/tickets_dq.sql", + ) + + load_ods_airports = PostgresOperator( + task_id="load_ods_airports", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/airports_load.sql", + ) + + dq_ods_airports = PostgresOperator( + task_id="dq_ods_airports", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/airports_dq.sql", + ) + + load_ods_airplanes = PostgresOperator( + task_id="load_ods_airplanes", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/airplanes_load.sql", + ) + + dq_ods_airplanes = PostgresOperator( + task_id="dq_ods_airplanes", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/airplanes_dq.sql", + ) + + load_ods_routes = PostgresOperator( + task_id="load_ods_routes", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/routes_load.sql", + ) + + dq_ods_routes = PostgresOperator( + task_id="dq_ods_routes", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/routes_dq.sql", + ) + + load_ods_seats = PostgresOperator( + task_id="load_ods_seats", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/seats_load.sql", + ) + + dq_ods_seats = PostgresOperator( + task_id="dq_ods_seats", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/seats_dq.sql", + ) + + load_ods_flights = PostgresOperator( + task_id="load_ods_flights", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/flights_load.sql", + ) + + dq_ods_flights = PostgresOperator( + task_id="dq_ods_flights", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/flights_dq.sql", + ) + + load_ods_segments = PostgresOperator( + task_id="load_ods_segments", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/segments_load.sql", + ) + + dq_ods_segments = PostgresOperator( + task_id="dq_ods_segments", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/segments_dq.sql", + ) + + load_ods_boarding_passes = PostgresOperator( + task_id="load_ods_boarding_passes", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/boarding_passes_load.sql", + ) + + dq_ods_boarding_passes = PostgresOperator( + task_id="dq_ods_boarding_passes", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="ods/boarding_passes_dq.sql", + ) + + finish_ods_summary = PythonOperator( + task_id="finish_ods_summary", + python_callable=_finish_summary, + ) + + # stg_batch_id нужен всем загрузочным веткам. + ( + resolve_stg_batch_id + >> load_ods_bookings + >> dq_ods_bookings + >> load_ods_tickets + >> dq_ods_tickets + ) + resolve_stg_batch_id >> load_ods_airports >> dq_ods_airports + resolve_stg_batch_id >> load_ods_airplanes >> dq_ods_airplanes + + [dq_ods_airports, dq_ods_airplanes] >> load_ods_routes >> dq_ods_routes + dq_ods_airplanes >> load_ods_seats >> dq_ods_seats + + dq_ods_routes >> load_ods_flights >> dq_ods_flights + [dq_ods_flights, dq_ods_tickets] >> load_ods_segments >> dq_ods_segments + dq_ods_segments >> load_ods_boarding_passes >> dq_ods_boarding_passes + + [dq_ods_boarding_passes, dq_ods_seats] >> finish_ods_summary diff --git a/docs/README.md b/docs/README.md index 2e727a5..8fa0215 100644 --- a/docs/README.md +++ b/docs/README.md @@ -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) diff --git a/docs/bookings_to_gp_ods.md b/docs/bookings_to_gp_ods.md new file mode 100644 index 0000000..e1db94f --- /dev/null +++ b/docs/bookings_to_gp_ods.md @@ -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 задач. diff --git a/docs/internal/bookings_ods_design.md b/docs/internal/bookings_ods_design.md index 06489c4..b3c77c6 100644 --- a/docs/internal/bookings_ods_design.md +++ b/docs/internal/bookings_ods_design.md @@ -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 ``` diff --git a/sql/ddl_gp.sql b/sql/ddl_gp.sql index 40da916..9df705b 100644 --- a/sql/ddl_gp.sql +++ b/sql/ddl_gp.sql @@ -37,3 +37,14 @@ FORMAT 'CUSTOM' (formatter='pxfwritable_import'); \i stg/flights_ddl.sql \i stg/segments_ddl.sql \i stg/boarding_passes_ddl.sql + +-- DDL для ODS-слоя (текущее состояние, SCD1). +\i ods/airports_ddl.sql +\i ods/airplanes_ddl.sql +\i ods/routes_ddl.sql +\i ods/seats_ddl.sql +\i ods/bookings_ddl.sql +\i ods/tickets_ddl.sql +\i ods/flights_ddl.sql +\i ods/segments_ddl.sql +\i ods/boarding_passes_ddl.sql diff --git a/sql/ods/airplanes_ddl.sql b/sql/ods/airplanes_ddl.sql new file mode 100644 index 0000000..56aa6dc --- /dev/null +++ b/sql/ods/airplanes_ddl.sql @@ -0,0 +1,13 @@ +-- DDL для ODS-слоя по таблице airplanes (текущее состояние, SCD1). + +CREATE SCHEMA IF NOT EXISTS ods; + +CREATE TABLE IF NOT EXISTS ods.airplanes ( + airplane_code TEXT NOT NULL, + model TEXT NOT NULL, + range_km INTEGER, + speed_kmh INTEGER, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (airplane_code); diff --git a/sql/ods/airplanes_dq.sql b/sql/ods/airplanes_dq.sql new file mode 100644 index 0000000..100c487 --- /dev/null +++ b/sql/ods/airplanes_dq.sql @@ -0,0 +1,99 @@ +-- DQ для ODS airplanes. + +DO $$ +DECLARE + v_batch_id TEXT := '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text; + v_stg_batch_count BIGINT; + v_dup_count BIGINT; + v_missing_keys_count BIGINT; + v_extra_keys_count BIGINT; + v_null_count BIGINT; +BEGIN + -- Для snapshot-справочников пустой батч — ошибка. + SELECT COUNT(*) + INTO v_stg_batch_count + FROM stg.airplanes + WHERE batch_id = v_batch_id; + + IF v_stg_batch_count = 0 THEN + RAISE EXCEPTION + 'DQ FAILED: batch_id=% для stg.airplanes пустой. Проверьте загрузку STG и PXF.', + v_batch_id; + END IF; + + -- В ODS не должно быть дублей по бизнес-ключу. + SELECT COUNT(*) - COUNT(DISTINCT airplane_code) + INTO v_dup_count + FROM ods.airplanes; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.airplanes найдены дубликаты airplane_code: %', + v_dup_count; + END IF; + + -- Все ключи из STG текущего батча должны присутствовать в ODS. + SELECT COUNT(*) + INTO v_missing_keys_count + FROM ( + SELECT DISTINCT airplane_code + FROM stg.airplanes + WHERE batch_id = v_batch_id + ) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.airplanes AS o + WHERE o.airplane_code = s.airplane_code + ); + + IF v_missing_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.airplanes отсутствуют ключи из stg.airplanes (batch_id=%): %', + v_batch_id, + v_missing_keys_count; + END IF; + + -- В ODS не должно быть лишних ключей, которых нет в snapshot текущего батча. + SELECT COUNT(*) + INTO v_extra_keys_count + FROM ods.airplanes AS o + WHERE NOT EXISTS ( + SELECT 1 + FROM ( + SELECT DISTINCT airplane_code + FROM stg.airplanes + WHERE batch_id = v_batch_id + ) AS s + WHERE s.airplane_code = o.airplane_code + ); + + IF v_extra_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.airplanes найдены лишние ключи вне stg batch_id=%: %', + v_batch_id, + v_extra_keys_count; + END IF; + + -- Обязательные поля в ODS. + SELECT COUNT(*) + INTO v_null_count + FROM ods.airplanes + WHERE airplane_code IS NULL + OR airplane_code = '' + OR model IS NULL + OR model = '' + OR _load_id IS NULL + OR _load_id = '' + OR _load_ts IS NULL; + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.airplanes найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: ods.airplanes ок (batch_id=%): stg_batch_rows=%', + v_batch_id, + v_stg_batch_count; +END $$; diff --git a/sql/ods/airplanes_load.sql b/sql/ods/airplanes_load.sql new file mode 100644 index 0000000..99a0a12 --- /dev/null +++ b/sql/ods/airplanes_load.sql @@ -0,0 +1,92 @@ +-- Загрузка ODS по airplanes: SCD1 (UPDATE изменившихся + INSERT новых). + +-- Statement 1: UPDATE существующих строк. +WITH src AS ( + SELECT + s.airplane_code, + s.model, + NULLIF(s.range, '')::INTEGER AS range_km, + NULLIF(s.speed, '')::INTEGER AS speed_kmh, + ROW_NUMBER() OVER ( + PARTITION BY s.airplane_code + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.airplanes AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +UPDATE ods.airplanes AS o +SET model = s.model, + range_km = s.range_km, + speed_kmh = s.speed_kmh, + _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + _load_ts = now() +FROM src AS s +WHERE s.rn = 1 + AND o.airplane_code = s.airplane_code + AND ( + o.model IS DISTINCT FROM s.model + OR o.range_km IS DISTINCT FROM s.range_km + OR o.speed_kmh IS DISTINCT FROM s.speed_kmh + ); + +-- Statement 2: INSERT новых строк. +WITH src AS ( + SELECT + s.airplane_code, + s.model, + NULLIF(s.range, '')::INTEGER AS range_km, + NULLIF(s.speed, '')::INTEGER AS speed_kmh, + ROW_NUMBER() OVER ( + PARTITION BY s.airplane_code + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.airplanes AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +INSERT INTO ods.airplanes ( + airplane_code, + model, + range_km, + speed_kmh, + _load_id, + _load_ts +) +SELECT + s.airplane_code, + s.model, + s.range_km, + s.speed_kmh, + '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + now() +FROM src AS s +WHERE s.rn = 1 + AND NOT EXISTS ( + SELECT 1 + FROM ods.airplanes AS o + WHERE o.airplane_code = s.airplane_code + ); + +-- Statement 3: DELETE ключей, которых нет в snapshot текущего батча. +-- Это делает ODS для справочника действительно "current state". +WITH src_keys AS ( + SELECT d.airplane_code + FROM ( + SELECT + s.airplane_code, + ROW_NUMBER() OVER ( + PARTITION BY s.airplane_code + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.airplanes AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + ) AS d + WHERE d.rn = 1 +) +DELETE FROM ods.airplanes AS o +WHERE NOT EXISTS ( + SELECT 1 + FROM src_keys AS s + WHERE s.airplane_code = o.airplane_code +); + +ANALYZE ods.airplanes; diff --git a/sql/ods/airports_ddl.sql b/sql/ods/airports_ddl.sql new file mode 100644 index 0000000..a6663c8 --- /dev/null +++ b/sql/ods/airports_ddl.sql @@ -0,0 +1,15 @@ +-- DDL для ODS-слоя по таблице airports (текущее состояние, SCD1). + +CREATE SCHEMA IF NOT EXISTS ods; + +CREATE TABLE IF NOT EXISTS ods.airports ( + airport_code TEXT NOT NULL, + airport_name TEXT NOT NULL, + city TEXT NOT NULL, + country TEXT NOT NULL, + coordinates TEXT, + timezone TEXT NOT NULL, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (airport_code); diff --git a/sql/ods/airports_dq.sql b/sql/ods/airports_dq.sql new file mode 100644 index 0000000..46b3c56 --- /dev/null +++ b/sql/ods/airports_dq.sql @@ -0,0 +1,105 @@ +-- DQ для ODS airports. + +DO $$ +DECLARE + v_batch_id TEXT := '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text; + v_stg_batch_count BIGINT; + v_dup_count BIGINT; + v_missing_keys_count BIGINT; + v_extra_keys_count BIGINT; + v_null_count BIGINT; +BEGIN + -- Для snapshot-справочников пустой батч — ошибка. + SELECT COUNT(*) + INTO v_stg_batch_count + FROM stg.airports + WHERE batch_id = v_batch_id; + + IF v_stg_batch_count = 0 THEN + RAISE EXCEPTION + 'DQ FAILED: batch_id=% для stg.airports пустой. Проверьте загрузку STG и PXF.', + v_batch_id; + END IF; + + -- В ODS не должно быть дублей по бизнес-ключу. + SELECT COUNT(*) - COUNT(DISTINCT airport_code) + INTO v_dup_count + FROM ods.airports; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.airports найдены дубликаты airport_code: %', + v_dup_count; + END IF; + + -- Все ключи из STG текущего батча должны присутствовать в ODS. + SELECT COUNT(*) + INTO v_missing_keys_count + FROM ( + SELECT DISTINCT airport_code + FROM stg.airports + WHERE batch_id = v_batch_id + ) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.airports AS o + WHERE o.airport_code = s.airport_code + ); + + IF v_missing_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.airports отсутствуют ключи из stg.airports (batch_id=%): %', + v_batch_id, + v_missing_keys_count; + END IF; + + -- В ODS не должно быть лишних ключей, которых нет в snapshot текущего батча. + SELECT COUNT(*) + INTO v_extra_keys_count + FROM ods.airports AS o + WHERE NOT EXISTS ( + SELECT 1 + FROM ( + SELECT DISTINCT airport_code + FROM stg.airports + WHERE batch_id = v_batch_id + ) AS s + WHERE s.airport_code = o.airport_code + ); + + IF v_extra_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.airports найдены лишние ключи вне stg batch_id=%: %', + v_batch_id, + v_extra_keys_count; + END IF; + + -- Обязательные поля в ODS. + SELECT COUNT(*) + INTO v_null_count + FROM ods.airports + WHERE airport_code IS NULL + OR airport_code = '' + OR airport_name IS NULL + OR airport_name = '' + OR city IS NULL + OR city = '' + OR country IS NULL + OR country = '' + OR timezone IS NULL + OR timezone = '' + OR _load_id IS NULL + OR _load_id = '' + OR _load_ts IS NULL; + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.airports найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: ods.airports ок (batch_id=%): stg_batch_rows=%', + v_batch_id, + v_stg_batch_count; +END $$; diff --git a/sql/ods/airports_load.sql b/sql/ods/airports_load.sql new file mode 100644 index 0000000..8a62893 --- /dev/null +++ b/sql/ods/airports_load.sql @@ -0,0 +1,104 @@ +-- Загрузка ODS по airports: SCD1 (UPDATE изменившихся + INSERT новых). + +-- Statement 1: UPDATE существующих строк. +WITH src AS ( + SELECT + s.airport_code, + s.airport_name, + s.city, + s.country, + s.coordinates, + s.timezone, + ROW_NUMBER() OVER ( + PARTITION BY s.airport_code + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.airports AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +UPDATE ods.airports AS o +SET airport_name = s.airport_name, + city = s.city, + country = s.country, + coordinates = s.coordinates, + timezone = s.timezone, + _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + _load_ts = now() +FROM src AS s +WHERE s.rn = 1 + AND o.airport_code = s.airport_code + AND ( + o.airport_name IS DISTINCT FROM s.airport_name + OR o.city IS DISTINCT FROM s.city + OR o.country IS DISTINCT FROM s.country + OR o.coordinates IS DISTINCT FROM s.coordinates + OR o.timezone IS DISTINCT FROM s.timezone + ); + +-- Statement 2: INSERT новых строк. +WITH src AS ( + SELECT + s.airport_code, + s.airport_name, + s.city, + s.country, + s.coordinates, + s.timezone, + ROW_NUMBER() OVER ( + PARTITION BY s.airport_code + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.airports AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +INSERT INTO ods.airports ( + airport_code, + airport_name, + city, + country, + coordinates, + timezone, + _load_id, + _load_ts +) +SELECT + s.airport_code, + s.airport_name, + s.city, + s.country, + s.coordinates, + s.timezone, + '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + now() +FROM src AS s +WHERE s.rn = 1 + AND NOT EXISTS ( + SELECT 1 + FROM ods.airports AS o + WHERE o.airport_code = s.airport_code + ); + +-- Statement 3: DELETE ключей, которых нет в snapshot текущего батча. +-- Это делает ODS для справочника действительно "current state". +WITH src_keys AS ( + SELECT d.airport_code + FROM ( + SELECT + s.airport_code, + ROW_NUMBER() OVER ( + PARTITION BY s.airport_code + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.airports AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + ) AS d + WHERE d.rn = 1 +) +DELETE FROM ods.airports AS o +WHERE NOT EXISTS ( + SELECT 1 + FROM src_keys AS s + WHERE s.airport_code = o.airport_code +); + +ANALYZE ods.airports; diff --git a/sql/ods/boarding_passes_ddl.sql b/sql/ods/boarding_passes_ddl.sql new file mode 100644 index 0000000..8a8640e --- /dev/null +++ b/sql/ods/boarding_passes_ddl.sql @@ -0,0 +1,15 @@ +-- DDL для ODS-слоя по таблице boarding_passes (текущее состояние, SCD1). + +CREATE SCHEMA IF NOT EXISTS ods; + +CREATE TABLE IF NOT EXISTS ods.boarding_passes ( + ticket_no TEXT NOT NULL, + flight_id INTEGER NOT NULL, + seat_no TEXT NOT NULL, + boarding_no INTEGER, + boarding_time TIMESTAMP WITH TIME ZONE, + event_ts TIMESTAMP, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (ticket_no); diff --git a/sql/ods/boarding_passes_dq.sql b/sql/ods/boarding_passes_dq.sql new file mode 100644 index 0000000..aeed460 --- /dev/null +++ b/sql/ods/boarding_passes_dq.sql @@ -0,0 +1,96 @@ +-- DQ для ODS boarding_passes. + +DO $$ +DECLARE + v_batch_id TEXT := '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text; + v_stg_batch_count BIGINT; + v_dup_count BIGINT; + v_missing_keys_count BIGINT; + v_null_count BIGINT; + v_orphan_segment_count BIGINT; +BEGIN + -- Для этой таблицы пустой батч допустим. + SELECT COUNT(*) + INTO v_stg_batch_count + FROM stg.boarding_passes + WHERE batch_id = v_batch_id; + + -- В ODS не должно быть дублей по составному бизнес-ключу. + SELECT COUNT(*) + INTO v_dup_count + FROM ( + SELECT ticket_no, flight_id + FROM ods.boarding_passes + GROUP BY ticket_no, flight_id + HAVING COUNT(*) > 1 + ) AS d; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.boarding_passes найдены дубликаты (ticket_no, flight_id): %', + v_dup_count; + END IF; + + -- Все ключи из STG текущего батча должны присутствовать в ODS. + SELECT COUNT(*) + INTO v_missing_keys_count + FROM ( + SELECT DISTINCT ticket_no, NULLIF(flight_id, '')::INTEGER AS flight_id + FROM stg.boarding_passes + WHERE batch_id = v_batch_id + ) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.boarding_passes AS o + WHERE o.ticket_no = s.ticket_no + AND o.flight_id = s.flight_id + ); + + IF v_missing_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.boarding_passes отсутствуют ключи из stg.boarding_passes (batch_id=%): %', + v_batch_id, + v_missing_keys_count; + END IF; + + -- Обязательные поля в ODS. + SELECT COUNT(*) + INTO v_null_count + FROM ods.boarding_passes + WHERE ticket_no IS NULL + OR ticket_no = '' + OR flight_id IS NULL + OR seat_no IS NULL + OR seat_no = '' + OR _load_id IS NULL + OR _load_id = '' + OR _load_ts IS NULL; + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.boarding_passes найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + -- Ссылочная целостность: boarding_passes (ticket_no, flight_id) -> segments. + SELECT COUNT(*) + INTO v_orphan_segment_count + FROM ods.boarding_passes AS bp + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.segments AS s + WHERE s.ticket_no = bp.ticket_no + AND s.flight_id = bp.flight_id + ); + + IF v_orphan_segment_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.boarding_passes найдены строки без соответствующего ods.segments: %', + v_orphan_segment_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: ods.boarding_passes ок (batch_id=%): stg_batch_rows=%', + v_batch_id, + v_stg_batch_count; +END $$; diff --git a/sql/ods/boarding_passes_load.sql b/sql/ods/boarding_passes_load.sql new file mode 100644 index 0000000..08d7359 --- /dev/null +++ b/sql/ods/boarding_passes_load.sql @@ -0,0 +1,81 @@ +-- Загрузка ODS по boarding_passes: SCD1 (UPDATE изменившихся + INSERT новых). + +-- Statement 1: UPDATE существующих строк. +WITH src AS ( + SELECT + s.ticket_no, + NULLIF(s.flight_id, '')::INTEGER AS flight_id, + s.seat_no, + NULLIF(s.boarding_no, '')::INTEGER AS boarding_no, + NULLIF(s.boarding_time, '')::TIMESTAMP WITH TIME ZONE AS boarding_time, + s.src_created_at_ts AS event_ts, + ROW_NUMBER() OVER ( + PARTITION BY s.ticket_no, s.flight_id + ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC + ) AS rn + FROM stg.boarding_passes AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +UPDATE ods.boarding_passes AS o +SET seat_no = s.seat_no, + boarding_no = s.boarding_no, + boarding_time = s.boarding_time, + event_ts = s.event_ts, + _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + _load_ts = now() +FROM src AS s +WHERE s.rn = 1 + AND o.ticket_no = s.ticket_no + AND o.flight_id = s.flight_id + AND ( + o.seat_no IS DISTINCT FROM s.seat_no + OR o.boarding_no IS DISTINCT FROM s.boarding_no + OR o.boarding_time IS DISTINCT FROM s.boarding_time + OR o.event_ts IS DISTINCT FROM s.event_ts + ); + +-- Statement 2: INSERT новых строк. +WITH src AS ( + SELECT + s.ticket_no, + NULLIF(s.flight_id, '')::INTEGER AS flight_id, + s.seat_no, + NULLIF(s.boarding_no, '')::INTEGER AS boarding_no, + NULLIF(s.boarding_time, '')::TIMESTAMP WITH TIME ZONE AS boarding_time, + s.src_created_at_ts AS event_ts, + ROW_NUMBER() OVER ( + PARTITION BY s.ticket_no, s.flight_id + ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC + ) AS rn + FROM stg.boarding_passes AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +INSERT INTO ods.boarding_passes ( + ticket_no, + flight_id, + seat_no, + boarding_no, + boarding_time, + event_ts, + _load_id, + _load_ts +) +SELECT + s.ticket_no, + s.flight_id, + s.seat_no, + s.boarding_no, + s.boarding_time, + s.event_ts, + '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + now() +FROM src AS s +WHERE s.rn = 1 + AND NOT EXISTS ( + SELECT 1 + FROM ods.boarding_passes AS o + WHERE o.ticket_no = s.ticket_no + AND o.flight_id = s.flight_id + ); + +ANALYZE ods.boarding_passes; diff --git a/sql/ods/bookings_ddl.sql b/sql/ods/bookings_ddl.sql new file mode 100644 index 0000000..1ee839f --- /dev/null +++ b/sql/ods/bookings_ddl.sql @@ -0,0 +1,13 @@ +-- DDL для ODS-слоя по таблице bookings (текущее состояние, SCD1). + +CREATE SCHEMA IF NOT EXISTS ods; + +CREATE TABLE IF NOT EXISTS ods.bookings ( + book_ref TEXT NOT NULL, + book_date TIMESTAMP WITH TIME ZONE NOT NULL, + total_amount NUMERIC(10,2) NOT NULL, + event_ts TIMESTAMP, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (book_ref); diff --git a/sql/ods/bookings_dq.sql b/sql/ods/bookings_dq.sql new file mode 100644 index 0000000..1709423 --- /dev/null +++ b/sql/ods/bookings_dq.sql @@ -0,0 +1,71 @@ +-- DQ для ODS bookings. + +DO $$ +DECLARE + v_batch_id TEXT := '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text; + v_stg_batch_count BIGINT; + v_dup_count BIGINT; + v_missing_keys_count BIGINT; + v_null_count BIGINT; +BEGIN + -- Для инкрементальных таблиц пустой батч допустим. + SELECT COUNT(*) + INTO v_stg_batch_count + FROM stg.bookings + WHERE batch_id = v_batch_id; + + -- В ODS не должно быть дублей по бизнес-ключу. + SELECT COUNT(*) - COUNT(DISTINCT book_ref) + INTO v_dup_count + FROM ods.bookings; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.bookings найдены дубликаты book_ref: %', + v_dup_count; + END IF; + + -- Все ключи из STG текущего батча должны присутствовать в ODS. + SELECT COUNT(*) + INTO v_missing_keys_count + FROM ( + SELECT DISTINCT book_ref + FROM stg.bookings + WHERE batch_id = v_batch_id + ) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.bookings AS o + WHERE o.book_ref = s.book_ref + ); + + IF v_missing_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.bookings отсутствуют ключи из stg.bookings (batch_id=%): %', + v_batch_id, + v_missing_keys_count; + END IF; + + -- Обязательные поля в ODS. + SELECT COUNT(*) + INTO v_null_count + FROM ods.bookings + WHERE book_ref IS NULL + OR book_ref = '' + OR book_date IS NULL + OR total_amount IS NULL + OR _load_id IS NULL + OR _load_id = '' + OR _load_ts IS NULL; + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.bookings найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: ods.bookings ок (batch_id=%): stg_batch_rows=%', + v_batch_id, + v_stg_batch_count; +END $$; diff --git a/sql/ods/bookings_load.sql b/sql/ods/bookings_load.sql new file mode 100644 index 0000000..45db3de --- /dev/null +++ b/sql/ods/bookings_load.sql @@ -0,0 +1,69 @@ +-- Загрузка ODS по bookings: SCD1 (UPDATE изменившихся + INSERT новых). + +-- Statement 1: UPDATE существующих строк. +WITH src AS ( + SELECT + s.book_ref, + NULLIF(s.book_date, '')::TIMESTAMP WITH TIME ZONE AS book_date, + NULLIF(s.total_amount, '')::NUMERIC(10,2) AS total_amount, + s.src_created_at_ts AS event_ts, + ROW_NUMBER() OVER ( + PARTITION BY s.book_ref + ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC + ) AS rn + FROM stg.bookings AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +UPDATE ods.bookings AS o +SET book_date = s.book_date, + total_amount = s.total_amount, + event_ts = s.event_ts, + _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + _load_ts = now() +FROM src AS s +WHERE s.rn = 1 + AND o.book_ref = s.book_ref + AND ( + o.book_date IS DISTINCT FROM s.book_date + OR o.total_amount IS DISTINCT FROM s.total_amount + OR o.event_ts IS DISTINCT FROM s.event_ts + ); + +-- Statement 2: INSERT новых строк. +WITH src AS ( + SELECT + s.book_ref, + NULLIF(s.book_date, '')::TIMESTAMP WITH TIME ZONE AS book_date, + NULLIF(s.total_amount, '')::NUMERIC(10,2) AS total_amount, + s.src_created_at_ts AS event_ts, + ROW_NUMBER() OVER ( + PARTITION BY s.book_ref + ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC + ) AS rn + FROM stg.bookings AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +INSERT INTO ods.bookings ( + book_ref, + book_date, + total_amount, + event_ts, + _load_id, + _load_ts +) +SELECT + s.book_ref, + s.book_date, + s.total_amount, + s.event_ts, + '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + now() +FROM src AS s +WHERE s.rn = 1 + AND NOT EXISTS ( + SELECT 1 + FROM ods.bookings AS o + WHERE o.book_ref = s.book_ref + ); + +ANALYZE ods.bookings; diff --git a/sql/ods/flights_ddl.sql b/sql/ods/flights_ddl.sql new file mode 100644 index 0000000..5b49996 --- /dev/null +++ b/sql/ods/flights_ddl.sql @@ -0,0 +1,17 @@ +-- DDL для ODS-слоя по таблице flights (текущее состояние, SCD1). + +CREATE SCHEMA IF NOT EXISTS ods; + +CREATE TABLE IF NOT EXISTS ods.flights ( + flight_id INTEGER NOT NULL, + route_no TEXT NOT NULL, + status TEXT NOT NULL, + scheduled_departure TIMESTAMP WITH TIME ZONE, + scheduled_arrival TIMESTAMP WITH TIME ZONE, + actual_departure TIMESTAMP WITH TIME ZONE, + actual_arrival TIMESTAMP WITH TIME ZONE, + event_ts TIMESTAMP, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (flight_id); diff --git a/sql/ods/flights_dq.sql b/sql/ods/flights_dq.sql new file mode 100644 index 0000000..912e340 --- /dev/null +++ b/sql/ods/flights_dq.sql @@ -0,0 +1,89 @@ +-- DQ для ODS flights. + +DO $$ +DECLARE + v_batch_id TEXT := '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text; + v_stg_batch_count BIGINT; + v_dup_count BIGINT; + v_missing_keys_count BIGINT; + v_null_count BIGINT; + v_orphan_route_count BIGINT; +BEGIN + -- Для инкрементальных таблиц пустой батч допустим. + SELECT COUNT(*) + INTO v_stg_batch_count + FROM stg.flights + WHERE batch_id = v_batch_id; + + -- В ODS не должно быть дублей по бизнес-ключу. + SELECT COUNT(*) - COUNT(DISTINCT flight_id) + INTO v_dup_count + FROM ods.flights; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.flights найдены дубликаты flight_id: %', + v_dup_count; + END IF; + + -- Все ключи из STG текущего батча должны присутствовать в ODS. + SELECT COUNT(*) + INTO v_missing_keys_count + FROM ( + SELECT DISTINCT NULLIF(flight_id, '')::INTEGER AS flight_id + FROM stg.flights + WHERE batch_id = v_batch_id + ) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.flights AS o + WHERE o.flight_id = s.flight_id + ); + + IF v_missing_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.flights отсутствуют ключи из stg.flights (batch_id=%): %', + v_batch_id, + v_missing_keys_count; + END IF; + + -- Обязательные поля в ODS. + SELECT COUNT(*) + INTO v_null_count + FROM ods.flights + WHERE flight_id IS NULL + OR route_no IS NULL + OR route_no = '' + OR status IS NULL + OR status = '' + OR _load_id IS NULL + OR _load_id = '' + OR _load_ts IS NULL; + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.flights найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + -- Ссылочная целостность: flights.route_no -> routes.route_no (упрощённо, без validity). + SELECT COUNT(*) + INTO v_orphan_route_count + FROM ods.flights AS f + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.routes AS r + WHERE r.route_no = f.route_no + ); + + IF v_orphan_route_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.flights найдены строки без соответствующего routes.route_no: %', + v_orphan_route_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: ods.flights ок (batch_id=%): stg_batch_rows=%', + v_batch_id, + v_stg_batch_count; +END $$; diff --git a/sql/ods/flights_load.sql b/sql/ods/flights_load.sql new file mode 100644 index 0000000..ff33580 --- /dev/null +++ b/sql/ods/flights_load.sql @@ -0,0 +1,93 @@ +-- Загрузка ODS по flights: SCD1 (UPDATE изменившихся + INSERT новых). + +-- Statement 1: UPDATE существующих строк. +WITH src AS ( + SELECT + NULLIF(s.flight_id, '')::INTEGER AS flight_id, + s.route_no, + s.status, + NULLIF(s.scheduled_departure, '')::TIMESTAMP WITH TIME ZONE AS scheduled_departure, + NULLIF(s.scheduled_arrival, '')::TIMESTAMP WITH TIME ZONE AS scheduled_arrival, + NULLIF(s.actual_departure, '')::TIMESTAMP WITH TIME ZONE AS actual_departure, + NULLIF(s.actual_arrival, '')::TIMESTAMP WITH TIME ZONE AS actual_arrival, + s.src_created_at_ts AS event_ts, + ROW_NUMBER() OVER ( + PARTITION BY s.flight_id + ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC + ) AS rn + FROM stg.flights AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +UPDATE ods.flights AS o +SET route_no = s.route_no, + status = s.status, + scheduled_departure = s.scheduled_departure, + scheduled_arrival = s.scheduled_arrival, + actual_departure = s.actual_departure, + actual_arrival = s.actual_arrival, + event_ts = s.event_ts, + _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + _load_ts = now() +FROM src AS s +WHERE s.rn = 1 + AND o.flight_id = s.flight_id + AND ( + o.route_no IS DISTINCT FROM s.route_no + OR o.status IS DISTINCT FROM s.status + OR o.scheduled_departure IS DISTINCT FROM s.scheduled_departure + OR o.scheduled_arrival IS DISTINCT FROM s.scheduled_arrival + OR o.actual_departure IS DISTINCT FROM s.actual_departure + OR o.actual_arrival IS DISTINCT FROM s.actual_arrival + OR o.event_ts IS DISTINCT FROM s.event_ts + ); + +-- Statement 2: INSERT новых строк. +WITH src AS ( + SELECT + NULLIF(s.flight_id, '')::INTEGER AS flight_id, + s.route_no, + s.status, + NULLIF(s.scheduled_departure, '')::TIMESTAMP WITH TIME ZONE AS scheduled_departure, + NULLIF(s.scheduled_arrival, '')::TIMESTAMP WITH TIME ZONE AS scheduled_arrival, + NULLIF(s.actual_departure, '')::TIMESTAMP WITH TIME ZONE AS actual_departure, + NULLIF(s.actual_arrival, '')::TIMESTAMP WITH TIME ZONE AS actual_arrival, + s.src_created_at_ts AS event_ts, + ROW_NUMBER() OVER ( + PARTITION BY s.flight_id + ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC + ) AS rn + FROM stg.flights AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +INSERT INTO ods.flights ( + flight_id, + route_no, + status, + scheduled_departure, + scheduled_arrival, + actual_departure, + actual_arrival, + event_ts, + _load_id, + _load_ts +) +SELECT + s.flight_id, + s.route_no, + s.status, + s.scheduled_departure, + s.scheduled_arrival, + s.actual_departure, + s.actual_arrival, + s.event_ts, + '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + now() +FROM src AS s +WHERE s.rn = 1 + AND NOT EXISTS ( + SELECT 1 + FROM ods.flights AS o + WHERE o.flight_id = s.flight_id + ); + +ANALYZE ods.flights; diff --git a/sql/ods/routes_ddl.sql b/sql/ods/routes_ddl.sql new file mode 100644 index 0000000..5f3d486 --- /dev/null +++ b/sql/ods/routes_ddl.sql @@ -0,0 +1,17 @@ +-- DDL для ODS-слоя по таблице routes (текущее состояние, SCD1). + +CREATE SCHEMA IF NOT EXISTS ods; + +CREATE TABLE IF NOT EXISTS ods.routes ( + route_no TEXT NOT NULL, + validity TEXT NOT NULL, + departure_airport TEXT NOT NULL, + arrival_airport TEXT NOT NULL, + airplane_code TEXT NOT NULL, + days_of_week TEXT, + departure_time TIME, + duration INTERVAL, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (route_no); diff --git a/sql/ods/routes_dq.sql b/sql/ods/routes_dq.sql new file mode 100644 index 0000000..2b36969 --- /dev/null +++ b/sql/ods/routes_dq.sql @@ -0,0 +1,163 @@ +-- DQ для ODS routes. + +DO $$ +DECLARE + v_batch_id TEXT := '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text; + v_stg_batch_count BIGINT; + v_dup_count BIGINT; + v_missing_keys_count BIGINT; + v_extra_keys_count BIGINT; + v_null_count BIGINT; + v_orphan_departure_count BIGINT; + v_orphan_arrival_count BIGINT; + v_orphan_airplane_count BIGINT; +BEGIN + -- Для snapshot-справочников пустой батч — ошибка. + SELECT COUNT(*) + INTO v_stg_batch_count + FROM stg.routes + WHERE batch_id = v_batch_id; + + IF v_stg_batch_count = 0 THEN + RAISE EXCEPTION + 'DQ FAILED: batch_id=% для stg.routes пустой. Проверьте загрузку STG и PXF.', + v_batch_id; + END IF; + + -- В ODS не должно быть дублей по составному бизнес-ключу. + SELECT COUNT(*) + INTO v_dup_count + FROM ( + SELECT route_no, validity + FROM ods.routes + GROUP BY route_no, validity + HAVING COUNT(*) > 1 + ) AS d; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.routes найдены дубликаты (route_no, validity): %', + v_dup_count; + END IF; + + -- Все ключи из STG текущего батча должны присутствовать в ODS. + SELECT COUNT(*) + INTO v_missing_keys_count + FROM ( + SELECT DISTINCT route_no, validity + FROM stg.routes + WHERE batch_id = v_batch_id + ) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.routes AS o + WHERE o.route_no = s.route_no + AND o.validity = s.validity + ); + + IF v_missing_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.routes отсутствуют ключи из stg.routes (batch_id=%): %', + v_batch_id, + v_missing_keys_count; + END IF; + + -- В ODS не должно быть лишних ключей, которых нет в snapshot текущего батча. + SELECT COUNT(*) + INTO v_extra_keys_count + FROM ods.routes AS o + WHERE NOT EXISTS ( + SELECT 1 + FROM ( + SELECT DISTINCT route_no, validity + FROM stg.routes + WHERE batch_id = v_batch_id + ) AS s + WHERE s.route_no = o.route_no + AND s.validity = o.validity + ); + + IF v_extra_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.routes найдены лишние ключи вне stg batch_id=%: %', + v_batch_id, + v_extra_keys_count; + END IF; + + -- Обязательные поля в ODS. + SELECT COUNT(*) + INTO v_null_count + FROM ods.routes + WHERE route_no IS NULL + OR route_no = '' + OR validity IS NULL + OR validity = '' + OR departure_airport IS NULL + OR departure_airport = '' + OR arrival_airport IS NULL + OR arrival_airport = '' + OR airplane_code IS NULL + OR airplane_code = '' + OR _load_id IS NULL + OR _load_id = '' + OR _load_ts IS NULL; + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.routes найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + -- Ссылочная целостность: departure_airport должен существовать в ods.airports. + SELECT COUNT(*) + INTO v_orphan_departure_count + FROM ods.routes AS r + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.airports AS a + WHERE a.airport_code = r.departure_airport + ); + + IF v_orphan_departure_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.routes найдены строки с невалидным departure_airport: %', + v_orphan_departure_count; + END IF; + + -- Ссылочная целостность: arrival_airport должен существовать в ods.airports. + SELECT COUNT(*) + INTO v_orphan_arrival_count + FROM ods.routes AS r + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.airports AS a + WHERE a.airport_code = r.arrival_airport + ); + + IF v_orphan_arrival_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.routes найдены строки с невалидным arrival_airport: %', + v_orphan_arrival_count; + END IF; + + -- Ссылочная целостность: airplane_code должен существовать в ods.airplanes. + SELECT COUNT(*) + INTO v_orphan_airplane_count + FROM ods.routes AS r + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.airplanes AS a + WHERE a.airplane_code = r.airplane_code + ); + + IF v_orphan_airplane_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.routes найдены строки с невалидным airplane_code: %', + v_orphan_airplane_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: ods.routes ок (batch_id=%): stg_batch_rows=%', + v_batch_id, + v_stg_batch_count; +END $$; diff --git a/sql/ods/routes_load.sql b/sql/ods/routes_load.sql new file mode 100644 index 0000000..c51a160 --- /dev/null +++ b/sql/ods/routes_load.sql @@ -0,0 +1,118 @@ +-- Загрузка ODS по routes: SCD1 (UPDATE изменившихся + INSERT новых). + +-- Statement 1: UPDATE существующих строк. +WITH src AS ( + SELECT + s.route_no, + s.validity, + s.departure_airport, + s.arrival_airport, + s.airplane_code, + s.days_of_week, + NULLIF(s.scheduled_time, '')::TIME AS departure_time, + NULLIF(s.duration, '')::INTERVAL AS duration, + ROW_NUMBER() OVER ( + PARTITION BY s.route_no, s.validity + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.routes AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +UPDATE ods.routes AS o +SET departure_airport = s.departure_airport, + arrival_airport = s.arrival_airport, + airplane_code = s.airplane_code, + days_of_week = s.days_of_week, + departure_time = s.departure_time, + duration = s.duration, + _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + _load_ts = now() +FROM src AS s +WHERE s.rn = 1 + AND o.route_no = s.route_no + AND o.validity = s.validity + AND ( + o.departure_airport IS DISTINCT FROM s.departure_airport + OR o.arrival_airport IS DISTINCT FROM s.arrival_airport + OR o.airplane_code IS DISTINCT FROM s.airplane_code + OR o.days_of_week IS DISTINCT FROM s.days_of_week + OR o.departure_time IS DISTINCT FROM s.departure_time + OR o.duration IS DISTINCT FROM s.duration + ); + +-- Statement 2: INSERT новых строк. +WITH src AS ( + SELECT + s.route_no, + s.validity, + s.departure_airport, + s.arrival_airport, + s.airplane_code, + s.days_of_week, + NULLIF(s.scheduled_time, '')::TIME AS departure_time, + NULLIF(s.duration, '')::INTERVAL AS duration, + ROW_NUMBER() OVER ( + PARTITION BY s.route_no, s.validity + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.routes AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +INSERT INTO ods.routes ( + route_no, + validity, + departure_airport, + arrival_airport, + airplane_code, + days_of_week, + departure_time, + duration, + _load_id, + _load_ts +) +SELECT + s.route_no, + s.validity, + s.departure_airport, + s.arrival_airport, + s.airplane_code, + s.days_of_week, + s.departure_time, + s.duration, + '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + now() +FROM src AS s +WHERE s.rn = 1 + AND NOT EXISTS ( + SELECT 1 + FROM ods.routes AS o + WHERE o.route_no = s.route_no + AND o.validity = s.validity + ); + +-- Statement 3: DELETE ключей, которых нет в snapshot текущего батча. +-- Это делает ODS для справочника действительно "current state". +WITH src_keys AS ( + SELECT d.route_no, d.validity + FROM ( + SELECT + s.route_no, + s.validity, + ROW_NUMBER() OVER ( + PARTITION BY s.route_no, s.validity + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.routes AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + ) AS d + WHERE d.rn = 1 +) +DELETE FROM ods.routes AS o +WHERE NOT EXISTS ( + SELECT 1 + FROM src_keys AS s + WHERE s.route_no = o.route_no + AND s.validity = o.validity +); + +ANALYZE ods.routes; diff --git a/sql/ods/seats_ddl.sql b/sql/ods/seats_ddl.sql new file mode 100644 index 0000000..08eaceb --- /dev/null +++ b/sql/ods/seats_ddl.sql @@ -0,0 +1,12 @@ +-- DDL для ODS-слоя по таблице seats (текущее состояние, SCD1). + +CREATE SCHEMA IF NOT EXISTS ods; + +CREATE TABLE IF NOT EXISTS ods.seats ( + airplane_code TEXT NOT NULL, + seat_no TEXT NOT NULL, + fare_conditions TEXT NOT NULL, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (airplane_code); diff --git a/sql/ods/seats_dq.sql b/sql/ods/seats_dq.sql new file mode 100644 index 0000000..ea2b61e --- /dev/null +++ b/sql/ods/seats_dq.sql @@ -0,0 +1,125 @@ +-- DQ для ODS seats. + +DO $$ +DECLARE + v_batch_id TEXT := '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text; + v_stg_batch_count BIGINT; + v_dup_count BIGINT; + v_missing_keys_count BIGINT; + v_extra_keys_count BIGINT; + v_null_count BIGINT; + v_orphan_airplane_count BIGINT; +BEGIN + -- Для snapshot-справочников пустой батч — ошибка. + SELECT COUNT(*) + INTO v_stg_batch_count + FROM stg.seats + WHERE batch_id = v_batch_id; + + IF v_stg_batch_count = 0 THEN + RAISE EXCEPTION + 'DQ FAILED: batch_id=% для stg.seats пустой. Проверьте загрузку STG и PXF.', + v_batch_id; + END IF; + + -- В ODS не должно быть дублей по составному бизнес-ключу. + SELECT COUNT(*) + INTO v_dup_count + FROM ( + SELECT airplane_code, seat_no + FROM ods.seats + GROUP BY airplane_code, seat_no + HAVING COUNT(*) > 1 + ) AS d; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.seats найдены дубликаты (airplane_code, seat_no): %', + v_dup_count; + END IF; + + -- Все ключи из STG текущего батча должны присутствовать в ODS. + SELECT COUNT(*) + INTO v_missing_keys_count + FROM ( + SELECT DISTINCT airplane_code, seat_no + FROM stg.seats + WHERE batch_id = v_batch_id + ) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.seats AS o + WHERE o.airplane_code = s.airplane_code + AND o.seat_no = s.seat_no + ); + + IF v_missing_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.seats отсутствуют ключи из stg.seats (batch_id=%): %', + v_batch_id, + v_missing_keys_count; + END IF; + + -- В ODS не должно быть лишних ключей, которых нет в snapshot текущего батча. + SELECT COUNT(*) + INTO v_extra_keys_count + FROM ods.seats AS o + WHERE NOT EXISTS ( + SELECT 1 + FROM ( + SELECT DISTINCT airplane_code, seat_no + FROM stg.seats + WHERE batch_id = v_batch_id + ) AS s + WHERE s.airplane_code = o.airplane_code + AND s.seat_no = o.seat_no + ); + + IF v_extra_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.seats найдены лишние ключи вне stg batch_id=%: %', + v_batch_id, + v_extra_keys_count; + END IF; + + -- Обязательные поля в ODS. + SELECT COUNT(*) + INTO v_null_count + FROM ods.seats + WHERE airplane_code IS NULL + OR airplane_code = '' + OR seat_no IS NULL + OR seat_no = '' + OR fare_conditions IS NULL + OR fare_conditions = '' + OR _load_id IS NULL + OR _load_id = '' + OR _load_ts IS NULL; + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.seats найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + -- Ссылочная целостность: airplane_code должен существовать в ods.airplanes. + SELECT COUNT(*) + INTO v_orphan_airplane_count + FROM ods.seats AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.airplanes AS a + WHERE a.airplane_code = s.airplane_code + ); + + IF v_orphan_airplane_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.seats найдены строки с невалидным airplane_code: %', + v_orphan_airplane_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: ods.seats ок (batch_id=%): stg_batch_rows=%', + v_batch_id, + v_stg_batch_count; +END $$; diff --git a/sql/ods/seats_load.sql b/sql/ods/seats_load.sql new file mode 100644 index 0000000..82fbf9f --- /dev/null +++ b/sql/ods/seats_load.sql @@ -0,0 +1,86 @@ +-- Загрузка ODS по seats: SCD1 (UPDATE изменившихся + INSERT новых). + +-- Statement 1: UPDATE существующих строк. +WITH src AS ( + SELECT + s.airplane_code, + s.seat_no, + s.fare_conditions, + ROW_NUMBER() OVER ( + PARTITION BY s.airplane_code, s.seat_no + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.seats AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +UPDATE ods.seats AS o +SET fare_conditions = s.fare_conditions, + _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + _load_ts = now() +FROM src AS s +WHERE s.rn = 1 + AND o.airplane_code = s.airplane_code + AND o.seat_no = s.seat_no + AND o.fare_conditions IS DISTINCT FROM s.fare_conditions; + +-- Statement 2: INSERT новых строк. +WITH src AS ( + SELECT + s.airplane_code, + s.seat_no, + s.fare_conditions, + ROW_NUMBER() OVER ( + PARTITION BY s.airplane_code, s.seat_no + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.seats AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +INSERT INTO ods.seats ( + airplane_code, + seat_no, + fare_conditions, + _load_id, + _load_ts +) +SELECT + s.airplane_code, + s.seat_no, + s.fare_conditions, + '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + now() +FROM src AS s +WHERE s.rn = 1 + AND NOT EXISTS ( + SELECT 1 + FROM ods.seats AS o + WHERE o.airplane_code = s.airplane_code + AND o.seat_no = s.seat_no + ); + +-- Statement 3: DELETE ключей, которых нет в snapshot текущего батча. +-- Это делает ODS для справочника действительно "current state". +WITH src_keys AS ( + SELECT d.airplane_code, d.seat_no + FROM ( + SELECT + s.airplane_code, + s.seat_no, + ROW_NUMBER() OVER ( + PARTITION BY s.airplane_code, s.seat_no + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.seats AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + ) AS d + WHERE d.rn = 1 +) +DELETE FROM ods.seats AS o +WHERE NOT EXISTS ( + SELECT 1 + FROM src_keys AS s + WHERE s.airplane_code = o.airplane_code + AND s.seat_no = o.seat_no +); + +ANALYZE ods.seats; diff --git a/sql/ods/segments_ddl.sql b/sql/ods/segments_ddl.sql new file mode 100644 index 0000000..b35034a --- /dev/null +++ b/sql/ods/segments_ddl.sql @@ -0,0 +1,14 @@ +-- DDL для ODS-слоя по таблице segments (текущее состояние, SCD1). + +CREATE SCHEMA IF NOT EXISTS ods; + +CREATE TABLE IF NOT EXISTS ods.segments ( + ticket_no TEXT NOT NULL, + flight_id INTEGER NOT NULL, + fare_conditions TEXT NOT NULL, + segment_amount NUMERIC(10,2), + event_ts TIMESTAMP, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (ticket_no); diff --git a/sql/ods/segments_dq.sql b/sql/ods/segments_dq.sql new file mode 100644 index 0000000..f1a05bc --- /dev/null +++ b/sql/ods/segments_dq.sql @@ -0,0 +1,112 @@ +-- DQ для ODS segments. + +DO $$ +DECLARE + v_batch_id TEXT := '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text; + v_stg_batch_count BIGINT; + v_dup_count BIGINT; + v_missing_keys_count BIGINT; + v_null_count BIGINT; + v_orphan_ticket_count BIGINT; + v_orphan_flight_count BIGINT; +BEGIN + -- Для инкрементальных таблиц пустой батч допустим. + SELECT COUNT(*) + INTO v_stg_batch_count + FROM stg.segments + WHERE batch_id = v_batch_id; + + -- В ODS не должно быть дублей по составному бизнес-ключу. + SELECT COUNT(*) + INTO v_dup_count + FROM ( + SELECT ticket_no, flight_id + FROM ods.segments + GROUP BY ticket_no, flight_id + HAVING COUNT(*) > 1 + ) AS d; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.segments найдены дубликаты (ticket_no, flight_id): %', + v_dup_count; + END IF; + + -- Все ключи из STG текущего батча должны присутствовать в ODS. + SELECT COUNT(*) + INTO v_missing_keys_count + FROM ( + SELECT DISTINCT ticket_no, NULLIF(flight_id, '')::INTEGER AS flight_id + FROM stg.segments + WHERE batch_id = v_batch_id + ) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.segments AS o + WHERE o.ticket_no = s.ticket_no + AND o.flight_id = s.flight_id + ); + + IF v_missing_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.segments отсутствуют ключи из stg.segments (batch_id=%): %', + v_batch_id, + v_missing_keys_count; + END IF; + + -- Обязательные поля в ODS. + SELECT COUNT(*) + INTO v_null_count + FROM ods.segments + WHERE ticket_no IS NULL + OR ticket_no = '' + OR flight_id IS NULL + OR fare_conditions IS NULL + OR fare_conditions = '' + OR _load_id IS NULL + OR _load_id = '' + OR _load_ts IS NULL; + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.segments найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + -- Ссылочная целостность: segments.ticket_no -> tickets.ticket_no. + SELECT COUNT(*) + INTO v_orphan_ticket_count + FROM ods.segments AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.tickets AS t + WHERE t.ticket_no = s.ticket_no + ); + + IF v_orphan_ticket_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.segments найдены строки без соответствующего tickets.ticket_no: %', + v_orphan_ticket_count; + END IF; + + -- Ссылочная целостность: segments.flight_id -> flights.flight_id. + SELECT COUNT(*) + INTO v_orphan_flight_count + FROM ods.segments AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.flights AS f + WHERE f.flight_id = s.flight_id + ); + + IF v_orphan_flight_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.segments найдены строки без соответствующего flights.flight_id: %', + v_orphan_flight_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: ods.segments ок (batch_id=%): stg_batch_rows=%', + v_batch_id, + v_stg_batch_count; +END $$; diff --git a/sql/ods/segments_load.sql b/sql/ods/segments_load.sql new file mode 100644 index 0000000..50e0e18 --- /dev/null +++ b/sql/ods/segments_load.sql @@ -0,0 +1,75 @@ +-- Загрузка ODS по segments: SCD1 (UPDATE изменившихся + INSERT новых). + +-- Statement 1: UPDATE существующих строк. +WITH src AS ( + SELECT + s.ticket_no, + NULLIF(s.flight_id, '')::INTEGER AS flight_id, + s.fare_conditions, + NULLIF(s.price, '')::NUMERIC(10,2) AS segment_amount, + s.src_created_at_ts AS event_ts, + ROW_NUMBER() OVER ( + PARTITION BY s.ticket_no, s.flight_id + ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC + ) AS rn + FROM stg.segments AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +UPDATE ods.segments AS o +SET fare_conditions = s.fare_conditions, + segment_amount = s.segment_amount, + event_ts = s.event_ts, + _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + _load_ts = now() +FROM src AS s +WHERE s.rn = 1 + AND o.ticket_no = s.ticket_no + AND o.flight_id = s.flight_id + AND ( + o.fare_conditions IS DISTINCT FROM s.fare_conditions + OR o.segment_amount IS DISTINCT FROM s.segment_amount + OR o.event_ts IS DISTINCT FROM s.event_ts + ); + +-- Statement 2: INSERT новых строк. +WITH src AS ( + SELECT + s.ticket_no, + NULLIF(s.flight_id, '')::INTEGER AS flight_id, + s.fare_conditions, + NULLIF(s.price, '')::NUMERIC(10,2) AS segment_amount, + s.src_created_at_ts AS event_ts, + ROW_NUMBER() OVER ( + PARTITION BY s.ticket_no, s.flight_id + ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC + ) AS rn + FROM stg.segments AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +INSERT INTO ods.segments ( + ticket_no, + flight_id, + fare_conditions, + segment_amount, + event_ts, + _load_id, + _load_ts +) +SELECT + s.ticket_no, + s.flight_id, + s.fare_conditions, + s.segment_amount, + s.event_ts, + '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + now() +FROM src AS s +WHERE s.rn = 1 + AND NOT EXISTS ( + SELECT 1 + FROM ods.segments AS o + WHERE o.ticket_no = s.ticket_no + AND o.flight_id = s.flight_id + ); + +ANALYZE ods.segments; diff --git a/sql/ods/tickets_ddl.sql b/sql/ods/tickets_ddl.sql new file mode 100644 index 0000000..a4a70ee --- /dev/null +++ b/sql/ods/tickets_ddl.sql @@ -0,0 +1,15 @@ +-- DDL для ODS-слоя по таблице tickets (текущее состояние, SCD1). + +CREATE SCHEMA IF NOT EXISTS ods; + +CREATE TABLE IF NOT EXISTS ods.tickets ( + ticket_no TEXT NOT NULL, + book_ref TEXT NOT NULL, + passenger_id TEXT NOT NULL, + passenger_name TEXT NOT NULL, + is_outbound BOOLEAN, + event_ts TIMESTAMP, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (ticket_no); diff --git a/sql/ods/tickets_dq.sql b/sql/ods/tickets_dq.sql new file mode 100644 index 0000000..3ad1007 --- /dev/null +++ b/sql/ods/tickets_dq.sql @@ -0,0 +1,92 @@ +-- DQ для ODS tickets. + +DO $$ +DECLARE + v_batch_id TEXT := '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text; + v_stg_batch_count BIGINT; + v_dup_count BIGINT; + v_missing_keys_count BIGINT; + v_null_count BIGINT; + v_orphan_booking_count BIGINT; +BEGIN + -- Для инкрементальных таблиц пустой батч допустим. + SELECT COUNT(*) + INTO v_stg_batch_count + FROM stg.tickets + WHERE batch_id = v_batch_id; + + -- В ODS не должно быть дублей по бизнес-ключу. + SELECT COUNT(*) - COUNT(DISTINCT ticket_no) + INTO v_dup_count + FROM ods.tickets; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.tickets найдены дубликаты ticket_no: %', + v_dup_count; + END IF; + + -- Все ключи из STG текущего батча должны присутствовать в ODS. + SELECT COUNT(*) + INTO v_missing_keys_count + FROM ( + SELECT DISTINCT ticket_no + FROM stg.tickets + WHERE batch_id = v_batch_id + ) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.tickets AS o + WHERE o.ticket_no = s.ticket_no + ); + + IF v_missing_keys_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.tickets отсутствуют ключи из stg.tickets (batch_id=%): %', + v_batch_id, + v_missing_keys_count; + END IF; + + -- Обязательные поля в ODS. + SELECT COUNT(*) + INTO v_null_count + FROM ods.tickets + WHERE ticket_no IS NULL + OR ticket_no = '' + OR book_ref IS NULL + OR book_ref = '' + OR passenger_id IS NULL + OR passenger_id = '' + OR passenger_name IS NULL + OR passenger_name = '' + OR _load_id IS NULL + OR _load_id = '' + OR _load_ts IS NULL; + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.tickets найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + -- Ссылочная целостность: tickets.book_ref -> bookings.book_ref. + SELECT COUNT(*) + INTO v_orphan_booking_count + FROM ods.tickets AS t + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.bookings AS b + WHERE b.book_ref = t.book_ref + ); + + IF v_orphan_booking_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в ods.tickets найдены строки без соответствующего bookings.book_ref: %', + v_orphan_booking_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: ods.tickets ок (batch_id=%): stg_batch_rows=%', + v_batch_id, + v_stg_batch_count; +END $$; diff --git a/sql/ods/tickets_load.sql b/sql/ods/tickets_load.sql new file mode 100644 index 0000000..b4c214e --- /dev/null +++ b/sql/ods/tickets_load.sql @@ -0,0 +1,81 @@ +-- Загрузка ODS по tickets: SCD1 (UPDATE изменившихся + INSERT новых). + +-- Statement 1: UPDATE существующих строк. +WITH src AS ( + SELECT + s.ticket_no, + s.book_ref, + s.passenger_id, + s.passenger_name, + NULLIF(s.outbound, '')::BOOLEAN AS is_outbound, + s.src_created_at_ts AS event_ts, + ROW_NUMBER() OVER ( + PARTITION BY s.ticket_no + ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC + ) AS rn + FROM stg.tickets AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +UPDATE ods.tickets AS o +SET book_ref = s.book_ref, + passenger_id = s.passenger_id, + passenger_name = s.passenger_name, + is_outbound = s.is_outbound, + event_ts = s.event_ts, + _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + _load_ts = now() +FROM src AS s +WHERE s.rn = 1 + AND o.ticket_no = s.ticket_no + AND ( + o.book_ref IS DISTINCT FROM s.book_ref + OR o.passenger_id IS DISTINCT FROM s.passenger_id + OR o.passenger_name IS DISTINCT FROM s.passenger_name + OR o.is_outbound IS DISTINCT FROM s.is_outbound + OR o.event_ts IS DISTINCT FROM s.event_ts + ); + +-- Statement 2: INSERT новых строк. +WITH src AS ( + SELECT + s.ticket_no, + s.book_ref, + s.passenger_id, + s.passenger_name, + NULLIF(s.outbound, '')::BOOLEAN AS is_outbound, + s.src_created_at_ts AS event_ts, + ROW_NUMBER() OVER ( + PARTITION BY s.ticket_no + ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC + ) AS rn + FROM stg.tickets AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) +INSERT INTO ods.tickets ( + ticket_no, + book_ref, + passenger_id, + passenger_name, + is_outbound, + event_ts, + _load_id, + _load_ts +) +SELECT + s.ticket_no, + s.book_ref, + s.passenger_id, + s.passenger_name, + s.is_outbound, + s.event_ts, + '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + now() +FROM src AS s +WHERE s.rn = 1 + AND NOT EXISTS ( + SELECT 1 + FROM ods.tickets AS o + WHERE o.ticket_no = s.ticket_no + ); + +ANALYZE ods.tickets; diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index e422d0e..2117dd9 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -199,3 +199,101 @@ def test_bookings_to_gp_stage_dag_structure(): assert airports not in airplanes.get_flat_relatives( upstream=False ), "airplanes не должен быть upstream для airports" + + +def test_bookings_ods_ddl_dag_structure(): + """Проверка структуры DAG bookings_ods_ddl.""" + dag = _load_dag("airflow.dags.bookings_ods_ddl") + + expected_tasks = { + "apply_ods_airports_ddl", + "apply_ods_airplanes_ddl", + "apply_ods_routes_ddl", + "apply_ods_seats_ddl", + "apply_ods_bookings_ddl", + "apply_ods_tickets_ddl", + "apply_ods_flights_ddl", + "apply_ods_segments_ddl", + "apply_ods_boarding_passes_ddl", + } + assert expected_tasks.issubset(dag.task_dict.keys()) + + _assert_reachable(dag, "apply_ods_airports_ddl", "apply_ods_airplanes_ddl") + + for task_id in expected_tasks - {"apply_ods_airports_ddl"}: + _assert_reachable(dag, "apply_ods_airports_ddl", task_id) + + +def test_bookings_to_gp_ods_dag_structure(): + """Проверка структуры DAG bookings_to_gp_ods.""" + dag = _load_dag("airflow.dags.bookings_to_gp_ods") + + expected_tasks = { + "resolve_stg_batch_id", + "load_ods_bookings", + "dq_ods_bookings", + "load_ods_tickets", + "dq_ods_tickets", + "load_ods_airports", + "dq_ods_airports", + "load_ods_airplanes", + "dq_ods_airplanes", + "load_ods_routes", + "dq_ods_routes", + "load_ods_seats", + "dq_ods_seats", + "load_ods_flights", + "dq_ods_flights", + "load_ods_segments", + "dq_ods_segments", + "load_ods_boarding_passes", + "dq_ods_boarding_passes", + "finish_ods_summary", + } + assert expected_tasks.issubset(dag.task_dict.keys()) + + # Базовая цепочка транзакций. + _assert_reachable(dag, "resolve_stg_batch_id", "load_ods_bookings") + _assert_reachable(dag, "load_ods_bookings", "dq_ods_bookings") + _assert_reachable(dag, "dq_ods_bookings", "load_ods_tickets") + _assert_reachable(dag, "load_ods_tickets", "dq_ods_tickets") + + # Инвариант "load -> dq" для каждой таблицы. + load_to_dq = [ + ("load_ods_bookings", "dq_ods_bookings"), + ("load_ods_tickets", "dq_ods_tickets"), + ("load_ods_airports", "dq_ods_airports"), + ("load_ods_airplanes", "dq_ods_airplanes"), + ("load_ods_routes", "dq_ods_routes"), + ("load_ods_seats", "dq_ods_seats"), + ("load_ods_flights", "dq_ods_flights"), + ("load_ods_segments", "dq_ods_segments"), + ("load_ods_boarding_passes", "dq_ods_boarding_passes"), + ] + for load_task_id, dq_task_id in load_to_dq: + _assert_direct_edge(dag, load_task_id, dq_task_id) + + # Справочники стартуют параллельно и не зависят друг от друга. + _assert_reachable(dag, "resolve_stg_batch_id", "load_ods_airports") + _assert_reachable(dag, "resolve_stg_batch_id", "load_ods_airplanes") + airports = dag.get_task("load_ods_airports") + airplanes = dag.get_task("load_ods_airplanes") + assert airplanes not in airports.get_flat_relatives( + upstream=False + ), "airports не должен быть upstream для airplanes" + assert airports not in airplanes.get_flat_relatives( + upstream=False + ), "airplanes не должен быть upstream для airports" + + # Барьеры по данным. + _assert_reachable(dag, "dq_ods_airports", "dq_ods_routes") + _assert_reachable(dag, "dq_ods_airplanes", "dq_ods_routes") + _assert_reachable(dag, "dq_ods_airplanes", "dq_ods_seats") + _assert_reachable(dag, "dq_ods_routes", "dq_ods_flights") + _assert_reachable(dag, "dq_ods_flights", "dq_ods_segments") + _assert_reachable(dag, "dq_ods_tickets", "dq_ods_segments") + _assert_reachable(dag, "dq_ods_segments", "dq_ods_boarding_passes") + + # Финальная сводка должна ждать обе ветки. + _assert_reachable(dag, "dq_ods_boarding_passes", "finish_ods_summary") + _assert_reachable(dag, "dq_ods_seats", "finish_ods_summary") diff --git a/tests/test_ods_snapshot_integration.py b/tests/test_ods_snapshot_integration.py new file mode 100644 index 0000000..3c227f6 --- /dev/null +++ b/tests/test_ods_snapshot_integration.py @@ -0,0 +1,150 @@ +from __future__ import annotations + +import os +import subprocess +from pathlib import Path + +import pytest + +PROJECT_ROOT = Path(__file__).resolve().parents[1] +RUN_ODS_INTEGRATION = os.getenv("RUN_ODS_INTEGRATION") == "1" +BATCH_TOKEN = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}' + +pytestmark = pytest.mark.skipif( + not RUN_ODS_INTEGRATION, + reason="Set RUN_ODS_INTEGRATION=1 to run ODS integration tests", +) + + +def _psql(sql: str) -> str: + """Выполняет SQL в Greenplum-контейнере и возвращает stdout psql.""" + docker_bin = os.getenv("DOCKER_BIN", "docker") + cmd = [ + docker_bin, + "compose", + "-f", + "docker-compose.yml", + "exec", + "-T", + "greenplum", + "bash", + "-lc", + "su - gpadmin -c '/usr/local/greenplum-db/bin/psql -v ON_ERROR_STOP=1 -d gp_dwh -At -f -'", + ] + result = subprocess.run( + cmd, + cwd=PROJECT_ROOT, + input=sql, + text=True, + capture_output=True, + check=True, + ) + return result.stdout + + +def _render_airports_sql( + path: str, batch_id: str, stg_table: str, ods_table: str +) -> str: + sql = (PROJECT_ROOT / path).read_text(encoding="utf-8") + sql = sql.replace(BATCH_TOKEN, batch_id) + sql = sql.replace("stg.airports", stg_table) + sql = sql.replace("ods.airports", ods_table) + return sql + + +def test_snapshot_airports_contract_upsert_delete_and_dq() -> None: + """ + Интеграционный тест контракта snapshot-таблицы: + - UPSERT обновляет и вставляет; + - DELETE синхронизирует current state по выбранному батчу; + - DQ проходит на корректном состоянии. + """ + stg_table = "public.it_stg_airports_ods" + ods_table = "public.it_ods_airports_ods" + + setup_sql = f""" + DROP TABLE IF EXISTS {stg_table}; + DROP TABLE IF EXISTS {ods_table}; + + CREATE TABLE {stg_table} ( + airport_code TEXT, + airport_name TEXT, + city TEXT, + country TEXT, + coordinates TEXT, + timezone TEXT, + src_created_at_ts TIMESTAMP, + load_dttm TIMESTAMP, + batch_id TEXT + ); + + CREATE TABLE {ods_table} ( + airport_code TEXT NOT NULL, + airport_name TEXT NOT NULL, + city TEXT NOT NULL, + country TEXT NOT NULL, + coordinates TEXT, + timezone TEXT NOT NULL, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() + ); + """ + _psql(setup_sql) + + try: + _psql( + f""" + INSERT INTO {stg_table} ( + airport_code, airport_name, city, country, coordinates, timezone, + src_created_at_ts, load_dttm, batch_id + ) VALUES + ('AAA', 'Airport A', 'City A', 'Country A', '(0,0)', 'UTC', now(), now(), 'batch_1'), + ('BBB', 'Airport B', 'City B', 'Country B', '(1,1)', 'UTC', now(), now(), 'batch_1'); + """ + ) + + _psql( + _render_airports_sql( + "sql/ods/airports_load.sql", "batch_1", stg_table, ods_table + ) + ) + + codes_batch_1 = _psql( + f"SELECT COALESCE(string_agg(airport_code, ',' ORDER BY airport_code), '') FROM {ods_table};" + ).strip() + assert codes_batch_1 == "AAA,BBB" + + _psql( + f""" + INSERT INTO {stg_table} ( + airport_code, airport_name, city, country, coordinates, timezone, + src_created_at_ts, load_dttm, batch_id + ) VALUES + ('AAA', 'Airport A v2', 'City A', 'Country A', '(0,0)', 'UTC', now(), now(), 'batch_2'), + ('CCC', 'Airport C', 'City C', 'Country C', '(2,2)', 'UTC', now(), now(), 'batch_2'); + """ + ) + + _psql( + _render_airports_sql( + "sql/ods/airports_load.sql", "batch_2", stg_table, ods_table + ) + ) + + codes_batch_2 = _psql( + f"SELECT COALESCE(string_agg(airport_code, ',' ORDER BY airport_code), '') FROM {ods_table};" + ).strip() + assert codes_batch_2 == "AAA,CCC" + + airport_a_name = _psql( + f"SELECT airport_name FROM {ods_table} WHERE airport_code = 'AAA';" + ).strip() + assert airport_a_name == "Airport A v2" + + _psql( + _render_airports_sql( + "sql/ods/airports_dq.sql", "batch_2", stg_table, ods_table + ) + ) + finally: + _psql(f"DROP TABLE IF EXISTS {stg_table}; DROP TABLE IF EXISTS {ods_table};") diff --git a/tests/test_ods_sql_contract.py b/tests/test_ods_sql_contract.py new file mode 100644 index 0000000..0471bda --- /dev/null +++ b/tests/test_ods_sql_contract.py @@ -0,0 +1,38 @@ +from __future__ import annotations + +from pathlib import Path + +PROJECT_ROOT = Path(__file__).resolve().parents[1] +SNAPSHOT_ENTITIES = ("airports", "airplanes", "routes", "seats") + + +def _read(path: str) -> str: + return (PROJECT_ROOT / path).read_text(encoding="utf-8") + + +def test_snapshot_load_scripts_sync_deleted_keys() -> None: + """Snapshot-таблицы в ODS должны удалять ключи, отсутствующие в текущем батче.""" + for entity in SNAPSHOT_ENTITIES: + sql = _read(f"sql/ods/{entity}_load.sql") + assert f"DELETE FROM ods.{entity} AS o" in sql + assert "WITH src_keys AS (" in sql + assert "WHERE NOT EXISTS (" in sql + + +def test_snapshot_dq_checks_extra_keys() -> None: + """DQ snapshot-таблиц должен ловить лишние ключи в ODS относительно текущего батча STG.""" + for entity in SNAPSHOT_ENTITIES: + sql = _read(f"sql/ods/{entity}_dq.sql") + assert "v_extra_keys_count" in sql + assert "найдены лишние ключи" in sql + + +def test_ods_batch_resolver_uses_consistent_snapshot_batches() -> None: + """Резолвер батча должен искать batch_id, общий для всех snapshot-таблиц STG.""" + dag_code = _read("airflow/dags/bookings_to_gp_ods.py") + + for table_name in ("stg.airports", "stg.airplanes", "stg.routes", "stg.seats"): + assert table_name in dag_code + + assert "INTERSECT" in dag_code + assert "candidate_batches" in dag_code