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,