From cfd20328d69559576089e858bbec9634344ea153 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 1 Mar 2026 18:15:04 +0300 Subject: [PATCH] =?UTF-8?q?fix(ods):=20=D0=B8=D0=B7=D0=BC=D0=B5=D0=BD?= =?UTF-8?q?=D0=B5=D0=BD=20=D0=B1=D0=B0=D1=82=D1=87=D0=B5=D0=B2=D1=8B=D0=B9?= =?UTF-8?q?=20=D1=80=D0=B5=D0=B7=D0=BE=D0=BB=D0=B2=D0=B5=D1=80=20=D0=B4?= =?UTF-8?q?=D0=BB=D1=8F=20=D1=82=D1=80=D0=B0=D0=BD=D0=B7=D0=B0=D0=BA=D1=86?= =?UTF-8?q?=D0=B8=D0=BE=D0=BD=D0=BD=D1=8B=D1=85=20=D1=82=D0=B0=D0=B1=D0=BB?= =?UTF-8?q?=D0=B8=D1=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - текущая реализация ODS batch_id теряла данные транзакционных таблиц, если между запусками ODS STG успевал отработать дважды (брался только последний батч). - Что: - изменены скрипты загрузки транзакционных таблиц (bookings, tickets, flights, segments, boarding_passes) для использования паттерна HWM по _load_ts вместо фильтрации по конкретному батчу. - обновлен комментарий в DAG bookings_to_gp_ods, объясняющий разное поведение для справочников и транзакционных данных. - сохранено использование оригинального batch_id из STG для поля _load_id в слое ODS для сквозного трассирования. - Проверка: - запуск пайплайнов и проверка, что все батчи загружаются из STG в ODS без потерь. --- airflow/dags/bookings_to_gp_ods.py | 9 +++++++++ sql/ods/boarding_passes_load.sql | 6 +++--- sql/ods/bookings_load.sql | 10 ++++++---- sql/ods/flights_load.sql | 10 +++++----- sql/ods/segments_load.sql | 6 +++--- sql/ods/tickets_load.sql | 6 +++--- 6 files changed, 29 insertions(+), 18 deletions(-) diff --git a/airflow/dags/bookings_to_gp_ods.py b/airflow/dags/bookings_to_gp_ods.py index 2fb7bbe..1dfc1fe 100644 --- a/airflow/dags/bookings_to_gp_ods.py +++ b/airflow/dags/bookings_to_gp_ods.py @@ -34,6 +34,15 @@ def _resolve_stg_batch_id(**context) -> str: """ Возвращает stg_batch_id из dag_run.conf или вычисляет последний согласованный батч. + ВНИМАНИЕ: Это значение используется ТОЛЬКО для загрузки snapshot-справочников + (airports, airplanes, routes, seats). + + Для инкрементальных транзакционных таблиц (bookings, tickets, flights, segments, + boarding_passes) этот батч НЕ используется. Вместо этого они грузят все новые + записи по HWM: WHERE load_dttm > (SELECT MAX(_load_ts) FROM ods.table). + Это сделано для того, чтобы не потерять инкременты, если STG-DAG запускался + несколько раз до запуска ODS-DAG'а. + Зачем нужна согласованность (INTERSECT по всем справочникам)? Чтобы ODS загружал только те данные, для которых уже приехали ВСЕ связанные справочники. Это защищает от рассинхрона данных, когда часть измерений в текущем diff --git a/sql/ods/boarding_passes_load.sql b/sql/ods/boarding_passes_load.sql index 08d7359..bf60499 100644 --- a/sql/ods/boarding_passes_load.sql +++ b/sql/ods/boarding_passes_load.sql @@ -14,7 +14,7 @@ WITH src AS ( ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC ) AS rn FROM stg.boarding_passes AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.boarding_passes) ) UPDATE ods.boarding_passes AS o SET seat_no = s.seat_no, @@ -48,7 +48,7 @@ WITH src AS ( ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC ) AS rn FROM stg.boarding_passes AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.boarding_passes) ) INSERT INTO ods.boarding_passes ( ticket_no, @@ -67,7 +67,7 @@ SELECT s.boarding_no, s.boarding_time, s.event_ts, - '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + s.batch_id, now() FROM src AS s WHERE s.rn = 1 diff --git a/sql/ods/bookings_load.sql b/sql/ods/bookings_load.sql index 45db3de..43f4f13 100644 --- a/sql/ods/bookings_load.sql +++ b/sql/ods/bookings_load.sql @@ -7,18 +7,19 @@ WITH src AS ( NULLIF(s.book_date, '')::TIMESTAMP WITH TIME ZONE AS book_date, NULLIF(s.total_amount, '')::NUMERIC(10,2) AS total_amount, s.src_created_at_ts AS event_ts, + s.batch_id, ROW_NUMBER() OVER ( PARTITION BY s.book_ref ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC ) AS rn FROM stg.bookings AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.bookings) ) UPDATE ods.bookings AS o SET book_date = s.book_date, total_amount = s.total_amount, event_ts = s.event_ts, - _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + _load_id = s.batch_id, _load_ts = now() FROM src AS s WHERE s.rn = 1 @@ -36,12 +37,13 @@ WITH src AS ( NULLIF(s.book_date, '')::TIMESTAMP WITH TIME ZONE AS book_date, NULLIF(s.total_amount, '')::NUMERIC(10,2) AS total_amount, s.src_created_at_ts AS event_ts, + s.batch_id, ROW_NUMBER() OVER ( PARTITION BY s.book_ref ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC ) AS rn FROM stg.bookings AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.bookings) ) INSERT INTO ods.bookings ( book_ref, @@ -56,7 +58,7 @@ SELECT s.book_date, s.total_amount, s.event_ts, - '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + s.batch_id, now() FROM src AS s WHERE s.rn = 1 diff --git a/sql/ods/flights_load.sql b/sql/ods/flights_load.sql index bd30210..2b3d47e 100644 --- a/sql/ods/flights_load.sql +++ b/sql/ods/flights_load.sql @@ -7,7 +7,7 @@ WITH segment_flights AS ( SELECT DISTINCT NULLIF(s.flight_id, '')::INTEGER AS flight_id FROM stg.segments AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.segments) AND s.flight_id IS NOT NULL AND s.flight_id <> '' ), @@ -40,7 +40,7 @@ src_union AS ( f.event_ts, f.load_dttm FROM stg_flights_typed AS f - WHERE f.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE f.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.flights) UNION ALL @@ -103,7 +103,7 @@ WITH segment_flights AS ( SELECT DISTINCT NULLIF(s.flight_id, '')::INTEGER AS flight_id FROM stg.segments AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.segments) AND s.flight_id IS NOT NULL AND s.flight_id <> '' ), @@ -135,7 +135,7 @@ src_union AS ( f.event_ts, f.load_dttm FROM stg_flights_typed AS f - WHERE f.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE f.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.flights) UNION ALL @@ -190,7 +190,7 @@ SELECT s.actual_departure, s.actual_arrival, s.event_ts, - '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + s.batch_id, now() FROM src AS s WHERE s.rn = 1 diff --git a/sql/ods/segments_load.sql b/sql/ods/segments_load.sql index 50e0e18..4fce77a 100644 --- a/sql/ods/segments_load.sql +++ b/sql/ods/segments_load.sql @@ -13,7 +13,7 @@ WITH src AS ( ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC ) AS rn FROM stg.segments AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.segments) ) UPDATE ods.segments AS o SET fare_conditions = s.fare_conditions, @@ -44,7 +44,7 @@ WITH src AS ( ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC ) AS rn FROM stg.segments AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.segments) ) INSERT INTO ods.segments ( ticket_no, @@ -61,7 +61,7 @@ SELECT s.fare_conditions, s.segment_amount, s.event_ts, - '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + s.batch_id, now() FROM src AS s WHERE s.rn = 1 diff --git a/sql/ods/tickets_load.sql b/sql/ods/tickets_load.sql index b4c214e..18b19c7 100644 --- a/sql/ods/tickets_load.sql +++ b/sql/ods/tickets_load.sql @@ -14,7 +14,7 @@ WITH src AS ( ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC ) AS rn FROM stg.tickets AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.tickets) ) UPDATE ods.tickets AS o SET book_ref = s.book_ref, @@ -49,7 +49,7 @@ WITH src AS ( ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC ) AS rn FROM stg.tickets AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.tickets) ) INSERT INTO ods.tickets ( ticket_no, @@ -68,7 +68,7 @@ SELECT s.passenger_name, s.is_outbound, s.event_ts, - '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, + s.batch_id, now() FROM src AS s WHERE s.rn = 1