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:
@@ -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
|
||||||
|
|||||||
@@ -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(
|
||||||
|
|||||||
@@ -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`);
|
||||||
- логирует краткую сводку в конце запуска.
|
- логирует краткую сводку в конце запуска.
|
||||||
|
|
||||||
## Как проверить результат
|
## Как проверить результат
|
||||||
|
|||||||
@@ -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"
|
||||||
|
|||||||
Reference in New Issue
Block a user