diff --git a/docs/bookings_to_gp_ods.md b/docs/bookings_to_gp_ods.md index e1db94f..53c7259 100644 --- a/docs/bookings_to_gp_ods.md +++ b/docs/bookings_to_gp_ods.md @@ -15,6 +15,9 @@ - На загрузке использует дедупликацию внутри батча + UPSERT (SCD1). - Для snapshot-справочников (`airports`, `airplanes`, `routes`, `seats`) дополнительно синхронизирует ключи (удаляет из ODS записи, отсутствующие в выбранном STG-батче). +- Для `flights` дополнительно добирает рейсы из истории `stg.flights`, если на них + ссылаются `stg.segments` выбранного батча (чтобы сохранить ссылочную целостность + `segments.flight_id -> flights.flight_id`). ## Что должно быть готово перед запуском diff --git a/sql/ods/flights_load.sql b/sql/ods/flights_load.sql index ff33580..bd30210 100644 --- a/sql/ods/flights_load.sql +++ b/sql/ods/flights_load.sql @@ -1,22 +1,79 @@ -- Загрузка ODS по flights: SCD1 (UPDATE изменившихся + INSERT новых). -- Statement 1: UPDATE существующих строк. -WITH src AS ( +WITH segment_flights AS ( + -- В segments текущего batch могут быть flight_id не только из stg.flights этого batch. + -- Поэтому заранее собираем список flight_id из segments текущего batch. + 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 + AND s.flight_id IS NOT NULL + AND s.flight_id <> '' +), +stg_flights_typed AS ( SELECT - NULLIF(s.flight_id, '')::INTEGER AS flight_id, + NULLIF(s.flight_id, '')::INTEGER AS flight_id, s.route_no, s.status, NULLIF(s.scheduled_departure, '')::TIMESTAMP WITH TIME ZONE AS scheduled_departure, NULLIF(s.scheduled_arrival, '')::TIMESTAMP WITH TIME ZONE AS scheduled_arrival, NULLIF(s.actual_departure, '')::TIMESTAMP WITH TIME ZONE AS actual_departure, NULLIF(s.actual_arrival, '')::TIMESTAMP WITH TIME ZONE AS actual_arrival, - s.src_created_at_ts AS event_ts, - ROW_NUMBER() OVER ( - PARTITION BY s.flight_id - ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC - ) AS rn + s.src_created_at_ts AS event_ts, + s.load_dttm, + s.batch_id FROM stg.flights AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE s.flight_id IS NOT NULL + AND s.flight_id <> '' +), +src_union AS ( + -- 1) Рейсы из текущего batch. + SELECT + f.flight_id, + f.route_no, + f.status, + f.scheduled_departure, + f.scheduled_arrival, + f.actual_departure, + f.actual_arrival, + 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 + + UNION ALL + + -- 2) Добираем из истории STG рейсы, на которые ссылаются segments текущего batch. + SELECT + f.flight_id, + f.route_no, + f.status, + f.scheduled_departure, + f.scheduled_arrival, + f.actual_departure, + f.actual_arrival, + f.event_ts, + f.load_dttm + FROM stg_flights_typed AS f + JOIN segment_flights AS sf + ON sf.flight_id = f.flight_id +), +src AS ( + SELECT + u.flight_id, + u.route_no, + u.status, + u.scheduled_departure, + u.scheduled_arrival, + u.actual_departure, + u.actual_arrival, + u.event_ts, + ROW_NUMBER() OVER ( + PARTITION BY u.flight_id + ORDER BY u.event_ts DESC NULLS LAST, u.load_dttm DESC + ) AS rn + FROM src_union AS u ) UPDATE ods.flights AS o SET route_no = s.route_no, @@ -42,22 +99,75 @@ WHERE s.rn = 1 ); -- Statement 2: INSERT новых строк. -WITH src AS ( +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 + AND s.flight_id IS NOT NULL + AND s.flight_id <> '' +), +stg_flights_typed AS ( SELECT - NULLIF(s.flight_id, '')::INTEGER AS flight_id, + NULLIF(s.flight_id, '')::INTEGER AS flight_id, s.route_no, s.status, NULLIF(s.scheduled_departure, '')::TIMESTAMP WITH TIME ZONE AS scheduled_departure, NULLIF(s.scheduled_arrival, '')::TIMESTAMP WITH TIME ZONE AS scheduled_arrival, NULLIF(s.actual_departure, '')::TIMESTAMP WITH TIME ZONE AS actual_departure, NULLIF(s.actual_arrival, '')::TIMESTAMP WITH TIME ZONE AS actual_arrival, - s.src_created_at_ts AS event_ts, - ROW_NUMBER() OVER ( - PARTITION BY s.flight_id - ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC - ) AS rn + s.src_created_at_ts AS event_ts, + s.load_dttm, + s.batch_id FROM stg.flights AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text + WHERE s.flight_id IS NOT NULL + AND s.flight_id <> '' +), +src_union AS ( + SELECT + f.flight_id, + f.route_no, + f.status, + f.scheduled_departure, + f.scheduled_arrival, + f.actual_departure, + f.actual_arrival, + 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 + + UNION ALL + + SELECT + f.flight_id, + f.route_no, + f.status, + f.scheduled_departure, + f.scheduled_arrival, + f.actual_departure, + f.actual_arrival, + f.event_ts, + f.load_dttm + FROM stg_flights_typed AS f + JOIN segment_flights AS sf + ON sf.flight_id = f.flight_id +), +src AS ( + SELECT + u.flight_id, + u.route_no, + u.status, + u.scheduled_departure, + u.scheduled_arrival, + u.actual_departure, + u.actual_arrival, + u.event_ts, + ROW_NUMBER() OVER ( + PARTITION BY u.flight_id + ORDER BY u.event_ts DESC NULLS LAST, u.load_dttm DESC + ) AS rn + FROM src_union AS u ) INSERT INTO ods.flights ( flight_id, diff --git a/tests/test_ods_sql_contract.py b/tests/test_ods_sql_contract.py index 0471bda..c93596d 100644 --- a/tests/test_ods_sql_contract.py +++ b/tests/test_ods_sql_contract.py @@ -36,3 +36,15 @@ def test_ods_batch_resolver_uses_consistent_snapshot_batches() -> None: assert "INTERSECT" in dag_code assert "candidate_batches" in dag_code + + +def test_flights_load_covers_segment_flight_ids_from_stg_history() -> None: + """ + Загрузка flights должна подтягивать рейсы из истории STG, + если на них ссылаются segments текущего батча. + """ + sql = _read("sql/ods/flights_load.sql") + + assert "segment_flights" in sql + assert "FROM stg.segments AS s" in sql + assert "JOIN segment_flights AS sf" in sql