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-23 01:51:40 +03:00
parent dc1dc9c66a
commit 03765d2808
37 changed files with 2704 additions and 25 deletions
+16 -10
View File
@@ -12,9 +12,9 @@
- построения ETL/ELT; - построения ETL/ELT;
- работы с Airflow и Greenplum. - работы с Airflow и Greenplum.
В курсовой у нас один источник данных — демо‑БД **bookings**. В стенде уже есть готовый учебный пример В курсовой у нас один источник данных — демо‑БД **bookings**. В стенде уже есть готовые учебные примеры
загрузки **bookings → stg в Greenplum**, чтобы вы могли сфокусироваться на DWH‑части (ODS/DDS/DM) и загрузки **bookings → stg** и **stg -> ods** в Greenplum, чтобы вы могли сфокусироваться на DWH‑части
не тратить время на инфраструктуру. (ODS/DDS/DM) и не тратить время на инфраструктуру.
## Что внутри ## Что внутри
@@ -44,7 +44,7 @@
Про PXF и технические детали стенда: [docs/stack.md](docs/stack.md). Про PXF и технические детали стенда: [docs/stack.md](docs/stack.md).
## Быстрый старт (основной сценарий: bookings stg) ## Быстрый старт (основной сценарий: bookings -> stg -> ods)
1) Скопируйте настройки: 1) Скопируйте настройки:
@@ -65,14 +65,16 @@ make up
make bookings-init make bookings-init
``` ```
4) Подготовьте STGобъекты в Greenplum (выберите один вариант): 4) Подготовьте STG/ODS-объекты в Greenplum (выберите один вариант):
- Учебный вариант: в Airflow UI запустите DAG `bookings_stg_ddl`; - Учебный вариант: в Airflow UI запустите DAG `bookings_stg_ddl`, затем `bookings_ods_ddl`;
- Технический шорткат: `make ddl-gp` (применяет все DDL разом вручную). - Технический шорткат: `make ddl-gp` (применяет DDL для STG и ODS разом вручную).
5) Запустите основной DAG `bookings_to_gp_stage`. 5) Запустите основной DAG `bookings_to_gp_stage`.
6) Проверьте результат в Greenplum: 6) Запустите DAG `bookings_to_gp_ods`.
7) Проверьте результат в Greenplum:
```bash ```bash
make gp-psql make gp-psql
@@ -80,6 +82,8 @@ make gp-psql
SELECT COUNT(*) FROM stg.bookings; SELECT COUNT(*) FROM stg.bookings;
SELECT COUNT(*) FROM stg.tickets; SELECT COUNT(*) FROM stg.tickets;
SELECT * FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10; 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`. Подробнее про логику DAG и проверки — `docs/bookings_to_gp_stage.md`.
@@ -88,10 +92,11 @@ SELECT * FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10;
Основные (для потока bookings → DWH): Основные (для потока 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, - `bookings_stg_ddl` — создаёт/обновляет весь STG слой для bookings (9 таблиц: bookings, tickets, airports, airplanes,
routes, seats, flights, segments, boarding_passes; включая внешние `*_ext` через PXF); routes, seats, flights, segments, boarding_passes; включая внешние `*_ext` через PXF);
- `bookings_to_gp_stage` — генерирует учебный день в `bookings-db`, затем загружает данные в STG и выполняет DQ‑проверки. - `bookings_to_gp_stage` — генерирует учебный день в `bookings-db`, затем загружает данные в STG и выполняет DQ‑проверки.
- `bookings_ods_ddl` — создаёт/обновляет ODS-таблицы по домену bookings.
- `bookings_to_gp_ods` — загружает данные из STG в ODS (SCD1 UPSERT) и выполняет DQ‑проверки.
Вспомогательные (побочный трек с CSV): Вспомогательные (побочный трек с CSV):
@@ -106,7 +111,7 @@ make up # поднять стек
make logs # логи airflow-webserver и airflow-scheduler make logs # логи airflow-webserver и airflow-scheduler
make gp-psql # psql в Greenplum make gp-psql # psql в Greenplum
make bookings-psql # psql в демо-БД bookings (Postgres) 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 down # остановить и удалить контейнеры/сети (volumes сохраняются)
make clean # полный reset: удалить контейнеры/сети и volumes (данные будут потеряны) make clean # полный reset: удалить контейнеры/сети и volumes (данные будут потеряны)
``` ```
@@ -136,6 +141,7 @@ make clean # полный reset: удалить контейнер
- Учебные задания: `educational-tasks.md`). - Учебные задания: `educational-tasks.md`).
- План тестирования/проверок и негативные кейсы: `TESTING.md`. - План тестирования/проверок и негативные кейсы: `TESTING.md`.
- Дополнительные заметки и технические детали: `docs/README.md`. - Дополнительные заметки и технические детали: `docs/README.md`.
- Детали по ODS DAG: `docs/bookings_to_gp_ods.md`.
## Типичные проблемы и решения ## Типичные проблемы и решения
+98
View File
@@ -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
)
+261
View File
@@ -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
+2
View File
@@ -8,6 +8,7 @@
- [Учебные задания](../educational-tasks.md) - [Учебные задания](../educational-tasks.md)
- [План тестирования и проверки](../TESTING.md) - [План тестирования и проверки](../TESTING.md)
- [Главный учебный DAG: bookings → stg](bookings_to_gp_stage.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) - [Единые конвенции нейминга DWH (служебные поля и SCD)](internal/naming_conventions.md)
- [PXF в этом проекте (проектная реализация)](internal/pxf_bookings.md) - [PXF в этом проекте (проектная реализация)](internal/pxf_bookings.md)
- [Дизайн stg для bookings (черновик)](internal/bookings_stg_design.md) - [Дизайн stg для bookings (черновик)](internal/bookings_stg_design.md)
- [Дизайн ods для bookings (черновик)](internal/bookings_ods_design.md)
- [Про время/UTC в bookings (черновик)](internal/bookings_tz.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/...`); - приведение типов (`TEXT -> TIMESTAMPTZ/NUMERIC/INT/BOOLEAN/...`);
- дедупликацию внутри батча; - дедупликацию внутри батча;
- `UPSERT` (SCD Type 1): обновляем текущую запись при изменении, вставляем новые. - `UPSERT` (SCD Type 1): обновляем текущую запись при изменении, вставляем новые.
- для snapshot-справочников (`airports`, `airplanes`, `routes`, `seats`) синхронизацию ключей:
удаляем из ODS записи, которых нет в выбранном `stg_batch_id`.
Не делаем в ODS (в базовом эталоне): Не делаем в ODS (в базовом эталоне):
- SCD Type 2 с периодами действия; - SCD Type 2 с периодами действия;
@@ -264,23 +266,47 @@ log = logging.getLogger(__name__)
GREENPLUM_CONN_ID = "greenplum_conn" GREENPLUM_CONN_ID = "greenplum_conn"
def _resolve_stg_batch_id(**context): 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 {} conf = context["dag_run"].conf or {}
stg_batch_id = conf.get("stg_batch_id") stg_batch_id = conf.get("stg_batch_id")
if not stg_batch_id: if not stg_batch_id:
# Берём batch_id с самым свежим load_dttm (TIMESTAMP, монотонно растёт). # Берём batch_id, который присутствует во всех snapshot-таблицах STG:
# MAX(batch_id) ненадёжен: run_id — строка вида "manual__2024-...", # airports, airplanes, routes, seats. Это защищает от частично успешных запусков.
# лексикографическая сортировка не гарантирует хронологический порядок.
hook = PostgresHook(postgres_conn_id=GREENPLUM_CONN_ID) hook = PostgresHook(postgres_conn_id=GREENPLUM_CONN_ID)
result = hook.get_first( 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 stg_batch_id = result[0] if result and result[0] else None
if not stg_batch_id: if not stg_batch_id:
raise ValueError( raise ValueError(
"stg_batch_id не найден: передайте в conf или сначала загрузите STG" "stg_batch_id не найден: передайте в conf или сначала выполните bookings_to_gp_stage"
) )
log.info("Используем stg_batch_id = %s", stg_batch_id) log.info("Используем stg_batch_id = %s", stg_batch_id)
@@ -382,6 +408,19 @@ WHERE NOT EXISTS (
WHERE o.airport_code = s.airport_code 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 -- Обновляем статистику для оптимизатора запросов Greenplum
ANALYZE ods.airports; ANALYZE ods.airports;
``` ```
@@ -454,7 +493,10 @@ ANALYZE ods.bookings;
### 6.3. Идемпотентность паттерна ### 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. Поведение при пустом батче ### 6.4. Поведение при пустом батче
@@ -541,14 +583,14 @@ sql/ods/
├── boarding_passes_load.sql ├── boarding_passes_load.sql
└── boarding_passes_dq.sql └── boarding_passes_dq.sql
sql/ddl_gp_ods.sql sql/ddl_gp.sql (+ подключение sql/ods/*_ddl.sql)
airflow/dags/ airflow/dags/
├── bookings_ods_ddl.py ├── bookings_ods_ddl.py
└── bookings_to_gp_ods.py └── bookings_to_gp_ods.py
docs/bookings_to_gp_ods.md docs/bookings_to_gp_ods.md
Makefile (+ ddl-gp-ods) Makefile (ddl-gp включает ODS DDL)
tests/test_dags_smoke.py (+ smoke для 2 новых DAG) tests/test_dags_smoke.py (+ smoke для 2 новых DAG)
``` ```
@@ -591,8 +633,8 @@ resolve_stg_batch_id
## 10) Порядок реализации ## 10) Порядок реализации
1. Подготовить DDL в `sql/ods/*_ddl.sql`. 1. Подготовить DDL в `sql/ods/*_ddl.sql`.
2. Сделать мастер-скрипт `sql/ddl_gp_ods.sql`. 2. Подключить `sql/ods/*_ddl.sql` в общий `sql/ddl_gp.sql`.
3. Добавить `Makefile`-таргет `ddl-gp-ods`. 3. Использовать существующий `Makefile`-таргет `ddl-gp` для STG+ODS.
4. Создать DAG `bookings_ods_ddl.py`. 4. Создать DAG `bookings_ods_ddl.py`.
5. Реализовать `sql/ods/*_load.sql` (SCD1 UPSERT). 5. Реализовать `sql/ods/*_load.sql` (SCD1 UPSERT).
6. Реализовать `sql/ods/*_dq.sql`. 6. Реализовать `sql/ods/*_dq.sql`.
@@ -607,7 +649,7 @@ resolve_stg_batch_id
Готово, если: Готово, если:
1. Оба новых DAG парсятся и проходят smoke-тесты (`make test`). 1. Оба новых DAG парсятся и проходят smoke-тесты (`make test`).
2. `make ddl-gp-ods` создаёт объекты без ошибок. 2. `make ddl-gp` создаёт объекты STG+ODS без ошибок.
3. Для тестового `stg_batch_id` ODS-загрузка завершается успешно. 3. Для тестового `stg_batch_id` ODS-загрузка завершается успешно.
4. Все DQ-задачи зелёные и реально валят DAG при искусственной ошибке. 4. Все DQ-задачи зелёные и реально валят DAG при искусственной ошибке.
5. В ODS нет дублей по бизнес-ключам. 5. В ODS нет дублей по бизнес-ключам.
@@ -619,10 +661,10 @@ resolve_stg_batch_id
```bash ```bash
make up make up
make ddl-gp # создать STG-объекты make ddl-gp # создать STG+ODS-объекты
make ddl-gp-ods # создать ODS-объекты
# Trigger bookings_to_gp_stage (загрузить STG) # 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": "<значение>"} # Trigger bookings_to_gp_ods с conf: {"stg_batch_id": "<значение>"}
make gp-psql make gp-psql
``` ```
+11
View File
@@ -37,3 +37,14 @@ FORMAT 'CUSTOM' (formatter='pxfwritable_import');
\i stg/flights_ddl.sql \i stg/flights_ddl.sql
\i stg/segments_ddl.sql \i stg/segments_ddl.sql
\i stg/boarding_passes_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
+13
View File
@@ -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);
+99
View File
@@ -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 $$;
+92
View File
@@ -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;
+15
View File
@@ -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);
+105
View File
@@ -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 $$;
+104
View File
@@ -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;
+15
View File
@@ -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);
+96
View File
@@ -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 $$;
+81
View File
@@ -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;
+13
View File
@@ -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);
+71
View File
@@ -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 $$;
+69
View File
@@ -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;
+17
View File
@@ -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);
+89
View File
@@ -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 $$;
+93
View File
@@ -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;
+17
View File
@@ -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);
+163
View File
@@ -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 $$;
+118
View File
@@ -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;
+12
View File
@@ -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);
+125
View File
@@ -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 $$;
+86
View File
@@ -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;
+14
View File
@@ -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);
+112
View File
@@ -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 $$;
+75
View File
@@ -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;
+15
View File
@@ -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);
+92
View File
@@ -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 $$;
+81
View File
@@ -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;
+98
View File
@@ -199,3 +199,101 @@ def test_bookings_to_gp_stage_dag_structure():
assert airports not in airplanes.get_flat_relatives( assert airports not in airplanes.get_flat_relatives(
upstream=False upstream=False
), "airplanes не должен быть upstream для airports" ), "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")
+150
View File
@@ -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};")
+38
View File
@@ -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