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