Files
airflow-greenplum/airflow/dags/bookings_to_gp_ods.py
T
ddadminandClaude Opus 4.6 5ca4e6b2ca feat(main): заглушки студенческих файлов, ослабление DQ, адаптация тестов
- fact_flight_sales_dq: student SK (passenger/route/airplane) → NOTICE,
  airport_sk → порог 1% (эталон через ods.routes)
- routes_dq: закомментирован RI airplane_code → ods.airplanes
- ODS DAG: убрана зависимость dq_ods_airplanes → load_ods_routes
- 18 студенческих файлов заменены заглушками (ODS/DDS/DM load+dq)
- test_ods_sql_contract: SNAPSHOT_ENTITIES без airplanes/seats
- test_dags_smoke: убран assert dq_ods_airplanes → dq_ods_routes

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-12 23:07:14 +03:00

276 lines
9.7 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
from __future__ import annotations
"""
Учебный DAG: загрузка из STG в ODS (Greenplum) по домену bookings.
Ключевая идея:
- весь запуск ODS работает с одним stg_batch_id;
- для каждой сущности выполняем пару задач load -> dq;
- загрузка реализована как SCD1 UPSERT (UPDATE изменившихся + INSERT новых).
"""
from datetime import timedelta
from logging import getLogger
import pendulum
from airflow.operators.python import PythonOperator
from airflow.providers.postgres.hooks.postgres import PostgresHook
from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow import DAG
GREENPLUM_CONN_ID = "greenplum_conn"
log = getLogger(__name__)
default_args = {
"owner": "airflow",
"retries": 1,
"retry_delay": timedelta(seconds=30),
}
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_ts > (SELECT MAX(_load_ts) FROM ods.table).
Это сделано для того, чтобы не потерять инкременты, если STG-DAG запускался
несколько раз до запуска ODS-DAG'а.
Зачем нужна согласованность (INTERSECT по всем справочникам)?
Чтобы ODS загружал только те данные, для которых уже приехали ВСЕ связанные
справочники. Это защищает от рассинхрона данных, когда часть измерений в текущем
батче обновилась, а часть — упала или не доехала, что могло бы привести к
потере ссылочной целостности при сборке витрин.
"""
conf = context["dag_run"].conf or {}
stg_batch_id = conf.get("stg_batch_id")
if not stg_batch_id:
hook = PostgresHook(postgres_conn_id=GREENPLUM_CONN_ID)
result = hook.get_first(
"""
WITH candidate_batches AS (
SELECT _load_id
FROM stg.airports
WHERE _load_id IS NOT NULL AND _load_id <> ''
GROUP BY _load_id
INTERSECT
SELECT _load_id
FROM stg.airplanes
WHERE _load_id IS NOT NULL AND _load_id <> ''
GROUP BY _load_id
INTERSECT
SELECT _load_id
FROM stg.routes
WHERE _load_id IS NOT NULL AND _load_id <> ''
GROUP BY _load_id
INTERSECT
SELECT _load_id
FROM stg.seats
WHERE _load_id IS NOT NULL AND _load_id <> ''
GROUP BY _load_id
),
batch_ready AS (
SELECT
c._load_id,
GREATEST(
COALESCE(
(SELECT MAX(_load_ts) FROM stg.airports a WHERE a._load_id = c._load_id),
TIMESTAMP '1900-01-01 00:00:00'
),
COALESCE(
(SELECT MAX(_load_ts) FROM stg.airplanes a WHERE a._load_id = c._load_id),
TIMESTAMP '1900-01-01 00:00:00'
),
COALESCE(
(SELECT MAX(_load_ts) FROM stg.routes r WHERE r._load_id = c._load_id),
TIMESTAMP '1900-01-01 00:00:00'
),
COALESCE(
(SELECT MAX(_load_ts) FROM stg.seats s WHERE s._load_id = c._load_id),
TIMESTAMP '1900-01-01 00:00:00'
)
) AS ready_dttm
FROM candidate_batches c
)
SELECT _load_id
FROM batch_ready
ORDER BY ready_dttm DESC
LIMIT 1
"""
)
stg_batch_id = result[0] if result and result[0] else None
if not stg_batch_id:
raise ValueError(
"stg_batch_id не найден: передайте stg_batch_id в conf или "
"сначала выполните bookings_to_gp_stage для snapshot-справочников"
)
log.info("Используем stg_batch_id=%s", stg_batch_id)
return stg_batch_id
def _finish_summary() -> None:
"""Логирует краткий итог выполнения ODS-ветки."""
log.info("DAG bookings_to_gp_ods завершён. Подробности смотрите в логах задач.")
with DAG(
dag_id="bookings_to_gp_ods",
start_date=pendulum.datetime(2017, 1, 1, tz="UTC"),
schedule=None,
catchup=False,
max_active_runs=1,
template_searchpath="/sql",
default_args=default_args,
tags=["demo", "bookings", "greenplum", "ods"],
description="Учебный DAG: загрузка STG -> ODS (SCD1 UPSERT) + DQ проверки",
) as dag:
resolve_stg_batch_id = PythonOperator(
task_id="resolve_stg_batch_id",
python_callable=_resolve_stg_batch_id,
)
load_ods_bookings = PostgresOperator(
task_id="load_ods_bookings",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/bookings_load.sql",
)
dq_ods_bookings = PostgresOperator(
task_id="dq_ods_bookings",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/bookings_dq.sql",
)
load_ods_tickets = PostgresOperator(
task_id="load_ods_tickets",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/tickets_load.sql",
)
dq_ods_tickets = PostgresOperator(
task_id="dq_ods_tickets",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/tickets_dq.sql",
)
load_ods_airports = PostgresOperator(
task_id="load_ods_airports",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/airports_load.sql",
)
dq_ods_airports = PostgresOperator(
task_id="dq_ods_airports",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/airports_dq.sql",
)
load_ods_airplanes = PostgresOperator(
task_id="load_ods_airplanes",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/airplanes_load.sql",
)
dq_ods_airplanes = PostgresOperator(
task_id="dq_ods_airplanes",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/airplanes_dq.sql",
)
load_ods_routes = PostgresOperator(
task_id="load_ods_routes",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/routes_load.sql",
)
dq_ods_routes = PostgresOperator(
task_id="dq_ods_routes",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/routes_dq.sql",
)
load_ods_seats = PostgresOperator(
task_id="load_ods_seats",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/seats_load.sql",
)
dq_ods_seats = PostgresOperator(
task_id="dq_ods_seats",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/seats_dq.sql",
)
load_ods_flights = PostgresOperator(
task_id="load_ods_flights",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/flights_load.sql",
)
dq_ods_flights = PostgresOperator(
task_id="dq_ods_flights",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/flights_dq.sql",
)
load_ods_segments = PostgresOperator(
task_id="load_ods_segments",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/segments_load.sql",
)
dq_ods_segments = PostgresOperator(
task_id="dq_ods_segments",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/segments_dq.sql",
)
load_ods_boarding_passes = PostgresOperator(
task_id="load_ods_boarding_passes",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/boarding_passes_load.sql",
)
dq_ods_boarding_passes = PostgresOperator(
task_id="dq_ods_boarding_passes",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="ods/boarding_passes_dq.sql",
)
finish_ods_summary = PythonOperator(
task_id="finish_ods_summary",
python_callable=_finish_summary,
)
# stg_batch_id нужен всем загрузочным веткам.
(
resolve_stg_batch_id
>> load_ods_bookings
>> dq_ods_bookings
>> load_ods_tickets
>> dq_ods_tickets
)
resolve_stg_batch_id >> load_ods_airports >> dq_ods_airports
resolve_stg_batch_id >> load_ods_airplanes >> dq_ods_airplanes
# На main routes не зависит от airplanes (ods.airplanes — студенческая заглушка).
# На ветке solution: [dq_ods_airports, dq_ods_airplanes] >> load_ods_routes
dq_ods_airports >> load_ods_routes >> dq_ods_routes
dq_ods_airplanes >> load_ods_seats >> dq_ods_seats
dq_ods_routes >> load_ods_flights >> dq_ods_flights
[dq_ods_flights, dq_ods_tickets] >> load_ods_segments >> dq_ods_segments
dq_ods_segments >> load_ods_boarding_passes >> dq_ods_boarding_passes
[dq_ods_boarding_passes, dq_ods_seats] >> finish_ods_summary