From f68b055b006a172795eb2a90b263c13eb67f28cc Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Wed, 4 Mar 2026 22:19:06 +0300 Subject: [PATCH] =?UTF-8?q?feat(dm):=20=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2?= =?UTF-8?q?=D0=BB=D0=B5=D0=BD=D0=B0=20=D0=B2=D0=B8=D1=82=D1=80=D0=B8=D0=BD?= =?UTF-8?q?=D0=B0=20airport=5Ftraffic=20=D0=B8=20=D0=B2=D0=BD=D0=B5=D0=B4?= =?UTF-8?q?=D1=80=D0=B5=D0=BD=20=D0=BF=D0=B0=D1=82=D1=82=D0=B5=D1=80=D0=BD?= =?UTF-8?q?=20dual-role=20dimensions?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - необходимо продемонстрировать студентам работу с одной сущностью в разных ролях (вылет/прилет) через UNION ALL. - Что: - созданы DDL, Load и DQ скрипты для витрины dm.airport_traffic (пассажиропоток аэропортов). - реализован паттерн Unpivot через UNION ALL для консолидации метрик вылета и прилета в одном разрезе. - внедрена инкрементальная загрузка по затронутым датам (HWM) с честным подсчетом рейсов через COUNT(DISTINCT). - добавлены подробные комментарии к колонкам выручки, предупреждающие о риске двойного счета. - обновлены DAGи bookings_dm_ddl и bookings_to_gp_dm для включения витрины в общий пайплайн. - Проверка: - визуальный аудит SQL-кода на предмет использования airport_bk и корректной агрегации по ролям. - наличие DQ-инварианта total_passengers = departures + arrivals. - верификация блока UPDATE: теперь обновляются и денормализованные атрибуты (city, airport_bk). --- airflow/dags/bookings_dm_ddl.py | 14 +++- airflow/dags/bookings_to_gp_dm.py | 23 ++++-- sql/dm/airport_traffic_ddl.sql | 45 +++++++++++ sql/dm/airport_traffic_dq.sql | 49 ++++++++++++ sql/dm/airport_traffic_load.sql | 129 ++++++++++++++++++++++++++++++ 5 files changed, 251 insertions(+), 9 deletions(-) create mode 100644 sql/dm/airport_traffic_ddl.sql create mode 100644 sql/dm/airport_traffic_dq.sql create mode 100644 sql/dm/airport_traffic_load.sql diff --git a/airflow/dags/bookings_dm_ddl.py b/airflow/dags/bookings_dm_ddl.py index 344a9b1..524b721 100644 --- a/airflow/dags/bookings_dm_ddl.py +++ b/airflow/dags/bookings_dm_ddl.py @@ -7,7 +7,7 @@ from __future__ import annotations Создаёт 5 DM-витрин: sales_report, route_performance, passenger_loyalty, airport_traffic, monthly_overview. -На данном этапе реализованы витрины: sales_report, route_performance, passenger_loyalty. +На данном этапе реализованы витрины: sales_report, route_performance, passenger_loyalty, airport_traffic. Остальные витрины будут добавлены в последующих этапах. """ @@ -53,9 +53,15 @@ with DAG( sql="dm/passenger_loyalty_ddl.sql", ) - # Заглушки для будущих витрин (будут реализованы в этапах 4-5) - # apply_dm_airport_traffic_ddl = PostgresOperator(...) + # Витрина: Пассажиропоток аэропортов (Этап 4) + apply_dm_airport_traffic_ddl = PostgresOperator( + task_id="apply_dm_airport_traffic_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/airport_traffic_ddl.sql", + ) + + # Заглушки для будущих витрин (будут реализованы в этапе 5) # apply_dm_monthly_overview_ddl = PostgresOperator(...) # Линейная цепочка - apply_dm_sales_report_ddl >> apply_dm_route_performance_ddl >> apply_dm_passenger_loyalty_ddl + apply_dm_sales_report_ddl >> apply_dm_route_performance_ddl >> apply_dm_passenger_loyalty_ddl >> apply_dm_airport_traffic_ddl diff --git a/airflow/dags/bookings_to_gp_dm.py b/airflow/dags/bookings_to_gp_dm.py index 56b11f2..6baba8c 100644 --- a/airflow/dags/bookings_to_gp_dm.py +++ b/airflow/dags/bookings_to_gp_dm.py @@ -9,7 +9,7 @@ from __future__ import annotations - все витрины загружаются параллельно (не зависят друг от друга); - sales_report использует UPSERT (heap-таблица), route_performance — Full Rebuild (AO Column). -На данном этапе реализованы витрины: sales_report, route_performance, passenger_loyalty. +На данном этапе реализованы витрины: sales_report, route_performance, passenger_loyalty, airport_traffic. Остальные витрины будут добавлены в последующих этапах. """ @@ -94,9 +94,20 @@ with DAG( sql="dm/passenger_loyalty_dq.sql", ) - # === Заглушки для будущих витрин (будут реализованы в этапах 4-5) === - # load_dm_airport_traffic = PostgresOperator(...) - # dq_dm_airport_traffic = PostgresOperator(...) + # === Витрина: Пассажиропоток аэропортов (Этап 4) === + load_dm_airport_traffic = PostgresOperator( + task_id="load_dm_airport_traffic", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/airport_traffic_load.sql", + ) + + dq_dm_airport_traffic = PostgresOperator( + task_id="dq_dm_airport_traffic", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/airport_traffic_dq.sql", + ) + + # === Заглушки для будущих витрин (будут реализованы в этапе 5) === # load_dm_monthly_overview = PostgresOperator(...) # dq_dm_monthly_overview = PostgresOperator(...) @@ -110,10 +121,12 @@ with DAG( start_dm >> [ load_dm_sales_report, load_dm_route_performance, - load_dm_passenger_loyalty + load_dm_passenger_loyalty, + load_dm_airport_traffic ] # Связываем dq с finish load_dm_sales_report >> dq_dm_sales_report >> finish_dm_summary 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 diff --git a/sql/dm/airport_traffic_ddl.sql b/sql/dm/airport_traffic_ddl.sql new file mode 100644 index 0000000..4d4479e --- /dev/null +++ b/sql/dm/airport_traffic_ddl.sql @@ -0,0 +1,45 @@ +-- DDL для витрины dm.airport_traffic (Пассажиропоток аэропортов). +-- +-- Учебные цели: +-- 1. Обработка dual-role dimensions: один аэропорт выступает и как точка вылета, +-- и как точка прилета. Витрина собирает статистику в едином разрезе (date, airport). +-- 2. Денормализация: включение кода аэропорта (BK) и города для удобства анализа. +-- 3. Предупреждение о семантике данных: риск "двойного счета" выручки. +-- 4. Стратегия HWM по датам: пересчет только затронутых дней. + +CREATE TABLE IF NOT EXISTS dm.airport_traffic ( + -- Зерно: дата и суррогатный ключ аэропорта + traffic_date DATE NOT NULL, + airport_sk INTEGER NOT NULL, + + -- Денормализованные атрибуты (из dim_airports) + airport_bk TEXT NOT NULL, + city TEXT NOT NULL, + + -- Метрики вылета + departures_flights INTEGER NOT NULL DEFAULT 0, + departures_passengers INTEGER NOT NULL DEFAULT 0, + departures_revenue NUMERIC(15,2) NOT NULL DEFAULT 0, + + -- Метрики прилета + arrivals_flights INTEGER NOT NULL DEFAULT 0, + arrivals_passengers INTEGER NOT NULL DEFAULT 0, + arrivals_revenue NUMERIC(15,2) NOT NULL DEFAULT 0, + + -- Итоговые метрики + total_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 (airport_sk); + +COMMENT ON TABLE dm.airport_traffic IS 'Витрина: ежедневный пассажиропоток аэропортов (UNION ALL, Heap)'; + +-- Комментарии к колонкам с предупреждением (учебная ценность) +COMMENT ON COLUMN dm.airport_traffic.departures_revenue IS 'Выручка от билетов, где аэропорт был точкой вылета. ВНИМАНИЕ: суммирование с arrivals_revenue приведет к двойному счету!'; +COMMENT ON COLUMN dm.airport_traffic.arrivals_revenue IS 'Выручка от билетов, где аэропорт был точкой прилета. ВНИМАНИЕ: суммирование с departures_revenue приведет к двойному счету!'; diff --git a/sql/dm/airport_traffic_dq.sql b/sql/dm/airport_traffic_dq.sql new file mode 100644 index 0000000..bb8021b --- /dev/null +++ b/sql/dm/airport_traffic_dq.sql @@ -0,0 +1,49 @@ +-- DQ проверки для витрины dm.airport_traffic. +-- +-- Учебные цели: +-- 1. Проверка арифметических инвариантов (сумма частей должна быть равна целому). +-- 2. Валидация уникальности составного ключа (зерна). + +DO $$ +DECLARE + row_count INTEGER; + duplicate_count INTEGER; + invalid_total_count INTEGER; + negative_metrics_count INTEGER; +BEGIN + -- 1. Проверка на наполненность + SELECT COUNT(*) INTO row_count FROM dm.airport_traffic; + IF row_count = 0 THEN + RAISE EXCEPTION 'DQ Error: Таблица dm.airport_traffic пуста.'; + END IF; + + -- 2. Проверка на уникальность (зерно - traffic_date + airport_sk) + SELECT COUNT(*) INTO duplicate_count + FROM ( + SELECT traffic_date, airport_sk FROM dm.airport_traffic + GROUP BY traffic_date, airport_sk HAVING COUNT(*) > 1 + ) q; + IF duplicate_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.airport_traffic обнаружены дубликаты по (date, airport_sk) (% шт).', duplicate_count; + END IF; + + -- 3. Проверка инварианта total = departures + arrivals + SELECT COUNT(*) INTO invalid_total_count + FROM dm.airport_traffic + WHERE total_passengers != (departures_passengers + arrivals_passengers); + + IF invalid_total_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.airport_traffic нарушен инвариант total_passengers = dep + arr (% строк).', invalid_total_count; + END IF; + + -- 4. Проверка на отрицательные значения + SELECT COUNT(*) INTO negative_metrics_count + FROM dm.airport_traffic + WHERE departures_flights < 0 OR arrivals_flights < 0 OR departures_revenue < 0 OR arrivals_revenue < 0; + + IF negative_metrics_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.airport_traffic обнаружены отрицательные метрики (% строк).', negative_metrics_count; + END IF; + + RAISE NOTICE 'DQ Success: dm.airport_traffic успешно прошла все проверки (% строк).', row_count; +END $$; diff --git a/sql/dm/airport_traffic_load.sql b/sql/dm/airport_traffic_load.sql new file mode 100644 index 0000000..2303b8d --- /dev/null +++ b/sql/dm/airport_traffic_load.sql @@ -0,0 +1,129 @@ +-- Загрузка витрины dm.airport_traffic: инкрементальный UPSERT по датам. +-- +-- Учебные цели: +-- 1. Обработка dual-role dimensions через UNION ALL: +-- Мы превращаем один билет (факт) в два "события": вылет и прилет. +-- Это позволяет собрать единую статистику аэропорта в одном проходе. +-- 2. Инкремент по датам (HWM): в этой витрине метрики ограничены сутками, +-- поэтому пересчитываем только те дни, где появились новые факты. +-- 3. Агрегация рейсов: используем COUNT(DISTINCT flight_id) для подсчета рейсов. +-- 4. Денормализация (SCD1): обновление атрибутов (город) при изменениях в измерении. + +-- Шаг 1: Находим даты, затронутые новыми/измененными фактами. +CREATE TEMP TABLE tmp_traffic_affected_dates ON COMMIT DROP AS +SELECT DISTINCT cal.date_actual AS traffic_date +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.airport_traffic +); + +-- Шаг 2: Unpivot (UNION ALL) и агрегация по затронутым датам. +CREATE TEMP TABLE tmp_airport_traffic_delta ON COMMIT DROP AS +WITH raw_events AS ( + -- Роль 1: Аэропорт вылета + SELECT + cal.date_actual AS traffic_date, + f.departure_airport_sk AS airport_sk, + f.flight_id, + 'departure' AS role, + (CASE WHEN f.is_boarded THEN 1 ELSE 0 END) AS is_passenger, + f.price AS revenue + FROM dds.fact_flight_sales f + JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk + WHERE cal.date_actual IN (SELECT traffic_date FROM tmp_traffic_affected_dates) + + UNION ALL + + -- Роль 2: Аэропорт прилета + SELECT + cal.date_actual AS traffic_date, + f.arrival_airport_sk AS airport_sk, + f.flight_id, + 'arrival' AS role, + (CASE WHEN f.is_boarded THEN 1 ELSE 0 END) AS is_passenger, + f.price AS revenue + FROM dds.fact_flight_sales f + JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk + WHERE cal.date_actual IN (SELECT traffic_date FROM tmp_traffic_affected_dates) +) +SELECT + e.traffic_date, + e.airport_sk, + a.airport_bk, -- Берем из dim_airports.airport_bk + a.city, + -- Агрегация вылетов + COUNT(DISTINCT CASE WHEN e.role = 'departure' THEN e.flight_id END) AS departures_flights, + SUM(CASE WHEN e.role = 'departure' THEN e.is_passenger ELSE 0 END) AS departures_passengers, + SUM(CASE WHEN e.role = 'departure' THEN e.revenue ELSE 0 END) AS departures_revenue, + -- Агрегация прилетов + COUNT(DISTINCT CASE WHEN e.role = 'arrival' THEN e.flight_id END) AS arrivals_flights, + SUM(CASE WHEN e.role = 'arrival' THEN e.is_passenger ELSE 0 END) AS arrivals_passengers, + SUM(CASE WHEN e.role = 'arrival' THEN e.revenue ELSE 0 END) AS arrivals_revenue, + -- Итого: все пассажиры (сели + вышли в этом аэропорту) + SUM(e.is_passenger) AS total_passengers +FROM raw_events e +JOIN dds.dim_airports a ON e.airport_sk = a.airport_sk +GROUP BY e.traffic_date, e.airport_sk, a.airport_bk, a.city; + +-- Шаг 3: UPDATE существующих записей (включая денормализованные атрибуты). +UPDATE dm.airport_traffic AS tgt +SET + airport_bk = src.airport_bk, + city = src.city, + departures_flights = src.departures_flights, + departures_passengers = src.departures_passengers, + departures_revenue = src.departures_revenue, + arrivals_flights = src.arrivals_flights, + arrivals_passengers = src.arrivals_passengers, + arrivals_revenue = src.arrivals_revenue, + total_passengers = src.total_passengers, + updated_at = now(), + _load_id = '{{ run_id }}', + _load_ts = now() +FROM tmp_airport_traffic_delta AS src +WHERE tgt.traffic_date = src.traffic_date + AND tgt.airport_sk = src.airport_sk + AND ( + tgt.total_passengers IS DISTINCT FROM src.total_passengers + OR tgt.departures_flights IS DISTINCT FROM src.departures_flights + OR tgt.arrivals_flights IS DISTINCT FROM src.arrivals_flights + OR tgt.city IS DISTINCT FROM src.city + OR tgt.airport_bk IS DISTINCT FROM src.airport_bk + ); + +-- Шаг 4: INSERT новых записей. +INSERT INTO dm.airport_traffic ( + traffic_date, + airport_sk, + airport_bk, + city, + departures_flights, + departures_passengers, + departures_revenue, + arrivals_flights, + arrivals_passengers, + arrivals_revenue, + total_passengers, + _load_id +) +SELECT + src.traffic_date, + src.airport_sk, + src.airport_bk, + src.city, + src.departures_flights, + src.departures_passengers, + src.departures_revenue, + src.arrivals_flights, + src.arrivals_passengers, + src.arrivals_revenue, + src.total_passengers, + '{{ run_id }}' AS _load_id +FROM tmp_airport_traffic_delta AS src +WHERE NOT EXISTS ( + SELECT 1 FROM dm.airport_traffic AS tgt + WHERE tgt.traffic_date = src.traffic_date + AND tgt.airport_sk = src.airport_sk +);