From 491f1e68d6908776a96e302d47b2b5310ae48db4 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Wed, 4 Mar 2026 21:52:10 +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=20route=5Fperformance=20=D0=B8=20=D0=BE=D0=B1=D0=BD?= =?UTF-8?q?=D0=BE=D0=B2=D0=BB=D0=B5=D0=BD=D1=8B=20DM=20DAG=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - необходимо продемонстрировать студентам альтернативный паттерн загрузки (Full Rebuild) и использование AO Column Store в Greenplum. - Что: - созданы DDL, Load и DQ скрипты для витрины dm.route_performance (эффективность маршрутов). - реализована агрегация по бизнес-ключу route_bk для корректной обработки SCD2-измерений. - настроен формат хранения AO Column Store с компрессией zstd (уровень 1). - обновлены DAGи bookings_dm_ddl и bookings_to_gp_dm для параллельной оркестрации новой витрины. - Проверка: - визуальный аудит SQL-кода на соответствие naming_conventions.md. - проверка структуры DAG в Airflow (параллельные ветки load -> dq). - наличие бизнес-инварианта total_boarded <= total_tickets в DQ-скрипте. --- airflow/dags/bookings_dm_ddl.py | 18 ++++--- airflow/dags/bookings_to_gp_dm.py | 31 +++++++++--- sql/dm/route_performance_ddl.sql | 48 ++++++++++++++++++ sql/dm/route_performance_dq.sql | 59 ++++++++++++++++++++++ sql/dm/route_performance_load.sql | 82 +++++++++++++++++++++++++++++++ 5 files changed, 224 insertions(+), 14 deletions(-) create mode 100644 sql/dm/route_performance_ddl.sql create mode 100644 sql/dm/route_performance_dq.sql create mode 100644 sql/dm/route_performance_load.sql diff --git a/airflow/dags/bookings_dm_ddl.py b/airflow/dags/bookings_dm_ddl.py index 540100e..aedb383 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. +На данном этапе реализованы витрины: sales_report, route_performance. Остальные витрины будут добавлены в последующих этапах. """ @@ -39,13 +39,17 @@ with DAG( sql="dm/sales_report_ddl.sql", ) - # Заглушки для будущих витрин (будут реализованы в этапах 2-5) - # apply_dm_route_performance_ddl = PostgresOperator(...) + # Витрина: Эффективность маршрутов (Этап 2) + apply_dm_route_performance_ddl = PostgresOperator( + task_id="apply_dm_route_performance_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/route_performance_ddl.sql", + ) + + # Заглушки для будущих витрин (будут реализованы в этапах 3-5) # apply_dm_passenger_loyalty_ddl = PostgresOperator(...) # apply_dm_airport_traffic_ddl = PostgresOperator(...) # apply_dm_monthly_overview_ddl = PostgresOperator(...) - # Линейная цепочка (пока только одна витрина) - # В будущем: sales_report >> route_performance >> passenger_loyalty - # >> airport_traffic >> monthly_overview - apply_dm_sales_report_ddl + # Линейная цепочка + apply_dm_sales_report_ddl >> apply_dm_route_performance_ddl diff --git a/airflow/dags/bookings_to_gp_dm.py b/airflow/dags/bookings_to_gp_dm.py index e35bd2f..80f7fab 100644 --- a/airflow/dags/bookings_to_gp_dm.py +++ b/airflow/dags/bookings_to_gp_dm.py @@ -7,9 +7,9 @@ from __future__ import annotations - DM читает из DDS (Star Schema); - для каждой витрины выполняем пару задач load -> dq; - все витрины загружаются параллельно (не зависят друг от друга); -- sales_report использует UPSERT (heap-таблица с UPDATE). +- sales_report использует UPSERT (heap-таблица), route_performance — Full Rebuild (AO Column). -На данном этапе реализована только эталонная витрина sales_report. +На данном этапе реализованы витрины: sales_report, route_performance. Остальные витрины будут добавлены в последующих этапах. """ @@ -68,9 +68,20 @@ with DAG( sql="dm/sales_report_dq.sql", ) - # === Заглушки для будущих витрин (будут реализованы в этапах 2-5) === - # load_dm_route_performance = PostgresOperator(...) - # dq_dm_route_performance = PostgresOperator(...) + # === Витрина: Эффективность маршрутов (Этап 2) === + load_dm_route_performance = PostgresOperator( + task_id="load_dm_route_performance", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/route_performance_load.sql", + ) + + dq_dm_route_performance = PostgresOperator( + task_id="dq_dm_route_performance", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/route_performance_dq.sql", + ) + + # === Заглушки для будущих витрин (будут реализованы в этапах 3-5) === # load_dm_passenger_loyalty = PostgresOperator(...) # dq_dm_passenger_loyalty = PostgresOperator(...) # load_dm_airport_traffic = PostgresOperator(...) @@ -85,5 +96,11 @@ with DAG( ) # Зависимости: параллельные ветки load -> dq - # В будущем: start_dm >> [load_sales >> dq_sales, load_route >> dq_route, ...] >> finish - start_dm >> load_dm_sales_report >> dq_dm_sales_report >> finish_dm_summary + start_dm >> [ + load_dm_sales_report, + load_dm_route_performance + ] + + # Связываем dq с finish + load_dm_sales_report >> dq_dm_sales_report >> finish_dm_summary + load_dm_route_performance >> dq_dm_route_performance >> finish_dm_summary diff --git a/sql/dm/route_performance_ddl.sql b/sql/dm/route_performance_ddl.sql new file mode 100644 index 0000000..eb1e895 --- /dev/null +++ b/sql/dm/route_performance_ddl.sql @@ -0,0 +1,48 @@ +-- DDL для витрины dm.route_performance (Эффективность маршрутов). +-- +-- Учебные цели: +-- 1. Демонстрация AO Column Store (Append-Only Columnar Storage). +-- В Greenplum формат AO Column идеален для аналитики (сжатие, чтение только нужных колонок). +-- 2. Демонстрация стратегии Full Rebuild (TRUNCATE + INSERT). +-- Для небольших справочных витрин это проще и надёжнее, чем сложный инкремент. +-- 3. Выбор служебных полей: +-- Для Full Rebuild таблиц created_at/updated_at не имеют смысла, т.к. строки +-- каждый раз пересоздаются. Достаточно _load_ts. + +CREATE TABLE IF NOT EXISTS dm.route_performance ( + -- Ключ: бизнес-код маршрута (напр. 'SVO-LED') + route_bk TEXT NOT NULL, + + -- SK текущей (актуальной) версии маршрута для связи с измерениями + route_sk INTEGER NOT NULL, + + -- Денормализованные атрибуты (из актуальной версии dim_routes) + departure_airport_bk TEXT NOT NULL, + departure_city TEXT NOT NULL, + arrival_airport_bk TEXT NOT NULL, + arrival_city TEXT NOT NULL, + airplane_bk TEXT NOT NULL, + airplane_model TEXT NOT NULL, + total_seats INTEGER NOT NULL, + + -- Метрики (агрегированы по всем версиям данного маршрута) + total_flights INTEGER NOT NULL, + total_tickets INTEGER NOT NULL, + total_boarded INTEGER NOT NULL, + total_revenue NUMERIC(15,2) NOT NULL, + avg_ticket_price NUMERIC(10,2), + avg_boarding_rate NUMERIC(5,4) NOT NULL, + avg_load_factor NUMERIC(5,4), -- средняя заполняемость кресел + first_flight_date DATE, + last_flight_date DATE, + + -- Служебные поля. + -- created_at/updated_at здесь не нужны: при Full Rebuild все строки пересоздаются, + -- поэтому _load_ts достаточно для отслеживания момента загрузки. + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() +) +WITH (appendonly=true, orientation=column, compresstype=zstd, compresslevel=1) +DISTRIBUTED BY (route_bk); + +COMMENT ON TABLE dm.route_performance IS 'Витрина: эффективность авиамаршрутов (Full Rebuild, AO Column, zstd)'; diff --git a/sql/dm/route_performance_dq.sql b/sql/dm/route_performance_dq.sql new file mode 100644 index 0000000..6c6c509 --- /dev/null +++ b/sql/dm/route_performance_dq.sql @@ -0,0 +1,59 @@ +-- DQ проверки для витрины dm.route_performance. +-- +-- Учебные цели: +-- 1. Использование PL/pgSQL блоков (DO $$ ...) для кастомных проверок. +-- 2. Валидация бизнес-инвариантов после загрузки Full Rebuild. + +DO $$ +DECLARE + row_count INTEGER; + duplicate_count INTEGER; + invalid_metrics_count INTEGER; + null_attributes_count INTEGER; +BEGIN + -- 1. Проверка на наполненность + SELECT COUNT(*) INTO row_count FROM dm.route_performance; + IF row_count = 0 THEN + RAISE EXCEPTION 'DQ Error: Таблица dm.route_performance пуста после загрузки.'; + END IF; + + -- 2. Проверка на уникальность бизнес-ключа + SELECT COUNT(*) INTO duplicate_count + FROM ( + SELECT route_bk FROM dm.route_performance + GROUP BY route_bk HAVING COUNT(*) > 1 + ) AS q; + IF duplicate_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.route_performance обнаружены дубликаты по route_bk (% шт).', duplicate_count; + END IF; + + -- 3. Проверка бизнес-метрик (инварианты) + -- Load factor не может быть отрицательным или физически невозможным (например, более 200%) + SELECT COUNT(*) INTO invalid_metrics_count + FROM dm.route_performance + WHERE avg_load_factor < 0 OR avg_load_factor > 2.0; + + IF invalid_metrics_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.route_performance обнаружены некорректные показатели load factor (% строк).', invalid_metrics_count; + END IF; + + -- 4. Проверка обязательных атрибутов на NULL + SELECT COUNT(*) INTO null_attributes_count + FROM dm.route_performance + WHERE departure_city IS NULL OR arrival_city IS NULL OR airplane_model IS NULL; + + IF null_attributes_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.route_performance обнаружены NULL-значения в обязательных полях (% строк).', null_attributes_count; + END IF; + + -- 5. Проверка: посаженных не может быть больше, чем купивших билет + SELECT COUNT(*) INTO invalid_metrics_count + FROM dm.route_performance + WHERE total_boarded > total_tickets; + + IF invalid_metrics_count > 0 THEN + RAISE EXCEPTION 'DQ Error: В dm.route_performance total_boarded > total_tickets (% строк).', invalid_metrics_count; + END IF; + + RAISE NOTICE 'DQ Success: dm.route_performance успешно прошла все проверки (% строк).', row_count; +END $$; diff --git a/sql/dm/route_performance_load.sql b/sql/dm/route_performance_load.sql new file mode 100644 index 0000000..2e441e8 --- /dev/null +++ b/sql/dm/route_performance_load.sql @@ -0,0 +1,82 @@ +-- Загрузка витрины dm.route_performance: Full Rebuild. +-- +-- Учебные цели: +-- 1. Демонстрация TRUNCATE + INSERT. Таблица маленькая (~1000 строк), +-- пересоздать её с нуля дешевле, чем вычислять дельту. +-- 2. Обработка SCD2: агрегируем факты по бизнес-ключу route_bk, +-- чтобы собрать данные со всех исторических версий маршрута. +-- 3. Денормализация атрибутов из ТЕКУЩЕЙ (актуальной) версии измерения. + +-- Шаг 1: Очистка таблицы (AO Column Store не поддерживает эффективный DELETE/UPDATE). +TRUNCATE dm.route_performance; + +-- Шаг 2: Агрегация всех фактов по бизнес-ключу маршрута. +-- Используем временную таблицу для удобства и читаемости. +CREATE TEMP TABLE tmp_route_metrics ON COMMIT DROP AS +SELECT + r.route_bk, + COUNT(DISTINCT f.flight_id) AS total_flights, + COUNT(*) AS total_tickets, + SUM(CASE WHEN f.is_boarded THEN 1 ELSE 0 END) AS total_boarded, + SUM(f.price) AS total_revenue, + MIN(cal.date_actual) AS first_flight_date, + MAX(cal.date_actual) AS last_flight_date, + -- Доля посаженных пассажиров среди всех проданных билетов маршрута + -- (упрощение: не взвешено по рейсам) + AVG(CASE WHEN f.is_boarded THEN 1 ELSE 0 END::NUMERIC) AS avg_boarding_rate +FROM dds.fact_flight_sales f +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: Объединяем агрегаты с АКТУАЛЬНЫМИ атрибутами маршрута. +INSERT INTO dm.route_performance ( + route_bk, + route_sk, + departure_airport_bk, + departure_city, + arrival_airport_bk, + arrival_city, + airplane_bk, + airplane_model, + total_seats, + total_flights, + total_tickets, + total_boarded, + total_revenue, + avg_ticket_price, + avg_boarding_rate, + avg_load_factor, + first_flight_date, + last_flight_date, + _load_id +) +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, + m.total_flights, + m.total_tickets, + m.total_boarded, + m.total_revenue, + -- Средний чек + (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), + 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; -- Берем только АКТУАЛЬНУЮ версию для денормализации