From 81ff39640203e404db4a81ef72db2c74cc96e5bf Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 1 Mar 2026 17:09:42 +0300 Subject: [PATCH] =?UTF-8?q?refactor(dm):=20=D0=BF=D0=B5=D1=80=D0=B5=D1=85?= =?UTF-8?q?=D0=BE=D0=B4=20sales=5Freport=20=D0=BD=D0=B0=20HWM-=D0=B8=D0=BD?= =?UTF-8?q?=D0=BA=D1=80=D0=B5=D0=BC=D0=B5=D0=BD=D1=82=D0=B0=D0=BB=D1=8C?= =?UTF-8?q?=D0=BD=D0=BE=D1=81=D1=82=D1=8C=20=D1=87=D0=B5=D1=80=D0=B5=D0=B7?= =?UTF-8?q?=20=5Fload=5Fts?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - фильтр `_load_id = '{{ run_id }}'` использовал run_id DM-DAG-а, который не совпадает с run_id DDS-DAG-а, записанным в факты — витрина не находила дельту. - Что: - load: заменён _load_id-фильтр на HWM-подзапрос `_load_ts > MAX(_load_ts)` из dm.sales_report. - dq: источник затронутых дат переключён с DDS на саму витрину (где _load_id уже корректный). - Проверка: - `make test` — smoke-тесты зелёные. - запуск `bookings_to_gp_dm` в Airflow после загрузки DDS. --- sql/dm/sales_report_dq.sql | 21 ++++++++++++--------- sql/dm/sales_report_load.sql | 18 +++++++++++++----- 2 files changed, 25 insertions(+), 14 deletions(-) diff --git a/sql/dm/sales_report_dq.sql b/sql/dm/sales_report_dq.sql index 52bc664..f88e6ba 100644 --- a/sql/dm/sales_report_dq.sql +++ b/sql/dm/sales_report_dq.sql @@ -7,8 +7,10 @@ -- 4. Обязательные поля не NULL -- -- ВАЖНО (Паттерн Data Quality): --- Проверки выполняются ИНКРЕМЕНТАЛЬНО. Мы находим все даты, измененные текущим батчем, --- и валидируем витрину только для этих дат. +-- Проверки выполняются ИНКРЕМЕНТАЛЬНО. Мы спрашиваем саму витрину, какие даты +-- были обновлены текущим запуском DAG-а (по _load_id = run_id), и валидируем +-- только эти даты. Это корректно, потому что на шаге load мы записали +-- текущий run_id DM-DAG-а в обновлённые строки витрины. DO $$ DECLARE @@ -18,16 +20,17 @@ DECLARE v_invalid_boarding BIGINT; v_null_required BIGINT; BEGIN - -- Создаем временную таблицу с датами, которые мы затронули в этом батче. + -- Создаем временную таблицу с датами, которые МЫ обновили в этом запуске DAG-а. + -- Спрашиваем саму витрину (не DDS!), потому что на шаге load мы записали + -- _load_id = '{{ run_id }}' текущего DM-DAG-а в обновлённые строки. CREATE TEMP TABLE tmp_dq_affected_dates ON COMMIT DROP AS - SELECT DISTINCT cal.date_actual - FROM dds.fact_flight_sales AS f - JOIN dds.dim_calendar AS cal ON f.calendar_sk = cal.calendar_sk - WHERE f._load_id = '{{ run_id }}'; + SELECT DISTINCT flight_date AS date_actual + FROM dm.sales_report + WHERE _load_id = '{{ run_id }}'; - -- Если батч пуст (нет новых или измененных фактов), нам нечего валидировать. + -- Если батч пуст (ничего не обновлялось), нам нечего валидировать. IF NOT EXISTS (SELECT 1 FROM tmp_dq_affected_dates) THEN - RAISE NOTICE 'DQ PASSED: Источник не содержит новых данных для витрины в батче %.', '{{ run_id }}'; + RAISE NOTICE 'DQ PASSED: Витрина не обновлялась в этом запуске (%). Новых данных в DDS нет.', '{{ run_id }}'; RETURN; END IF; diff --git a/sql/dm/sales_report_load.sql b/sql/dm/sales_report_load.sql index a1cfd6c..88fe9ce 100644 --- a/sql/dm/sales_report_load.sql +++ b/sql/dm/sales_report_load.sql @@ -2,13 +2,15 @@ -- -- Паттерны для студентов: -- - Использование TEMP TABLE: агрегация считается только ОДИН раз (канон для MPP) --- - Инкрементальность (Batch-driven): мы читаем факты ТОЛЬКО за те дни, которые были затронуты --- в текущем батче загрузки (по `_load_id`), что решает проблему late-arriving facts. +-- - Инкрементальность (HWM — High-Water Mark): витрина сравнивает свой MAX(_load_ts) +-- с _load_ts фактов в DDS и пересчитывает агрегаты только для затронутых дат. +-- Это делает конвейер самовосстанавливающимся: если DAG не запускался несколько дней, +-- при следующем запуске витрина автоматически «догонит» всю накопленную дельту. -- - Денормализация: города и названия тарифов тащим в витрину -- - UPSERT: UPDATE изменившихся + INSERT новых (нужен heap для UPDATE) -- - IS DISTINCT FROM для корректного сравнения NULL --- Шаг 1: Считаем агрегаты для дельты (текущего батча) и кладем во временную таблицу. +-- Шаг 1: Считаем агрегаты для дельты (новых данных с момента последнего обновления витрины). -- ON COMMIT DROP гарантирует, что таблица исчезнет после завершения транзакции (Airflow сессии). CREATE TEMP TABLE tmp_sales_report_delta ON COMMIT DROP AS SELECT @@ -52,11 +54,17 @@ JOIN dds.dim_airports AS arr JOIN dds.dim_tariffs AS tar ON tar.tariff_sk = f.tariff_sk WHERE cal.date_actual IN ( - -- ИНКРЕМЕНТАЛЬНЫЙ ФИЛЬТР: Выбираем только те даты, факты за которые обновились в текущем запуске + -- ИНКРЕМЕНТАЛЬНЫЙ ФИЛЬТР (HWM — High-Water Mark): + -- Ищем даты полётов, в которых появились новые или изменённые факты + -- с момента последнего обновления витрины (MAX(_load_ts) в dm.sales_report). + -- Если витрина пуста — '1900-01-01' заберёт всю историю (первичная загрузка). SELECT DISTINCT cal_sq.date_actual FROM dds.fact_flight_sales AS f_sq JOIN dds.dim_calendar AS cal_sq ON f_sq.calendar_sk = cal_sq.calendar_sk - WHERE f_sq._load_id = '{{ run_id }}' + WHERE f_sq._load_ts > ( + SELECT COALESCE(MAX(_load_ts), '1900-01-01'::TIMESTAMP) + FROM dm.sales_report + ) ) GROUP BY cal.date_actual,