From 84ec949ff938bcb9acbafc9024297191ae060c31 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 1 Mar 2026 01:54:57 +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=20=D0=BD=D0=B0=20batch-driven=20=D0=B8=D0=BD=D0=BA?= =?UTF-8?q?=D1=80=D0=B5=D0=BC=D0=B5=D0=BD=D1=82=D0=B0=D0=BB=D1=8C=D0=BD?= =?UTF-8?q?=D0=BE=D1=81=D1=82=D1=8C=20=D0=B4=D0=BB=D1=8F=20sales=5Freport?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - жесткая привязка инкремента к логической дате Airflow ({{ ds }}) приводила к пустой витрине при обработке исторических и "опоздавших" (late-arriving) данных. - Что: - изменена фильтрация в скрипте загрузки витрины: теперь динамически определяются даты, затронутые текущим батчем (через _load_id). - обновлены DQ-проверки для валидации только тех дат, которые были изменены в рамках запущенного батча. - в дизайн-документ добавлено описание паттерна работы с late-arriving facts для студентов. - Проверка: - запуск пайплайна "с нуля" за логическую дату 2024-01-01 приводит к корректному расчету агрегатов для исторических данных 2017 года (>8000 строк). --- docs/internal/bookings_dm_design.md | 5 +- sql/dm/sales_report_dq.sql | 74 +++++++++++++++-------------- sql/dm/sales_report_load.sql | 11 ++++- 3 files changed, 50 insertions(+), 40 deletions(-) diff --git a/docs/internal/bookings_dm_design.md b/docs/internal/bookings_dm_design.md index b5dd18f..3bb0fa3 100644 --- a/docs/internal/bookings_dm_design.md +++ b/docs/internal/bookings_dm_design.md @@ -49,11 +49,12 @@ created_at, updated_at, _load_id, _load_ts **Источники**: `fact_flight_sales` JOIN `dim_calendar`, `dim_airports` (x2), `dim_tariffs` -**Загрузка**: инкрементальный UPSERT (UPDATE изменившихся + INSERT новых по ключу) +**Загрузка**: Batch-driven инкрементальный UPSERT (UPDATE изменившихся + INSERT новых по ключу). +*Архитектурный нюанс:* Вместо жесткой фильтрации по дате запуска Airflow (`{{ ds }}`), витрина динамически определяет, какие исторические даты были затронуты в текущем загружаемом батче фактов (по `_load_id`), и пересчитывает агрегаты только для этих дат. Это решает проблему "опоздавших данных" (late-arriving facts). **Хранение**: `DISTRIBUTED BY (flight_date)`, heap (нужен UPDATE) -**Учит**: денормализация измерений, GROUP BY + агрегация, UPSERT по составному ключу, IS DISTINCT FROM +**Учит**: Batch-driven инкрементальность, ограничение радиуса обновления, денормализация измерений, GROUP BY + агрегация, UPSERT по составному ключу, IS DISTINCT FROM --- diff --git a/sql/dm/sales_report_dq.sql b/sql/dm/sales_report_dq.sql index e2e06c4..52bc664 100644 --- a/sql/dm/sales_report_dq.sql +++ b/sql/dm/sales_report_dq.sql @@ -1,14 +1,14 @@ -- DQ для DM витрины sales_report. -- -- Проверки: --- 1. Таблица не пуста (если источник за день {{ ds }} не пуст) --- 2. Нет дублей по составному ключу (flight_date, departure_airport_sk, arrival_airport_sk, tariff_sk) +-- 1. Таблица не пуста (если источник за затронутые батчем дни не пуст) +-- 2. Нет дублей по составному ключу -- 3. Бизнес-инварианты: tickets_sold >= passengers_boarded, boarding_rate BETWEEN 0 AND 1 -- 4. Обязательные поля не NULL -- -- ВАЖНО (Паттерн Data Quality): --- Все проверки выполняются ИНКРЕМЕНТАЛЬНО (только для данных за '{{ ds }}'::date). --- Валидировать терабайты исторических данных при каждой загрузке недопустимо. +-- Проверки выполняются ИНКРЕМЕНТАЛЬНО. Мы находим все даты, измененные текущим батчем, +-- и валидируем витрину только для этих дат. DO $$ DECLARE @@ -18,69 +18,71 @@ DECLARE v_invalid_boarding BIGINT; v_null_required BIGINT; BEGIN - -- Проверка 1: Таблица не пуста при непустом источнике (инкрементально) + -- Создаем временную таблицу с датами, которые мы затронули в этом батче. + 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 }}'; + + -- Если батч пуст (нет новых или измененных фактов), нам нечего валидировать. + IF NOT EXISTS (SELECT 1 FROM tmp_dq_affected_dates) THEN + RAISE NOTICE 'DQ PASSED: Источник не содержит новых данных для витрины в батче %.', '{{ run_id }}'; + RETURN; + END IF; + + -- Проверка 1: Таблица не пуста при непустом источнике (для затронутых дат) SELECT COUNT(*) INTO v_src_count 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; + WHERE cal.date_actual IN (SELECT date_actual FROM tmp_dq_affected_dates); SELECT COUNT(*) INTO v_row_count FROM dm.sales_report - WHERE flight_date = '{{ ds }}'::date; + WHERE flight_date IN (SELECT date_actual FROM tmp_dq_affected_dates); - -- Если источник за день пуст, допускаем пустую витрину - IF v_src_count = 0 THEN - IF v_row_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: fact_flight_sales за {{ ds }} пуст, но dm.sales_report содержит % строк', - v_row_count; - END IF; - RAISE NOTICE 'DQ PASSED: Источник и витрина за {{ ds }} пусты (нет данных для обработки).'; - RETURN; - END IF; - - IF v_row_count = 0 THEN + IF v_src_count > 0 AND v_row_count = 0 THEN RAISE EXCEPTION - 'DQ FAILED: dm.sales_report за {{ ds }} пуста при непустом источнике (% строк в fact_flight_sales)', + 'DQ FAILED: dm.sales_report пуста для дат батча при непустом источнике (% строк в fact_flight_sales)', v_src_count; END IF; - RAISE NOTICE 'DQ INFO: dm.sales_report за {{ ds }} содержит % строк', v_row_count; + RAISE NOTICE 'DQ INFO: dm.sales_report содержит % строк для дат текущего батча', v_row_count; - -- Проверка 2: Нет дублей по составному ключу (в рамках партиции {{ ds }}) + -- Проверка 2: Нет дублей по составному ключу (в рамках затронутых дат) 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 + WHERE flight_date IN (SELECT date_actual FROM tmp_dq_affected_dates) 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 за {{ ds }}', + 'DQ FAILED: Найдено % дублирующихся комбинаций ключа в dm.sales_report', v_dup_count; END IF; - RAISE NOTICE 'DQ PASSED: Дублей по составному ключу за день нет'; + RAISE NOTICE 'DQ PASSED: Дублей по составному ключу для дат батча нет'; - -- Проверка 3: Бизнес-инварианты (в рамках партиции {{ ds }}) + -- Проверка 3: Бизнес-инварианты (в рамках затронутых дат) -- tickets_sold >= passengers_boarded SELECT COUNT(*) INTO v_invalid_boarding FROM dm.sales_report - WHERE flight_date = '{{ ds }}'::date + WHERE flight_date IN (SELECT date_actual FROM tmp_dq_affected_dates) AND tickets_sold < passengers_boarded; IF v_invalid_boarding <> 0 THEN RAISE EXCEPTION - 'DQ FAILED: % строк с tickets_sold < passengers_boarded за {{ ds }}', + 'DQ FAILED: % строк с tickets_sold < passengers_boarded', v_invalid_boarding; END IF; @@ -88,12 +90,12 @@ BEGIN SELECT COUNT(*) INTO v_invalid_boarding FROM dm.sales_report - WHERE flight_date = '{{ ds }}'::date + WHERE flight_date IN (SELECT date_actual FROM tmp_dq_affected_dates) AND (boarding_rate < 0 OR boarding_rate > 1); IF v_invalid_boarding <> 0 THEN RAISE EXCEPTION - 'DQ FAILED: % строк с boarding_rate вне диапазона [0, 1] за {{ ds }}', + 'DQ FAILED: % строк с boarding_rate вне диапазона [0, 1]', v_invalid_boarding; END IF; @@ -101,22 +103,22 @@ BEGIN SELECT COUNT(*) INTO v_invalid_boarding FROM dm.sales_report - WHERE flight_date = '{{ ds }}'::date + WHERE flight_date IN (SELECT date_actual FROM tmp_dq_affected_dates) AND total_revenue < 0; IF v_invalid_boarding <> 0 THEN RAISE EXCEPTION - 'DQ FAILED: % строк с отрицательной total_revenue за {{ ds }}', + 'DQ FAILED: % строк с отрицательной total_revenue', v_invalid_boarding; END IF; RAISE NOTICE 'DQ PASSED: Бизнес-инварианты соблюдены'; - -- Проверка 4: Обязательные поля не NULL (в рамках партиции {{ ds }}) + -- Проверка 4: Обязательные поля не NULL (в рамках затронутых дат) SELECT COUNT(*) INTO v_null_required FROM dm.sales_report - WHERE flight_date = '{{ ds }}'::date + WHERE flight_date IN (SELECT date_actual FROM tmp_dq_affected_dates) AND ( departure_airport_sk IS NULL OR arrival_airport_sk IS NULL @@ -128,12 +130,12 @@ BEGIN IF v_null_required <> 0 THEN RAISE EXCEPTION - 'DQ FAILED: % строк с NULL в обязательных полях за {{ ds }}', + 'DQ FAILED: % строк с NULL в обязательных полях', v_null_required; END IF; RAISE NOTICE 'DQ PASSED: Обязательные поля заполнены'; -- Итог - RAISE NOTICE 'DQ COMPLETE: dm.sales_report прошла все проверки за {{ ds }}'; + RAISE NOTICE 'DQ COMPLETE: dm.sales_report прошла все проверки для текущего батча'; END $$; diff --git a/sql/dm/sales_report_load.sql b/sql/dm/sales_report_load.sql index b4439e8..a1cfd6c 100644 --- a/sql/dm/sales_report_load.sql +++ b/sql/dm/sales_report_load.sql @@ -2,7 +2,8 @@ -- -- Паттерны для студентов: -- - Использование TEMP TABLE: агрегация считается только ОДИН раз (канон для MPP) --- - Инкрементальность: мы читаем факты ТОЛЬКО за день запуска {{ ds }}, избегая Full Scan +-- - Инкрементальность (Batch-driven): мы читаем факты ТОЛЬКО за те дни, которые были затронуты +-- в текущем батче загрузки (по `_load_id`), что решает проблему late-arriving facts. -- - Денормализация: города и названия тарифов тащим в витрину -- - UPSERT: UPDATE изменившихся + INSERT новых (нужен heap для UPDATE) -- - IS DISTINCT FROM для корректного сравнения NULL @@ -50,7 +51,13 @@ 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 -- ИНКРЕМЕНТАЛЬНЫЙ ФИЛЬТР: читаем только нужный день! +WHERE cal.date_actual IN ( + -- ИНКРЕМЕНТАЛЬНЫЙ ФИЛЬТР: Выбираем только те даты, факты за которые обновились в текущем запуске + 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 }}' +) GROUP BY cal.date_actual, dep.airport_sk, dep.city, dep.airport_bk,