From a06b52c38d3269611e59883c381a532d65336142 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Wed, 4 Mar 2026 22:03:26 +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=20passenger=5Floyalty=20=D0=B8=20=D0=B2=D0=BD=D0=B5?= =?UTF-8?q?=D0=B4=D1=80=D0=B5=D0=BD=20=D0=BF=D0=B0=D1=82=D1=82=D0=B5=D1=80?= =?UTF-8?q?=D0=BD=20HWM=20=D0=BF=D0=BE=20=D0=BA=D0=BB=D1=8E=D1=87=D1=83?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - необходимо продемонстрировать студентам метод инкрементального пересчета "затронутых ключей" для больших справочных витрин. - Что: - созданы DDL, Load и DQ скрипты для витрины dm.passenger_loyalty (лояльность пассажиров). - реализован расчет моды (самый частый тариф) через PostgreSQL-специфику DISTINCT ON. - внедрена корректная агрегация SCD2-измерений (unique_routes) по бизнес-ключу route_bk. - настроено Heap-хранилище (WITH appendonly=false) для эффективного выполнения UPSERT. - обновлены DAGи bookings_dm_ddl и bookings_to_gp_dm для включения новой витрины в конвейер. - Проверка: - визуальный аудит SQL на предмет использования p.passenger_id (BK) и r.route_bk. - наличие учебной DQ-проверки ссылочной целостности и инварианта дат (first <= last). - проверка параллельности задач в Airflow DAG. --- airflow/dags/bookings_dm_ddl.py | 14 +++- airflow/dags/bookings_to_gp_dm.py | 23 ++++-- sql/dm/passenger_loyalty_ddl.sql | 41 ++++++++++ sql/dm/passenger_loyalty_dq.sql | 51 ++++++++++++ sql/dm/passenger_loyalty_load.sql | 129 ++++++++++++++++++++++++++++++ 5 files changed, 249 insertions(+), 9 deletions(-) create mode 100644 sql/dm/passenger_loyalty_ddl.sql create mode 100644 sql/dm/passenger_loyalty_dq.sql create mode 100644 sql/dm/passenger_loyalty_load.sql diff --git a/airflow/dags/bookings_dm_ddl.py b/airflow/dags/bookings_dm_ddl.py index aedb383..344a9b1 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. +На данном этапе реализованы витрины: sales_report, route_performance, passenger_loyalty. Остальные витрины будут добавлены в последующих этапах. """ @@ -46,10 +46,16 @@ with DAG( sql="dm/route_performance_ddl.sql", ) - # Заглушки для будущих витрин (будут реализованы в этапах 3-5) - # apply_dm_passenger_loyalty_ddl = PostgresOperator(...) + # Витрина: Лояльность пассажиров (Этап 3) + apply_dm_passenger_loyalty_ddl = PostgresOperator( + task_id="apply_dm_passenger_loyalty_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/passenger_loyalty_ddl.sql", + ) + + # Заглушки для будущих витрин (будут реализованы в этапах 4-5) # apply_dm_airport_traffic_ddl = PostgresOperator(...) # apply_dm_monthly_overview_ddl = PostgresOperator(...) # Линейная цепочка - apply_dm_sales_report_ddl >> apply_dm_route_performance_ddl + apply_dm_sales_report_ddl >> apply_dm_route_performance_ddl >> apply_dm_passenger_loyalty_ddl diff --git a/airflow/dags/bookings_to_gp_dm.py b/airflow/dags/bookings_to_gp_dm.py index 80f7fab..56b11f2 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. +На данном этапе реализованы витрины: sales_report, route_performance, passenger_loyalty. Остальные витрины будут добавлены в последующих этапах. """ @@ -81,9 +81,20 @@ with DAG( sql="dm/route_performance_dq.sql", ) - # === Заглушки для будущих витрин (будут реализованы в этапах 3-5) === - # load_dm_passenger_loyalty = PostgresOperator(...) - # dq_dm_passenger_loyalty = PostgresOperator(...) + # === Витрина: Лояльность пассажиров (Этап 3) === + load_dm_passenger_loyalty = PostgresOperator( + task_id="load_dm_passenger_loyalty", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/passenger_loyalty_load.sql", + ) + + dq_dm_passenger_loyalty = PostgresOperator( + task_id="dq_dm_passenger_loyalty", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/passenger_loyalty_dq.sql", + ) + + # === Заглушки для будущих витрин (будут реализованы в этапах 4-5) === # load_dm_airport_traffic = PostgresOperator(...) # dq_dm_airport_traffic = PostgresOperator(...) # load_dm_monthly_overview = PostgresOperator(...) @@ -98,9 +109,11 @@ with DAG( # Зависимости: параллельные ветки load -> dq start_dm >> [ load_dm_sales_report, - load_dm_route_performance + load_dm_route_performance, + load_dm_passenger_loyalty ] # Связываем 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 diff --git a/sql/dm/passenger_loyalty_ddl.sql b/sql/dm/passenger_loyalty_ddl.sql new file mode 100644 index 0000000..e16ce5c --- /dev/null +++ b/sql/dm/passenger_loyalty_ddl.sql @@ -0,0 +1,41 @@ +-- DDL для витрины dm.passenger_loyalty (Лояльность пассажиров). +-- +-- Учебные цели: +-- 1. Выбор зерна для SCD1-измерений: +-- В SCD2-витринах (напр. route_performance) зерно — бизнес-ключ (BK), т.к. версий много. +-- В SCD1 (dim_passengers) BK и SK связаны 1:1, поэтому зерном выступает суррогатный ключ (SK). +-- Это упрощает JOIN с фактами и ускоряет запросы в Greenplum. +-- 2. Использование Heap-таблицы для UPSERT: +-- В отличие от AO Column, Heap поддерживает эффективный UPDATE, что критично для +-- накопительных витрин с большим количеством строк. +-- 3. Служебные поля жизненного цикла: +-- Т.к. мы будем обновлять метрики пассажиров, нам нужны и created_at, и updated_at. + +CREATE TABLE IF NOT EXISTS dm.passenger_loyalty ( + -- Ключ: суррогатный ключ пассажира (SCD1 гарантирует 1:1 к бизнес-ключу) + passenger_sk INTEGER NOT NULL, + passenger_bk TEXT NOT NULL, + passenger_name TEXT NOT NULL, + + -- Метрики лояльности + total_bookings INTEGER NOT NULL, -- кол-во бронирований (book_ref) + total_flights INTEGER NOT NULL, -- кол-во перелетов + total_boarded INTEGER NOT NULL, -- кол-во успешных посадок + total_spent NUMERIC(15,2) NOT NULL, + avg_ticket_price NUMERIC(10,2), + favorite_fare_conditions TEXT, -- самый частый класс обслуживания + unique_routes INTEGER NOT NULL, -- кол-во уникальных маршрутов + first_flight_date DATE, + last_flight_date DATE, + days_as_customer 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() +) +WITH (appendonly=false) -- Явно указываем Heap для поддержки эффективного UPDATE +DISTRIBUTED BY (passenger_sk); + +COMMENT ON TABLE dm.passenger_loyalty IS 'Витрина: лояльность и активность пассажиров (Incremental UPSERT, Heap)'; diff --git a/sql/dm/passenger_loyalty_dq.sql b/sql/dm/passenger_loyalty_dq.sql new file mode 100644 index 0000000..6e57d32 --- /dev/null +++ b/sql/dm/passenger_loyalty_dq.sql @@ -0,0 +1,51 @@ +-- DQ проверки для витрины dm.passenger_loyalty. +-- +-- Учебные цели: +-- 1. Валидация бизнес-логики дат (стаж не может быть отрицательным). +-- 2. Учебная проверка ссылочной целостности (FK Integrity). + +DO $$ +DECLARE + row_count INTEGER; + duplicate_count INTEGER; + invalid_dates_count INTEGER; + orphan_keys_count INTEGER; +BEGIN + -- 1. Проверка на наполненность + SELECT COUNT(*) INTO row_count FROM dm.passenger_loyalty; + IF row_count = 0 THEN + RAISE EXCEPTION 'DQ Error: Таблица dm.passenger_loyalty пуста.'; + END IF; + + -- 2. Проверка на уникальность (зерно - passenger_sk) + SELECT COUNT(*) INTO duplicate_count + FROM ( + SELECT passenger_sk FROM dm.passenger_loyalty + GROUP BY passenger_sk HAVING COUNT(*) > 1 + ) q; + IF duplicate_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.passenger_loyalty обнаружены дубликаты по passenger_sk (% шт).', duplicate_count; + END IF; + + -- 3. Проверка логики дат + SELECT COUNT(*) INTO invalid_dates_count + FROM dm.passenger_loyalty + WHERE first_flight_date > last_flight_date; + + IF invalid_dates_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.passenger_loyalty обнаружены записи с first_date > last_date (% строк).', invalid_dates_count; + END IF; + + -- 4. Учебная проверка ссылочной целостности (FK Check) + -- В продакшене это обычно гарантируется JOIN при загрузке, но здесь мы показываем саму возможность проверки. + SELECT COUNT(*) INTO orphan_keys_count + FROM dm.passenger_loyalty tgt + LEFT JOIN dds.dim_passengers p ON tgt.passenger_sk = p.passenger_sk + WHERE p.passenger_sk IS NULL; + + IF orphan_keys_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.passenger_loyalty обнаружены пассажиры, отсутствующие в dim_passengers (% строк).', orphan_keys_count; + END IF; + + RAISE NOTICE 'DQ Success: dm.passenger_loyalty успешно прошла все проверки (% строк).', row_count; +END $$; diff --git a/sql/dm/passenger_loyalty_load.sql b/sql/dm/passenger_loyalty_load.sql new file mode 100644 index 0000000..273030f --- /dev/null +++ b/sql/dm/passenger_loyalty_load.sql @@ -0,0 +1,129 @@ +-- Загрузка витрины dm.passenger_loyalty: инкрементальный UPSERT по затронутым ключам. +-- +-- Учебные цели: +-- 1. Метод "затронутых ключей": HWM по _load_ts находит затронутые ID пассажиров, +-- а затем мы ПЕРЕСЧИТЫВАЕМ всю историю именно для этого круга лиц. +-- Это гарантирует точность накопительных агрегатов (total_spent, dates). +-- 2. Использование DISTINCT ON (PostgreSQL-специфика): самый лаконичный способ +-- найти "самое частое" (моду) в рамках группы. +-- 3. Обработка NULL в фактах: фильтрация (passenger_sk IS NOT NULL). +-- 4. Агрегация SCD2-измерений: при подсчете уникальных маршрутов (dim_routes) +-- нужно агрегировать по BK (route_bk), т.к. один маршрут может иметь несколько SK (версий). + +-- Шаг 1: Находим ID пассажиров, чьи данные изменились или добавились в фактах. +CREATE TEMP TABLE tmp_loyalty_affected_keys ON COMMIT DROP AS +SELECT DISTINCT passenger_sk +FROM dds.fact_flight_sales +WHERE _load_ts > ( + SELECT COALESCE(MAX(_load_ts), '1900-01-01'::TIMESTAMP) + FROM dm.passenger_loyalty +) +AND passenger_sk IS NOT NULL; + +-- Шаг 2: Для затронутых лиц считаем ИТОГОВЫЕ агрегаты по ВСЕЙ истории фактов. +CREATE TEMP TABLE tmp_passenger_delta ON COMMIT DROP AS +WITH base_metrics AS ( + -- Агрегируем количественные метрики. + -- JOIN dds.dim_routes нужен для подсчета УНИКАЛЬНЫХ маршрутов по BK (т.к. dim_routes - SCD2). + SELECT + f.passenger_sk, + COUNT(DISTINCT f.book_ref) AS total_bookings, + COUNT(*) AS total_flights, + SUM(CASE WHEN f.is_boarded THEN 1 ELSE 0 END) AS total_boarded, + SUM(f.price) AS total_spent, + COUNT(DISTINCT r.route_bk) AS unique_routes, -- Агрегация по BK (бизнес-ключу) маршрута + MIN(cal.date_actual) AS first_flight_date, + MAX(cal.date_actual) AS last_flight_date + 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 f.passenger_sk IN (SELECT passenger_sk FROM tmp_loyalty_affected_keys) + GROUP BY f.passenger_sk +), +fare_modes AS ( + -- Находим самый частый класс обслуживания для каждого пассажира. + -- DISTINCT ON при ничьей выбирает произвольный вариант из топ-результатов. + SELECT DISTINCT ON (f.passenger_sk) + f.passenger_sk, + tar.fare_conditions AS favorite_fare_conditions + FROM dds.fact_flight_sales f + JOIN dds.dim_tariffs tar ON f.tariff_sk = tar.tariff_sk + WHERE f.passenger_sk IN (SELECT passenger_sk FROM tmp_loyalty_affected_keys) + GROUP BY f.passenger_sk, tar.fare_conditions + ORDER BY f.passenger_sk, COUNT(*) DESC +) +SELECT + bm.*, + fm.favorite_fare_conditions, + p.passenger_id AS passenger_bk, -- Исправлено: в dim_passengers BK называется passenger_id + p.passenger_name, + (bm.last_flight_date - bm.first_flight_date) AS days_as_customer +FROM base_metrics bm +JOIN dds.dim_passengers p ON bm.passenger_sk = p.passenger_sk +JOIN fare_modes fm ON bm.passenger_sk = fm.passenger_sk; + +-- Шаг 3: UPDATE существующих записей. +UPDATE dm.passenger_loyalty AS tgt +SET + passenger_name = src.passenger_name, + total_bookings = src.total_bookings, + total_flights = src.total_flights, + total_boarded = src.total_boarded, + total_spent = src.total_spent, + avg_ticket_price = (src.total_spent / NULLIF(src.total_flights, 0))::NUMERIC(10,2), + favorite_fare_conditions = src.favorite_fare_conditions, + unique_routes = src.unique_routes, + first_flight_date = src.first_flight_date, + last_flight_date = src.last_flight_date, + days_as_customer = src.days_as_customer, + updated_at = now(), + _load_id = '{{ run_id }}', + _load_ts = now() +FROM tmp_passenger_delta AS src +WHERE tgt.passenger_sk = src.passenger_sk + AND ( + -- Обновляем только если что-то реально изменилось + tgt.total_flights IS DISTINCT FROM src.total_flights + OR tgt.total_boarded IS DISTINCT FROM src.total_boarded + OR tgt.total_spent IS DISTINCT FROM src.total_spent + OR tgt.last_flight_date IS DISTINCT FROM src.last_flight_date + OR tgt.passenger_name IS DISTINCT FROM src.passenger_name + ); + +-- Шаг 4: INSERT новых пассажиров. +INSERT INTO dm.passenger_loyalty ( + passenger_sk, + passenger_bk, + passenger_name, + total_bookings, + total_flights, + total_boarded, + total_spent, + avg_ticket_price, + favorite_fare_conditions, + unique_routes, + first_flight_date, + last_flight_date, + days_as_customer, + _load_id +) +SELECT + src.passenger_sk, + src.passenger_bk, + src.passenger_name, + src.total_bookings, + src.total_flights, + src.total_boarded, + src.total_spent, + (src.total_spent / NULLIF(src.total_flights, 0))::NUMERIC(10,2), + src.favorite_fare_conditions, + src.unique_routes, + src.first_flight_date, + src.last_flight_date, + src.days_as_customer, + '{{ run_id }}' AS _load_id +FROM tmp_passenger_delta AS src +WHERE NOT EXISTS ( + SELECT 1 FROM dm.passenger_loyalty AS tgt + WHERE tgt.passenger_sk = src.passenger_sk +);