refactor(dm): исправление архитектуры витрины sales_report
- Зачем:
- исходная реализация содержала критические ошибки MPP (Load/Processing Skew) и Full Scan.
- требуется демонстрация студентам эталонного инкрементального UPSERT.
- Что:
- изменен ключ дистрибуции с flight_date на (departure_airport_sk, arrival_airport_sk).
- внедрен каноничный UPSERT через TEMP TABLE (агрегация выполняется один раз).
- добавлен инкрементальный фильтр по {{ ds }} для предотвращения Full Scan dds.fact_flight_sales.
- проверки DQ переведены в инкрементальный режим (валидация только текущего батча).
- Проверка:
- airflow dags test bookings_to_gp_dm 2026-02-28.
This commit is contained in:
@@ -1,13 +1,66 @@
|
||||
-- Загрузка DM витрины sales_report: инкрементальный UPSERT.
|
||||
--
|
||||
-- Паттерны для студентов:
|
||||
-- - GROUP BY агрегация фактов перед JOIN с измерениями
|
||||
-- - Использование TEMP TABLE: агрегация считается только ОДИН раз (канон для MPP)
|
||||
-- - Инкрементальность: мы читаем факты ТОЛЬКО за день запуска {{ ds }}, избегая Full Scan
|
||||
-- - Денормализация: города и названия тарифов тащим в витрину
|
||||
-- - UPSERT: UPDATE изменившихся + INSERT новых (нужен heap для UPDATE)
|
||||
-- - IS DISTINCT FROM для корректного сравнения NULL
|
||||
|
||||
-- Statement 1: UPDATE существующих строк.
|
||||
-- Обновляем метрики и денормализованные атрибуты (на случай изменений в справочниках).
|
||||
-- Шаг 1: Считаем агрегаты для дельты (текущего батча) и кладем во временную таблицу.
|
||||
-- ON COMMIT DROP гарантирует, что таблица исчезнет после завершения транзакции (Airflow сессии).
|
||||
CREATE TEMP TABLE tmp_sales_report_delta ON COMMIT DROP AS
|
||||
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
|
||||
WHERE cal.date_actual = '{{ ds }}'::date -- ИНКРЕМЕНТАЛЬНЫЙ ФИЛЬТР: читаем только нужный день!
|
||||
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
|
||||
DISTRIBUTED BY (departure_airport_sk, arrival_airport_sk);
|
||||
|
||||
-- Шаг 2: UPDATE существующих строк.
|
||||
-- Обновляем метрики и денормализованные атрибуты.
|
||||
UPDATE dm.sales_report AS tgt
|
||||
SET
|
||||
departure_city = src.departure_city,
|
||||
@@ -26,55 +79,7 @@ SET
|
||||
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
|
||||
FROM tmp_sales_report_delta 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
|
||||
@@ -89,7 +94,7 @@ WHERE tgt.flight_date = src.flight_date
|
||||
OR tgt.arrival_city IS DISTINCT FROM src.arrival_city
|
||||
);
|
||||
|
||||
-- Statement 2: INSERT новых строк (те, которых нет по составному ключу).
|
||||
-- Шаг 3: INSERT новых строк (те, которых нет по составному ключу).
|
||||
INSERT INTO dm.sales_report (
|
||||
flight_date,
|
||||
departure_airport_sk,
|
||||
@@ -133,51 +138,7 @@ SELECT
|
||||
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
|
||||
FROM tmp_sales_report_delta AS src
|
||||
WHERE NOT EXISTS (
|
||||
SELECT 1
|
||||
FROM dm.sales_report AS tgt
|
||||
|
||||
Reference in New Issue
Block a user