From 7b6136b0fa3c2f44b04a3e06cb4d5b600d021fd6 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Thu, 12 Mar 2026 13:11:40 +0300 Subject: [PATCH] =?UTF-8?q?fix(stg):=20=D0=BF=D0=B5=D1=80=D0=B5=D0=B2?= =?UTF-8?q?=D0=B5=D0=B4=D1=91=D0=BD=20boarding=5Fpasses=20=D0=BD=D0=B0=20?= =?UTF-8?q?=D0=B8=D0=BD=D0=BA=D1=80=D0=B5=D0=BC=D0=B5=D0=BD=D1=82=D0=B0?= =?UTF-8?q?=D0=BB=D1=8C=D0=BD=D1=83=D1=8E=20=D0=B7=D0=B0=D0=B3=D1=80=D1=83?= =?UTF-8?q?=D0=B7=D0=BA=D1=83?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - full snapshot boarding_passes при повторных запусках создавал orphan-записи без соответствующих tickets/segments (инкрементальных), DQ корректно падал. - Что: - boarding_passes_load.sql: HWM через book_date (JOIN tickets_ext → bookings_ext), аналогично segments_load.sql. - boarding_passes_dq.sql: подсчёт источника с фильтром по окну инкремента, обработка пустого окна (NOTICE + RETURN). - обновлена документация (4 файла): db_schema, bookings_to_gp_stage, bookings_ods_design, inline-комментарий в DAG. - Проверка: - make test (4 passed), статический ревью Codex CLI (0 замечаний по SQL). Co-Authored-By: Claude Opus 4.6 --- airflow/dags/bookings_to_gp_stage.py | 2 +- docs/bookings_to_gp_stage.md | 4 +-- docs/design/bookings_ods_design.md | 2 +- docs/design/db_schema.md | 2 +- sql/stg/boarding_passes_dq.sql | 51 +++++++++++++++++++++------- sql/stg/boarding_passes_load.sql | 30 ++++++++++++---- 6 files changed, 66 insertions(+), 25 deletions(-) diff --git a/airflow/dags/bookings_to_gp_stage.py b/airflow/dags/bookings_to_gp_stage.py index ad7b8bb..6097b9d 100644 --- a/airflow/dags/bookings_to_gp_stage.py +++ b/airflow/dags/bookings_to_gp_stage.py @@ -162,7 +162,7 @@ with DAG( sql="stg/seats_dq.sql", ) - # Загрузка транзакций (инкремент/full snapshot) + # Загрузка транзакций (инкремент) load_flights_to_stg = PostgresOperator( task_id="load_flights_to_stg", postgres_conn_id=GREENPLUM_CONN_ID, diff --git a/docs/bookings_to_gp_stage.md b/docs/bookings_to_gp_stage.md index 1b7f267..9cbb3e2 100644 --- a/docs/bookings_to_gp_stage.md +++ b/docs/bookings_to_gp_stage.md @@ -13,7 +13,7 @@ - Загружает инкремент в `stg.tickets` через внешнюю таблицу `stg.tickets_ext`, используя PXF. - Запускает DQ‑проверки для `stg.tickets` (количество, ссылочная целостность, обязательные поля). - Загружает справочники (full load): `stg.airports`, `stg.airplanes`, `stg.routes`, `stg.seats` + DQ. -- Загружает транзакции: `stg.flights` (инкремент), `stg.segments` (инкремент), `stg.boarding_passes` (full snapshot) + DQ. +- Загружает транзакции: `stg.flights` (инкремент), `stg.segments` (инкремент), `stg.boarding_passes` (инкремент через tickets/bookings) + DQ. ## Что должно быть готово перед запуском @@ -110,7 +110,7 @@ check_tickets_dq - `load_flights_to_stg` → `check_flights_dq` (инкремент по `scheduled_departure`) — зависит от **routes** - `load_segments_to_stg` → `check_segments_dq` (инкремент по `book_date` через tickets/bookings) — зависит от **flights** -- `load_boarding_passes_to_stg` → `check_boarding_passes_dq` (full snapshot) — зависит от **segments** +- `load_boarding_passes_to_stg` → `check_boarding_passes_dq` (инкремент по `book_date` через tickets/bookings) — зависит от **segments** Ветка `seats` работает параллельно с веткой `routes → flights → segments → boarding_passes`. Обе ветки сходятся на `finish_summary`. diff --git a/docs/design/bookings_ods_design.md b/docs/design/bookings_ods_design.md index 065e845..139382c 100644 --- a/docs/design/bookings_ods_design.md +++ b/docs/design/bookings_ods_design.md @@ -511,7 +511,7 @@ ANALYZE ods.bookings; ### 6.4. Поведение при пустом батче -- для инкрементальных таблиц (`bookings`, `tickets`, `flights`, `segments`) пустой батч допустим; +- для инкрементальных таблиц (`bookings`, `tickets`, `flights`, `segments`, `boarding_passes`) пустой батч допустим; - для snapshot-справочников (`airports`, `airplanes`, `routes`, `seats`) пустой батч считаем ошибкой. --- diff --git a/docs/design/db_schema.md b/docs/design/db_schema.md index ffee180..400290a 100644 --- a/docs/design/db_schema.md +++ b/docs/design/db_schema.md @@ -172,7 +172,7 @@ graph LR | `stg.tickets` | `bookings.tickets` | `ticket_no` | Инкремент (через bookings) | `book_ref` | | `stg.flights` | `bookings.flights` | `flight_id` | Инкремент (scheduled_departure) | `flight_id` | | `stg.segments` | `bookings.segments` | `(ticket_no, flight_id)` | Инкремент (через tickets) | `ticket_no` | -| `stg.boarding_passes` | `bookings.boarding_passes` | `(ticket_no, flight_id)` | Full snapshot | `ticket_no` | +| `stg.boarding_passes` | `bookings.boarding_passes` | `(ticket_no, flight_id)` | Инкремент (через tickets/bookings) | `ticket_no` | | `stg.airports` | `bookings.airports_data` | `airport_code` | Full snapshot | `airport_code` | | `stg.airplanes` | `bookings.airplanes_data` | `airplane_code` | Full snapshot | `airplane_code` | | `stg.routes` | `bookings.routes` | `(route_no, validity)` | Full snapshot | `route_no` | diff --git a/sql/stg/boarding_passes_dq.sql b/sql/stg/boarding_passes_dq.sql index e6cd7a2..b9c43cf 100644 --- a/sql/stg/boarding_passes_dq.sql +++ b/sql/stg/boarding_passes_dq.sql @@ -1,23 +1,48 @@ --- Проверки качества данных для boarding_passes +-- Проверки качества данных для boarding_passes (инкрементальная загрузка). DO $$ DECLARE - v_batch_id TEXT := '{{ run_id }}'::text; - v_src_count BIGINT; - v_stg_count BIGINT; - v_dup_count BIGINT; - v_null_count BIGINT; - v_orphan_ticket_count BIGINT; + v_batch_id TEXT := '{{ run_id }}'::text; + v_prev_ts TIMESTAMP; + v_src_count BIGINT; + v_stg_count BIGINT; + v_dup_count BIGINT; + v_null_count BIGINT; + v_orphan_ticket_count BIGINT; v_orphan_segment_count BIGINT; BEGIN - -- Источник: считаем все строки во внешней таблице + -- Опорная метка: максимум event_ts среди предыдущих батчей + SELECT MAX(event_ts) + INTO v_prev_ts + FROM stg.boarding_passes + WHERE _load_id <> v_batch_id + OR _load_id IS NULL; + + -- Количество в источнике (boarding_passes в том же окне инкремента, что и загрузка) SELECT COUNT(*) INTO v_src_count - FROM stg.boarding_passes_ext; + FROM stg.boarding_passes_ext AS ext + JOIN stg.tickets_ext AS t ON ext.ticket_no = t.ticket_no + JOIN stg.bookings_ext AS b ON t.book_ref = b.book_ref + WHERE b.book_date > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); IF v_src_count = 0 THEN + -- Пустое окно инкремента допустимо: новых данных может не быть. + SELECT COUNT(*) + INTO v_stg_count + FROM stg.boarding_passes + WHERE _load_id = v_batch_id; + + IF v_stg_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: источник boarding_passes_ext за окно инкремента пустой, но в stg.boarding_passes есть строки текущего _load_id (_load_id=%): %', + v_batch_id, + v_stg_count; + END IF; + RAISE NOTICE - 'В источнике boarding_passes_ext нет строк - пропускаем DQ проверки (_load_id=%).', + 'В источнике boarding_passes_ext нет строк для окна инкремента (book_date > %). Пропускаем DQ проверки (_load_id=%).', + COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'), v_batch_id; RETURN; END IF; @@ -30,13 +55,13 @@ BEGIN IF v_src_count <> v_stg_count THEN RAISE EXCEPTION - 'DQ FAILED: несовпадение количества строк. Источник: %, STG: %', + 'DQ FAILED: несовпадение количества строк. Источник (окно инкремента): %, STG (_load_id=%): %', v_src_count, + v_batch_id, v_stg_count; END IF; - -- Проверка на дубликаты (ticket_no, flight_id) - -- Используем md5 от ROW, чтобы избежать коллизий при склейке строк. + -- Проверка на дубликаты (ticket_no, flight_id) в текущем батче SELECT COUNT(*) - COUNT(DISTINCT md5(ROW(ticket_no, flight_id)::text)) INTO v_dup_count FROM stg.boarding_passes AS bp diff --git a/sql/stg/boarding_passes_load.sql b/sql/stg/boarding_passes_load.sql index 8afe640..b46391e 100644 --- a/sql/stg/boarding_passes_load.sql +++ b/sql/stg/boarding_passes_load.sql @@ -1,7 +1,20 @@ --- Загрузка всех строк из stg.boarding_passes_ext в stg.boarding_passes. --- Используем full snapshot: все строки при каждом запуске. --- Используем _load_id для отслеживания загрузки. +-- Загрузка инкремента из stg.boarding_passes_ext в stg.boarding_passes. +-- Инкремент определяется по дате бронирования (book_date из bookings.bookings) +-- через JOIN с таблицами tickets и bookings (аналогично segments_load.sql). +-- +-- Учебный комментарий: boarding_passes привязаны к билетам, а билеты — к бронированиям. +-- Чтобы инкремент boarding_passes совпадал с инкрементом tickets и segments, +-- используем ту же точку отсечения — book_date бронирования. Иначе при повторных +-- запусках full snapshot загрузит все boarding_passes, а tickets/segments — только +-- новые, и DQ обнаружит «сиротские» boarding_passes без соответствующих tickets. +-- CTE для определения максимальной даты загрузки предыдущего батча +WITH max_batch_ts AS ( + SELECT COALESCE(MAX(event_ts), TIMESTAMP '1900-01-01 00:00:00') AS max_ts + FROM stg.boarding_passes + WHERE _load_id <> '{{ run_id }}'::text + OR _load_id IS NULL +) INSERT INTO stg.boarding_passes ( ticket_no, flight_id, @@ -18,12 +31,16 @@ SELECT ext.seat_no, ext.boarding_no::text, ext.boarding_time::text, - now()::timestamp, + b.book_date::timestamp, -- временная метка из бронирования now(), '{{ run_id }}'::text FROM stg.boarding_passes_ext AS ext -WHERE NOT EXISTS ( - -- Идемпотентность: при повторном запуске/ретрае не вставляем повторно те же строки в рамках текущего _load_id. +JOIN stg.tickets_ext AS t ON ext.ticket_no = t.ticket_no +JOIN stg.bookings_ext AS b ON t.book_ref = b.book_ref +CROSS JOIN max_batch_ts AS mb +WHERE b.book_date > mb.max_ts +AND NOT EXISTS ( + -- Идемпотентность: при повторном запуске/ретрае не вставляем повторно те же строки. -- Считаем ключом строки (ticket_no, flight_id). SELECT 1 FROM stg.boarding_passes AS bp @@ -33,5 +50,4 @@ WHERE NOT EXISTS ( ); -- Обновляем статистику для оптимизатора Greenplum --- Это критично для корректной работы оптимизатора и выбора оптимального плана выполнения ANALYZE stg.boarding_passes;