From 2228312b5549c350ea7a193b552b02e733b235c4 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Wed, 4 Mar 2026 22:28:52 +0300 Subject: [PATCH] =?UTF-8?q?feat(dm):=20=D0=B7=D0=B0=D0=B2=D0=B5=D1=80?= =?UTF-8?q?=D1=88=D0=B5=D0=BD=D0=BE=20=D0=BF=D0=BE=D1=81=D1=82=D1=80=D0=BE?= =?UTF-8?q?=D0=B5=D0=BD=D0=B8=D0=B5=20DM-=D1=81=D0=BB=D0=BE=D1=8F=20(5=20?= =?UTF-8?q?=D0=B2=D0=B8=D1=82=D1=80=D0=B8=D0=BD)=20=D0=B8=20=D0=B8=D1=81?= =?UTF-8?q?=D0=BF=D1=80=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD=D1=8B=20=D0=B1=D0=B0?= =?UTF-8?q?=D0=B3=D0=B8=20=D0=B4=D0=B5=D0=BD=D0=BE=D1=80=D0=BC=D0=B0=D0=BB?= =?UTF-8?q?=D0=B8=D0=B7=D0=B0=D1=86=D0=B8=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - предоставить студентам полный набор аналитических витрин с примерами различных паттернов (UPSERT, Full Rebuild, UNION ALL, двухуровневая агрегация). - Что: - реализованы витрины: sales_report, route_performance, passenger_loyalty, airport_traffic, monthly_overview. - исправлен баг в route_performance_load.sql: добавлены JOIN к dim_airports и dim_airplanes для корректной денормализации атрибутов. - обновлен скрипт e2e_etl.sh: добавлена верификация всех 5 витрин и проверка бизнес-логики (load factor). - обновлены DAGи, тесты и главный DDL скрипт. - Проверка: - автоматизированный прогон e2e_etl.sh через REST API Airflow. --- airflow/dags/bookings_dm_ddl.py | 13 ++- airflow/dags/bookings_to_gp_dm.py | 22 ++-- scripts/e2e_etl.sh | 15 ++- sql/ddl_gp.sql | 4 + sql/dm/monthly_overview_ddl.sql | 40 ++++++++ sql/dm/monthly_overview_dq.sql | 52 ++++++++++ sql/dm/monthly_overview_load.sql | 165 ++++++++++++++++++++++++++++++ sql/dm/route_performance_load.sql | 25 ++--- tests/test_dags_smoke.py | 38 +++++-- 9 files changed, 342 insertions(+), 32 deletions(-) create mode 100644 sql/dm/monthly_overview_ddl.sql create mode 100644 sql/dm/monthly_overview_dq.sql create mode 100644 sql/dm/monthly_overview_load.sql diff --git a/airflow/dags/bookings_dm_ddl.py b/airflow/dags/bookings_dm_ddl.py index 524b721..7d9b7b7 100644 --- a/airflow/dags/bookings_dm_ddl.py +++ b/airflow/dags/bookings_dm_ddl.py @@ -7,8 +7,7 @@ from __future__ import annotations Создаёт 5 DM-витрин: sales_report, route_performance, passenger_loyalty, airport_traffic, monthly_overview. -На данном этапе реализованы витрины: sales_report, route_performance, passenger_loyalty, airport_traffic. -Остальные витрины будут добавлены в последующих этапах. +На данном этапе реализованы все 5 витрин: sales_report, route_performance, passenger_loyalty, airport_traffic, monthly_overview. """ from datetime import timedelta @@ -60,8 +59,12 @@ with DAG( sql="dm/airport_traffic_ddl.sql", ) - # Заглушки для будущих витрин (будут реализованы в этапе 5) - # apply_dm_monthly_overview_ddl = PostgresOperator(...) + # Витрина: Помесячная сводка (Этап 5) + apply_dm_monthly_overview_ddl = PostgresOperator( + task_id="apply_dm_monthly_overview_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/monthly_overview_ddl.sql", + ) # Линейная цепочка - apply_dm_sales_report_ddl >> apply_dm_route_performance_ddl >> apply_dm_passenger_loyalty_ddl >> apply_dm_airport_traffic_ddl + apply_dm_sales_report_ddl >> apply_dm_route_performance_ddl >> apply_dm_passenger_loyalty_ddl >> apply_dm_airport_traffic_ddl >> apply_dm_monthly_overview_ddl diff --git a/airflow/dags/bookings_to_gp_dm.py b/airflow/dags/bookings_to_gp_dm.py index 6baba8c..0297594 100644 --- a/airflow/dags/bookings_to_gp_dm.py +++ b/airflow/dags/bookings_to_gp_dm.py @@ -9,8 +9,7 @@ from __future__ import annotations - все витрины загружаются параллельно (не зависят друг от друга); - sales_report использует UPSERT (heap-таблица), route_performance — Full Rebuild (AO Column). -На данном этапе реализованы витрины: sales_report, route_performance, passenger_loyalty, airport_traffic. -Остальные витрины будут добавлены в последующих этапах. +На данном этапе реализованы все 5 витрин: sales_report, route_performance, passenger_loyalty, airport_traffic, monthly_overview. """ from datetime import timedelta @@ -107,9 +106,18 @@ with DAG( sql="dm/airport_traffic_dq.sql", ) - # === Заглушки для будущих витрин (будут реализованы в этапе 5) === - # load_dm_monthly_overview = PostgresOperator(...) - # dq_dm_monthly_overview = PostgresOperator(...) + # === Витрина: Помесячная сводка (Этап 5) === + load_dm_monthly_overview = PostgresOperator( + task_id="load_dm_monthly_overview", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/monthly_overview_load.sql", + ) + + dq_dm_monthly_overview = PostgresOperator( + task_id="dq_dm_monthly_overview", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/monthly_overview_dq.sql", + ) # finish-задача finish_dm_summary = PythonOperator( @@ -122,7 +130,8 @@ with DAG( load_dm_sales_report, load_dm_route_performance, load_dm_passenger_loyalty, - load_dm_airport_traffic + load_dm_airport_traffic, + load_dm_monthly_overview ] # Связываем dq с finish @@ -130,3 +139,4 @@ with DAG( load_dm_route_performance >> dq_dm_route_performance >> finish_dm_summary load_dm_passenger_loyalty >> dq_dm_passenger_loyalty >> finish_dm_summary load_dm_airport_traffic >> dq_dm_airport_traffic >> finish_dm_summary + load_dm_monthly_overview >> dq_dm_monthly_overview >> finish_dm_summary diff --git a/scripts/e2e_etl.sh b/scripts/e2e_etl.sh index 907698e..2a840b7 100755 --- a/scripts/e2e_etl.sh +++ b/scripts/e2e_etl.sh @@ -93,12 +93,23 @@ trigger_and_wait "bookings_to_gp_dds" trigger_and_wait "bookings_to_gp_dm" warn "=== Верификация данных ===" -# Выполняем простые проверки строк, чтобы убедиться, что данные дошли до витрины +# Выполняем проверки строк и бизнес-логики, чтобы убедиться, что все витрины DM-слоя заполнены docker compose -f docker-compose.yml exec greenplum bash -c "su - gpadmin -c \"psql -d gp_dwh -c \\\" SELECT 'STG bookings' as layer, COUNT(*) FROM stg.bookings UNION ALL SELECT 'ODS bookings', COUNT(*) FROM ods.bookings UNION ALL SELECT 'DDS fact', COUNT(*) FROM dds.fact_flight_sales UNION ALL -SELECT 'DM report', COUNT(*) FROM dm.sales_report; +SELECT 'DM sales_report', COUNT(*) FROM dm.sales_report UNION ALL +SELECT 'DM route_performance', COUNT(*) FROM dm.route_performance UNION ALL +SELECT 'DM passenger_loyalty', COUNT(*) FROM dm.passenger_loyalty UNION ALL +SELECT 'DM airport_traffic', COUNT(*) FROM dm.airport_traffic UNION ALL +SELECT 'DM monthly_overview', COUNT(*) FROM dm.monthly_overview; + +-- Дополнительная проверка бизнес-логики (load factor не должен быть NULL и должен быть в пределах разумного) +SELECT + 'Check LF' as check, + COUNT(*) as total_rows, + SUM(CASE WHEN avg_load_factor IS NOT NULL AND avg_load_factor >= 0 THEN 1 ELSE 0 END) as valid_lf +FROM dm.monthly_overview; \\\"\"" warn "E2E ETL тест успешно завершен!" diff --git a/sql/ddl_gp.sql b/sql/ddl_gp.sql index 5e10504..893d895 100644 --- a/sql/ddl_gp.sql +++ b/sql/ddl_gp.sql @@ -60,3 +60,7 @@ FORMAT 'CUSTOM' (formatter='pxfwritable_import'); -- DDL для DM-слоя (Data Mart). \i dm/sales_report_ddl.sql +\i dm/route_performance_ddl.sql +\i dm/passenger_loyalty_ddl.sql +\i dm/airport_traffic_ddl.sql +\i dm/monthly_overview_ddl.sql diff --git a/sql/dm/monthly_overview_ddl.sql b/sql/dm/monthly_overview_ddl.sql new file mode 100644 index 0000000..80f6304 --- /dev/null +++ b/sql/dm/monthly_overview_ddl.sql @@ -0,0 +1,40 @@ +-- DDL для витрины dm.monthly_overview (Помесячная сводка). +-- +-- Учебные цели: +-- 1. Двухуровневая агрегация (до рейса, затем до месяца) для точного расчета +-- средних показателей (load factor), избегая парадокса Симпсона. +-- 2. Выбор ключа распределения (Distribution Key): +-- Мы используем airplane_sk, а НЕ (year, month). Распределение по датам в MPP — +-- это антипаттерн, приводящий к Data Skew (весь месяц на одном сегменте). + +CREATE TABLE IF NOT EXISTS dm.monthly_overview ( + -- Зерно: Месяц и самолет + year_actual INTEGER NOT NULL, + month_actual INTEGER NOT NULL, + airplane_sk INTEGER NOT NULL, + + -- Денормализованные атрибуты + airplane_bk TEXT NOT NULL, + airplane_model TEXT NOT NULL, + total_seats INTEGER NOT NULL, + + -- Метрики + total_flights INTEGER NOT NULL DEFAULT 0, + total_tickets INTEGER NOT NULL DEFAULT 0, + total_boarded INTEGER NOT NULL DEFAULT 0, + total_revenue NUMERIC(15,2) NOT NULL DEFAULT 0, + avg_ticket_price NUMERIC(10,2), + avg_load_factor NUMERIC(5,4), -- Точный load factor (сначала по рейсу, потом среднее) + unique_routes INTEGER NOT NULL DEFAULT 0, + unique_passengers INTEGER NOT NULL DEFAULT 0, + + -- Служебные + 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() +) +WITH (appendonly=false) -- Heap для UPSERT +DISTRIBUTED BY (airplane_sk); -- Избегаем распределения по (year, month) + +COMMENT ON TABLE dm.monthly_overview IS 'Витрина: помесячная аналитика с точным load_factor (Двухуровневая агрегация, Heap)'; diff --git a/sql/dm/monthly_overview_dq.sql b/sql/dm/monthly_overview_dq.sql new file mode 100644 index 0000000..7ca896c --- /dev/null +++ b/sql/dm/monthly_overview_dq.sql @@ -0,0 +1,52 @@ +-- DQ проверки для витрины dm.monthly_overview. +-- +-- Учебные цели: +-- 1. Валидация логики двухуровневой агрегации: load factor не должен превышать 1.0 +-- (в отличие от route_performance, где допускалась погрешность из-за упрощенного расчета). +-- 2. Базовые санити-проверки для дат (month BETWEEN 1 AND 12). + +DO $$ +DECLARE + row_count INTEGER; + duplicate_count INTEGER; + invalid_month_count INTEGER; + invalid_load_factor_count INTEGER; +BEGIN + -- 1. Проверка на наполненность + SELECT COUNT(*) INTO row_count FROM dm.monthly_overview; + IF row_count = 0 THEN + RAISE EXCEPTION 'DQ Error: Таблица dm.monthly_overview пуста.'; + END IF; + + -- 2. Проверка на уникальность (зерно - год + месяц + самолет) + SELECT COUNT(*) INTO duplicate_count + FROM ( + SELECT year_actual, month_actual, airplane_sk FROM dm.monthly_overview + GROUP BY year_actual, month_actual, airplane_sk HAVING COUNT(*) > 1 + ) q; + IF duplicate_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.monthly_overview обнаружены дубликаты по (year, month, airplane_sk) (% шт).', duplicate_count; + END IF; + + -- 3. Санити-проверка месяца + SELECT COUNT(*) INTO invalid_month_count + FROM dm.monthly_overview + WHERE month_actual < 1 OR month_actual > 12; + + IF invalid_month_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.monthly_overview обнаружены некорректные месяцы (% строк).', invalid_month_count; + END IF; + + -- 4. Проверка точного load_factor + -- В этой витрине мы считаем его честно (по рейсам), поэтому он строго <= 1.0 + -- (исключая экзотические случаи овербукинга, но для учебного стенда ставим жесткий лимит 1.0). + SELECT COUNT(*) INTO invalid_load_factor_count + FROM dm.monthly_overview + WHERE avg_load_factor < 0 OR avg_load_factor > 1.0; + + IF invalid_load_factor_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.monthly_overview обнаружен avg_load_factor вне диапазона 0..1 (% строк).', invalid_load_factor_count; + END IF; + + RAISE NOTICE 'DQ Success: dm.monthly_overview успешно прошла все проверки (% строк).', row_count; +END $$; diff --git a/sql/dm/monthly_overview_load.sql b/sql/dm/monthly_overview_load.sql new file mode 100644 index 0000000..f17dd08 --- /dev/null +++ b/sql/dm/monthly_overview_load.sql @@ -0,0 +1,165 @@ +-- Загрузка витрины dm.monthly_overview: инкрементальный UPSERT по месяцам. +-- +-- Учебные цели: +-- 1. Двухуровневая агрегация: чтобы честно посчитать среднюю заполняемость (avg_load_factor), +-- мы сначала считаем её для КАЖДОГО рейса, а затем берем среднее по месяцу. +-- 2. Ограничения SCD1: dim_airplanes — это SCD1-измерение. Мы берем total_seats из +-- его текущего состояния. Если бы самолет переоборудовали (изменили число мест) в прошлом, +-- для точного исторического расчета нам потребовалось бы SCD2-измерение. + +-- Шаг 1: Находим затронутые месяцы (HWM). +CREATE TEMP TABLE tmp_monthly_affected_dates ON COMMIT DROP AS +SELECT DISTINCT cal.year_actual, cal.month_actual +FROM dds.fact_flight_sales f +JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk +WHERE f._load_ts > ( + SELECT COALESCE(MAX(_load_ts), '1900-01-01'::TIMESTAMP) + FROM dm.monthly_overview +); + +-- Шаг 2: Уровень 1 - Агрегация фактов до рейса (flight_id). +-- Считаем точный load_factor для каждого отдельного перелета. +CREATE TEMP TABLE tmp_flight_level_metrics ON COMMIT DROP AS +SELECT + cal.year_actual, + cal.month_actual, + f.flight_id, + f.airplane_sk, + a.total_seats, -- ВНИМАНИЕ: берется текущее значение (SCD1) + COUNT(*) AS tickets_sold_per_flight, + SUM(CASE WHEN f.is_boarded THEN 1 ELSE 0 END) AS boarded_per_flight, + SUM(f.price) AS revenue_per_flight, + -- Точный load factor конкретного рейса: посаженные пассажиры / кол-во мест + (SUM(CASE WHEN f.is_boarded THEN 1 ELSE 0 END)::NUMERIC / NULLIF(a.total_seats, 0)) AS flight_load_factor +FROM dds.fact_flight_sales f +JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk +JOIN dds.dim_airplanes a ON f.airplane_sk = a.airplane_sk +-- Фильтруем только затронутые месяцы +WHERE EXISTS ( + SELECT 1 FROM tmp_monthly_affected_dates tad + WHERE tad.year_actual = cal.year_actual AND tad.month_actual = cal.month_actual +) +GROUP BY cal.year_actual, cal.month_actual, f.flight_id, f.airplane_sk, a.total_seats; + +-- Шаг 3: Уровень 2 - Агрегация от рейсов до месяцев. +-- Плюс собираем уникальные метрики (пассажиры, маршруты) напрямую из фактов. +CREATE TEMP TABLE tmp_monthly_overview_delta ON COMMIT DROP AS +WITH flight_aggs AS ( + SELECT + year_actual, + month_actual, + airplane_sk, + COUNT(DISTINCT flight_id) AS total_flights, + SUM(tickets_sold_per_flight) AS total_tickets, + SUM(boarded_per_flight) AS total_boarded, + SUM(revenue_per_flight) AS total_revenue, + -- Честное среднее от рейсовых показателей + AVG(flight_load_factor) AS avg_load_factor + FROM tmp_flight_level_metrics + GROUP BY year_actual, month_actual, airplane_sk +), +unique_aggs AS ( + SELECT + cal.year_actual, + cal.month_actual, + f.airplane_sk, + COUNT(DISTINCT r.route_bk) AS unique_routes, -- Считаем по BK (SCD2) + COUNT(DISTINCT f.passenger_sk) AS unique_passengers + FROM dds.fact_flight_sales f + JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk + JOIN dds.dim_routes r ON f.route_sk = r.route_sk + WHERE EXISTS ( + SELECT 1 FROM tmp_monthly_affected_dates tad + WHERE tad.year_actual = cal.year_actual AND tad.month_actual = cal.month_actual + ) + GROUP BY cal.year_actual, cal.month_actual, f.airplane_sk +) +SELECT + fa.year_actual, + fa.month_actual, + fa.airplane_sk, + a.airplane_bk, + a.model AS airplane_model, + a.total_seats, + fa.total_flights, + fa.total_tickets, + fa.total_boarded, + fa.total_revenue, + (fa.total_revenue / NULLIF(fa.total_tickets, 0))::NUMERIC(10,2) AS avg_ticket_price, + fa.avg_load_factor::NUMERIC(5,4) AS avg_load_factor, + ua.unique_routes, + ua.unique_passengers +FROM flight_aggs fa +JOIN unique_aggs ua ON fa.year_actual = ua.year_actual AND fa.month_actual = ua.month_actual AND fa.airplane_sk = ua.airplane_sk +JOIN dds.dim_airplanes a ON fa.airplane_sk = a.airplane_sk; + +-- Шаг 4: UPDATE существующих записей. +UPDATE dm.monthly_overview AS tgt +SET + airplane_bk = src.airplane_bk, + airplane_model = src.airplane_model, + total_seats = src.total_seats, + total_flights = src.total_flights, + total_tickets = src.total_tickets, + total_boarded = src.total_boarded, + total_revenue = src.total_revenue, + avg_ticket_price = src.avg_ticket_price, + avg_load_factor = src.avg_load_factor, + unique_routes = src.unique_routes, + unique_passengers = src.unique_passengers, + updated_at = now(), + _load_id = '{{ run_id }}', + _load_ts = now() +FROM tmp_monthly_overview_delta AS src +WHERE tgt.year_actual = src.year_actual + AND tgt.month_actual = src.month_actual + AND tgt.airplane_sk = src.airplane_sk + AND ( + tgt.total_flights IS DISTINCT FROM src.total_flights + OR tgt.total_revenue IS DISTINCT FROM src.total_revenue + OR tgt.total_boarded IS DISTINCT FROM src.total_boarded + OR tgt.unique_passengers IS DISTINCT FROM src.unique_passengers + OR tgt.airplane_model IS DISTINCT FROM src.airplane_model + ); + +-- Шаг 5: INSERT новых записей. +INSERT INTO dm.monthly_overview ( + year_actual, + month_actual, + airplane_sk, + airplane_bk, + airplane_model, + total_seats, + total_flights, + total_tickets, + total_boarded, + total_revenue, + avg_ticket_price, + avg_load_factor, + unique_routes, + unique_passengers, + _load_id +) +SELECT + src.year_actual, + src.month_actual, + src.airplane_sk, + src.airplane_bk, + src.airplane_model, + src.total_seats, + src.total_flights, + src.total_tickets, + src.total_boarded, + src.total_revenue, + src.avg_ticket_price, + src.avg_load_factor, + src.unique_routes, + src.unique_passengers, + '{{ run_id }}' AS _load_id +FROM tmp_monthly_overview_delta AS src +WHERE NOT EXISTS ( + SELECT 1 FROM dm.monthly_overview AS tgt + WHERE tgt.year_actual = src.year_actual + AND tgt.month_actual = src.month_actual + AND tgt.airplane_sk = src.airplane_sk +); diff --git a/sql/dm/route_performance_load.sql b/sql/dm/route_performance_load.sql index 2e441e8..5bb5dc6 100644 --- a/sql/dm/route_performance_load.sql +++ b/sql/dm/route_performance_load.sql @@ -29,7 +29,7 @@ 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: Объединяем агрегаты с АКТУАЛЬНЫМИ атрибутами маршрута, аэропортов и самолетов. INSERT INTO dm.route_performance ( route_bk, route_sk, @@ -54,13 +54,13 @@ INSERT INTO dm.route_performance ( SELECT m.route_bk, r_curr.route_sk, - r_curr.departure_airport_bk, - r_curr.departure_city, - r_curr.arrival_airport_bk, - r_curr.arrival_city, - r_curr.airplane_bk, - r_curr.airplane_model, - r_curr.total_seats, + 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, m.total_flights, m.total_tickets, m.total_boarded, @@ -69,14 +69,15 @@ SELECT (m.total_revenue / NULLIF(m.total_tickets, 0))::NUMERIC(10,2), m.avg_boarding_rate::NUMERIC(5,4), -- Load Factor: общее кол-во посадок / (кол-во рейсов * мест в самолете) - -- Это учебное допущение (упрощение), т.к. самолет мог меняться. ROUND( - m.total_boarded::NUMERIC / NULLIF(m.total_flights * r_curr.total_seats, 0), + m.total_boarded::NUMERIC / NULLIF(m.total_flights * air.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 -WHERE r_curr.valid_to IS NULL; -- Берем только АКТУАЛЬНУЮ версию для денормализации +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; diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index 8c7bc32..03424aa 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -390,11 +390,18 @@ def test_bookings_dm_ddl_dag_structure(): expected_tasks = { "apply_dm_sales_report_ddl", + "apply_dm_route_performance_ddl", + "apply_dm_passenger_loyalty_ddl", + "apply_dm_airport_traffic_ddl", + "apply_dm_monthly_overview_ddl", } assert expected_tasks.issubset(dag.task_dict.keys()) - # На данном этапе только одна витрина - # В будущем: линейная цепочка из 5 задач + # Проверяем линейную цепочку + _assert_direct_edge(dag, "apply_dm_sales_report_ddl", "apply_dm_route_performance_ddl") + _assert_direct_edge(dag, "apply_dm_route_performance_ddl", "apply_dm_passenger_loyalty_ddl") + _assert_direct_edge(dag, "apply_dm_passenger_loyalty_ddl", "apply_dm_airport_traffic_ddl") + _assert_direct_edge(dag, "apply_dm_airport_traffic_ddl", "apply_dm_monthly_overview_ddl") def test_bookings_to_gp_dm_dag_structure(): @@ -405,13 +412,30 @@ def test_bookings_to_gp_dm_dag_structure(): "start_dm", "load_dm_sales_report", "dq_dm_sales_report", + "load_dm_route_performance", + "dq_dm_route_performance", + "load_dm_passenger_loyalty", + "dq_dm_passenger_loyalty", + "load_dm_airport_traffic", + "dq_dm_airport_traffic", + "load_dm_monthly_overview", + "dq_dm_monthly_overview", "finish_dm_summary", } assert expected_tasks.issubset(dag.task_dict.keys()) - # Проверяем зависимости load -> dq - _assert_direct_edge(dag, "load_dm_sales_report", "dq_dm_sales_report") + marts = [ + "sales_report", + "route_performance", + "passenger_loyalty", + "airport_traffic", + "monthly_overview", + ] - # Проверяем, что start -> load -> dq -> finish - _assert_reachable(dag, "start_dm", "load_dm_sales_report") - _assert_reachable(dag, "dq_dm_sales_report", "finish_dm_summary") + for mart in marts: + # load -> dq + _assert_direct_edge(dag, f"load_dm_{mart}", f"dq_dm_{mart}") + # start -> load + _assert_reachable(dag, "start_dm", f"load_dm_{mart}") + # dq -> finish + _assert_reachable(dag, f"dq_dm_{mart}", "finish_dm_summary")