fix(ods): исправлена загрузка flights для ссылок из segments
- Зачем:
- dq_ods_segments падал на непустых батчах из-за orphan flight_id в ods.segments.
- Что:
- доработан sql/ods/flights_load.sql: добавлено добирание рейсов из истории stg.flights для flight_id из stg.segments текущего batch.
- добавлен контрактный тест в tests/test_ods_sql_contract.py на покрытие flight_id из segments.
- обновлена документация DAG в docs/bookings_to_gp_ods.md.
- Проверка:
- make test.
- airflow dags trigger bookings_to_gp_ods -c '{"stg_batch_id":"manual__2026-01-18T18:47:18.316091+00:00"}'.
This commit is contained in:
@@ -15,6 +15,9 @@
|
||||
- На загрузке использует дедупликацию внутри батча + UPSERT (SCD1).
|
||||
- Для snapshot-справочников (`airports`, `airplanes`, `routes`, `seats`) дополнительно
|
||||
синхронизирует ключи (удаляет из ODS записи, отсутствующие в выбранном STG-батче).
|
||||
- Для `flights` дополнительно добирает рейсы из истории `stg.flights`, если на них
|
||||
ссылаются `stg.segments` выбранного батча (чтобы сохранить ссылочную целостность
|
||||
`segments.flight_id -> flights.flight_id`).
|
||||
|
||||
## Что должно быть готово перед запуском
|
||||
|
||||
|
||||
+122
-12
@@ -1,7 +1,17 @@
|
||||
-- Загрузка 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,
|
||||
s.route_no,
|
||||
@@ -11,12 +21,59 @@ WITH src AS (
|
||||
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.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,7 +99,15 @@ 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,
|
||||
s.route_no,
|
||||
@@ -52,12 +117,57 @@ WITH src AS (
|
||||
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.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,
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user