feat(dm): добавлена эталонная витрина dm.sales_report
- Зачем: - требуется эталонная витрина для обучения паттернам DM-слоя - демонстрация UPSERT-логики с IS DISTINCT FROM для идемпотентности - Что: - DDL: heap-таблица dm.sales_report с 18 полями, DISTRIBUTED BY (flight_date) - Load: UPSERT (UPDATE + INSERT) с JOIN dim_calendar, dim_airports (x2), dim_tariffs, fact_flight_sales - DQ: PL/pgSQL DO $$ с проверками непустоты, уникальности, tickets_sold >= passengers_boarded, boarding_rate BETWEEN 0 AND 1 - DAGs: bookings_dm_ddl (DDL), bookings_to_gp_dm (ETL + DQ с параллельными ветками) - Tests: smoke-тесты для обоих DAG - sql/ddl_gp.sql: добавлен \i dm/sales_report_ddl.sql - Проверка: - make fmt && make lint — passed - make test — 15 passed, 11 skipped - make ddl-gp — DDL applied successfully - airflow dags test bookings_to_gp_dm 2026-02-28T13:00:00 — 4 tasks SUCCESS - 9243 rows loaded, _load_id подставлен корректно (Jinja2 templating works) - DQ checks passed
This commit is contained in:
@@ -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
|
||||||
@@ -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
|
||||||
@@ -57,3 +57,6 @@ FORMAT 'CUSTOM' (formatter='pxfwritable_import');
|
|||||||
\i dds/dim_passengers_ddl.sql
|
\i dds/dim_passengers_ddl.sql
|
||||||
\i dds/dim_routes_ddl.sql
|
\i dds/dim_routes_ddl.sql
|
||||||
\i dds/fact_flight_sales_ddl.sql
|
\i dds/fact_flight_sales_ddl.sql
|
||||||
|
|
||||||
|
-- DDL для DM-слоя (Data Mart).
|
||||||
|
\i 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)';
|
||||||
@@ -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 $$;
|
||||||
@@ -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
|
||||||
|
);
|
||||||
@@ -382,3 +382,36 @@ def test_bookings_to_gp_dds_dag_structure():
|
|||||||
|
|
||||||
# Финальная сводка должна ждать DQ факта.
|
# Финальная сводка должна ждать DQ факта.
|
||||||
_assert_reachable(dag, "dq_dds_fact_flight_sales", "finish_dds_summary")
|
_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")
|
||||||
|
|||||||
Reference in New Issue
Block a user