diff --git a/airflow/dags/bookings_dm_ddl.py b/airflow/dags/bookings_dm_ddl.py new file mode 100644 index 0000000..45c6b5b --- /dev/null +++ b/airflow/dags/bookings_dm_ddl.py @@ -0,0 +1,51 @@ +from __future__ import annotations + +""" +Учебный DAG: создаёт/обновляет DM-слой (Data Mart) в Greenplum для домена bookings. + +Запускается вручную перед DAG загрузки `bookings_to_gp_dm` или после изменения DM DDL. +Создаёт 5 DM-витрин: sales_report, route_performance, passenger_loyalty, +airport_traffic, monthly_overview. + +На данном этапе реализована только эталонная витрина sales_report. +Остальные витрины будут добавлены в последующих этапах. +""" + +from datetime import timedelta + +import pendulum +from airflow.providers.postgres.operators.postgres import PostgresOperator + +from airflow import DAG + +GREENPLUM_CONN_ID = "greenplum_conn" + +default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)} + +with DAG( + dag_id="bookings_dm_ddl", + start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), + schedule=None, + catchup=False, + template_searchpath="/sql", + default_args=default_args, + tags=["demo", "greenplum", "ddl", "bookings", "dm"], + description="Учебный DDL DAG: создаёт/обновляет dm.* для bookings", +) as dag: + # Эталонная витрина: sales_report + apply_dm_sales_report_ddl = PostgresOperator( + task_id="apply_dm_sales_report_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/sales_report_ddl.sql", + ) + + # Заглушки для будущих витрин (будут реализованы в этапах 2-5) + # apply_dm_route_performance_ddl = PostgresOperator(...) + # 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 diff --git a/airflow/dags/bookings_to_gp_dm.py b/airflow/dags/bookings_to_gp_dm.py new file mode 100644 index 0000000..1683fa9 --- /dev/null +++ b/airflow/dags/bookings_to_gp_dm.py @@ -0,0 +1,89 @@ +from __future__ import annotations + +""" +Учебный DAG: загрузка из DDS в DM (Greenplum) по домену bookings. + +Ключевая идея: +- DM читает из DDS (Star Schema); +- для каждой витрины выполняем пару задач load -> dq; +- все витрины загружаются параллельно (не зависят друг от друга); +- sales_report использует UPSERT (heap-таблица с UPDATE). + +На данном этапе реализована только эталонная витрина sales_report. +Остальные витрины будут добавлены в последующих этапах. +""" + +from datetime import timedelta +from logging import getLogger + +import pendulum +from airflow.operators.python import PythonOperator +from airflow.providers.postgres.operators.postgres import PostgresOperator + +from airflow import DAG + +GREENPLUM_CONN_ID = "greenplum_conn" + +log = getLogger(__name__) + +default_args = { + "owner": "airflow", + "retries": 1, + "retry_delay": timedelta(seconds=30), +} + + +def _finish_summary() -> None: + """Логирует краткий итог выполнения DM-ветки.""" + log.info("DAG bookings_to_gp_dm завершён. Подробности смотрите в логах задач.") + + +with DAG( + dag_id="bookings_to_gp_dm", + start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), + schedule=None, + catchup=False, + max_active_runs=1, + template_searchpath="/sql", + default_args=default_args, + tags=["demo", "bookings", "greenplum", "dm"], + description="Учебный DAG: загрузка DDS -> DM (Data Mart) + DQ проверки", +) as dag: + # start-задача (для структуры, вдруг понадобятся pre-checks) + start_dm = PythonOperator( + task_id="start_dm", + python_callable=lambda: log.info("Начало загрузки DM-слоя..."), + ) + + # === Эталонная витрина: sales_report === + load_dm_sales_report = PostgresOperator( + task_id="load_dm_sales_report", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/sales_report_load.sql", + ) + + dq_dm_sales_report = PostgresOperator( + task_id="dq_dm_sales_report", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dm/sales_report_dq.sql", + ) + + # === Заглушки для будущих витрин (будут реализованы в этапах 2-5) === + # load_dm_route_performance = PostgresOperator(...) + # dq_dm_route_performance = PostgresOperator(...) + # load_dm_passenger_loyalty = PostgresOperator(...) + # dq_dm_passenger_loyalty = PostgresOperator(...) + # load_dm_airport_traffic = PostgresOperator(...) + # dq_dm_airport_traffic = PostgresOperator(...) + # load_dm_monthly_overview = PostgresOperator(...) + # dq_dm_monthly_overview = PostgresOperator(...) + + # finish-задача + finish_dm_summary = PythonOperator( + task_id="finish_dm_summary", + python_callable=_finish_summary, + ) + + # Зависимости: параллельные ветки 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 diff --git a/sql/ddl_gp.sql b/sql/ddl_gp.sql index a270a0b..5e10504 100644 --- a/sql/ddl_gp.sql +++ b/sql/ddl_gp.sql @@ -57,3 +57,6 @@ FORMAT 'CUSTOM' (formatter='pxfwritable_import'); \i dds/dim_passengers_ddl.sql \i dds/dim_routes_ddl.sql \i dds/fact_flight_sales_ddl.sql + +-- DDL для DM-слоя (Data Mart). +\i dm/sales_report_ddl.sql diff --git a/sql/dm/sales_report_ddl.sql b/sql/dm/sales_report_ddl.sql new file mode 100644 index 0000000..ebab2aa --- /dev/null +++ b/sql/dm/sales_report_ddl.sql @@ -0,0 +1,54 @@ +-- DDL для DM-слоя: витрина sales_report. +-- +-- Бизнес-вопрос: "Какова выручка, кол-во билетов и boarding rate +-- по направлениям/тарифам за каждый день?" +-- +-- Паттерны для студентов: +-- - Денормализация измерений (города, аэропорты, тарифы) для удобства аналитики +-- - Служебные поля календаря (day_of_week, day_name, is_weekend) +-- - Heap-таблица с UPDATE (нужен для UPSERT) + +CREATE SCHEMA IF NOT EXISTS dm; + +CREATE TABLE IF NOT EXISTS dm.sales_report ( + -- Ключ (зерно витрины) + flight_date DATE NOT NULL, + departure_airport_sk INTEGER NOT NULL, + arrival_airport_sk INTEGER NOT NULL, + tariff_sk INTEGER NOT NULL, + + -- Денормализованные атрибуты (для удобства аналитики) + departure_city TEXT NOT NULL, + departure_airport_bk TEXT NOT NULL, + arrival_city TEXT NOT NULL, + arrival_airport_bk TEXT NOT NULL, + fare_conditions TEXT NOT NULL, + + -- Атрибуты календаря + day_of_week INTEGER NOT NULL, + day_name TEXT NOT NULL, + is_weekend BOOLEAN NOT NULL, + + -- Метрики + tickets_sold INTEGER NOT NULL, + passengers_boarded INTEGER NOT NULL, + total_revenue NUMERIC(15,2) NOT NULL, + avg_price NUMERIC(10,2) NOT NULL, + min_price NUMERIC(10,2), + max_price NUMERIC(10,2), + boarding_rate NUMERIC(5,4) NOT NULL, -- boarded / sold + + -- Служебные поля (канон из naming_conventions.md) + 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() +) +DISTRIBUTED BY (flight_date); + +-- Комментарии для документирования +COMMENT ON TABLE dm.sales_report IS + 'Витрина продаж: выручка, билеты и boarding rate по направлениям/тарифам/дням'; + +COMMENT ON COLUMN dm.sales_report.boarding_rate IS + 'Доля пассажиров, прошедших посадку (passengers_boarded / tickets_sold)'; diff --git a/sql/dm/sales_report_dq.sql b/sql/dm/sales_report_dq.sql new file mode 100644 index 0000000..f8b35f7 --- /dev/null +++ b/sql/dm/sales_report_dq.sql @@ -0,0 +1,125 @@ +-- DQ для DM витрины sales_report. +-- +-- Проверки: +-- 1. Таблица не пуста (если источник не пуст) +-- 2. Нет дублей по составному ключу (flight_date, departure_airport_sk, arrival_airport_sk, tariff_sk) +-- 3. Бизнес-инварианты: tickets_sold >= passengers_boarded, boarding_rate BETWEEN 0 AND 1 +-- 4. Обязательные поля не NULL + +DO $$ +DECLARE + v_row_count BIGINT; + v_src_count BIGINT; + v_dup_count BIGINT; + v_invalid_boarding BIGINT; + v_null_required BIGINT; +BEGIN + -- Проверка 1: Таблица не пуста при непустом источнике + SELECT COUNT(*) + INTO v_src_count + FROM dds.fact_flight_sales; + + SELECT COUNT(*) + INTO v_row_count + FROM dm.sales_report; + + -- Если источник пуст, допускаем пустую витрину (инкрементальное окно без данных) + IF v_src_count = 0 THEN + IF v_row_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: fact_flight_sales пуст, но dm.sales_report содержит % строк', + v_row_count; + END IF; + RAISE NOTICE 'DQ PASSED: Источник и витрина пусты (нет данных для обработки).'; + RETURN; + END IF; + + IF v_row_count = 0 THEN + RAISE EXCEPTION + 'DQ FAILED: dm.sales_report пуста при непустом источнике (% строк в fact_flight_sales)', + v_src_count; + END IF; + + RAISE NOTICE 'DQ INFO: dm.sales_report содержит % строк', v_row_count; + + -- Проверка 2: Нет дублей по составному ключу + SELECT COUNT(*) + INTO v_dup_count + FROM ( + SELECT flight_date, departure_airport_sk, arrival_airport_sk, tariff_sk + FROM dm.sales_report + GROUP BY flight_date, departure_airport_sk, arrival_airport_sk, tariff_sk + HAVING COUNT(*) > 1 + ) AS dups; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: Найдено % дублирующихся комбинаций ключа в dm.sales_report', + v_dup_count; + END IF; + + RAISE NOTICE 'DQ PASSED: Дублей по составному ключу нет'; + + -- Проверка 3: Бизнес-инварианты + + -- tickets_sold >= passengers_boarded + SELECT COUNT(*) + INTO v_invalid_boarding + FROM dm.sales_report + WHERE tickets_sold < passengers_boarded; + + IF v_invalid_boarding <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: % строк с tickets_sold < passengers_boarded', + v_invalid_boarding; + END IF; + + -- boarding_rate BETWEEN 0 AND 1 + SELECT COUNT(*) + INTO v_invalid_boarding + FROM dm.sales_report + WHERE boarding_rate < 0 OR boarding_rate > 1; + + IF v_invalid_boarding <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: % строк с boarding_rate вне диапазона [0, 1]', + v_invalid_boarding; + END IF; + + -- total_revenue >= 0 + SELECT COUNT(*) + INTO v_invalid_boarding + FROM dm.sales_report + WHERE total_revenue < 0; + + IF v_invalid_boarding <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: % строк с отрицательной total_revenue', + v_invalid_boarding; + END IF; + + RAISE NOTICE 'DQ PASSED: Бизнес-инварианты соблюдены'; + + -- Проверка 4: Обязательные поля не NULL + SELECT COUNT(*) + INTO v_null_required + FROM dm.sales_report + WHERE flight_date IS NULL + OR departure_airport_sk IS NULL + OR arrival_airport_sk IS NULL + OR tariff_sk IS NULL + OR tickets_sold IS NULL + OR total_revenue IS NULL + OR boarding_rate IS NULL; + + IF v_null_required <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: % строк с NULL в обязательных полях', + v_null_required; + END IF; + + RAISE NOTICE 'DQ PASSED: Обязательные поля заполнены'; + + -- Итог + RAISE NOTICE 'DQ COMPLETE: dm.sales_report прошла все проверки'; +END $$; diff --git a/sql/dm/sales_report_load.sql b/sql/dm/sales_report_load.sql new file mode 100644 index 0000000..6f6d284 --- /dev/null +++ b/sql/dm/sales_report_load.sql @@ -0,0 +1,188 @@ +-- Загрузка DM витрины sales_report: инкрементальный UPSERT. +-- +-- Паттерны для студентов: +-- - GROUP BY агрегация фактов перед JOIN с измерениями +-- - Денормализация: города и названия тарифов тащим в витрину +-- - UPSERT: UPDATE изменившихся + INSERT новых (нужен heap для UPDATE) +-- - IS DISTINCT FROM для корректного сравнения NULL + +-- Statement 1: UPDATE существующих строк. +-- Обновляем метрики и денормализованные атрибуты (на случай изменений в справочниках). +UPDATE dm.sales_report AS tgt +SET + departure_city = src.departure_city, + arrival_city = src.arrival_city, + fare_conditions = src.fare_conditions, + day_of_week = src.day_of_week, + day_name = src.day_name, + is_weekend = src.is_weekend, + tickets_sold = src.tickets_sold, + passengers_boarded = src.passengers_boarded, + total_revenue = src.total_revenue, + avg_price = src.avg_price, + min_price = src.min_price, + max_price = src.max_price, + boarding_rate = src.boarding_rate, + updated_at = now(), + _load_id = '{{ run_id }}', + _load_ts = now() +FROM ( + -- Агрегация фактов по зерну витрины + SELECT + cal.date_actual AS flight_date, + dep.airport_sk AS departure_airport_sk, + arr.airport_sk AS arrival_airport_sk, + tar.tariff_sk, + + -- Денормализованные атрибуты + dep.city AS departure_city, + dep.airport_bk AS departure_airport_bk, + arr.city AS arrival_city, + arr.airport_bk AS arrival_airport_bk, + tar.fare_conditions, + + -- Атрибуты календаря + cal.day_of_week, + cal.day_name, + cal.is_weekend, + + -- Метрики + COUNT(*) AS tickets_sold, + SUM(CASE WHEN f.is_boarded THEN 1 ELSE 0 END) AS passengers_boarded, + SUM(f.price) AS total_revenue, + AVG(f.price) AS avg_price, + MIN(f.price) AS min_price, + MAX(f.price) AS max_price, + -- boarding_rate: делим boarded на sold с защитой от деления на 0 + ROUND( + SUM(CASE WHEN f.is_boarded THEN 1 ELSE 0 END)::NUMERIC / NULLIF(COUNT(*), 0), + 4 + ) AS boarding_rate + + FROM dds.fact_flight_sales AS f + JOIN dds.dim_calendar AS cal + ON cal.calendar_sk = f.calendar_sk + JOIN dds.dim_airports AS dep + ON dep.airport_sk = f.departure_airport_sk + JOIN dds.dim_airports AS arr + ON arr.airport_sk = f.arrival_airport_sk + JOIN dds.dim_tariffs AS tar + ON tar.tariff_sk = f.tariff_sk + GROUP BY + cal.date_actual, + dep.airport_sk, dep.city, dep.airport_bk, + arr.airport_sk, arr.city, arr.airport_bk, + tar.tariff_sk, tar.fare_conditions, + cal.day_of_week, cal.day_name, cal.is_weekend +) AS src +WHERE tgt.flight_date = src.flight_date + AND tgt.departure_airport_sk = src.departure_airport_sk + AND tgt.arrival_airport_sk = src.arrival_airport_sk + AND tgt.tariff_sk = src.tariff_sk + AND ( + -- Обновляем только если что-то реально изменилось + tgt.tickets_sold IS DISTINCT FROM src.tickets_sold + OR tgt.passengers_boarded IS DISTINCT FROM src.passengers_boarded + OR tgt.total_revenue IS DISTINCT FROM src.total_revenue + OR tgt.boarding_rate IS DISTINCT FROM src.boarding_rate + OR tgt.departure_city IS DISTINCT FROM src.departure_city + OR tgt.arrival_city IS DISTINCT FROM src.arrival_city + ); + +-- Statement 2: INSERT новых строк (те, которых нет по составному ключу). +INSERT INTO dm.sales_report ( + flight_date, + departure_airport_sk, + arrival_airport_sk, + tariff_sk, + departure_city, + departure_airport_bk, + arrival_city, + arrival_airport_bk, + fare_conditions, + day_of_week, + day_name, + is_weekend, + tickets_sold, + passengers_boarded, + total_revenue, + avg_price, + min_price, + max_price, + boarding_rate, + _load_id +) +SELECT + src.flight_date, + src.departure_airport_sk, + src.arrival_airport_sk, + src.tariff_sk, + src.departure_city, + src.departure_airport_bk, + src.arrival_city, + src.arrival_airport_bk, + src.fare_conditions, + src.day_of_week, + src.day_name, + src.is_weekend, + src.tickets_sold, + src.passengers_boarded, + src.total_revenue, + src.avg_price, + src.min_price, + src.max_price, + src.boarding_rate, + '{{ run_id }}' AS _load_id +FROM ( + -- Агрегация фактов (тот же CTE, что и в UPDATE) + SELECT + cal.date_actual AS flight_date, + dep.airport_sk AS departure_airport_sk, + arr.airport_sk AS arrival_airport_sk, + tar.tariff_sk, + + dep.city AS departure_city, + dep.airport_bk AS departure_airport_bk, + arr.city AS arrival_city, + arr.airport_bk AS arrival_airport_bk, + tar.fare_conditions, + + cal.day_of_week, + cal.day_name, + cal.is_weekend, + + COUNT(*) AS tickets_sold, + SUM(CASE WHEN f.is_boarded THEN 1 ELSE 0 END) AS passengers_boarded, + SUM(f.price) AS total_revenue, + AVG(f.price) AS avg_price, + MIN(f.price) AS min_price, + MAX(f.price) AS max_price, + ROUND( + SUM(CASE WHEN f.is_boarded THEN 1 ELSE 0 END)::NUMERIC / NULLIF(COUNT(*), 0), + 4 + ) AS boarding_rate + + FROM dds.fact_flight_sales AS f + JOIN dds.dim_calendar AS cal + ON cal.calendar_sk = f.calendar_sk + JOIN dds.dim_airports AS dep + ON dep.airport_sk = f.departure_airport_sk + JOIN dds.dim_airports AS arr + ON arr.airport_sk = f.arrival_airport_sk + JOIN dds.dim_tariffs AS tar + ON tar.tariff_sk = f.tariff_sk + GROUP BY + cal.date_actual, + dep.airport_sk, dep.city, dep.airport_bk, + arr.airport_sk, arr.city, arr.airport_bk, + tar.tariff_sk, tar.fare_conditions, + cal.day_of_week, cal.day_name, cal.is_weekend +) AS src +WHERE NOT EXISTS ( + SELECT 1 + FROM dm.sales_report AS tgt + WHERE tgt.flight_date = src.flight_date + AND tgt.departure_airport_sk = src.departure_airport_sk + AND tgt.arrival_airport_sk = src.arrival_airport_sk + AND tgt.tariff_sk = src.tariff_sk +); diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index 2fbe4f3..8c7bc32 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -382,3 +382,36 @@ def test_bookings_to_gp_dds_dag_structure(): # Финальная сводка должна ждать DQ факта. _assert_reachable(dag, "dq_dds_fact_flight_sales", "finish_dds_summary") + + +def test_bookings_dm_ddl_dag_structure(): + """Проверка структуры DAG bookings_dm_ddl.""" + dag = _load_dag("airflow.dags.bookings_dm_ddl") + + expected_tasks = { + "apply_dm_sales_report_ddl", + } + assert expected_tasks.issubset(dag.task_dict.keys()) + + # На данном этапе только одна витрина + # В будущем: линейная цепочка из 5 задач + + +def test_bookings_to_gp_dm_dag_structure(): + """Проверка структуры DAG bookings_to_gp_dm.""" + dag = _load_dag("airflow.dags.bookings_to_gp_dm") + + expected_tasks = { + "start_dm", + "load_dm_sales_report", + "dq_dm_sales_report", + "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") + + # Проверяем, что start -> load -> dq -> finish + _assert_reachable(dag, "start_dm", "load_dm_sales_report") + _assert_reachable(dag, "dq_dm_sales_report", "finish_dm_summary")