feat(dm): добавлена витрина airport_traffic и внедрен паттерн dual-role dimensions
- Зачем: - необходимо продемонстрировать студентам работу с одной сущностью в разных ролях (вылет/прилет) через 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).
This commit is contained in:
@@ -7,7 +7,7 @@ from __future__ import annotations
|
|||||||
Создаёт 5 DM-витрин: sales_report, route_performance, passenger_loyalty,
|
Создаёт 5 DM-витрин: sales_report, route_performance, passenger_loyalty,
|
||||||
airport_traffic, monthly_overview.
|
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",
|
sql="dm/passenger_loyalty_ddl.sql",
|
||||||
)
|
)
|
||||||
|
|
||||||
# Заглушки для будущих витрин (будут реализованы в этапах 4-5)
|
# Витрина: Пассажиропоток аэропортов (Этап 4)
|
||||||
# apply_dm_airport_traffic_ddl = PostgresOperator(...)
|
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_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
|
||||||
|
|||||||
@@ -9,7 +9,7 @@ from __future__ import annotations
|
|||||||
- все витрины загружаются параллельно (не зависят друг от друга);
|
- все витрины загружаются параллельно (не зависят друг от друга);
|
||||||
- sales_report использует UPSERT (heap-таблица), route_performance — Full Rebuild (AO Column).
|
- 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",
|
sql="dm/passenger_loyalty_dq.sql",
|
||||||
)
|
)
|
||||||
|
|
||||||
# === Заглушки для будущих витрин (будут реализованы в этапах 4-5) ===
|
# === Витрина: Пассажиропоток аэропортов (Этап 4) ===
|
||||||
# load_dm_airport_traffic = PostgresOperator(...)
|
load_dm_airport_traffic = PostgresOperator(
|
||||||
# dq_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(...)
|
# load_dm_monthly_overview = PostgresOperator(...)
|
||||||
# dq_dm_monthly_overview = PostgresOperator(...)
|
# dq_dm_monthly_overview = PostgresOperator(...)
|
||||||
|
|
||||||
@@ -110,10 +121,12 @@ with DAG(
|
|||||||
start_dm >> [
|
start_dm >> [
|
||||||
load_dm_sales_report,
|
load_dm_sales_report,
|
||||||
load_dm_route_performance,
|
load_dm_route_performance,
|
||||||
load_dm_passenger_loyalty
|
load_dm_passenger_loyalty,
|
||||||
|
load_dm_airport_traffic
|
||||||
]
|
]
|
||||||
|
|
||||||
# Связываем dq с finish
|
# Связываем dq с finish
|
||||||
load_dm_sales_report >> dq_dm_sales_report >> finish_dm_summary
|
load_dm_sales_report >> dq_dm_sales_report >> finish_dm_summary
|
||||||
load_dm_route_performance >> dq_dm_route_performance >> 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_passenger_loyalty >> dq_dm_passenger_loyalty >> finish_dm_summary
|
||||||
|
load_dm_airport_traffic >> dq_dm_airport_traffic >> finish_dm_summary
|
||||||
|
|||||||
@@ -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 приведет к двойному счету!';
|
||||||
@@ -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 $$;
|
||||||
@@ -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
|
||||||
|
);
|
||||||
Reference in New Issue
Block a user