From a8cce6cd6beb5911b796b1d2ab3dd2e0b2b9e1b8 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 28 Feb 2026 23:56:31 +0300 Subject: [PATCH] =?UTF-8?q?refactor(dm):=20=D0=B8=D1=81=D0=BF=D1=80=D0=B0?= =?UTF-8?q?=D0=B2=D0=BB=D0=B5=D0=BD=D0=B8=D0=B5=20=D0=B0=D1=80=D1=85=D0=B8?= =?UTF-8?q?=D1=82=D0=B5=D0=BA=D1=82=D1=83=D1=80=D1=8B=20=D0=B2=D0=B8=D1=82?= =?UTF-8?q?=D1=80=D0=B8=D0=BD=D1=8B=20sales=5Freport?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - исходная реализация содержала критические ошибки 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. --- sql/dm/sales_report_ddl.sql | 11 ++- sql/dm/sales_report_dq.sql | 72 +++++++++------- sql/dm/sales_report_load.sql | 157 +++++++++++++---------------------- 3 files changed, 112 insertions(+), 128 deletions(-) diff --git a/sql/dm/sales_report_ddl.sql b/sql/dm/sales_report_ddl.sql index ebab2aa..3ce6db0 100644 --- a/sql/dm/sales_report_ddl.sql +++ b/sql/dm/sales_report_ddl.sql @@ -7,6 +7,15 @@ -- - Денормализация измерений (города, аэропорты, тарифы) для удобства аналитики -- - Служебные поля календаря (day_of_week, day_name, is_weekend) -- - Heap-таблица с UPDATE (нужен для UPSERT) +-- +-- ВНИМАНИЕ (Антипаттерн распределения): +-- Никогда не распределяйте таблицы по дате (DISTRIBUTED BY flight_date) в MPP-системах! +-- 1. Load Skew: При инкрементальной загрузке весь батч за один день запишется на ОДИН сегмент, +-- а остальные будут простаивать. Кластер превратится в одиночный сервер. +-- 2. Processing Skew: Запросы аналитиков за конкретный день будут читаться только с одного сегмента. +-- +-- ПРАВИЛЬНЫЙ ВЫБОР: Распределение по полям с высокой кардинальностью (направления). +-- В нашем случае это комбинация departure_airport_sk и arrival_airport_sk. CREATE SCHEMA IF NOT EXISTS dm; @@ -44,7 +53,7 @@ CREATE TABLE IF NOT EXISTS dm.sales_report ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) -DISTRIBUTED BY (flight_date); +DISTRIBUTED BY (departure_airport_sk, arrival_airport_sk); -- Комментарии для документирования COMMENT ON TABLE dm.sales_report IS diff --git a/sql/dm/sales_report_dq.sql b/sql/dm/sales_report_dq.sql index f8b35f7..e2e06c4 100644 --- a/sql/dm/sales_report_dq.sql +++ b/sql/dm/sales_report_dq.sql @@ -1,10 +1,14 @@ -- DQ для DM витрины sales_report. -- -- Проверки: --- 1. Таблица не пуста (если источник не пуст) +-- 1. Таблица не пуста (если источник за день {{ ds }} не пуст) -- 2. Нет дублей по составному ключу (flight_date, departure_airport_sk, arrival_airport_sk, tariff_sk) -- 3. Бизнес-инварианты: tickets_sold >= passengers_boarded, boarding_rate BETWEEN 0 AND 1 -- 4. Обязательные поля не NULL +-- +-- ВАЖНО (Паттерн Data Quality): +-- Все проверки выполняются ИНКРЕМЕНТАЛЬНО (только для данных за '{{ ds }}'::date). +-- Валидировать терабайты исторических данных при каждой загрузке недопустимо. DO $$ DECLARE @@ -14,63 +18,69 @@ DECLARE v_invalid_boarding BIGINT; v_null_required BIGINT; BEGIN - -- Проверка 1: Таблица не пуста при непустом источнике + -- Проверка 1: Таблица не пуста при непустом источнике (инкрементально) SELECT COUNT(*) INTO v_src_count - FROM dds.fact_flight_sales; + FROM dds.fact_flight_sales AS f + JOIN dds.dim_calendar AS cal + ON cal.calendar_sk = f.calendar_sk + WHERE cal.date_actual = '{{ ds }}'::date; SELECT COUNT(*) INTO v_row_count - FROM dm.sales_report; + FROM dm.sales_report + WHERE flight_date = '{{ ds }}'::date; - -- Если источник пуст, допускаем пустую витрину (инкрементальное окно без данных) + -- Если источник за день пуст, допускаем пустую витрину IF v_src_count = 0 THEN IF v_row_count <> 0 THEN RAISE EXCEPTION - 'DQ FAILED: fact_flight_sales пуст, но dm.sales_report содержит % строк', + 'DQ FAILED: fact_flight_sales за {{ ds }} пуст, но dm.sales_report содержит % строк', v_row_count; END IF; - RAISE NOTICE 'DQ PASSED: Источник и витрина пусты (нет данных для обработки).'; + RAISE NOTICE 'DQ PASSED: Источник и витрина за {{ ds }} пусты (нет данных для обработки).'; RETURN; END IF; IF v_row_count = 0 THEN RAISE EXCEPTION - 'DQ FAILED: dm.sales_report пуста при непустом источнике (% строк в fact_flight_sales)', + 'DQ FAILED: dm.sales_report за {{ ds }} пуста при непустом источнике (% строк в fact_flight_sales)', v_src_count; END IF; - RAISE NOTICE 'DQ INFO: dm.sales_report содержит % строк', v_row_count; + RAISE NOTICE 'DQ INFO: dm.sales_report за {{ ds }} содержит % строк', v_row_count; - -- Проверка 2: Нет дублей по составному ключу + -- Проверка 2: Нет дублей по составному ключу (в рамках партиции {{ ds }}) SELECT COUNT(*) INTO v_dup_count FROM ( SELECT flight_date, departure_airport_sk, arrival_airport_sk, tariff_sk FROM dm.sales_report + WHERE flight_date = '{{ ds }}'::date 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', + 'DQ FAILED: Найдено % дублирующихся комбинаций ключа в dm.sales_report за {{ ds }}', v_dup_count; END IF; - RAISE NOTICE 'DQ PASSED: Дублей по составному ключу нет'; + RAISE NOTICE 'DQ PASSED: Дублей по составному ключу за день нет'; - -- Проверка 3: Бизнес-инварианты + -- Проверка 3: Бизнес-инварианты (в рамках партиции {{ ds }}) -- tickets_sold >= passengers_boarded SELECT COUNT(*) INTO v_invalid_boarding FROM dm.sales_report - WHERE tickets_sold < passengers_boarded; + WHERE flight_date = '{{ ds }}'::date + AND tickets_sold < passengers_boarded; IF v_invalid_boarding <> 0 THEN RAISE EXCEPTION - 'DQ FAILED: % строк с tickets_sold < passengers_boarded', + 'DQ FAILED: % строк с tickets_sold < passengers_boarded за {{ ds }}', v_invalid_boarding; END IF; @@ -78,11 +88,12 @@ BEGIN SELECT COUNT(*) INTO v_invalid_boarding FROM dm.sales_report - WHERE boarding_rate < 0 OR boarding_rate > 1; + WHERE flight_date = '{{ ds }}'::date + AND (boarding_rate < 0 OR boarding_rate > 1); IF v_invalid_boarding <> 0 THEN RAISE EXCEPTION - 'DQ FAILED: % строк с boarding_rate вне диапазона [0, 1]', + 'DQ FAILED: % строк с boarding_rate вне диапазона [0, 1] за {{ ds }}', v_invalid_boarding; END IF; @@ -90,36 +101,39 @@ BEGIN SELECT COUNT(*) INTO v_invalid_boarding FROM dm.sales_report - WHERE total_revenue < 0; + WHERE flight_date = '{{ ds }}'::date + AND total_revenue < 0; IF v_invalid_boarding <> 0 THEN RAISE EXCEPTION - 'DQ FAILED: % строк с отрицательной total_revenue', + 'DQ FAILED: % строк с отрицательной total_revenue за {{ ds }}', v_invalid_boarding; END IF; RAISE NOTICE 'DQ PASSED: Бизнес-инварианты соблюдены'; - -- Проверка 4: Обязательные поля не NULL + -- Проверка 4: Обязательные поля не NULL (в рамках партиции {{ ds }}) 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; + WHERE flight_date = '{{ ds }}'::date + AND ( + 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 в обязательных полях', + 'DQ FAILED: % строк с NULL в обязательных полях за {{ ds }}', v_null_required; END IF; RAISE NOTICE 'DQ PASSED: Обязательные поля заполнены'; -- Итог - RAISE NOTICE 'DQ COMPLETE: dm.sales_report прошла все проверки'; + RAISE NOTICE 'DQ COMPLETE: dm.sales_report прошла все проверки за {{ ds }}'; END $$; diff --git a/sql/dm/sales_report_load.sql b/sql/dm/sales_report_load.sql index 6f6d284..b4439e8 100644 --- a/sql/dm/sales_report_load.sql +++ b/sql/dm/sales_report_load.sql @@ -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