refactor(dm): переход на batch-driven инкрементальность для sales_report

- Зачем:
  - жесткая привязка инкремента к логической дате Airflow ({{ ds }}) приводила к пустой витрине при обработке исторических и "опоздавших" (late-arriving) данных.
- Что:
  - изменена фильтрация в скрипте загрузки витрины: теперь динамически определяются даты, затронутые текущим батчем (через _load_id).
  - обновлены DQ-проверки для валидации только тех дат, которые были изменены в рамках запущенного батча.
  - в дизайн-документ добавлено описание паттерна работы с late-arriving facts для студентов.
- Проверка:
  - запуск пайплайна "с нуля" за логическую дату 2024-01-01 приводит к корректному расчету агрегатов для исторических данных 2017 года (>8000 строк).
This commit is contained in:
2026-03-01 01:54:57 +03:00
parent a8cce6cd6b
commit 84ec949ff9
3 changed files with 50 additions and 40 deletions
+3 -2
View File
@@ -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
---
+38 -36
View File
@@ -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
IF v_src_count > 0 AND 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
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 $$;
+9 -2
View File
@@ -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,