From f2f0468cb190a3286e72fa77d449cb7774491e4f Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 1 Mar 2026 17:09:55 +0300 Subject: [PATCH] =?UTF-8?q?docs(dm):=20=D0=BE=D0=BF=D0=B8=D1=81=D0=B0?= =?UTF-8?q?=D0=BD=D0=B8=D0=B5=20=D0=BF=D0=B0=D1=82=D1=82=D0=B5=D1=80=D0=BD?= =?UTF-8?q?=D0=BE=D0=B2=20HWM=20=D0=B8=20TEMP=20TABLE=20=D0=B2=20=D0=B4?= =?UTF-8?q?=D0=B8=D0=B7=D0=B0=D0=B9=D0=BD-=D0=B4=D0=BE=D0=BA=D1=83=D0=BC?= =?UTF-8?q?=D0=B5=D0=BD=D1=82=D0=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - дизайн-документ описывал устаревший batch-driven подход и не фиксировал паттерн TEMP TABLE, используемый в реализации. - Что: - заменено описание загрузки всех UPSERT-витрин на HWM-инкрементальность через _load_ts. - добавлена секция «Общий паттерн загрузки UPSERT-витрин» с SQL-скелетом и таблицей применимости. - в секции «Учит» каждой витрины добавлены паттерны HWM и TEMP TABLE. - Проверка: - визуальная проверка docs/internal/bookings_dm_design.md. --- docs/internal/bookings_dm_design.md | 75 +++++++++++++++++++++++++---- 1 file changed, 66 insertions(+), 9 deletions(-) diff --git a/docs/internal/bookings_dm_design.md b/docs/internal/bookings_dm_design.md index 3bb0fa3..5170fcb 100644 --- a/docs/internal/bookings_dm_design.md +++ b/docs/internal/bookings_dm_design.md @@ -49,12 +49,12 @@ created_at, updated_at, _load_id, _load_ts **Источники**: `fact_flight_sales` JOIN `dim_calendar`, `dim_airports` (x2), `dim_tariffs` -**Загрузка**: Batch-driven инкрементальный UPSERT (UPDATE изменившихся + INSERT новых по ключу). -*Архитектурный нюанс:* Вместо жесткой фильтрации по дате запуска Airflow (`{{ ds }}`), витрина динамически определяет, какие исторические даты были затронуты в текущем загружаемом батче фактов (по `_load_id`), и пересчитывает агрегаты только для этих дат. Это решает проблему "опоздавших данных" (late-arriving facts). +**Загрузка**: Инкрементальный UPSERT по High-Water Mark (`_load_ts`). +*Архитектурный нюанс (HWM-паттерн):* Витрина сравнивает свой `MAX(_load_ts)` с `_load_ts` фактов в DDS и пересчитывает агрегаты только для затронутых дат. Это делает конвейер самовосстанавливающимся: если DAG не запускался несколько дней, при следующем запуске витрина автоматически «догонит» всю накопленную дельту. При этом `_load_id` и `_load_ts` самой витрины фиксируют текущий DM-DAG run_id — для аудита и DQ-проверок. **Хранение**: `DISTRIBUTED BY (flight_date)`, heap (нужен UPDATE) -**Учит**: Batch-driven инкрементальность, ограничение радиуса обновления, денормализация измерений, GROUP BY + агрегация, UPSERT по составному ключу, IS DISTINCT FROM +**Учит**: HWM-инкрементальность через `_load_ts`, TEMP TABLE для однократной агрегации (канон MPP), ограничение радиуса обновления, денормализация измерений, GROUP BY + агрегация, UPSERT по составному ключу, IS DISTINCT FROM --- @@ -123,11 +123,11 @@ created_at, updated_at, _load_id, _load_ts **Источники**: `fact_flight_sales` JOIN `dim_passengers`, `dim_tariffs`, `dim_calendar` -**Загрузка**: UPSERT (664K пассажиров — full rebuild дорогой) +**Загрузка**: Инкрементальный UPSERT по HWM (`_load_ts`). Из дельты фактов определяем затронутых `passenger_sk`, пересчитываем агрегаты только для них. (664K пассажиров — full rebuild дорогой.) **Хранение**: `DISTRIBUTED BY (passenger_sk)`, heap -**Учит**: DISTINCT ON / ROW_NUMBER для "самого частого", COUNT(DISTINCT) по нескольким полям, RFM-подобные метрики, UPSERT на большой таблице +**Учит**: HWM-инкрементальность, TEMP TABLE для однократной агрегации, DISTINCT ON / ROW_NUMBER для "самого частого", COUNT(DISTINCT) по нескольким полям, RFM-подобные метрики, UPSERT на большой таблице --- @@ -176,11 +176,11 @@ SELECT airport_sk, date_actual, FROM traffic GROUP BY ... ``` -**Загрузка**: UPSERT по (traffic_date, airport_sk) +**Загрузка**: Инкрементальный UPSERT по HWM (`_load_ts`). Из дельты фактов определяем затронутые `(traffic_date, airport_sk)`, пересчитываем агрегаты только для них. **Хранение**: `DISTRIBUTED BY (traffic_date)`, heap -**Учит**: dual-role dimension join (UNION ALL), conditional aggregation (CASE WHEN + SUM), паттерн "unpivot → aggregate" +**Учит**: HWM-инкрементальность, TEMP TABLE для однократной агрегации, dual-role dimension join (UNION ALL), conditional aggregation (CASE WHEN + SUM), паттерн "unpivot → aggregate" --- @@ -213,11 +213,11 @@ created_at, updated_at, _load_id, _load_ts **Источники**: `fact_flight_sales` JOIN `dim_calendar`, `dim_airplanes` -**Загрузка**: UPSERT по (year_actual, month_actual, airplane_sk) +**Загрузка**: Инкрементальный UPSERT по HWM (`_load_ts`). Из дельты фактов определяем затронутые `(year_actual, month_actual, airplane_sk)`, пересчитываем агрегаты только для них. **Хранение**: `DISTRIBUTED BY (year_actual)`, heap -**Учит**: двухуровневая агрегация (сначала по рейсу для load factor, потом по месяцу), NULLIF для деления, COUNT(DISTINCT) на нескольких полях, executive-дашборд +**Учит**: HWM-инкрементальность, TEMP TABLE для однократной агрегации, двухуровневая агрегация (сначала по рейсу для load factor, потом по месяцу), NULLIF для деления, COUNT(DISTINCT) на нескольких полях, executive-дашборд --- @@ -314,6 +314,62 @@ start_dm ──>> load_dm_passenger_loyalty → dq_dm_passenger_loyalty --- +## Общий паттерн загрузки UPSERT-витрин (TEMP TABLE + HWM) + +Все инкрементальные витрины (кроме `route_performance` — full rebuild) строятся по единому скелету: + +```sql +-- Шаг 1: Агрегация дельты во временную таблицу. +-- Тяжёлый SELECT с JOIN-ами выполняется ОДИН раз — канон для MPP (Greenplum). +-- Без TEMP TABLE пришлось бы дублировать тот же SELECT в UPDATE и INSERT, +-- что означает двойной скан таблицы фактов. +CREATE TEMP TABLE tmp__delta ON COMMIT DROP AS +SELECT + , + , + +FROM dds.fact_flight_sales AS f +JOIN ... +WHERE IN ( + -- HWM-фильтр: берём только зёрна, затронутые новыми фактами + SELECT DISTINCT + FROM dds.fact_flight_sales AS f_sq + JOIN ... + WHERE f_sq._load_ts > ( + SELECT COALESCE(MAX(_load_ts), '1900-01-01'::TIMESTAMP) + FROM dm. + ) +) +GROUP BY , ; + +-- Шаг 2: UPDATE существующих строк (только если что-то изменилось). +UPDATE dm. AS tgt +SET = src., + updated_at = now(), + _load_id = '{{ run_id }}', + _load_ts = now() +FROM tmp__delta AS src +WHERE tgt. = src. + AND ( IS DISTINCT FROM ...); + +-- Шаг 3: INSERT новых строк. +INSERT INTO dm. (...) +SELECT ... FROM tmp__delta AS src +WHERE NOT EXISTS ( + SELECT 1 FROM dm. AS tgt + WHERE tgt. = src. +); +``` + +**Применимость по витринам:** +- `sales_report` — TEMP TABLE + HWM UPSERT (эталон, уже реализован) +- `passenger_loyalty` — TEMP TABLE + HWM UPSERT (затронутые `passenger_sk`) +- `airport_traffic` — TEMP TABLE + HWM UPSERT (затронутые `(traffic_date, airport_sk)`) +- `monthly_overview` — TEMP TABLE + HWM UPSERT (затронутые `(year_actual, month_actual, airplane_sk)`) +- `route_performance` — **не использует** (full rebuild: TRUNCATE + INSERT, один проход) + +--- + ## DQ-проверки (общий паттерн для всех витрин) PL/pgSQL `DO $$` блоки (как в DDS): @@ -332,6 +388,7 @@ PL/pgSQL `DO $$` блоки (как в DDS): | ETL DAG (PostgresOperator, зависимости) | `airflow/dags/bookings_to_gp_dds.py` | | DDL DAG (линейная цепочка) | `airflow/dags/bookings_dds_ddl.py` | | UPSERT SQL (UPDATE + INSERT + CTE) | `sql/dds/fact_flight_sales_load.sql` | +| HWM-инкремент (UPSERT по _load_ts) | `sql/dm/sales_report_load.sql` | | DQ PL/pgSQL (RAISE EXCEPTION/NOTICE) | `sql/dds/fact_flight_sales_dq.sql` | | DDL (CREATE TABLE IF NOT EXISTS) | `sql/dds/dim_airports_ddl.sql` | | Smoke-тесты DAG | `tests/test_dags_smoke.py` |