From 69651f98fff7e566e13201fa4cb19a670b8d4be3 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 22 Feb 2026 20:19:58 +0300 Subject: [PATCH] =?UTF-8?q?refactor(dags):=20=D0=BF=D0=B0=D1=80=D0=B0?= =?UTF-8?q?=D0=BB=D0=BB=D0=B5=D0=BB=D0=B8=D0=B7=D0=BE=D0=B2=D0=B0=D0=BD=20?= =?UTF-8?q?=D0=B3=D1=80=D0=B0=D1=84=20DAG=20bookings=5Fto=5Fgp=5Fstage?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - линейный граф маскировал реальные зависимости данных; для эталонного стенда важно показать менти параллельный граф там, где данные независимы. - Что: - airports и airplanes грузятся параллельно после check_tickets_dq. - routes ждёт обоих (DQ проверяет ссылочную целостность на оба справочника). - seats зависит только от airplanes и работает параллельно с веткой routes → flights → segments → boarding_passes. - finish_summary ждёт обе ветки (check_boarding_passes_dq + check_seats_dq). - datetime.utcnow() заменён на datetime.now(UTC) в csv_to_greenplum.py. - smoke-тесты дополнены проверкой параллельности и второй ветки. - документация обновлена с ASCII-схемой нового графа. - Проверка: - make test (11 passed, 4 skipped), make lint — чисто. Co-Authored-By: Claude Opus 4.6 --- airflow/dags/bookings_to_gp_stage.py | 47 +++++++++++++++++++++------- airflow/dags/csv_to_greenplum.py | 8 ++--- docs/bookings_to_gp_stage.md | 28 ++++++++++++----- tests/test_dags_smoke.py | 29 ++++++++++++----- 4 files changed, 83 insertions(+), 29 deletions(-) diff --git a/airflow/dags/bookings_to_gp_stage.py b/airflow/dags/bookings_to_gp_stage.py index 2cf5832..bcacfd4 100644 --- a/airflow/dags/bookings_to_gp_stage.py +++ b/airflow/dags/bookings_to_gp_stage.py @@ -15,6 +15,23 @@ from __future__ import annotations - `run_id` используется как метка запуска (в `batch_id`, в логах и DQ). Важно: для инкрементальных таблиц «пустое окно инкремента» допустимо (это не ошибка). + +Граф зависимостей (параллельный там, где данные независимы): + + generate_bookings_day → load_bookings → check_bookings_dq + → load_tickets → check_tickets_dq + ├─ load_airports → check_airports_dq ─┐ + │ ├─ load_routes → check_routes_dq + ├─ load_airplanes → check_airplanes_dq ┤ → load_flights → check_flights_dq + │ │ → load_segments → check_segments_dq + │ │ → load_boarding_passes → check_bp_dq ─┐ + │ └─ load_seats → check_seats_dq ──────────────────────┤ + │ ▼ + └──────────────────────────────────────────────────────────────────────────── finish_summary + +airports и airplanes грузятся параллельно (они не зависят друг от друга). +routes зависит от обоих (DQ проверяет ссылочную целостность на airports и airplanes). +seats зависит только от airplanes (DQ проверяет airplane_code → airplanes). """ from datetime import timedelta @@ -188,22 +205,30 @@ with DAG( python_callable=_finish_summary, ) - # Сначала загружаются и проверяются bookings и tickets + # === Этап 1. Транзакции: bookings → tickets (последовательно, т.к. tickets зависят от bookings) === generate_bookings_day >> load_bookings_to_stg >> check_row_counts check_row_counts >> load_tickets_to_stg >> check_tickets_dq - # Затем загружаются справочники. - # Для простоты (и более понятных логов для новичков) делаем это последовательно. - # Если позже понадобится ускорить DAG, эти шаги можно распараллелить, сохранив зависимости. + # === Этап 2. Справочники (параллельно, где данные независимы) === + # airports и airplanes не зависят друг от друга — грузим параллельно. check_tickets_dq >> load_airports_to_stg >> check_airports_dq - check_airports_dq >> load_airplanes_to_stg >> check_airplanes_dq - check_airplanes_dq >> load_routes_to_stg >> check_routes_dq - check_routes_dq >> load_seats_to_stg >> check_seats_dq + check_tickets_dq >> load_airplanes_to_stg >> check_airplanes_dq - # Затем загружаются транзакции (тоже последовательно, по тем же причинам). - check_seats_dq >> load_flights_to_stg >> check_flights_dq + # routes зависит от airports И airplanes (DQ проверяет ссылочную целостность). + [check_airports_dq, check_airplanes_dq] >> load_routes_to_stg >> check_routes_dq + + # seats зависит только от airplanes (DQ проверяет ссылочную целостность). + check_airplanes_dq >> load_seats_to_stg >> check_seats_dq + + # === Этап 3. Транзакции (последовательно, каждая зависит от предыдущей) === + # flights зависят от routes (DQ проверяет ссылочную целостность route_no → routes). + check_routes_dq >> load_flights_to_stg >> check_flights_dq + + # segments зависят от flights и tickets (DQ проверяет обе ссылки). check_flights_dq >> load_segments_to_stg >> check_segments_dq + + # boarding_passes зависят от segments и tickets (DQ проверяет обе ссылки). check_segments_dq >> load_boarding_passes_to_stg >> check_boarding_passes_dq - # В конце финальный лог - check_boarding_passes_dq >> finish_summary + # === Финал: ждём завершения ВСЕХ веток === + [check_boarding_passes_dq, check_seats_dq] >> finish_summary diff --git a/airflow/dags/csv_to_greenplum.py b/airflow/dags/csv_to_greenplum.py index 9289403..ad5273a 100644 --- a/airflow/dags/csv_to_greenplum.py +++ b/airflow/dags/csv_to_greenplum.py @@ -3,7 +3,7 @@ from __future__ import annotations import logging import os import random -from datetime import datetime, timedelta +from datetime import UTC, datetime, timedelta from pathlib import Path from typing import List @@ -37,11 +37,11 @@ def _create_table() -> None: def _generate_csv(rows: int, csv_dir: Path) -> str: """Генерирует CSV c заказами с помощью pandas и сохраняет на диск.""" csv_dir.mkdir(parents=True, exist_ok=True) - timestamp = datetime.utcnow().strftime("%Y%m%d_%H%M%S") + timestamp = datetime.now(UTC).strftime("%Y%m%d_%H%M%S") csv_path = csv_dir / f"orders_{timestamp}.csv" # Генерируем данные в pandas-стиле - base_order_id = int(datetime.utcnow().timestamp() * 1_000) + base_order_id = int(datetime.now(UTC).timestamp() * 1_000) # Создаём DataFrame с использованием pandas методов df = pd.DataFrame( @@ -52,7 +52,7 @@ def _generate_csv(rows: int, csv_dir: Path) -> str: ), # Временные метки с интервалом в 1 секунду в обратном порядке "order_ts": pd.date_range( - end=datetime.utcnow(), periods=rows, freq="1S" + end=datetime.now(UTC), periods=rows, freq="1S" ).sort_values(ascending=False), # Случайные customer_id от 1 до 1000 "customer_id": pd.Series( diff --git a/docs/bookings_to_gp_stage.md b/docs/bookings_to_gp_stage.md index 296493e..335a604 100644 --- a/docs/bookings_to_gp_stage.md +++ b/docs/bookings_to_gp_stage.md @@ -87,26 +87,40 @@ make bookings-init - проверяет количество строк в том же окне инкремента, а также ссылочную целостность и обязательные поля; - при проблемах делает `RAISE EXCEPTION`, чтобы DAG падал “красным”. -6) Справочники (full load) +6) Справочники (full load, параллельно где возможно) -Каждый справочник загружается “снэпшотом” (все строки) и затем проверяется DQ-скриптом: +Справочники загружаются “снэпшотом” (все строки) и затем проверяются DQ-скриптом. +Порядок определяется зависимостями данных — **airports** и **airplanes** грузятся **параллельно**, +потому что не зависят друг от друга: + +``` +check_tickets_dq + ├─ load_airports → check_airports_dq ─┐ + │ ├─ load_routes → check_routes_dq + └─ load_airplanes → check_airplanes_dq ─┤ + └─ load_seats → check_seats_dq +``` - `load_airports_to_stg` → `check_airports_dq` (`sql/stg/airports_load.sql`, `sql/stg/airports_dq.sql`) - `load_airplanes_to_stg` → `check_airplanes_dq` (`sql/stg/airplanes_load.sql`, `sql/stg/airplanes_dq.sql`) -- `load_routes_to_stg` → `check_routes_dq` (`sql/stg/routes_load.sql`, `sql/stg/routes_dq.sql`) -- `load_seats_to_stg` → `check_seats_dq` (`sql/stg/seats_load.sql`, `sql/stg/seats_dq.sql`) +- `load_routes_to_stg` → `check_routes_dq` (`sql/stg/routes_load.sql`, `sql/stg/routes_dq.sql`) — зависит от **airports** и **airplanes** (DQ проверяет ссылочную целостность) +- `load_seats_to_stg` → `check_seats_dq` (`sql/stg/seats_load.sql`, `sql/stg/seats_dq.sql`) — зависит от **airplanes** (DQ проверяет `airplane_code → airplanes`) 7) Транзакции -- `load_flights_to_stg` → `check_flights_dq` (инкремент по `scheduled_departure`) -- `load_segments_to_stg` → `check_segments_dq` (инкремент по `book_date` через tickets/bookings) -- `load_boarding_passes_to_stg` → `check_boarding_passes_dq` (full snapshot) +- `load_flights_to_stg` → `check_flights_dq` (инкремент по `scheduled_departure`) — зависит от **routes** +- `load_segments_to_stg` → `check_segments_dq` (инкремент по `book_date` через tickets/bookings) — зависит от **flights** +- `load_boarding_passes_to_stg` → `check_boarding_passes_dq` (full snapshot) — зависит от **segments** + +Ветка `seats` работает параллельно с веткой `routes → flights → segments → boarding_passes`. +Обе ветки сходятся на `finish_summary`. Важно: для инкрементальных таблиц “пустое окно инкремента” допустимо — загрузка и DQ логируют `NOTICE` и завершаются успешно. Для snapshot-справочников (airports/airplanes/routes/seats) пустой источник считается ошибкой (DQ делает `RAISE EXCEPTION`). 8) `finish_summary` +- ждёт завершения **обеих** параллельных веток (`check_boarding_passes_dq` и `check_seats_dq`); - логирует краткую сводку в конце запуска. ## Как проверить результат diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index 100b722..e422d0e 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -28,17 +28,17 @@ def _load_dag(module_name: str): def _assert_direct_edge(dag, upstream_task_id: str, downstream_task_id: str) -> None: upstream = dag.get_task(upstream_task_id) downstream = dag.get_task(downstream_task_id) - assert downstream in upstream.get_direct_relatives(upstream=False), ( - f"Expected direct edge {upstream_task_id} -> {downstream_task_id}" - ) + assert downstream in upstream.get_direct_relatives( + upstream=False + ), f"Expected direct edge {upstream_task_id} -> {downstream_task_id}" def _assert_reachable(dag, upstream_task_id: str, downstream_task_id: str) -> None: upstream = dag.get_task(upstream_task_id) downstream = dag.get_task(downstream_task_id) - assert downstream in upstream.get_flat_relatives(upstream=False), ( - f"Expected {downstream_task_id} to be downstream of {upstream_task_id}" - ) + assert downstream in upstream.get_flat_relatives( + upstream=False + ), f"Expected {downstream_task_id} to be downstream of {upstream_task_id}" def test_csv_to_greenplum_dag_structure(): @@ -182,5 +182,20 @@ def test_bookings_to_gp_stage_dag_structure(): # boarding_passes_dq проверяет наличие segments/tickets (STG-история); для первой загрузки segments должны быть до DQ. _assert_reachable(dag, "check_segments_dq", "check_boarding_passes_dq") - # Финальная сводка должна быть в конце графа. + # Финальная сводка должна быть в конце графа (обе ветки). _assert_reachable(dag, "check_boarding_passes_dq", "finish_summary") + _assert_reachable(dag, "check_seats_dq", "finish_summary") + + # Параллельность: airports и airplanes оба downstream от check_tickets_dq, + # но НЕ зависят друг от друга (ни прямо, ни транзитивно). + _assert_reachable(dag, "check_tickets_dq", "load_airports_to_stg") + _assert_reachable(dag, "check_tickets_dq", "load_airplanes_to_stg") + + airports = dag.get_task("load_airports_to_stg") + airplanes = dag.get_task("load_airplanes_to_stg") + assert airplanes not in airports.get_flat_relatives( + upstream=False + ), "airports не должен быть upstream для airplanes" + assert airports not in airplanes.get_flat_relatives( + upstream=False + ), "airplanes не должен быть upstream для airports"