From d89434c3a8b9f45916b2e291bbdd7053032d207c Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Thu, 5 Mar 2026 22:41:17 +0300 Subject: [PATCH] =?UTF-8?q?feat(dds):=20=D0=B4=D0=B5=D0=BD=D0=BE=D1=80?= =?UTF-8?q?=D0=BC=D0=B0=D0=BB=D0=B8=D0=B7=D0=BE=D0=B2=D0=B0=D0=BD=D0=BE=20?= =?UTF-8?q?=D0=B8=D0=B7=D0=BC=D0=B5=D1=80=D0=B5=D0=BD=D0=B8=D0=B5=20dim=5F?= =?UTF-8?q?routes=20=D0=B8=20=D1=83=D0=BF=D1=80=D0=BE=D1=89=D0=B5=D0=BD?= =?UTF-8?q?=D0=B0=20=D0=B2=D0=B8=D1=82=D1=80=D0=B8=D0=BD=D0=B0=20route=5Fp?= =?UTF-8?q?erformance?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - улучшение производительности аналитических запросов и упрощение витрин согласно принципам Kimball Star Schema. - Что: - в dds.dim_routes добавлены денормализованные поля городов и моделей самолетов. - в скрипт загрузки dim_routes_load.sql добавлена фаза refresh для актуализации атрибутов. - загрузка dm.route_performance упрощена до 1 JOIN к измерению маршрутов. - обновлен DAG bookings_to_gp_dds и smoke-тесты структуры графа. - в документации (db_schema.md) отражены денормализация и lineage версий. - в dim_routes_load.sql исправлено затирание _load_id при refresh исторических версий. - Проверка: - make test (smoke-тесты структуры DAG проходят успешно). --- TODO.md | 2 +- airflow/dags/bookings_to_gp_dds.py | 5 ++- docs/internal/db_schema.md | 4 ++- .../dim_routes_denormalization_plan.md | 2 +- sql/dds/dim_routes_ddl.sql | 12 +++++++ sql/dds/dim_routes_dq.sql | 8 +++++ sql/dds/dim_routes_load.sql | 35 +++++++++++++++++++ sql/dm/route_performance_load.sql | 24 ++++++------- tests/test_dags_smoke.py | 4 +++ 9 files changed, 79 insertions(+), 17 deletions(-) diff --git a/TODO.md b/TODO.md index 1fa6d1d..142ed11 100644 --- a/TODO.md +++ b/TODO.md @@ -53,7 +53,7 @@ - [x] Добавить в образ Airflow установку `psql`, чтобы тестировать загрузку CSV из CLI внутри контейнера (без root и дополнительных зависимостей на хосте). -- [ ] Денормализовать `dds.dim_routes` (добавить departure_city, arrival_city, airplane_model, total_seats): +- [x] Денормализовать `dds.dim_routes` (добавить departure_city, arrival_city, airplane_model, total_seats): - привести измерение в соответствие с принципом Кимбалла («самодостаточное измерение»); - упростить `dm.route_performance` с 4-JOIN до 1-JOIN; - обновить DAG-зависимости: airports+airplanes DQ → routes load; diff --git a/airflow/dags/bookings_to_gp_dds.py b/airflow/dags/bookings_to_gp_dds.py index f09096a..ee44efa 100644 --- a/airflow/dags/bookings_to_gp_dds.py +++ b/airflow/dags/bookings_to_gp_dds.py @@ -142,11 +142,14 @@ with DAG( load_dds_dim_airplanes, load_dds_dim_tariffs, load_dds_dim_passengers, - load_dds_dim_routes, ] + # dim_routes зависит от airports и airplanes (денормализация). + [dq_dds_dim_airports, dq_dds_dim_airplanes] >> 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 diff --git a/docs/internal/db_schema.md b/docs/internal/db_schema.md index 1a63323..d53c343 100644 --- a/docs/internal/db_schema.md +++ b/docs/internal/db_schema.md @@ -70,7 +70,7 @@ | `dds.dim_airplanes` | `airplane_code` (`airplane_bk`) | `airplane_sk` | `model`, `range_km`, `speed_kmh`, `total_seats` | | `dds.dim_tariffs` | `fare_conditions` | `tariff_sk` | `fare_conditions` | | `dds.dim_passengers` | `passenger_id` (`passenger_bk`) | `passenger_sk` | `passenger_name` (SCD1) | -| `dds.dim_routes` | `route_no` (`route_bk`) | `route_sk` | `departure_airport`, `arrival_airport`, `airplane_code`, `hashdiff`, `valid_from`, `valid_to` (SCD2) | +| `dds.dim_routes` | `route_no` (`route_bk`) | `route_sk` | `departure_airport`, `arrival_airport`, `airplane_code`, `departure_city` (денормализовано), `arrival_city` (денормализовано), `airplane_model` (денормализовано), `total_seats` (денормализовано), `hashdiff`, `valid_from`, `valid_to` (SCD2) | ### Факт DDS (Fact) @@ -325,6 +325,8 @@ graph LR ODS_Tickets -->|Extract Unique| DIM_Passengers ODS_Routes -->|SCD2 with hashdiff| DIM_Routes + DIM_Airports -.->|Enrich cities| DIM_Routes + DIM_Airplanes -.->|Enrich model and seats| DIM_Routes %% Fact assembly (Main process + route point-in-time) ODS_Segments -->|Main Stream| FACT_Sales diff --git a/docs/internal/dim_routes_denormalization_plan.md b/docs/internal/dim_routes_denormalization_plan.md index 48070d2..c72c6c8 100644 --- a/docs/internal/dim_routes_denormalization_plan.md +++ b/docs/internal/dim_routes_denormalization_plan.md @@ -1,6 +1,6 @@ # План: Денормализация dim_routes + упрощение route_performance -> **Статус:** запланировано, не реализовано. +> **Статус:** реализовано. ## Контекст diff --git a/sql/dds/dim_routes_ddl.sql b/sql/dds/dim_routes_ddl.sql index 0bd0021..dfd07bb 100644 --- a/sql/dds/dim_routes_ddl.sql +++ b/sql/dds/dim_routes_ddl.sql @@ -5,12 +5,24 @@ CREATE SCHEMA IF NOT EXISTS dds; -- Тип таблицы: Heap (стандартная). -- Обоснование: Необходим row-level UPDATE для реализации SCD2 (закрытие версий). -- Использование Append-Only при частых обновлениях приводит к раздуванию (bloat) таблицы. +-- +-- Учебный комментарий (Kimball Star Schema): +-- Измерение должно быть «самодостаточным»: один JOIN к dim_routes — +-- и аналитик видит маршрут, города, модель самолёта и кол-во мест. +-- Денормализованные атрибуты (departure_city, arrival_city, airplane_model, total_seats) +-- НЕ участвуют в hashdiff. Версия SCD2 фиксирует изменения атрибутов маршрута +-- (аэропорт, самолёт, расписание). Если изменится название города — +-- обновим отдельным refresh-шагом, не создавая новую версию. 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, + departure_city TEXT NOT NULL, + arrival_city TEXT NOT NULL, + airplane_model TEXT NOT NULL, + total_seats INTEGER NOT NULL, days_of_week TEXT, departure_time TIME, duration INTERVAL, diff --git a/sql/dds/dim_routes_dq.sql b/sql/dds/dim_routes_dq.sql index fe1e56f..6e1a6bd 100644 --- a/sql/dds/dim_routes_dq.sql +++ b/sql/dds/dim_routes_dq.sql @@ -128,6 +128,13 @@ BEGIN OR arrival_airport = '' OR airplane_code IS NULL OR airplane_code = '' + OR departure_city IS NULL + OR departure_city = '' + OR arrival_city IS NULL + OR arrival_city = '' + OR airplane_model IS NULL + OR airplane_model = '' + OR total_seats IS NULL OR hashdiff IS NULL OR hashdiff = '' OR valid_from IS NULL @@ -146,4 +153,5 @@ BEGIN 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 index fac9b52..c95266b 100644 --- a/sql/dds/dim_routes_load.sql +++ b/sql/dds/dim_routes_load.sql @@ -62,6 +62,10 @@ INSERT INTO dds.dim_routes ( departure_airport, arrival_airport, airplane_code, + departure_city, + arrival_city, + airplane_model, + total_seats, days_of_week, departure_time, duration, @@ -79,6 +83,10 @@ SELECT s.departure_airport, s.arrival_airport, s.airplane_code, + dep.city, + arr.city, + air.model, + air.total_seats, s.days_of_week, s.departure_time, s.duration, @@ -97,6 +105,9 @@ SELECT '{{ run_id }}', now() FROM tmp_routes_src AS s +LEFT JOIN dds.dim_airports AS dep ON s.departure_airport = dep.airport_bk +LEFT JOIN dds.dim_airports AS arr ON s.arrival_airport = arr.airport_bk +LEFT JOIN dds.dim_airplanes AS air ON s.airplane_code = air.airplane_bk WHERE s.rn = 1 AND NOT EXISTS ( SELECT 1 @@ -106,4 +117,28 @@ WHERE s.rn = 1 AND d.hashdiff = s.hashdiff ); +-- Фаза 3: Обновление денормализованных атрибутов (refresh). +-- Нужна для SCD1-изменений в dim_airports/dim_airplanes (напр. переименование города). +-- Обновляем все версии (и текущие, и исторические), т.к. измерения SCD1. +-- При этом мы не перезаписываем _load_id и _load_ts, чтобы не размывать lineage версий. +UPDATE dds.dim_routes AS d +SET departure_city = dep.city, + arrival_city = arr.city, + airplane_model = air.model, + total_seats = air.total_seats, + updated_at = now() +FROM dds.dim_airports AS dep, + dds.dim_airports AS arr, + dds.dim_airplanes AS air +WHERE dep.airport_bk = d.departure_airport + AND arr.airport_bk = d.arrival_airport + AND air.airplane_bk = d.airplane_code + AND ( + d.departure_city IS DISTINCT FROM dep.city + OR d.arrival_city IS DISTINCT FROM arr.city + OR d.airplane_model IS DISTINCT FROM air.model + OR d.total_seats IS DISTINCT FROM air.total_seats + ); + ANALYZE dds.dim_routes; + diff --git a/sql/dm/route_performance_load.sql b/sql/dm/route_performance_load.sql index 5bb5dc6..64e9353 100644 --- a/sql/dm/route_performance_load.sql +++ b/sql/dm/route_performance_load.sql @@ -29,7 +29,8 @@ JOIN dds.dim_routes r ON f.route_sk = r.route_sk JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk GROUP BY r.route_bk; --- Шаг 3: Объединяем агрегаты с АКТУАЛЬНЫМИ атрибутами маршрута, аэропортов и самолетов. +-- Шаг 3: Объединяем агрегаты с АКТУАЛЬНЫМИ атрибутами маршрута. +-- Учебный комментарий: Благодаря денормализации dim_routes — один JOIN вместо четырёх. INSERT INTO dm.route_performance ( route_bk, route_sk, @@ -54,13 +55,13 @@ INSERT INTO dm.route_performance ( SELECT m.route_bk, r_curr.route_sk, - dep.airport_bk AS departure_airport_bk, - dep.city AS departure_city, - arr.airport_bk AS arrival_airport_bk, - arr.city AS arrival_city, - air.airplane_bk, - air.model AS airplane_model, - air.total_seats, + r_curr.departure_airport AS departure_airport_bk, + r_curr.departure_city, + r_curr.arrival_airport AS arrival_airport_bk, + r_curr.arrival_city, + r_curr.airplane_code AS airplane_bk, + r_curr.airplane_model, + r_curr.total_seats, m.total_flights, m.total_tickets, m.total_boarded, @@ -70,14 +71,11 @@ SELECT m.avg_boarding_rate::NUMERIC(5,4), -- Load Factor: общее кол-во посадок / (кол-во рейсов * мест в самолете) ROUND( - m.total_boarded::NUMERIC / NULLIF(m.total_flights * air.total_seats, 0), + m.total_boarded::NUMERIC / NULLIF(m.total_flights * r_curr.total_seats, 0), 4 ) AS avg_load_factor, m.first_flight_date, m.last_flight_date, '{{ run_id }}' AS _load_id FROM tmp_route_metrics m -JOIN dds.dim_routes r_curr ON m.route_bk = r_curr.route_bk AND r_curr.valid_to IS NULL -JOIN dds.dim_airports dep ON r_curr.departure_airport = dep.airport_bk -JOIN dds.dim_airports arr ON r_curr.arrival_airport = arr.airport_bk -JOIN dds.dim_airplanes air ON r_curr.airplane_code = air.airplane_bk; +JOIN dds.dim_routes r_curr ON m.route_bk = r_curr.route_bk AND r_curr.valid_to IS NULL; diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index 03424aa..a1ca87e 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -372,6 +372,10 @@ def test_bookings_to_gp_dds_dag_structure(): upstream=False ), "dds airplanes не должен быть upstream для airports" + # dim_routes зависит от airports и airplanes (денормализация). + _assert_reachable(dag, "dq_dds_dim_airports", "load_dds_dim_routes") + _assert_reachable(dag, "dq_dds_dim_airplanes", "load_dds_dim_routes") + # Факт должен стартовать только после всех измерений. _assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_fact_flight_sales") _assert_reachable(dag, "dq_dds_dim_airports", "load_dds_fact_flight_sales")