From d25c753eb45409517d3baa8667bc513720bd049e Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Wed, 25 Feb 2026 23:05:43 +0300 Subject: [PATCH] =?UTF-8?q?feat(dds):=20=D1=80=D0=B5=D0=B0=D0=BB=D0=B8?= =?UTF-8?q?=D0=B7=D0=BE=D0=B2=D0=B0=D0=BD=20=D1=81=D0=BB=D0=BE=D0=B9=20dds?= =?UTF-8?q?=20=D0=B4=D0=BB=D1=8F=20bookings?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - подготовлен учебный Star Schema слой для перехода от ODS к аналитике и витринам. - Что: - добавлены 21 SQL-файл для DDS (DDL/LOAD/DQ) с SCD1/SCD2 и фактом `fact_flight_sales`. - добавлены DAG `bookings_dds_ddl` и `bookings_to_gp_dds`, а также smoke-тесты структуры DAG. - обновлены `sql/ddl_gp.sql` и документация (`README`, `docs/*`, `db_schema`) под поток `stg -> ods -> dds`. - Проверка: - make test. --- README.md | 23 ++-- TESTING.md | 4 +- airflow/dags/bookings_dds_ddl.py | 84 +++++++++++++++ airflow/dags/bookings_to_gp_dds.py | 162 +++++++++++++++++++++++++++++ docs/README.md | 2 + docs/bookings_to_gp_dds.md | 72 +++++++++++++ docs/internal/db_schema.md | 12 ++- sql/ddl_gp.sql | 9 ++ sql/dds/dim_airplanes_ddl.sql | 17 +++ sql/dds/dim_airplanes_dq.sql | 83 +++++++++++++++ sql/dds/dim_airplanes_load.sql | 92 ++++++++++++++++ sql/dds/dim_airports_ddl.sql | 18 ++++ sql/dds/dim_airports_dq.sql | 88 ++++++++++++++++ sql/dds/dim_airports_load.sql | 60 +++++++++++ sql/dds/dim_calendar_ddl.sql | 15 +++ sql/dds/dim_calendar_dq.sql | 87 ++++++++++++++++ sql/dds/dim_calendar_load.sql | 29 ++++++ sql/dds/dim_passengers_ddl.sql | 14 +++ sql/dds/dim_passengers_dq.sql | 100 ++++++++++++++++++ sql/dds/dim_passengers_load.sql | 82 +++++++++++++++ sql/dds/dim_routes_ddl.sql | 22 ++++ sql/dds/dim_routes_dq.sql | 149 ++++++++++++++++++++++++++ sql/dds/dim_routes_load.sql | 130 +++++++++++++++++++++++ sql/dds/dim_tariffs_ddl.sql | 13 +++ sql/dds/dim_tariffs_dq.sql | 98 +++++++++++++++++ sql/dds/dim_tariffs_load.sql | 36 +++++++ sql/dds/fact_flight_sales_ddl.sql | 23 ++++ sql/dds/fact_flight_sales_dq.sql | 151 +++++++++++++++++++++++++++ sql/dds/fact_flight_sales_load.sql | 111 ++++++++++++++++++++ tests/test_dags_smoke.py | 85 +++++++++++++++ 30 files changed, 1856 insertions(+), 15 deletions(-) create mode 100644 airflow/dags/bookings_dds_ddl.py create mode 100644 airflow/dags/bookings_to_gp_dds.py create mode 100644 docs/bookings_to_gp_dds.md create mode 100644 sql/dds/dim_airplanes_ddl.sql create mode 100644 sql/dds/dim_airplanes_dq.sql create mode 100644 sql/dds/dim_airplanes_load.sql create mode 100644 sql/dds/dim_airports_ddl.sql create mode 100644 sql/dds/dim_airports_dq.sql create mode 100644 sql/dds/dim_airports_load.sql create mode 100644 sql/dds/dim_calendar_ddl.sql create mode 100644 sql/dds/dim_calendar_dq.sql create mode 100644 sql/dds/dim_calendar_load.sql create mode 100644 sql/dds/dim_passengers_ddl.sql create mode 100644 sql/dds/dim_passengers_dq.sql create mode 100644 sql/dds/dim_passengers_load.sql create mode 100644 sql/dds/dim_routes_ddl.sql create mode 100644 sql/dds/dim_routes_dq.sql create mode 100644 sql/dds/dim_routes_load.sql create mode 100644 sql/dds/dim_tariffs_ddl.sql create mode 100644 sql/dds/dim_tariffs_dq.sql create mode 100644 sql/dds/dim_tariffs_load.sql create mode 100644 sql/dds/fact_flight_sales_ddl.sql create mode 100644 sql/dds/fact_flight_sales_dq.sql create mode 100644 sql/dds/fact_flight_sales_load.sql diff --git a/README.md b/README.md index a63ecac..ef25906 100644 --- a/README.md +++ b/README.md @@ -13,8 +13,8 @@ - работы с Airflow и Greenplum. В курсовой у нас один источник данных — демо‑БД **bookings**. В стенде уже есть готовые учебные примеры -загрузки **bookings → stg** и **stg -> ods** в Greenplum, чтобы вы могли сфокусироваться на DWH‑части -(ODS/DDS/DM) и не тратить время на инфраструктуру. +загрузки **bookings → stg**, **stg -> ods** и **ods -> dds** в Greenplum, чтобы вы могли сфокусироваться +на DWH‑части (ODS/DDS/DM) и не тратить время на инфраструктуру. ## Что внутри @@ -44,7 +44,7 @@ Про PXF и технические детали стенда: [docs/stack.md](docs/stack.md). -## Быстрый старт (основной сценарий: bookings -> stg -> ods) +## Быстрый старт (основной сценарий: bookings -> stg -> ods -> dds) 1) Скопируйте настройки: @@ -65,16 +65,18 @@ make up make bookings-init ``` -4) Подготовьте STG/ODS-объекты в Greenplum (выберите один вариант): +4) Подготовьте STG/ODS/DDS-объекты в Greenplum (выберите один вариант): -- Учебный вариант: в Airflow UI запустите DAG `bookings_stg_ddl`, затем `bookings_ods_ddl`; -- Технический шорткат: `make ddl-gp` (применяет DDL для STG и ODS разом вручную). +- Учебный вариант: в Airflow UI запустите DAG `bookings_stg_ddl`, затем `bookings_ods_ddl`, затем `bookings_dds_ddl`; +- Технический шорткат: `make ddl-gp` (применяет DDL для STG, ODS и DDS разом вручную). 5) Запустите основной DAG `bookings_to_gp_stage`. 6) Запустите DAG `bookings_to_gp_ods`. -7) Проверьте результат в Greenplum: +7) Запустите DAG `bookings_to_gp_dds`. + +8) Проверьте результат в Greenplum: ```bash make gp-psql @@ -84,6 +86,8 @@ 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; +SELECT COUNT(*) FROM dds.dim_routes; +SELECT COUNT(*) FROM dds.fact_flight_sales; ``` Подробнее про логику DAG и проверки — `docs/bookings_to_gp_stage.md`. @@ -97,6 +101,8 @@ SELECT COUNT(*) FROM ods.tickets; - `bookings_to_gp_stage` — генерирует учебный день в `bookings-db`, затем загружает данные в STG и выполняет DQ‑проверки. - `bookings_ods_ddl` — создаёт/обновляет ODS-таблицы по домену bookings. - `bookings_to_gp_ods` — загружает данные из STG в ODS (SCD1 UPSERT) и выполняет DQ‑проверки. +- `bookings_dds_ddl` — создаёт/обновляет DDS-таблицы (`dim_*`, `fact_flight_sales`) по домену bookings. +- `bookings_to_gp_dds` — загружает данные из ODS в DDS (SCD1/SCD2 + факт) и выполняет DQ‑проверки. Вспомогательные (побочный трек с CSV): @@ -111,7 +117,7 @@ make up # поднять стек make logs # логи airflow-webserver и airflow-scheduler make gp-psql # psql в Greenplum make bookings-psql # psql в демо-БД bookings (Postgres) -make ddl-gp # применить DDL STG+ODS к Greenplum вручную (вместо DDL-DAG) +make ddl-gp # применить DDL STG+ODS+DDS к Greenplum вручную (вместо DDL-DAG) make down # остановить и удалить контейнеры/сети (volumes сохраняются) make clean # полный reset: удалить контейнеры/сети и volumes (данные будут потеряны) ``` @@ -142,6 +148,7 @@ make clean # полный reset: удалить контейнер - План тестирования/проверок и негативные кейсы: `TESTING.md`. - Дополнительные заметки и технические детали: `docs/README.md`. - Детали по ODS DAG: `docs/bookings_to_gp_ods.md`. +- Детали по DDS DAG: `docs/bookings_to_gp_dds.md`. ## Типичные проблемы и решения diff --git a/TESTING.md b/TESTING.md index 3e58938..fb5f99d 100644 --- a/TESTING.md +++ b/TESTING.md @@ -38,7 +38,7 @@ - Проверить, что все 5 задач Success и логи содержат `Проверка пройдена`. - DAG `bookings_to_gp_stage` (полная проверка цепочки bookings → Greenplum STG): - - предварительно выполнить один раз: `make bookings-init` (установка демобазы `demo` в контейнере `bookings-db`) и `make ddl-gp` (создаёт STG слой в Greenplum, включая внешние `*_ext` через PXF); + - предварительно выполнить один раз: `make bookings-init` (установка демобазы `demo` в контейнере `bookings-db`) и `make ddl-gp` (создаёт STG/ODS/DDS слои в Greenplum, включая внешние `*_ext` через PXF); - перед Trigger проверить, что в source реально есть данные (все значения должны быть `> 0`): - `docker compose exec bookings-db psql -U bookings -d demo -At -c "SELECT COUNT(*) FROM bookings.bookings;"` - `docker compose exec bookings-db psql -U bookings -d demo -At -c "SELECT COUNT(*) FROM bookings.airports_data;"` @@ -90,6 +90,6 @@ - При необходимости сохранить данные: скопировать CSV из `data/` и сделать дампы до `make clean`. ## Текущий статус (пример успешного прогона) -- `uv run pytest -q` — 11 passed, 2 smoke-теста DAG пропущены (Airflow не установлен в venv). +- `uv run pytest -q` — 14 passed, 9 smoke-тестов DAG пропущены (Airflow не установлен в venv). - `make lint` — проходит (DAG‑файлы отформатированы black/isort). - Docker-стенд не запускался в рамках этой сессии; ожидается, что инструкции выше обеспечат полноценную проверку. diff --git a/airflow/dags/bookings_dds_ddl.py b/airflow/dags/bookings_dds_ddl.py new file mode 100644 index 0000000..f1e00b7 --- /dev/null +++ b/airflow/dags/bookings_dds_ddl.py @@ -0,0 +1,84 @@ +from __future__ import annotations + +""" +Учебный DAG: создаёт/обновляет слой dds в Greenplum для домена bookings. + +Запускается вручную перед DAG загрузки `bookings_to_gp_dds` или после изменения DDS DDL. +Создаёт 7 DDS-таблиц: dim_calendar, dim_airports, dim_airplanes, +dim_tariffs, dim_passengers, dim_routes, fact_flight_sales. +""" + +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_dds_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", "dds"], + description="Учебный DDL DAG: создаёт/обновляет dds.* для bookings", +) as dag: + apply_dds_dim_calendar_ddl = PostgresOperator( + task_id="apply_dds_dim_calendar_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_calendar_ddl.sql", + ) + + apply_dds_dim_airports_ddl = PostgresOperator( + task_id="apply_dds_dim_airports_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_airports_ddl.sql", + ) + + apply_dds_dim_airplanes_ddl = PostgresOperator( + task_id="apply_dds_dim_airplanes_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_airplanes_ddl.sql", + ) + + apply_dds_dim_tariffs_ddl = PostgresOperator( + task_id="apply_dds_dim_tariffs_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_tariffs_ddl.sql", + ) + + apply_dds_dim_passengers_ddl = PostgresOperator( + task_id="apply_dds_dim_passengers_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_passengers_ddl.sql", + ) + + apply_dds_dim_routes_ddl = PostgresOperator( + task_id="apply_dds_dim_routes_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_routes_ddl.sql", + ) + + apply_dds_fact_flight_sales_ddl = PostgresOperator( + task_id="apply_dds_fact_flight_sales_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/fact_flight_sales_ddl.sql", + ) + + # DDL применяем последовательно, чтобы порядок был понятным для новичков, + # а ошибки — воспроизводимыми (в логах сразу видно, на каком объекте упали). + ( + apply_dds_dim_calendar_ddl + >> apply_dds_dim_airports_ddl + >> apply_dds_dim_airplanes_ddl + >> apply_dds_dim_tariffs_ddl + >> apply_dds_dim_passengers_ddl + >> apply_dds_dim_routes_ddl + >> apply_dds_fact_flight_sales_ddl + ) diff --git a/airflow/dags/bookings_to_gp_dds.py b/airflow/dags/bookings_to_gp_dds.py new file mode 100644 index 0000000..a194b5f --- /dev/null +++ b/airflow/dags/bookings_to_gp_dds.py @@ -0,0 +1,162 @@ +from __future__ import annotations + +""" +Учебный DAG: загрузка из ODS в DDS (Greenplum) по домену bookings. + +Ключевая идея: +- DDS читает текущее состояние ODS; +- для каждой сущности выполняем пару задач load -> dq; +- для dim_routes применяем SCD2, для остальных измерений — SCD1 UPSERT; +- факт грузим инкрементальным UPSERT по зерну (ticket_no, flight_id). +""" + +from datetime import timedelta +from logging import getLogger + +import pendulum +from airflow.operators.python import PythonOperator +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 _finish_summary() -> None: + """Логирует краткий итог выполнения DDS-ветки.""" + log.info("DAG bookings_to_gp_dds завершён. Подробности смотрите в логах задач.") + + +with DAG( + dag_id="bookings_to_gp_dds", + 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", "dds"], + description="Учебный DAG: загрузка ODS -> DDS (Star Schema) + DQ проверки", +) as dag: + load_dds_dim_calendar = PostgresOperator( + task_id="load_dds_dim_calendar", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_calendar_load.sql", + ) + + dq_dds_dim_calendar = PostgresOperator( + task_id="dq_dds_dim_calendar", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_calendar_dq.sql", + ) + + load_dds_dim_airports = PostgresOperator( + task_id="load_dds_dim_airports", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_airports_load.sql", + ) + + dq_dds_dim_airports = PostgresOperator( + task_id="dq_dds_dim_airports", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_airports_dq.sql", + ) + + load_dds_dim_airplanes = PostgresOperator( + task_id="load_dds_dim_airplanes", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_airplanes_load.sql", + ) + + dq_dds_dim_airplanes = PostgresOperator( + task_id="dq_dds_dim_airplanes", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_airplanes_dq.sql", + ) + + load_dds_dim_tariffs = PostgresOperator( + task_id="load_dds_dim_tariffs", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_tariffs_load.sql", + ) + + dq_dds_dim_tariffs = PostgresOperator( + task_id="dq_dds_dim_tariffs", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_tariffs_dq.sql", + ) + + load_dds_dim_passengers = PostgresOperator( + task_id="load_dds_dim_passengers", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_passengers_load.sql", + ) + + dq_dds_dim_passengers = PostgresOperator( + task_id="dq_dds_dim_passengers", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_passengers_dq.sql", + ) + + load_dds_dim_routes = PostgresOperator( + task_id="load_dds_dim_routes", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_routes_load.sql", + ) + + dq_dds_dim_routes = PostgresOperator( + task_id="dq_dds_dim_routes", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_routes_dq.sql", + ) + + load_dds_fact_flight_sales = PostgresOperator( + task_id="load_dds_fact_flight_sales", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/fact_flight_sales_load.sql", + ) + + dq_dds_fact_flight_sales = PostgresOperator( + task_id="dq_dds_fact_flight_sales", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/fact_flight_sales_dq.sql", + ) + + finish_dds_summary = PythonOperator( + task_id="finish_dds_summary", + python_callable=_finish_summary, + ) + + load_dds_dim_calendar >> dq_dds_dim_calendar + + dq_dds_dim_calendar >> [ + load_dds_dim_airports, + load_dds_dim_airplanes, + load_dds_dim_tariffs, + load_dds_dim_passengers, + load_dds_dim_routes, + ] + + load_dds_dim_airports >> dq_dds_dim_airports + load_dds_dim_airplanes >> dq_dds_dim_airplanes + load_dds_dim_tariffs >> dq_dds_dim_tariffs + load_dds_dim_passengers >> dq_dds_dim_passengers + load_dds_dim_routes >> dq_dds_dim_routes + + [ + dq_dds_dim_airports, + dq_dds_dim_airplanes, + dq_dds_dim_tariffs, + dq_dds_dim_passengers, + dq_dds_dim_routes, + ] >> load_dds_fact_flight_sales + + load_dds_fact_flight_sales >> dq_dds_fact_flight_sales >> finish_dds_summary diff --git a/docs/README.md b/docs/README.md index 8fa0215..fc12a28 100644 --- a/docs/README.md +++ b/docs/README.md @@ -9,6 +9,7 @@ - [План тестирования и проверки](../TESTING.md) - [Главный учебный DAG: bookings → stg](bookings_to_gp_stage.md) - [Учебный DAG: stg -> ods](bookings_to_gp_ods.md) +- [Учебный DAG: ods -> dds](bookings_to_gp_dds.md) ## Технические детали (опционально) @@ -17,4 +18,5 @@ - [PXF в этом проекте (проектная реализация)](internal/pxf_bookings.md) - [Дизайн stg для bookings (черновик)](internal/bookings_stg_design.md) - [Дизайн ods для bookings (черновик)](internal/bookings_ods_design.md) +- [Дизайн dds для bookings (черновик)](internal/bookings_dds_design.md) - [Про время/UTC в bookings (черновик)](internal/bookings_tz.md) diff --git a/docs/bookings_to_gp_dds.md b/docs/bookings_to_gp_dds.md new file mode 100644 index 0000000..9cb0c75 --- /dev/null +++ b/docs/bookings_to_gp_dds.md @@ -0,0 +1,72 @@ +# DAG `bookings_to_gp_dds`: `ods` -> `dds` в Greenplum + +Этот DAG — учебный пример загрузки аналитического слоя **DDS** (Star Schema) из текущего состояния **ODS**. +Логика: измерения + факт, проверки качества данных после каждой загрузки. + +## Что делает DAG + +- Загружает измерения DDS: + - `dds.dim_calendar` (статическое измерение дат); + - `dds.dim_airports`, `dds.dim_airplanes`, `dds.dim_tariffs`, `dds.dim_passengers` (SCD1 UPSERT); + - `dds.dim_routes` (SCD2 с `hashdiff`, `valid_from`, `valid_to`). +- Загружает факт `dds.fact_flight_sales` инкрементальным UPSERT по зерну `(ticket_no, flight_id)`. +- Для каждой таблицы выполняет пару задач `load -> dq`. +- Использует `_load_id = {{ run_id }}` (DDS не требует `stg_batch_id`, потому что читает current state ODS). + +## Что должно быть готово перед запуском + +1) Стенд поднят: + +```bash +make up +``` + +2) STG и ODS уже загружены: + +- выполнены DAG-и `bookings_to_gp_stage` и `bookings_to_gp_ods`; +- DDL-объекты созданы (`bookings_dds_ddl` или `make ddl-gp`). + +## Как запустить + +1) Откройте Airflow UI: http://localhost:8080. +2) Если запускаете DDS впервые — выполните `bookings_dds_ddl`. +3) Запустите `bookings_to_gp_dds`. + +## Граф зависимостей (упрощённо) + +- `load_dds_dim_calendar -> dq_dds_dim_calendar` +- После calendar параллельно: + - `load_dds_dim_airports -> dq_dds_dim_airports` + - `load_dds_dim_airplanes -> dq_dds_dim_airplanes` + - `load_dds_dim_tariffs -> dq_dds_dim_tariffs` + - `load_dds_dim_passengers -> dq_dds_dim_passengers` + - `load_dds_dim_routes -> dq_dds_dim_routes` +- Факт: + - `load_dds_fact_flight_sales -> dq_dds_fact_flight_sales -> finish_dds_summary` + +## Как проверить результат + +```bash +make gp-psql +``` + +```sql +SELECT COUNT(*) FROM dds.dim_calendar; +SELECT COUNT(*) FROM dds.dim_routes; +SELECT COUNT(*) FROM dds.fact_flight_sales; + +SELECT + (SELECT COUNT(*) FROM dds.fact_flight_sales) AS fact_rows, + (SELECT COUNT(*) FROM ods.segments) AS ods_rows; +``` + +Ожидаемо: `fact_rows = ods_rows`. + +## Типичные ошибки + +- `relation "dds..." does not exist`: + - не применён DDS DDL (`bookings_dds_ddl` или `make ddl-gp`). +- DQ падает на `dim_routes`: + - проверьте согласованность `ods.routes` (дубли/аномальные версии) и перезапустите DAG. +- DQ падает на `fact_flight_sales` по coverage: + - проверьте, что ODS DAG завершился успешно без пропуска задач. diff --git a/docs/internal/db_schema.md b/docs/internal/db_schema.md index 4800b5c..1a63323 100644 --- a/docs/internal/db_schema.md +++ b/docs/internal/db_schema.md @@ -1,6 +1,6 @@ # Схема БД DWH (Bookings → Greenplum) -> **Статус:** Проект в разработке. Реализованы STG и ODS (по 9 таблиц). DDS зафиксирован как дизайн и готовится к реализации. +> **Статус:** Проект в разработке. Реализованы STG, ODS и DDS (bookings). ## Обзор @@ -32,7 +32,7 @@ | **Source** | ✅ Готово | Демо-БД bookings (Postgres) | | **STG** | ✅ Готово | 9 из 9 таблиц (bookings, tickets, airports, airplanes, routes, seats, flights, segments, boarding_passes) | | **ODS** | ✅ Готово | 9 из 9 таблиц + DAG `bookings_ods_ddl` и `bookings_to_gp_ods` | -| **DDS** | ⚙️ В проектировании | Подготовлен дизайн `docs/internal/bookings_dds_design.md`, реализация запланирована | +| **DDS** | ✅ Готово | 6 измерений + 1 факт + DAG `bookings_dds_ddl` и `bookings_to_gp_dds` | ### Архитектура слоёв @@ -59,7 +59,7 @@ - **Хранение**: Heap или AO-CO (Append-Only Column-oriented) для аналитических запросов - **Структура**: Измерения (Dimensions) + Факты (Facts) - **Ключи**: Суррогатные ключи (SK) для измерений, FK в фактах -- **Текущий статус**: Зафиксирован детальный дизайн (см. `docs/internal/bookings_dds_design.md`) +- **Текущий статус**: Реализован (6 измерений, 1 факт, SQL DQ, DAG загрузки) ### Измерения DDS (Dimensions) @@ -419,6 +419,7 @@ graph LR - [`docs/internal/bookings_stg_design.md`](bookings_stg_design.md) — Детальный дизайн STG слоя для bookings - [`docs/internal/bookings_ods_design.md`](bookings_ods_design.md) — Детальный дизайн ODS слоя (SCD1, batch contract, DQ) - [`docs/internal/bookings_dds_design.md`](bookings_dds_design.md) — План реализации DDS слоя (Star Schema, SCD2 для routes) +- [`docs/bookings_to_gp_dds.md`](../bookings_to_gp_dds.md) — Запуск и проверка DAG `bookings_to_gp_dds` - [`docs/internal/bookings_stg_code_review.md`](bookings_stg_code_review.md) — Ревью решения и рекомендации по улучшению - [`docs/internal/bookings_tz.md`](bookings_tz.md) — Работа с часовыми поясами в источнике - [`docs/internal/pxf_bookings.md`](pxf_bookings.md) — Настройка PXF для чтения из bookings-db @@ -430,6 +431,7 @@ graph LR | Дата | Версия | Описание изменений | |------|--------|-------------------| +| 2026-02-25 | 2.2 | Реализован DDS: добавлены `sql/dds/*` (DDL/LOAD/DQ), DAG `bookings_dds_ddl`, DAG `bookings_to_gp_dds`, обновлены smoke-тесты и документация. | | 2026-02-23 | 2.1 | Актуализирован статус: STG+ODS реализованы. Обновлены DDS-объекты (`dds.dim_*`, `dds.fact_flight_sales`), добавлен `dds.dim_routes` (SCD2), исправлены диаграмма и TODO. | | 2025-01-17 | 2.0 | Удалён слой DQ для упрощения учебного стенда. Добавлены спецификации для LLM и обучающие материалы для студентов. Добавлен глоссарий терминов. | | 2025-01-17 | 1.1 | Исправлены названия таблиц (`aircrafts_data` → `airplanes_data`, `ticket_flights` → `segments`), удалено `dim.bookings`, добавлены суррогатные ключи, добавлен слой DQ, исправлены связи | @@ -441,6 +443,6 @@ graph LR - [x] Реализовать STG слой полностью (все 9 таблиц) - [x] Реализовать ODS слой -- [ ] Реализовать DDS слой (измерения и факт) +- [x] Реализовать DDS слой (измерения и факт) - [x] Создать DAG для загрузки ODS -- [ ] Создать DAG для загрузки DDS +- [x] Создать DAG для загрузки DDS diff --git a/sql/ddl_gp.sql b/sql/ddl_gp.sql index 9df705b..a270a0b 100644 --- a/sql/ddl_gp.sql +++ b/sql/ddl_gp.sql @@ -48,3 +48,12 @@ FORMAT 'CUSTOM' (formatter='pxfwritable_import'); \i ods/flights_ddl.sql \i ods/segments_ddl.sql \i ods/boarding_passes_ddl.sql + +-- DDL для DDS-слоя (Star Schema, SCD1 + SCD2). +\i dds/dim_calendar_ddl.sql +\i dds/dim_airports_ddl.sql +\i dds/dim_airplanes_ddl.sql +\i dds/dim_tariffs_ddl.sql +\i dds/dim_passengers_ddl.sql +\i dds/dim_routes_ddl.sql +\i dds/fact_flight_sales_ddl.sql diff --git a/sql/dds/dim_airplanes_ddl.sql b/sql/dds/dim_airplanes_ddl.sql new file mode 100644 index 0000000..5376870 --- /dev/null +++ b/sql/dds/dim_airplanes_ddl.sql @@ -0,0 +1,17 @@ +-- DDL для DDS-слоя по таблице dim_airplanes (SCD1-измерение). + +CREATE SCHEMA IF NOT EXISTS dds; + +CREATE TABLE IF NOT EXISTS dds.dim_airplanes ( + airplane_sk INTEGER NOT NULL, + airplane_bk TEXT NOT NULL, + model TEXT NOT NULL, + range_km INTEGER, + speed_kmh INTEGER, + total_seats INTEGER, + created_at TIMESTAMP NOT NULL DEFAULT now(), + updated_at TIMESTAMP NOT NULL DEFAULT now(), + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (airplane_sk); diff --git a/sql/dds/dim_airplanes_dq.sql b/sql/dds/dim_airplanes_dq.sql new file mode 100644 index 0000000..c0d01dd --- /dev/null +++ b/sql/dds/dim_airplanes_dq.sql @@ -0,0 +1,83 @@ +-- DQ для DDS dim_airplanes. + +DO $$ +DECLARE + v_row_count BIGINT; + v_dup_sk BIGINT; + v_dup_bk BIGINT; + v_missing_bk BIGINT; + v_null_count BIGINT; +BEGIN + -- Таблица не пуста. + SELECT COUNT(*) + INTO v_row_count + FROM dds.dim_airplanes; + + IF v_row_count = 0 THEN + RAISE EXCEPTION 'DQ FAILED: dds.dim_airplanes пуста.'; + END IF; + + -- Нет дублей по SK. + SELECT COUNT(*) - COUNT(DISTINCT airplane_sk) + INTO v_dup_sk + FROM dds.dim_airplanes; + + IF v_dup_sk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_airplanes найдены дубликаты airplane_sk: %', + v_dup_sk; + END IF; + + -- Нет дублей по BK. + SELECT COUNT(*) - COUNT(DISTINCT airplane_bk) + INTO v_dup_bk + FROM dds.dim_airplanes; + + IF v_dup_bk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_airplanes найдены дубликаты airplane_bk: %', + v_dup_bk; + END IF; + + -- Покрытие ODS: все airplane_code из ODS есть в DDS. + SELECT COUNT(*) + INTO v_missing_bk + FROM (SELECT DISTINCT airplane_code FROM ods.airplanes) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM dds.dim_airplanes AS d + WHERE d.airplane_bk = s.airplane_code + ); + + IF v_missing_bk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_airplanes отсутствуют ключи из ods.airplanes: %', + v_missing_bk; + END IF; + + -- Обязательные поля. + SELECT COUNT(*) + INTO v_null_count + FROM dds.dim_airplanes + WHERE airplane_sk IS NULL + OR airplane_bk IS NULL + OR airplane_bk = '' + OR model IS NULL + OR model = '' + OR total_seats IS NULL + OR created_at IS NULL + OR updated_at 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: в dds.dim_airplanes найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: dds.dim_airplanes ок, строк=%', + v_row_count; +END $$; diff --git a/sql/dds/dim_airplanes_load.sql b/sql/dds/dim_airplanes_load.sql new file mode 100644 index 0000000..6eb1344 --- /dev/null +++ b/sql/dds/dim_airplanes_load.sql @@ -0,0 +1,92 @@ +-- Загрузка DDS dim_airplanes: SCD1 UPSERT (UPDATE изменившихся + INSERT новых). + +-- Statement 1: UPDATE существующих записей (если атрибуты изменились). +WITH seats_agg AS ( + SELECT + s.airplane_code, + COUNT(*)::INTEGER AS total_seats + FROM ods.seats AS s + GROUP BY s.airplane_code +), +src AS ( + SELECT + a.airplane_code, + a.model, + a.range_km, + a.speed_kmh, + COALESCE(sa.total_seats, 0) AS total_seats + FROM ods.airplanes AS a + LEFT JOIN seats_agg AS sa + ON sa.airplane_code = a.airplane_code +) +UPDATE dds.dim_airplanes AS d +SET model = s.model, + range_km = s.range_km, + speed_kmh = s.speed_kmh, + total_seats = s.total_seats, + updated_at = now(), + _load_id = '{{ run_id }}', + _load_ts = now() +FROM src AS s +WHERE d.airplane_bk = s.airplane_code + AND ( + d.model IS DISTINCT FROM s.model + OR d.range_km IS DISTINCT FROM s.range_km + OR d.speed_kmh IS DISTINCT FROM s.speed_kmh + OR d.total_seats IS DISTINCT FROM s.total_seats + ); + +-- Statement 2: INSERT новых записей (MAX(sk) + ROW_NUMBER()). +WITH seats_agg AS ( + SELECT + s.airplane_code, + COUNT(*)::INTEGER AS total_seats + FROM ods.seats AS s + GROUP BY s.airplane_code +), +src AS ( + SELECT + a.airplane_code, + a.model, + a.range_km, + a.speed_kmh, + COALESCE(sa.total_seats, 0) AS total_seats + FROM ods.airplanes AS a + LEFT JOIN seats_agg AS sa + ON sa.airplane_code = a.airplane_code +), +max_sk AS ( + SELECT COALESCE(MAX(airplane_sk), 0) AS v + FROM dds.dim_airplanes +) +INSERT INTO dds.dim_airplanes ( + airplane_sk, + airplane_bk, + model, + range_km, + speed_kmh, + total_seats, + created_at, + updated_at, + _load_id, + _load_ts +) +SELECT + (SELECT v FROM max_sk) + ROW_NUMBER() OVER (ORDER BY s.airplane_code)::INTEGER, + s.airplane_code, + s.model, + s.range_km, + s.speed_kmh, + s.total_seats, + now(), + now(), + '{{ run_id }}', + now() +FROM src AS s +WHERE NOT EXISTS ( + SELECT 1 + FROM dds.dim_airplanes AS d + WHERE d.airplane_bk = s.airplane_code +); + +ANALYZE dds.dim_airplanes; diff --git a/sql/dds/dim_airports_ddl.sql b/sql/dds/dim_airports_ddl.sql new file mode 100644 index 0000000..4495aad --- /dev/null +++ b/sql/dds/dim_airports_ddl.sql @@ -0,0 +1,18 @@ +-- DDL для DDS-слоя по таблице dim_airports (SCD1-измерение). + +CREATE SCHEMA IF NOT EXISTS dds; + +CREATE TABLE IF NOT EXISTS dds.dim_airports ( + airport_sk INTEGER NOT NULL, + airport_bk TEXT NOT NULL, + airport_name TEXT NOT NULL, + city TEXT NOT NULL, + country TEXT NOT NULL, + timezone TEXT NOT NULL, + coordinates TEXT, + created_at TIMESTAMP NOT NULL DEFAULT now(), + updated_at TIMESTAMP NOT NULL DEFAULT now(), + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (airport_sk); diff --git a/sql/dds/dim_airports_dq.sql b/sql/dds/dim_airports_dq.sql new file mode 100644 index 0000000..640951d --- /dev/null +++ b/sql/dds/dim_airports_dq.sql @@ -0,0 +1,88 @@ +-- DQ для DDS dim_airports. + +DO $$ +DECLARE + v_row_count BIGINT; + v_dup_sk BIGINT; + v_dup_bk BIGINT; + v_missing_bk BIGINT; + v_null_count BIGINT; +BEGIN + -- Таблица не пуста. + SELECT COUNT(*) + INTO v_row_count + FROM dds.dim_airports; + + IF v_row_count = 0 THEN + RAISE EXCEPTION 'DQ FAILED: dds.dim_airports пуста.'; + END IF; + + -- Нет дублей по SK. + SELECT COUNT(*) - COUNT(DISTINCT airport_sk) + INTO v_dup_sk + FROM dds.dim_airports; + + IF v_dup_sk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_airports найдены дубликаты airport_sk: %', + v_dup_sk; + END IF; + + -- Нет дублей по BK. + SELECT COUNT(*) - COUNT(DISTINCT airport_bk) + INTO v_dup_bk + FROM dds.dim_airports; + + IF v_dup_bk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_airports найдены дубликаты airport_bk: %', + v_dup_bk; + END IF; + + -- Покрытие ODS: все airport_code из ODS есть в DDS. + SELECT COUNT(*) + INTO v_missing_bk + FROM (SELECT DISTINCT airport_code FROM ods.airports) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM dds.dim_airports AS d + WHERE d.airport_bk = s.airport_code + ); + + IF v_missing_bk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_airports отсутствуют ключи из ods.airports: %', + v_missing_bk; + END IF; + + -- Обязательные поля. + SELECT COUNT(*) + INTO v_null_count + FROM dds.dim_airports + WHERE airport_sk IS NULL + OR airport_bk IS NULL + OR airport_bk = '' + 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 created_at IS NULL + OR updated_at 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: в dds.dim_airports найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: dds.dim_airports ок, строк=%', + v_row_count; +END $$; diff --git a/sql/dds/dim_airports_load.sql b/sql/dds/dim_airports_load.sql new file mode 100644 index 0000000..666e0bc --- /dev/null +++ b/sql/dds/dim_airports_load.sql @@ -0,0 +1,60 @@ +-- Загрузка DDS dim_airports: SCD1 UPSERT (UPDATE изменившихся + INSERT новых). + +-- Statement 1: UPDATE существующих записей (если атрибуты изменились). +UPDATE dds.dim_airports AS d +SET airport_name = s.airport_name, + city = s.city, + country = s.country, + timezone = s.timezone, + coordinates = s.coordinates, + updated_at = now(), + _load_id = '{{ run_id }}', + _load_ts = now() +FROM ods.airports AS s +WHERE d.airport_bk = s.airport_code + AND ( + d.airport_name IS DISTINCT FROM s.airport_name + OR d.city IS DISTINCT FROM s.city + OR d.country IS DISTINCT FROM s.country + OR d.timezone IS DISTINCT FROM s.timezone + OR d.coordinates IS DISTINCT FROM s.coordinates + ); + +-- Statement 2: INSERT новых записей (MAX(sk) + ROW_NUMBER()). +WITH max_sk AS ( + SELECT COALESCE(MAX(airport_sk), 0) AS v + FROM dds.dim_airports +) +INSERT INTO dds.dim_airports ( + airport_sk, + airport_bk, + airport_name, + city, + country, + timezone, + coordinates, + created_at, + updated_at, + _load_id, + _load_ts +) +SELECT + (SELECT v FROM max_sk) + ROW_NUMBER() OVER (ORDER BY s.airport_code)::INTEGER, + s.airport_code, + s.airport_name, + s.city, + s.country, + s.timezone, + s.coordinates, + now(), + now(), + '{{ run_id }}', + now() +FROM ods.airports AS s +WHERE NOT EXISTS ( + SELECT 1 + FROM dds.dim_airports AS d + WHERE d.airport_bk = s.airport_code +); + +ANALYZE dds.dim_airports; diff --git a/sql/dds/dim_calendar_ddl.sql b/sql/dds/dim_calendar_ddl.sql new file mode 100644 index 0000000..8611829 --- /dev/null +++ b/sql/dds/dim_calendar_ddl.sql @@ -0,0 +1,15 @@ +-- DDL для DDS-слоя по таблице dim_calendar (статическое измерение дат). + +CREATE SCHEMA IF NOT EXISTS dds; + +CREATE TABLE IF NOT EXISTS dds.dim_calendar ( + calendar_sk INTEGER NOT NULL, + date_actual DATE NOT NULL, + year_actual INTEGER NOT NULL, + month_actual INTEGER NOT NULL, + day_actual INTEGER NOT NULL, + day_of_week INTEGER NOT NULL, + day_name TEXT NOT NULL, + is_weekend BOOLEAN NOT NULL +) +DISTRIBUTED BY (calendar_sk); diff --git a/sql/dds/dim_calendar_dq.sql b/sql/dds/dim_calendar_dq.sql new file mode 100644 index 0000000..f0c2089 --- /dev/null +++ b/sql/dds/dim_calendar_dq.sql @@ -0,0 +1,87 @@ +-- DQ для DDS dim_calendar. + +DO $$ +DECLARE + v_row_count BIGINT; + v_dup_sk BIGINT; + v_dup_date BIGINT; + v_null_count BIGINT; + v_missing_flight_dates BIGINT; +BEGIN + -- Таблица должна быть достаточно заполнена. + SELECT COUNT(*) + INTO v_row_count + FROM dds.dim_calendar; + + IF v_row_count < 1000 THEN + RAISE EXCEPTION + 'DQ FAILED: dds.dim_calendar содержит слишком мало строк: %', + v_row_count; + END IF; + + -- Нет дублей по surrogate key. + SELECT COUNT(*) - COUNT(DISTINCT calendar_sk) + INTO v_dup_sk + FROM dds.dim_calendar; + + IF v_dup_sk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_calendar найдены дубликаты calendar_sk: %', + v_dup_sk; + END IF; + + -- Нет дублей по business key (date_actual). + SELECT COUNT(*) - COUNT(DISTINCT date_actual) + INTO v_dup_date + FROM dds.dim_calendar; + + IF v_dup_date <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_calendar найдены дубликаты date_actual: %', + v_dup_date; + END IF; + + -- Обязательные поля не NULL. + SELECT COUNT(*) + INTO v_null_count + FROM dds.dim_calendar + WHERE calendar_sk IS NULL + OR date_actual IS NULL + OR year_actual IS NULL + OR month_actual IS NULL + OR day_actual IS NULL + OR day_of_week IS NULL + OR day_name IS NULL + OR day_name = '' + OR is_weekend IS NULL; + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_calendar найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + -- Календарь покрывает даты вылета из ODS (где scheduled_departure не NULL). + SELECT COUNT(*) + INTO v_missing_flight_dates + FROM ( + SELECT DISTINCT f.scheduled_departure::DATE AS departure_date + FROM ods.flights AS f + WHERE f.scheduled_departure IS NOT NULL + ) AS src + WHERE NOT EXISTS ( + SELECT 1 + FROM dds.dim_calendar AS c + WHERE c.date_actual = src.departure_date + ); + + IF v_missing_flight_dates <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: dds.dim_calendar не покрывает даты вылета из ods.flights: %', + v_missing_flight_dates; + END IF; + + RAISE NOTICE + 'DQ PASSED: dds.dim_calendar ок, строк=%', + v_row_count; +END $$; diff --git a/sql/dds/dim_calendar_load.sql b/sql/dds/dim_calendar_load.sql new file mode 100644 index 0000000..1c9e31f --- /dev/null +++ b/sql/dds/dim_calendar_load.sql @@ -0,0 +1,29 @@ +-- Загрузка DDS dim_calendar: статическое измерение (генерация дат). +-- Заполняем только если таблица пуста (идемпотентно). + +INSERT INTO dds.dim_calendar ( + calendar_sk, + date_actual, + year_actual, + month_actual, + day_actual, + day_of_week, + day_name, + is_weekend +) +SELECT + ROW_NUMBER() OVER (ORDER BY d.date_actual)::INTEGER AS calendar_sk, + d.date_actual, + EXTRACT(YEAR FROM d.date_actual)::INTEGER AS year_actual, + EXTRACT(MONTH FROM d.date_actual)::INTEGER AS month_actual, + EXTRACT(DAY FROM d.date_actual)::INTEGER AS day_actual, + EXTRACT(ISODOW FROM d.date_actual)::INTEGER AS day_of_week, + TO_CHAR(d.date_actual, 'FMDay') AS day_name, + EXTRACT(ISODOW FROM d.date_actual) IN (6, 7) AS is_weekend +FROM ( + SELECT generate_series('2016-01-01'::DATE, '2030-12-31'::DATE, '1 day'::INTERVAL)::DATE + AS date_actual +) AS d +WHERE NOT EXISTS (SELECT 1 FROM dds.dim_calendar LIMIT 1); + +ANALYZE dds.dim_calendar; diff --git a/sql/dds/dim_passengers_ddl.sql b/sql/dds/dim_passengers_ddl.sql new file mode 100644 index 0000000..1d1de23 --- /dev/null +++ b/sql/dds/dim_passengers_ddl.sql @@ -0,0 +1,14 @@ +-- DDL для DDS-слоя по таблице dim_passengers (SCD1-измерение). + +CREATE SCHEMA IF NOT EXISTS dds; + +CREATE TABLE IF NOT EXISTS dds.dim_passengers ( + passenger_sk INTEGER NOT NULL, + passenger_bk TEXT NOT NULL, + passenger_name TEXT NOT NULL, + created_at TIMESTAMP NOT NULL DEFAULT now(), + updated_at TIMESTAMP NOT NULL DEFAULT now(), + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (passenger_sk); diff --git a/sql/dds/dim_passengers_dq.sql b/sql/dds/dim_passengers_dq.sql new file mode 100644 index 0000000..8cbcaa5 --- /dev/null +++ b/sql/dds/dim_passengers_dq.sql @@ -0,0 +1,100 @@ +-- DQ для DDS dim_passengers. + +DO $$ +DECLARE + v_row_count BIGINT; + v_src_count BIGINT; + v_dup_sk BIGINT; + v_dup_bk BIGINT; + v_missing_bk BIGINT; + v_null_count BIGINT; +BEGIN + -- Для инкрементальных периодов без новых билетов допускаем пустую dim_passengers. + SELECT COUNT(*) + INTO v_row_count + FROM dds.dim_passengers; + + SELECT COUNT(DISTINCT passenger_id) + INTO v_src_count + FROM ods.tickets + WHERE passenger_id IS NOT NULL + AND passenger_id <> ''; + + IF v_row_count = 0 AND v_src_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: dds.dim_passengers пуста при непустом источнике ods.tickets (passenger_id=%).', + v_src_count; + ELSIF v_row_count = 0 AND v_src_count = 0 THEN + RAISE NOTICE + 'DQ PASSED: dds.dim_passengers пуста, т.к. в ods.tickets нет passenger_id для загрузки.'; + RETURN; + END IF; + + -- Нет дублей по SK. + SELECT COUNT(*) - COUNT(DISTINCT passenger_sk) + INTO v_dup_sk + FROM dds.dim_passengers; + + IF v_dup_sk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_passengers найдены дубликаты passenger_sk: %', + v_dup_sk; + END IF; + + -- Нет дублей по BK. + SELECT COUNT(*) - COUNT(DISTINCT passenger_bk) + INTO v_dup_bk + FROM dds.dim_passengers; + + IF v_dup_bk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_passengers найдены дубликаты passenger_bk: %', + v_dup_bk; + END IF; + + -- Покрытие ODS: все passenger_id из ods.tickets есть в DDS. + SELECT COUNT(*) + INTO v_missing_bk + FROM ( + SELECT DISTINCT passenger_id + FROM ods.tickets + WHERE passenger_id IS NOT NULL + AND passenger_id <> '' + ) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM dds.dim_passengers AS d + WHERE d.passenger_bk = s.passenger_id + ); + + IF v_missing_bk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_passengers отсутствуют passenger_id из ods.tickets: %', + v_missing_bk; + END IF; + + -- Обязательные поля. + SELECT COUNT(*) + INTO v_null_count + FROM dds.dim_passengers + WHERE passenger_sk IS NULL + OR passenger_bk IS NULL + OR passenger_bk = '' + OR passenger_name IS NULL + OR passenger_name = '' + OR created_at IS NULL + OR updated_at 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: в dds.dim_passengers найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: dds.dim_passengers ок, строк=%', + v_row_count; +END $$; diff --git a/sql/dds/dim_passengers_load.sql b/sql/dds/dim_passengers_load.sql new file mode 100644 index 0000000..df48c0e --- /dev/null +++ b/sql/dds/dim_passengers_load.sql @@ -0,0 +1,82 @@ +-- Загрузка DDS dim_passengers: SCD1 UPSERT (UPDATE изменившихся + INSERT новых). + +-- Statement 1: UPDATE существующих записей (если атрибуты изменились). +WITH src AS ( + SELECT + d.passenger_id, + d.passenger_name + FROM ( + SELECT + t.passenger_id, + t.passenger_name, + ROW_NUMBER() OVER ( + PARTITION BY t.passenger_id + ORDER BY t.event_ts DESC NULLS LAST, t._load_ts DESC, t.ticket_no DESC + ) AS rn + FROM ods.tickets AS t + WHERE t.passenger_id IS NOT NULL + AND t.passenger_id <> '' + AND t.passenger_name IS NOT NULL + AND t.passenger_name <> '' + ) AS d + WHERE d.rn = 1 +) +UPDATE dds.dim_passengers AS d +SET passenger_name = s.passenger_name, + updated_at = now(), + _load_id = '{{ run_id }}', + _load_ts = now() +FROM src AS s +WHERE d.passenger_bk = s.passenger_id + AND d.passenger_name IS DISTINCT FROM s.passenger_name; + +-- Statement 2: INSERT новых записей (MAX(sk) + ROW_NUMBER()). +WITH src AS ( + SELECT + d.passenger_id, + d.passenger_name + FROM ( + SELECT + t.passenger_id, + t.passenger_name, + ROW_NUMBER() OVER ( + PARTITION BY t.passenger_id + ORDER BY t.event_ts DESC NULLS LAST, t._load_ts DESC, t.ticket_no DESC + ) AS rn + FROM ods.tickets AS t + WHERE t.passenger_id IS NOT NULL + AND t.passenger_id <> '' + AND t.passenger_name IS NOT NULL + AND t.passenger_name <> '' + ) AS d + WHERE d.rn = 1 +), +max_sk AS ( + SELECT COALESCE(MAX(passenger_sk), 0) AS v + FROM dds.dim_passengers +) +INSERT INTO dds.dim_passengers ( + passenger_sk, + passenger_bk, + passenger_name, + created_at, + updated_at, + _load_id, + _load_ts +) +SELECT + (SELECT v FROM max_sk) + ROW_NUMBER() OVER (ORDER BY s.passenger_id)::INTEGER, + s.passenger_id, + s.passenger_name, + now(), + now(), + '{{ run_id }}', + now() +FROM src AS s +WHERE NOT EXISTS ( + SELECT 1 + FROM dds.dim_passengers AS d + WHERE d.passenger_bk = s.passenger_id +); + +ANALYZE dds.dim_passengers; diff --git a/sql/dds/dim_routes_ddl.sql b/sql/dds/dim_routes_ddl.sql new file mode 100644 index 0000000..81892d4 --- /dev/null +++ b/sql/dds/dim_routes_ddl.sql @@ -0,0 +1,22 @@ +-- DDL для DDS-слоя по таблице dim_routes (SCD2-измерение). + +CREATE SCHEMA IF NOT EXISTS dds; + +CREATE TABLE IF NOT EXISTS dds.dim_routes ( + route_sk INTEGER NOT NULL, + route_bk 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, + hashdiff TEXT NOT NULL, + valid_from DATE NOT NULL, + valid_to DATE, + created_at TIMESTAMP NOT NULL DEFAULT now(), + updated_at TIMESTAMP NOT NULL DEFAULT now(), + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (route_sk); diff --git a/sql/dds/dim_routes_dq.sql b/sql/dds/dim_routes_dq.sql new file mode 100644 index 0000000..fe1e56f --- /dev/null +++ b/sql/dds/dim_routes_dq.sql @@ -0,0 +1,149 @@ +-- DQ для DDS dim_routes (SCD2). + +DO $$ +DECLARE + v_row_count BIGINT; + v_dup_sk BIGINT; + v_dup_current BIGINT; + v_overlap_count BIGINT; + v_missing_count BIGINT; + v_orphan_current BIGINT; + v_null_count BIGINT; +BEGIN + -- Таблица не пуста. + SELECT COUNT(*) + INTO v_row_count + FROM dds.dim_routes; + + IF v_row_count = 0 THEN + RAISE EXCEPTION 'DQ FAILED: dds.dim_routes пуста.'; + END IF; + + -- Нет дублей по SK. + SELECT COUNT(*) - COUNT(DISTINCT route_sk) + INTO v_dup_sk + FROM dds.dim_routes; + + IF v_dup_sk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_routes найдены дубликаты route_sk: %', + v_dup_sk; + END IF; + + -- Корректность интервалов (valid_from <= valid_to для закрытых версий). + SELECT COUNT(*) + INTO v_null_count + FROM dds.dim_routes + WHERE valid_to IS NOT NULL + AND valid_from > valid_to; + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_routes найдены версии с valid_from > valid_to: %', + v_null_count; + END IF; + + -- Нет перекрытий интервалов для одного route_bk. + SELECT COUNT(*) + INTO v_overlap_count + FROM ( + SELECT 1 + FROM dds.dim_routes AS d1 + JOIN dds.dim_routes AS d2 + ON d1.route_bk = d2.route_bk + AND d1.route_sk < d2.route_sk + AND d1.valid_from < COALESCE(d2.valid_to, DATE '9999-12-31') + AND d2.valid_from < COALESCE(d1.valid_to, DATE '9999-12-31') + ) AS overlap_rows; + + IF v_overlap_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_routes найдены перекрытия SCD2-интервалов: %', + v_overlap_count; + END IF; + + -- Не более одной текущей версии на route_bk. + SELECT COUNT(*) + INTO v_dup_current + FROM ( + SELECT route_bk + FROM dds.dim_routes + WHERE valid_to IS NULL + GROUP BY route_bk + HAVING COUNT(*) > 1 + ) AS d; + + IF v_dup_current <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_routes найдены route_bk с > 1 текущей версией: %', + v_dup_current; + END IF; + + -- Покрытие ODS: все route_no имеют хотя бы одну версию в DDS. + SELECT COUNT(*) + INTO v_missing_count + FROM (SELECT DISTINCT route_no FROM ods.routes) AS o + WHERE NOT EXISTS ( + SELECT 1 + FROM dds.dim_routes AS d + WHERE d.route_bk = o.route_no + ); + + IF v_missing_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_routes отсутствуют маршруты из ODS: %', + v_missing_count; + END IF; + + -- Current-срез DDS не содержит route_bk, которых нет в ODS. + SELECT COUNT(*) + INTO v_orphan_current + FROM ( + SELECT DISTINCT route_bk + FROM dds.dim_routes + WHERE valid_to IS NULL + ) AS d + WHERE NOT EXISTS ( + SELECT 1 + FROM ods.routes AS o + WHERE o.route_no = d.route_bk + ); + + IF v_orphan_current <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в current-срезе dds.dim_routes есть route_bk вне ODS: %', + v_orphan_current; + END IF; + + -- Обязательные поля. + SELECT COUNT(*) + INTO v_null_count + FROM dds.dim_routes + WHERE route_sk IS NULL + OR route_bk IS NULL + OR route_bk = '' + 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 hashdiff IS NULL + OR hashdiff = '' + OR valid_from IS NULL + OR created_at IS NULL + OR updated_at 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: в dds.dim_routes найдены NULL обязательные поля: %', + v_null_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: dds.dim_routes ок, строк=% (версий)', + v_row_count; +END $$; diff --git a/sql/dds/dim_routes_load.sql b/sql/dds/dim_routes_load.sql new file mode 100644 index 0000000..49e0585 --- /dev/null +++ b/sql/dds/dim_routes_load.sql @@ -0,0 +1,130 @@ +-- Загрузка DDS dim_routes: SCD2 с hashdiff. + +-- Statement 1: Закрыть устаревшие версии (valid_to = текущая дата). +WITH src AS ( + SELECT + route_no, + departure_airport, + arrival_airport, + airplane_code, + days_of_week, + departure_time, + duration, + md5( + COALESCE(departure_airport, '') || '|' || + COALESCE(arrival_airport, '') || '|' || + COALESCE(airplane_code, '') || '|' || + COALESCE(days_of_week, '') || '|' || + COALESCE(departure_time::TEXT, '') || '|' || + COALESCE(duration::TEXT, '') + ) AS hashdiff, + ROW_NUMBER() OVER (PARTITION BY route_no ORDER BY validity DESC) AS rn + FROM ods.routes +) +UPDATE dds.dim_routes AS d +SET valid_to = CURRENT_DATE, + updated_at = now(), + _load_id = '{{ run_id }}', + _load_ts = now() +FROM src AS s +WHERE s.rn = 1 + AND d.route_bk = s.route_no + AND d.valid_to IS NULL + AND d.hashdiff <> s.hashdiff; + +-- Statement 1.1: Закрыть "исчезнувшие" маршруты. +WITH src AS ( + SELECT + route_no, + ROW_NUMBER() OVER (PARTITION BY route_no ORDER BY validity DESC) AS rn + FROM ods.routes +) +UPDATE dds.dim_routes AS d +SET valid_to = CURRENT_DATE, + updated_at = now(), + _load_id = '{{ run_id }}', + _load_ts = now() +WHERE d.valid_to IS NULL + AND NOT EXISTS ( + SELECT 1 + FROM src AS s + WHERE s.rn = 1 + AND s.route_no = d.route_bk + ); + +-- Statement 2: Вставить новые версии (для изменённых и новых route_no). +WITH src AS ( + SELECT + route_no, + departure_airport, + arrival_airport, + airplane_code, + days_of_week, + departure_time, + duration, + md5( + COALESCE(departure_airport, '') || '|' || + COALESCE(arrival_airport, '') || '|' || + COALESCE(airplane_code, '') || '|' || + COALESCE(days_of_week, '') || '|' || + COALESCE(departure_time::TEXT, '') || '|' || + COALESCE(duration::TEXT, '') + ) AS hashdiff, + ROW_NUMBER() OVER (PARTITION BY route_no ORDER BY validity DESC) AS rn + FROM ods.routes +), +max_sk AS ( + SELECT COALESCE(MAX(route_sk), 0) AS v + FROM dds.dim_routes +) +INSERT INTO dds.dim_routes ( + route_sk, + route_bk, + departure_airport, + arrival_airport, + airplane_code, + days_of_week, + departure_time, + duration, + hashdiff, + valid_from, + valid_to, + created_at, + updated_at, + _load_id, + _load_ts +) +SELECT + (SELECT v FROM max_sk) + ROW_NUMBER() OVER (ORDER BY s.route_no)::INTEGER, + s.route_no, + s.departure_airport, + s.arrival_airport, + s.airplane_code, + s.days_of_week, + s.departure_time, + s.duration, + s.hashdiff, + CASE + WHEN EXISTS ( + SELECT 1 + FROM dds.dim_routes AS d2 + WHERE d2.route_bk = s.route_no + ) THEN CURRENT_DATE + ELSE '1900-01-01'::DATE + END AS valid_from, + NULL, + now(), + now(), + '{{ run_id }}', + now() +FROM src AS s +WHERE s.rn = 1 + AND NOT EXISTS ( + SELECT 1 + FROM dds.dim_routes AS d + WHERE d.route_bk = s.route_no + AND d.valid_to IS NULL + AND d.hashdiff = s.hashdiff + ); + +ANALYZE dds.dim_routes; diff --git a/sql/dds/dim_tariffs_ddl.sql b/sql/dds/dim_tariffs_ddl.sql new file mode 100644 index 0000000..bd93796 --- /dev/null +++ b/sql/dds/dim_tariffs_ddl.sql @@ -0,0 +1,13 @@ +-- DDL для DDS-слоя по таблице dim_tariffs (SCD1-измерение). + +CREATE SCHEMA IF NOT EXISTS dds; + +CREATE TABLE IF NOT EXISTS dds.dim_tariffs ( + tariff_sk INTEGER NOT NULL, + fare_conditions TEXT NOT NULL, + created_at TIMESTAMP NOT NULL DEFAULT now(), + updated_at TIMESTAMP NOT NULL DEFAULT now(), + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (tariff_sk); diff --git a/sql/dds/dim_tariffs_dq.sql b/sql/dds/dim_tariffs_dq.sql new file mode 100644 index 0000000..925bba9 --- /dev/null +++ b/sql/dds/dim_tariffs_dq.sql @@ -0,0 +1,98 @@ +-- DQ для DDS dim_tariffs. + +DO $$ +DECLARE + v_row_count BIGINT; + v_src_count BIGINT; + v_dup_sk BIGINT; + v_dup_bk BIGINT; + v_missing_bk BIGINT; + v_null_count BIGINT; +BEGIN + -- Для инкрементальных периодов без новых сегментов допускаем пустую dim_tariffs. + SELECT COUNT(*) + INTO v_row_count + FROM dds.dim_tariffs; + + SELECT COUNT(DISTINCT fare_conditions) + INTO v_src_count + FROM ods.segments + WHERE fare_conditions IS NOT NULL + AND fare_conditions <> ''; + + IF v_row_count = 0 AND v_src_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: dds.dim_tariffs пуста при непустом источнике ods.segments (fare_conditions=%).', + v_src_count; + ELSIF v_row_count = 0 AND v_src_count = 0 THEN + RAISE NOTICE + 'DQ PASSED: dds.dim_tariffs пуста, т.к. в ods.segments нет тарифов для загрузки.'; + RETURN; + END IF; + + -- Нет дублей по SK. + SELECT COUNT(*) - COUNT(DISTINCT tariff_sk) + INTO v_dup_sk + FROM dds.dim_tariffs; + + IF v_dup_sk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_tariffs найдены дубликаты tariff_sk: %', + v_dup_sk; + END IF; + + -- Нет дублей по BK. + SELECT COUNT(*) - COUNT(DISTINCT fare_conditions) + INTO v_dup_bk + FROM dds.dim_tariffs; + + IF v_dup_bk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_tariffs найдены дубликаты fare_conditions: %', + v_dup_bk; + END IF; + + -- Покрытие ODS: все fare_conditions из ods.segments есть в DDS. + SELECT COUNT(*) + INTO v_missing_bk + FROM ( + SELECT DISTINCT fare_conditions + FROM ods.segments + WHERE fare_conditions IS NOT NULL + AND fare_conditions <> '' + ) AS s + WHERE NOT EXISTS ( + SELECT 1 + FROM dds.dim_tariffs AS d + WHERE d.fare_conditions = s.fare_conditions + ); + + IF v_missing_bk <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.dim_tariffs отсутствуют значения fare_conditions из ods.segments: %', + v_missing_bk; + END IF; + + -- Обязательные поля. + SELECT COUNT(*) + INTO v_null_count + FROM dds.dim_tariffs + WHERE tariff_sk IS NULL + OR fare_conditions IS NULL + OR fare_conditions = '' + OR created_at IS NULL + OR updated_at 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: в dds.dim_tariffs найдены NULL/пустые обязательные поля: %', + v_null_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: dds.dim_tariffs ок, строк=%', + v_row_count; +END $$; diff --git a/sql/dds/dim_tariffs_load.sql b/sql/dds/dim_tariffs_load.sql new file mode 100644 index 0000000..ebe514f --- /dev/null +++ b/sql/dds/dim_tariffs_load.sql @@ -0,0 +1,36 @@ +-- Загрузка DDS dim_tariffs: SCD1 UPSERT (INSERT новых тарифов). + +WITH src AS ( + SELECT DISTINCT + s.fare_conditions + FROM ods.segments AS s + WHERE s.fare_conditions IS NOT NULL + AND s.fare_conditions <> '' +), +max_sk AS ( + SELECT COALESCE(MAX(tariff_sk), 0) AS v + FROM dds.dim_tariffs +) +INSERT INTO dds.dim_tariffs ( + tariff_sk, + fare_conditions, + created_at, + updated_at, + _load_id, + _load_ts +) +SELECT + (SELECT v FROM max_sk) + ROW_NUMBER() OVER (ORDER BY s.fare_conditions)::INTEGER, + s.fare_conditions, + now(), + now(), + '{{ run_id }}', + now() +FROM src AS s +WHERE NOT EXISTS ( + SELECT 1 + FROM dds.dim_tariffs AS d + WHERE d.fare_conditions = s.fare_conditions +); + +ANALYZE dds.dim_tariffs; diff --git a/sql/dds/fact_flight_sales_ddl.sql b/sql/dds/fact_flight_sales_ddl.sql new file mode 100644 index 0000000..1b087d4 --- /dev/null +++ b/sql/dds/fact_flight_sales_ddl.sql @@ -0,0 +1,23 @@ +-- DDL для DDS-слоя по таблице fact_flight_sales. + +CREATE SCHEMA IF NOT EXISTS dds; + +CREATE TABLE IF NOT EXISTS dds.fact_flight_sales ( + calendar_sk INTEGER, + departure_airport_sk INTEGER, + arrival_airport_sk INTEGER, + airplane_sk INTEGER, + tariff_sk INTEGER, + passenger_sk INTEGER, + route_sk INTEGER, + book_ref TEXT NOT NULL, + ticket_no TEXT NOT NULL, + flight_id INTEGER NOT NULL, + book_date DATE, + seat_no TEXT, + price NUMERIC(10,2), + is_boarded BOOLEAN NOT NULL, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +DISTRIBUTED BY (ticket_no); diff --git a/sql/dds/fact_flight_sales_dq.sql b/sql/dds/fact_flight_sales_dq.sql new file mode 100644 index 0000000..c4fa53c --- /dev/null +++ b/sql/dds/fact_flight_sales_dq.sql @@ -0,0 +1,151 @@ +-- DQ для DDS fact_flight_sales. + +DO $$ +DECLARE + v_row_count BIGINT; + v_ods_count BIGINT; + v_dup_count BIGINT; + v_null_passenger BIGINT; + v_null_tariff BIGINT; + v_null_route_related BIGINT; + v_null_calendar BIGINT; + v_null_required BIGINT; +BEGIN + -- Для пустого инкрементального окна (ods.segments) допускаем пустой факт. + SELECT COUNT(*) + INTO v_ods_count + FROM ods.segments; + + SELECT COUNT(*) + INTO v_row_count + FROM dds.fact_flight_sales; + + IF v_ods_count = 0 THEN + IF v_row_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: ods.segments пустая, но в dds.fact_flight_sales есть строки: %', + v_row_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: ods.segments и dds.fact_flight_sales пустые (инкрементальное окно без сегментов).'; + RETURN; + END IF; + + IF v_row_count = 0 THEN + RAISE EXCEPTION + 'DQ FAILED: dds.fact_flight_sales пуста при непустом источнике ods.segments (%).', + v_ods_count; + END IF; + + -- Покрытие: количество строк = ods.segments. + IF v_row_count <> v_ods_count THEN + RAISE EXCEPTION + 'DQ FAILED: dds.fact_flight_sales (%) <> ods.segments (%). Потеряны строки.', + v_row_count, + v_ods_count; + END IF; + + -- Нет дублей по зерну. + SELECT COUNT(*) + INTO v_dup_count + FROM ( + SELECT ticket_no, flight_id + FROM dds.fact_flight_sales + GROUP BY ticket_no, flight_id + HAVING COUNT(*) > 1 + ) AS d; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в dds.fact_flight_sales дубликаты (ticket_no, flight_id): %', + v_dup_count; + END IF; + + -- Ссылочная целостность: passenger_sk. + SELECT COUNT(*) + INTO v_null_passenger + FROM dds.fact_flight_sales + WHERE passenger_sk IS NULL; + + IF v_null_passenger <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в fact_flight_sales строки без passenger_sk: %', + v_null_passenger; + END IF; + + -- Ссылочная целостность: tariff_sk. + SELECT COUNT(*) + INTO v_null_tariff + FROM dds.fact_flight_sales + WHERE tariff_sk IS NULL; + + IF v_null_tariff <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в fact_flight_sales строки без tariff_sk: %', + v_null_tariff; + END IF; + + -- Route-related FK: допустимо при аномалиях, фейлим если > 1%. + SELECT COUNT(*) + INTO v_null_route_related + FROM dds.fact_flight_sales + WHERE route_sk IS NULL + OR departure_airport_sk IS NULL + OR arrival_airport_sk IS NULL + OR airplane_sk IS NULL; + + IF v_null_route_related > 0 THEN + IF v_null_route_related * 100.0 / NULLIF(v_row_count, 0) > 1.0 THEN + RAISE EXCEPTION + 'DQ FAILED: в fact_flight_sales слишком много строк с NULL в route-related FK: % (>1%%)', + v_null_route_related; + ELSE + RAISE NOTICE + 'DQ WARNING: в fact_flight_sales строк с NULL в route-related FK: % (<=1%%, допустимо)', + v_null_route_related; + END IF; + END IF; + + -- Calendar: допустимо если scheduled_departure IS NULL, фейлим если > 1%. + SELECT COUNT(*) + INTO v_null_calendar + FROM dds.fact_flight_sales + WHERE calendar_sk IS NULL; + + IF v_null_calendar > 0 THEN + IF v_null_calendar * 100.0 / NULLIF(v_row_count, 0) > 1.0 THEN + RAISE EXCEPTION + 'DQ FAILED: в fact_flight_sales слишком много строк без calendar_sk: % (>1%%)', + v_null_calendar; + ELSE + RAISE NOTICE + 'DQ WARNING: в fact_flight_sales строк без calendar_sk: % (<=1%%, допустимо)', + v_null_calendar; + END IF; + END IF; + + -- Обязательные поля. + SELECT COUNT(*) + INTO v_null_required + FROM dds.fact_flight_sales + WHERE book_ref IS NULL + OR book_ref = '' + OR ticket_no IS NULL + OR ticket_no = '' + OR flight_id IS NULL + OR is_boarded IS NULL + OR _load_id IS NULL + OR _load_id = '' + OR _load_ts IS NULL; + + IF v_null_required <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: в fact_flight_sales NULL обязательные поля: %', + v_null_required; + END IF; + + RAISE NOTICE + 'DQ PASSED: dds.fact_flight_sales ок, строк=%', + v_row_count; +END $$; diff --git a/sql/dds/fact_flight_sales_load.sql b/sql/dds/fact_flight_sales_load.sql new file mode 100644 index 0000000..053cfaf --- /dev/null +++ b/sql/dds/fact_flight_sales_load.sql @@ -0,0 +1,111 @@ +-- Загрузка DDS fact_flight_sales: инкрементальный UPSERT по (ticket_no, flight_id). + +-- Statement 1: UPDATE существующих строк факта. +-- Обновляем только мутабельные поля; SK измерений не перезаписываем. +UPDATE dds.fact_flight_sales AS f +SET seat_no = bp.seat_no, + price = seg.segment_amount, + is_boarded = (bp.ticket_no IS NOT NULL), + _load_id = '{{ run_id }}', + _load_ts = now() +FROM ods.segments AS seg +LEFT JOIN ods.boarding_passes AS bp + ON bp.ticket_no = seg.ticket_no + AND bp.flight_id = seg.flight_id +WHERE f.ticket_no = seg.ticket_no + AND f.flight_id = seg.flight_id + AND ( + f.is_boarded IS DISTINCT FROM (bp.ticket_no IS NOT NULL) + OR f.price IS DISTINCT FROM seg.segment_amount + OR f.seat_no IS DISTINCT FROM bp.seat_no + ); + +-- Statement 2: INSERT новых строк факта. +-- Dimension SK фиксируются на момент вставки (point-in-time для SCD2 routes). +WITH fact_src AS ( + SELECT + seg.ticket_no, + seg.flight_id, + cal.calendar_sk, + dep.airport_sk AS departure_airport_sk, + arr.airport_sk AS arrival_airport_sk, + ap.airplane_sk, + tar.tariff_sk, + pax.passenger_sk, + rte.route_sk, + tkt.book_ref, + bkg.book_date::DATE AS book_date, + bp.seat_no, + seg.segment_amount AS price, + (bp.ticket_no IS NOT NULL) AS is_boarded + FROM ods.segments AS seg + JOIN ods.tickets AS tkt + ON tkt.ticket_no = seg.ticket_no + JOIN ods.bookings AS bkg + ON bkg.book_ref = tkt.book_ref + JOIN ods.flights AS flt + ON flt.flight_id = seg.flight_id + LEFT JOIN dds.dim_routes AS rte + ON rte.route_bk = flt.route_no + AND flt.scheduled_departure::DATE >= rte.valid_from + AND (rte.valid_to IS NULL OR flt.scheduled_departure::DATE < rte.valid_to) + LEFT JOIN dds.dim_calendar AS cal + ON cal.date_actual = flt.scheduled_departure::DATE + LEFT JOIN dds.dim_airports AS dep + ON dep.airport_bk = rte.departure_airport + LEFT JOIN dds.dim_airports AS arr + ON arr.airport_bk = rte.arrival_airport + LEFT JOIN dds.dim_airplanes AS ap + ON ap.airplane_bk = rte.airplane_code + LEFT JOIN dds.dim_tariffs AS tar + ON tar.fare_conditions = seg.fare_conditions + LEFT JOIN dds.dim_passengers AS pax + ON pax.passenger_bk = tkt.passenger_id + LEFT JOIN ods.boarding_passes AS bp + ON bp.ticket_no = seg.ticket_no + AND bp.flight_id = seg.flight_id +) +INSERT INTO dds.fact_flight_sales ( + calendar_sk, + departure_airport_sk, + arrival_airport_sk, + airplane_sk, + tariff_sk, + passenger_sk, + route_sk, + book_ref, + ticket_no, + flight_id, + book_date, + seat_no, + price, + is_boarded, + _load_id, + _load_ts +) +SELECT + s.calendar_sk, + s.departure_airport_sk, + s.arrival_airport_sk, + s.airplane_sk, + s.tariff_sk, + s.passenger_sk, + s.route_sk, + s.book_ref, + s.ticket_no, + s.flight_id, + s.book_date, + s.seat_no, + s.price, + s.is_boarded, + '{{ run_id }}', + now() +FROM fact_src AS s +WHERE NOT EXISTS ( + SELECT 1 + FROM dds.fact_flight_sales AS f + WHERE f.ticket_no = s.ticket_no + AND f.flight_id = s.flight_id +); + +ANALYZE dds.fact_flight_sales; diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index 2117dd9..2fbe4f3 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -297,3 +297,88 @@ def test_bookings_to_gp_ods_dag_structure(): # Финальная сводка должна ждать обе ветки. _assert_reachable(dag, "dq_ods_boarding_passes", "finish_ods_summary") _assert_reachable(dag, "dq_ods_seats", "finish_ods_summary") + + +def test_bookings_dds_ddl_dag_structure(): + """Проверка структуры DAG bookings_dds_ddl.""" + dag = _load_dag("airflow.dags.bookings_dds_ddl") + + expected_tasks = { + "apply_dds_dim_calendar_ddl", + "apply_dds_dim_airports_ddl", + "apply_dds_dim_airplanes_ddl", + "apply_dds_dim_tariffs_ddl", + "apply_dds_dim_passengers_ddl", + "apply_dds_dim_routes_ddl", + "apply_dds_fact_flight_sales_ddl", + } + assert expected_tasks.issubset(dag.task_dict.keys()) + + _assert_reachable(dag, "apply_dds_dim_calendar_ddl", "apply_dds_dim_airports_ddl") + + for task_id in expected_tasks - {"apply_dds_dim_calendar_ddl"}: + _assert_reachable(dag, "apply_dds_dim_calendar_ddl", task_id) + + +def test_bookings_to_gp_dds_dag_structure(): + """Проверка структуры DAG bookings_to_gp_dds.""" + dag = _load_dag("airflow.dags.bookings_to_gp_dds") + + expected_tasks = { + "load_dds_dim_calendar", + "dq_dds_dim_calendar", + "load_dds_dim_airports", + "dq_dds_dim_airports", + "load_dds_dim_airplanes", + "dq_dds_dim_airplanes", + "load_dds_dim_tariffs", + "dq_dds_dim_tariffs", + "load_dds_dim_passengers", + "dq_dds_dim_passengers", + "load_dds_dim_routes", + "dq_dds_dim_routes", + "load_dds_fact_flight_sales", + "dq_dds_fact_flight_sales", + "finish_dds_summary", + } + assert expected_tasks.issubset(dag.task_dict.keys()) + + load_to_dq = [ + ("load_dds_dim_calendar", "dq_dds_dim_calendar"), + ("load_dds_dim_airports", "dq_dds_dim_airports"), + ("load_dds_dim_airplanes", "dq_dds_dim_airplanes"), + ("load_dds_dim_tariffs", "dq_dds_dim_tariffs"), + ("load_dds_dim_passengers", "dq_dds_dim_passengers"), + ("load_dds_dim_routes", "dq_dds_dim_routes"), + ("load_dds_fact_flight_sales", "dq_dds_fact_flight_sales"), + ] + for load_task_id, dq_task_id in load_to_dq: + _assert_direct_edge(dag, load_task_id, dq_task_id) + + # После calendar все остальные измерения должны быть reachable. + _assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_dim_airports") + _assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_dim_airplanes") + _assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_dim_tariffs") + _assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_dim_passengers") + _assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_dim_routes") + + # Параллельность: airports и airplanes не должны зависеть друг от друга. + airports = dag.get_task("load_dds_dim_airports") + airplanes = dag.get_task("load_dds_dim_airplanes") + assert airplanes not in airports.get_flat_relatives( + upstream=False + ), "dds airports не должен быть upstream для airplanes" + assert airports not in airplanes.get_flat_relatives( + upstream=False + ), "dds airplanes не должен быть upstream для airports" + + # Факт должен стартовать только после всех измерений. + _assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_fact_flight_sales") + _assert_reachable(dag, "dq_dds_dim_airports", "load_dds_fact_flight_sales") + _assert_reachable(dag, "dq_dds_dim_airplanes", "load_dds_fact_flight_sales") + _assert_reachable(dag, "dq_dds_dim_tariffs", "load_dds_fact_flight_sales") + _assert_reachable(dag, "dq_dds_dim_passengers", "load_dds_fact_flight_sales") + _assert_reachable(dag, "dq_dds_dim_routes", "load_dds_fact_flight_sales") + + # Финальная сводка должна ждать DQ факта. + _assert_reachable(dag, "dq_dds_fact_flight_sales", "finish_dds_summary")