refactor(dm): переход sales_report на HWM-инкрементальность через _load_ts
- Зачем:
- фильтр `_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.
This commit is contained in:
@@ -7,8 +7,10 @@
|
|||||||
-- 4. Обязательные поля не NULL
|
-- 4. Обязательные поля не NULL
|
||||||
--
|
--
|
||||||
-- ВАЖНО (Паттерн Data Quality):
|
-- ВАЖНО (Паттерн Data Quality):
|
||||||
-- Проверки выполняются ИНКРЕМЕНТАЛЬНО. Мы находим все даты, измененные текущим батчем,
|
-- Проверки выполняются ИНКРЕМЕНТАЛЬНО. Мы спрашиваем саму витрину, какие даты
|
||||||
-- и валидируем витрину только для этих дат.
|
-- были обновлены текущим запуском DAG-а (по _load_id = run_id), и валидируем
|
||||||
|
-- только эти даты. Это корректно, потому что на шаге load мы записали
|
||||||
|
-- текущий run_id DM-DAG-а в обновлённые строки витрины.
|
||||||
|
|
||||||
DO $$
|
DO $$
|
||||||
DECLARE
|
DECLARE
|
||||||
@@ -18,16 +20,17 @@ DECLARE
|
|||||||
v_invalid_boarding BIGINT;
|
v_invalid_boarding BIGINT;
|
||||||
v_null_required BIGINT;
|
v_null_required BIGINT;
|
||||||
BEGIN
|
BEGIN
|
||||||
-- Создаем временную таблицу с датами, которые мы затронули в этом батче.
|
-- Создаем временную таблицу с датами, которые МЫ обновили в этом запуске DAG-а.
|
||||||
|
-- Спрашиваем саму витрину (не DDS!), потому что на шаге load мы записали
|
||||||
|
-- _load_id = '{{ run_id }}' текущего DM-DAG-а в обновлённые строки.
|
||||||
CREATE TEMP TABLE tmp_dq_affected_dates ON COMMIT DROP AS
|
CREATE TEMP TABLE tmp_dq_affected_dates ON COMMIT DROP AS
|
||||||
SELECT DISTINCT cal.date_actual
|
SELECT DISTINCT flight_date AS date_actual
|
||||||
FROM dds.fact_flight_sales AS f
|
FROM dm.sales_report
|
||||||
JOIN dds.dim_calendar AS cal ON f.calendar_sk = cal.calendar_sk
|
WHERE _load_id = '{{ run_id }}';
|
||||||
WHERE f._load_id = '{{ run_id }}';
|
|
||||||
|
|
||||||
-- Если батч пуст (нет новых или измененных фактов), нам нечего валидировать.
|
-- Если батч пуст (ничего не обновлялось), нам нечего валидировать.
|
||||||
IF NOT EXISTS (SELECT 1 FROM tmp_dq_affected_dates) THEN
|
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;
|
RETURN;
|
||||||
END IF;
|
END IF;
|
||||||
|
|
||||||
|
|||||||
@@ -2,13 +2,15 @@
|
|||||||
--
|
--
|
||||||
-- Паттерны для студентов:
|
-- Паттерны для студентов:
|
||||||
-- - Использование TEMP TABLE: агрегация считается только ОДИН раз (канон для MPP)
|
-- - Использование TEMP TABLE: агрегация считается только ОДИН раз (канон для MPP)
|
||||||
-- - Инкрементальность (Batch-driven): мы читаем факты ТОЛЬКО за те дни, которые были затронуты
|
-- - Инкрементальность (HWM — High-Water Mark): витрина сравнивает свой MAX(_load_ts)
|
||||||
-- в текущем батче загрузки (по `_load_id`), что решает проблему late-arriving facts.
|
-- с _load_ts фактов в DDS и пересчитывает агрегаты только для затронутых дат.
|
||||||
|
-- Это делает конвейер самовосстанавливающимся: если DAG не запускался несколько дней,
|
||||||
|
-- при следующем запуске витрина автоматически «догонит» всю накопленную дельту.
|
||||||
-- - Денормализация: города и названия тарифов тащим в витрину
|
-- - Денормализация: города и названия тарифов тащим в витрину
|
||||||
-- - UPSERT: UPDATE изменившихся + INSERT новых (нужен heap для UPDATE)
|
-- - UPSERT: UPDATE изменившихся + INSERT новых (нужен heap для UPDATE)
|
||||||
-- - IS DISTINCT FROM для корректного сравнения NULL
|
-- - IS DISTINCT FROM для корректного сравнения NULL
|
||||||
|
|
||||||
-- Шаг 1: Считаем агрегаты для дельты (текущего батча) и кладем во временную таблицу.
|
-- Шаг 1: Считаем агрегаты для дельты (новых данных с момента последнего обновления витрины).
|
||||||
-- ON COMMIT DROP гарантирует, что таблица исчезнет после завершения транзакции (Airflow сессии).
|
-- ON COMMIT DROP гарантирует, что таблица исчезнет после завершения транзакции (Airflow сессии).
|
||||||
CREATE TEMP TABLE tmp_sales_report_delta ON COMMIT DROP AS
|
CREATE TEMP TABLE tmp_sales_report_delta ON COMMIT DROP AS
|
||||||
SELECT
|
SELECT
|
||||||
@@ -52,11 +54,17 @@ JOIN dds.dim_airports AS arr
|
|||||||
JOIN dds.dim_tariffs AS tar
|
JOIN dds.dim_tariffs AS tar
|
||||||
ON tar.tariff_sk = f.tariff_sk
|
ON tar.tariff_sk = f.tariff_sk
|
||||||
WHERE cal.date_actual IN (
|
WHERE cal.date_actual IN (
|
||||||
-- ИНКРЕМЕНТАЛЬНЫЙ ФИЛЬТР: Выбираем только те даты, факты за которые обновились в текущем запуске
|
-- ИНКРЕМЕНТАЛЬНЫЙ ФИЛЬТР (HWM — High-Water Mark):
|
||||||
|
-- Ищем даты полётов, в которых появились новые или изменённые факты
|
||||||
|
-- с момента последнего обновления витрины (MAX(_load_ts) в dm.sales_report).
|
||||||
|
-- Если витрина пуста — '1900-01-01' заберёт всю историю (первичная загрузка).
|
||||||
SELECT DISTINCT cal_sq.date_actual
|
SELECT DISTINCT cal_sq.date_actual
|
||||||
FROM dds.fact_flight_sales AS f_sq
|
FROM dds.fact_flight_sales AS f_sq
|
||||||
JOIN dds.dim_calendar AS cal_sq ON f_sq.calendar_sk = cal_sq.calendar_sk
|
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
|
GROUP BY
|
||||||
cal.date_actual,
|
cal.date_actual,
|
||||||
|
|||||||
Reference in New Issue
Block a user