refactor(dags): параллелизован граф DAG bookings_to_gp_stage

- Зачем:
  - линейный граф маскировал реальные зависимости данных; для эталонного
    стенда важно показать менти параллельный граф там, где данные независимы.
- Что:
  - 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 <noreply@anthropic.com>
This commit is contained in:
2026-02-22 20:19:58 +03:00
co-authored by Claude Opus 4.6
parent a0258fd63e
commit 69651f98ff
4 changed files with 83 additions and 29 deletions
+36 -11
View File
@@ -15,6 +15,23 @@ from __future__ import annotations
- `run_id` используется как метка запуска (в `batch_id`, в логах и DQ). - `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 from datetime import timedelta
@@ -188,22 +205,30 @@ with DAG(
python_callable=_finish_summary, python_callable=_finish_summary,
) )
# Сначала загружаются и проверяются bookings и tickets # === Этап 1. Транзакции: bookings tickets (последовательно, т.к. tickets зависят от bookings) ===
generate_bookings_day >> load_bookings_to_stg >> check_row_counts generate_bookings_day >> load_bookings_to_stg >> check_row_counts
check_row_counts >> load_tickets_to_stg >> check_tickets_dq check_row_counts >> load_tickets_to_stg >> check_tickets_dq
# Затем загружаются справочники. # === Этап 2. Справочники (параллельно, где данные независимы) ===
# Для простоты (и более понятных логов для новичков) делаем это последовательно. # airports и airplanes не зависят друг от друга — грузим параллельно.
# Если позже понадобится ускорить DAG, эти шаги можно распараллелить, сохранив зависимости.
check_tickets_dq >> load_airports_to_stg >> check_airports_dq check_tickets_dq >> load_airports_to_stg >> check_airports_dq
check_airports_dq >> load_airplanes_to_stg >> check_airplanes_dq check_tickets_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
# Затем загружаются транзакции (тоже последовательно, по тем же причинам). # routes зависит от airports И airplanes (DQ проверяет ссылочную целостность).
check_seats_dq >> load_flights_to_stg >> check_flights_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 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_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
+4 -4
View File
@@ -3,7 +3,7 @@ from __future__ import annotations
import logging import logging
import os import os
import random import random
from datetime import datetime, timedelta from datetime import UTC, datetime, timedelta
from pathlib import Path from pathlib import Path
from typing import List from typing import List
@@ -37,11 +37,11 @@ def _create_table() -> None:
def _generate_csv(rows: int, csv_dir: Path) -> str: def _generate_csv(rows: int, csv_dir: Path) -> str:
"""Генерирует CSV c заказами с помощью pandas и сохраняет на диск.""" """Генерирует CSV c заказами с помощью pandas и сохраняет на диск."""
csv_dir.mkdir(parents=True, exist_ok=True) 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" csv_path = csv_dir / f"orders_{timestamp}.csv"
# Генерируем данные в pandas-стиле # Генерируем данные в pandas-стиле
base_order_id = int(datetime.utcnow().timestamp() * 1_000) base_order_id = int(datetime.now(UTC).timestamp() * 1_000)
# Создаём DataFrame с использованием pandas методов # Создаём DataFrame с использованием pandas методов
df = pd.DataFrame( df = pd.DataFrame(
@@ -52,7 +52,7 @@ def _generate_csv(rows: int, csv_dir: Path) -> str:
), ),
# Временные метки с интервалом в 1 секунду в обратном порядке # Временные метки с интервалом в 1 секунду в обратном порядке
"order_ts": pd.date_range( "order_ts": pd.date_range(
end=datetime.utcnow(), periods=rows, freq="1S" end=datetime.now(UTC), periods=rows, freq="1S"
).sort_values(ascending=False), ).sort_values(ascending=False),
# Случайные customer_id от 1 до 1000 # Случайные customer_id от 1 до 1000
"customer_id": pd.Series( "customer_id": pd.Series(
+21 -7
View File
@@ -87,26 +87,40 @@ make bookings-init
- проверяет количество строк в том же окне инкремента, а также ссылочную целостность и обязательные поля; - проверяет количество строк в том же окне инкремента, а также ссылочную целостность и обязательные поля;
- при проблемах делает `RAISE EXCEPTION`, чтобы DAG падал “красным”. - при проблемах делает `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_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_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_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`) - `load_seats_to_stg``check_seats_dq` (`sql/stg/seats_load.sql`, `sql/stg/seats_dq.sql`) — зависит от **airplanes** (DQ проверяет `airplane_code → airplanes`)
7) Транзакции 7) Транзакции
- `load_flights_to_stg``check_flights_dq` (инкремент по `scheduled_departure`) - `load_flights_to_stg``check_flights_dq` (инкремент по `scheduled_departure`) — зависит от **routes**
- `load_segments_to_stg``check_segments_dq` (инкремент по `book_date` через tickets/bookings) - `load_segments_to_stg``check_segments_dq` (инкремент по `book_date` через tickets/bookings) — зависит от **flights**
- `load_boarding_passes_to_stg``check_boarding_passes_dq` (full snapshot) - `load_boarding_passes_to_stg``check_boarding_passes_dq` (full snapshot) — зависит от **segments**
Ветка `seats` работает параллельно с веткой `routes → flights → segments → boarding_passes`.
Обе ветки сходятся на `finish_summary`.
Важно: для инкрементальных таблиц “пустое окно инкремента” допустимо — загрузка и DQ логируют `NOTICE` и завершаются успешно. Важно: для инкрементальных таблиц “пустое окно инкремента” допустимо — загрузка и DQ логируют `NOTICE` и завершаются успешно.
Для snapshot-справочников (airports/airplanes/routes/seats) пустой источник считается ошибкой (DQ делает `RAISE EXCEPTION`). Для snapshot-справочников (airports/airplanes/routes/seats) пустой источник считается ошибкой (DQ делает `RAISE EXCEPTION`).
8) `finish_summary` 8) `finish_summary`
- ждёт завершения **обеих** параллельных веток (`check_boarding_passes_dq` и `check_seats_dq`);
- логирует краткую сводку в конце запуска. - логирует краткую сводку в конце запуска.
## Как проверить результат ## Как проверить результат
+22 -7
View File
@@ -28,17 +28,17 @@ def _load_dag(module_name: str):
def _assert_direct_edge(dag, upstream_task_id: str, downstream_task_id: str) -> None: def _assert_direct_edge(dag, upstream_task_id: str, downstream_task_id: str) -> None:
upstream = dag.get_task(upstream_task_id) upstream = dag.get_task(upstream_task_id)
downstream = dag.get_task(downstream_task_id) downstream = dag.get_task(downstream_task_id)
assert downstream in upstream.get_direct_relatives(upstream=False), ( assert downstream in upstream.get_direct_relatives(
f"Expected direct edge {upstream_task_id} -> {downstream_task_id}" 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: def _assert_reachable(dag, upstream_task_id: str, downstream_task_id: str) -> None:
upstream = dag.get_task(upstream_task_id) upstream = dag.get_task(upstream_task_id)
downstream = dag.get_task(downstream_task_id) downstream = dag.get_task(downstream_task_id)
assert downstream in upstream.get_flat_relatives(upstream=False), ( assert downstream in upstream.get_flat_relatives(
f"Expected {downstream_task_id} to be downstream of {upstream_task_id}" upstream=False
) ), f"Expected {downstream_task_id} to be downstream of {upstream_task_id}"
def test_csv_to_greenplum_dag_structure(): 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. # boarding_passes_dq проверяет наличие segments/tickets (STG-история); для первой загрузки segments должны быть до DQ.
_assert_reachable(dag, "check_segments_dq", "check_boarding_passes_dq") _assert_reachable(dag, "check_segments_dq", "check_boarding_passes_dq")
# Финальная сводка должна быть в конце графа. # Финальная сводка должна быть в конце графа (обе ветки).
_assert_reachable(dag, "check_boarding_passes_dq", "finish_summary") _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"