Files
airflow-greenplum/tests/test_dags_smoke.py
ddadmin a51e2a6243 feat(validate): добавлен валидационный DAG bookings_validate
- Зачем:
  - студент реализует 9 объектов DWH без автоматической обратной связи;
  - ошибки (NULL в PK, дубли BK, сломанный SCD2) обнаруживались только на слое DM
    через 2–3 слоя, где отладка многократно сложнее.
- Что:
  - создан DAG bookings_validate.py — три параллельных TaskGroup (validate_ods,
    validate_dds, validate_dm), schedule=None (только ручной trigger);
  - создано 17 SQL-скриптов в sql/validate/: ODS BK coverage vs STG-батч,
    согласованность _load_id между ods.airplanes и ods.seats, дубли BK, NULL в PK,
    покрытие ODS→DDS для измерений (dim_routes требует открытой версии valid_to IS NULL),
    SCD2 активный тест backup→mutate→student_load→check(5 инвариантов)→restore
    (trigger_rule=all_done), проверка «дыр» между SCD2-версиями,
    exists-чеки для 4 DM-витрин;
  - добавлен класс TestBookingsValidate (5 smoke-тестов) в tests/test_dags_smoke.py.
- Проверка:
  - make test — 4 passed, 14 skipped (DAG-тесты скипаются без Airflow, норма);
  - ручной trigger bookings_validate в Airflow UI на solution-ветке — все таски зелёные.
2026-03-12 12:12:44 +03:00

467 lines
19 KiB
Python

from __future__ import annotations
import importlib
import pytest
def _airflow_available() -> bool:
try:
af = importlib.import_module("airflow")
except Exception:
return False
# Real Airflow exposes DAG at top-level
return hasattr(af, "DAG")
pytestmark = pytest.mark.skipif(
not _airflow_available(), reason="Airflow is not installed for DAG smoke tests"
)
def _load_dag(module_name: str):
mod = importlib.import_module(module_name)
assert hasattr(mod, "dag"), f"{module_name} must expose variable 'dag'"
return getattr(mod, "dag")
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}"
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}"
def test_bookings_stg_ddl_dag_structure():
"""Проверка структуры DAG bookings_stg_ddl."""
dag = _load_dag("airflow.dags.bookings_stg_ddl")
expected_tasks = {
"apply_stg_bookings_ddl",
"apply_stg_tickets_ddl",
"apply_stg_airports_ddl",
"apply_stg_airplanes_ddl",
"apply_stg_routes_ddl",
"apply_stg_seats_ddl",
"apply_stg_flights_ddl",
"apply_stg_segments_ddl",
"apply_stg_boarding_passes_ddl",
}
assert expected_tasks.issubset(dag.task_dict.keys())
# Smoke-test графа: проверяем ключевые инварианты, не фиксируя линейный порядок.
# Это позволяет в будущем распараллеливать независимые DDL-задачи.
_assert_reachable(dag, "apply_stg_bookings_ddl", "apply_stg_tickets_ddl")
for task_id in expected_tasks - {"apply_stg_bookings_ddl"}:
_assert_reachable(dag, "apply_stg_bookings_ddl", task_id)
def test_bookings_to_gp_stage_dag_structure():
"""Проверка структуры DAG bookings_to_gp_stage."""
dag = _load_dag("airflow.dags.bookings_to_gp_stage")
expected_tasks = {
"generate_bookings_day",
"load_bookings_to_stg",
"check_row_counts",
"load_tickets_to_stg",
"check_tickets_dq",
"load_airports_to_stg",
"check_airports_dq",
"load_airplanes_to_stg",
"check_airplanes_dq",
"load_routes_to_stg",
"check_routes_dq",
"load_seats_to_stg",
"check_seats_dq",
"load_flights_to_stg",
"check_flights_dq",
"load_segments_to_stg",
"check_segments_dq",
"load_boarding_passes_to_stg",
"check_boarding_passes_dq",
"finish_summary",
}
assert expected_tasks.issubset(dag.task_dict.keys())
# Smoke-test графа: проверяем инварианты, не фиксируя линейный порядок.
# Это позволяет в будущем распараллеливать независимые загрузки справочников/транзакций.
# Базовая цепочка должна сохраниться: генерация → bookings → DQ → tickets → DQ.
_assert_reachable(dag, "generate_bookings_day", "load_bookings_to_stg")
_assert_reachable(dag, "load_bookings_to_stg", "check_row_counts")
_assert_reachable(dag, "check_row_counts", "load_tickets_to_stg")
_assert_reachable(dag, "load_tickets_to_stg", "check_tickets_dq")
# Инвариант "load → dq" для каждой таблицы.
load_to_dq = [
("load_bookings_to_stg", "check_row_counts"),
("load_tickets_to_stg", "check_tickets_dq"),
("load_airports_to_stg", "check_airports_dq"),
("load_airplanes_to_stg", "check_airplanes_dq"),
("load_routes_to_stg", "check_routes_dq"),
("load_seats_to_stg", "check_seats_dq"),
("load_flights_to_stg", "check_flights_dq"),
("load_segments_to_stg", "check_segments_dq"),
("load_boarding_passes_to_stg", "check_boarding_passes_dq"),
]
for load_task_id, dq_task_id in load_to_dq:
_assert_direct_edge(dag, load_task_id, dq_task_id)
# Барьеры по данным (не обязательно прямые рёбра).
# routes_dq использует airports/airplanes текущего _load_id.
_assert_reachable(dag, "check_airports_dq", "check_routes_dq")
_assert_reachable(dag, "check_airplanes_dq", "check_routes_dq")
# seats_dq использует airplanes текущего _load_id.
_assert_reachable(dag, "check_airplanes_dq", "check_seats_dq")
# flights_dq использует routes текущего _load_id.
_assert_reachable(dag, "check_routes_dq", "check_flights_dq")
# segments_dq проверяет наличие flights (STG-история); для первой загрузки flights должны быть до segments.
_assert_reachable(dag, "check_flights_dq", "check_segments_dq")
# 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"
def test_bookings_ods_ddl_dag_structure():
"""Проверка структуры DAG bookings_ods_ddl."""
dag = _load_dag("airflow.dags.bookings_ods_ddl")
expected_tasks = {
"apply_ods_airports_ddl",
"apply_ods_airplanes_ddl",
"apply_ods_routes_ddl",
"apply_ods_seats_ddl",
"apply_ods_bookings_ddl",
"apply_ods_tickets_ddl",
"apply_ods_flights_ddl",
"apply_ods_segments_ddl",
"apply_ods_boarding_passes_ddl",
}
assert expected_tasks.issubset(dag.task_dict.keys())
_assert_reachable(dag, "apply_ods_airports_ddl", "apply_ods_airplanes_ddl")
for task_id in expected_tasks - {"apply_ods_airports_ddl"}:
_assert_reachable(dag, "apply_ods_airports_ddl", task_id)
def test_bookings_to_gp_ods_dag_structure():
"""Проверка структуры DAG bookings_to_gp_ods."""
dag = _load_dag("airflow.dags.bookings_to_gp_ods")
expected_tasks = {
"resolve_stg_batch_id",
"load_ods_bookings",
"dq_ods_bookings",
"load_ods_tickets",
"dq_ods_tickets",
"load_ods_airports",
"dq_ods_airports",
"load_ods_airplanes",
"dq_ods_airplanes",
"load_ods_routes",
"dq_ods_routes",
"load_ods_seats",
"dq_ods_seats",
"load_ods_flights",
"dq_ods_flights",
"load_ods_segments",
"dq_ods_segments",
"load_ods_boarding_passes",
"dq_ods_boarding_passes",
"finish_ods_summary",
}
assert expected_tasks.issubset(dag.task_dict.keys())
# Базовая цепочка транзакций.
_assert_reachable(dag, "resolve_stg_batch_id", "load_ods_bookings")
_assert_reachable(dag, "load_ods_bookings", "dq_ods_bookings")
_assert_reachable(dag, "dq_ods_bookings", "load_ods_tickets")
_assert_reachable(dag, "load_ods_tickets", "dq_ods_tickets")
# Инвариант "load -> dq" для каждой таблицы.
load_to_dq = [
("load_ods_bookings", "dq_ods_bookings"),
("load_ods_tickets", "dq_ods_tickets"),
("load_ods_airports", "dq_ods_airports"),
("load_ods_airplanes", "dq_ods_airplanes"),
("load_ods_routes", "dq_ods_routes"),
("load_ods_seats", "dq_ods_seats"),
("load_ods_flights", "dq_ods_flights"),
("load_ods_segments", "dq_ods_segments"),
("load_ods_boarding_passes", "dq_ods_boarding_passes"),
]
for load_task_id, dq_task_id in load_to_dq:
_assert_direct_edge(dag, load_task_id, dq_task_id)
# Справочники стартуют параллельно и не зависят друг от друга.
_assert_reachable(dag, "resolve_stg_batch_id", "load_ods_airports")
_assert_reachable(dag, "resolve_stg_batch_id", "load_ods_airplanes")
airports = dag.get_task("load_ods_airports")
airplanes = dag.get_task("load_ods_airplanes")
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"
# Барьеры по данным.
_assert_reachable(dag, "dq_ods_airports", "dq_ods_routes")
_assert_reachable(dag, "dq_ods_airplanes", "dq_ods_routes")
_assert_reachable(dag, "dq_ods_airplanes", "dq_ods_seats")
_assert_reachable(dag, "dq_ods_routes", "dq_ods_flights")
_assert_reachable(dag, "dq_ods_flights", "dq_ods_segments")
_assert_reachable(dag, "dq_ods_tickets", "dq_ods_segments")
_assert_reachable(dag, "dq_ods_segments", "dq_ods_boarding_passes")
# Финальная сводка должна ждать обе ветки.
_assert_reachable(dag, "dq_ods_boarding_passes", "finish_ods_summary")
_assert_reachable(dag, "dq_ods_seats", "finish_ods_summary")
def test_bookings_dds_ddl_dag_structure():
"""Проверка структуры DAG bookings_dds_ddl."""
dag = _load_dag("airflow.dags.bookings_dds_ddl")
expected_tasks = {
"apply_dds_dim_calendar_ddl",
"apply_dds_dim_airports_ddl",
"apply_dds_dim_airplanes_ddl",
"apply_dds_dim_tariffs_ddl",
"apply_dds_dim_passengers_ddl",
"apply_dds_dim_routes_ddl",
"apply_dds_fact_flight_sales_ddl",
}
assert expected_tasks.issubset(dag.task_dict.keys())
_assert_reachable(dag, "apply_dds_dim_calendar_ddl", "apply_dds_dim_airports_ddl")
for task_id in expected_tasks - {"apply_dds_dim_calendar_ddl"}:
_assert_reachable(dag, "apply_dds_dim_calendar_ddl", task_id)
def test_bookings_to_gp_dds_dag_structure():
"""Проверка структуры DAG bookings_to_gp_dds."""
dag = _load_dag("airflow.dags.bookings_to_gp_dds")
expected_tasks = {
"load_dds_dim_calendar",
"dq_dds_dim_calendar",
"load_dds_dim_airports",
"dq_dds_dim_airports",
"load_dds_dim_airplanes",
"dq_dds_dim_airplanes",
"load_dds_dim_tariffs",
"dq_dds_dim_tariffs",
"load_dds_dim_passengers",
"dq_dds_dim_passengers",
"load_dds_dim_routes",
"dq_dds_dim_routes",
"load_dds_fact_flight_sales",
"dq_dds_fact_flight_sales",
"finish_dds_summary",
}
assert expected_tasks.issubset(dag.task_dict.keys())
load_to_dq = [
("load_dds_dim_calendar", "dq_dds_dim_calendar"),
("load_dds_dim_airports", "dq_dds_dim_airports"),
("load_dds_dim_airplanes", "dq_dds_dim_airplanes"),
("load_dds_dim_tariffs", "dq_dds_dim_tariffs"),
("load_dds_dim_passengers", "dq_dds_dim_passengers"),
("load_dds_dim_routes", "dq_dds_dim_routes"),
("load_dds_fact_flight_sales", "dq_dds_fact_flight_sales"),
]
for load_task_id, dq_task_id in load_to_dq:
_assert_direct_edge(dag, load_task_id, dq_task_id)
# После calendar все остальные измерения должны быть reachable.
_assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_dim_airports")
_assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_dim_airplanes")
_assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_dim_tariffs")
_assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_dim_passengers")
_assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_dim_routes")
# Параллельность: airports и airplanes не должны зависеть друг от друга.
airports = dag.get_task("load_dds_dim_airports")
airplanes = dag.get_task("load_dds_dim_airplanes")
assert airplanes not in airports.get_flat_relatives(
upstream=False
), "dds airports не должен быть upstream для airplanes"
assert airports not in airplanes.get_flat_relatives(
upstream=False
), "dds airplanes не должен быть upstream для airports"
# dim_routes зависит от airports и airplanes (денормализация).
_assert_reachable(dag, "dq_dds_dim_airports", "load_dds_dim_routes")
_assert_reachable(dag, "dq_dds_dim_airplanes", "load_dds_dim_routes")
# Факт должен стартовать только после всех измерений.
_assert_reachable(dag, "dq_dds_dim_calendar", "load_dds_fact_flight_sales")
_assert_reachable(dag, "dq_dds_dim_airports", "load_dds_fact_flight_sales")
_assert_reachable(dag, "dq_dds_dim_airplanes", "load_dds_fact_flight_sales")
_assert_reachable(dag, "dq_dds_dim_tariffs", "load_dds_fact_flight_sales")
_assert_reachable(dag, "dq_dds_dim_passengers", "load_dds_fact_flight_sales")
_assert_reachable(dag, "dq_dds_dim_routes", "load_dds_fact_flight_sales")
# Финальная сводка должна ждать DQ факта.
_assert_reachable(dag, "dq_dds_fact_flight_sales", "finish_dds_summary")
def test_bookings_dm_ddl_dag_structure():
"""Проверка структуры DAG bookings_dm_ddl."""
dag = _load_dag("airflow.dags.bookings_dm_ddl")
expected_tasks = {
"apply_dm_sales_report_ddl",
"apply_dm_route_performance_ddl",
"apply_dm_passenger_loyalty_ddl",
"apply_dm_airport_traffic_ddl",
"apply_dm_monthly_overview_ddl",
}
assert expected_tasks.issubset(dag.task_dict.keys())
# Проверяем линейную цепочку
_assert_direct_edge(
dag, "apply_dm_sales_report_ddl", "apply_dm_route_performance_ddl"
)
_assert_direct_edge(
dag, "apply_dm_route_performance_ddl", "apply_dm_passenger_loyalty_ddl"
)
_assert_direct_edge(
dag, "apply_dm_passenger_loyalty_ddl", "apply_dm_airport_traffic_ddl"
)
_assert_direct_edge(
dag, "apply_dm_airport_traffic_ddl", "apply_dm_monthly_overview_ddl"
)
class TestBookingsValidate:
"""Smoke-тесты DAG bookings_validate."""
def test_dag_loads(self):
dag = _load_dag("airflow.dags.bookings_validate")
assert dag is not None
def test_expected_tasks(self):
dag = _load_dag("airflow.dags.bookings_validate")
expected = {
# ODS
"validate_ods.check_ods_airplanes_rowcount",
"validate_ods.check_ods_seats_rowcount",
"validate_ods.check_ods_no_dup_bk",
"validate_ods.check_ods_no_null_pks",
# DDS
"validate_dds.check_dim_airplanes_exists",
"validate_dds.check_dim_passengers_exists",
"validate_dds.check_dim_passengers_no_dup_bk",
"validate_dds.check_dim_routes_exists",
"validate_dds.scd2_backup",
"validate_dds.scd2_mutate",
"validate_dds.scd2_run_student_load",
"validate_dds.scd2_check",
"validate_dds.scd2_restore",
"validate_dds.check_dim_routes_no_gaps",
# DM
"validate_dm.check_airport_traffic_exists",
"validate_dm.check_route_performance_exists",
"validate_dm.check_monthly_overview_exists",
"validate_dm.check_passenger_loyalty_exists",
}
assert expected.issubset(dag.task_dict.keys())
def test_scd2_chain(self):
"""SCD2 цепочка backup → mutate → load → check → restore."""
dag = _load_dag("airflow.dags.bookings_validate")
_assert_direct_edge(dag, "validate_dds.scd2_backup", "validate_dds.scd2_mutate")
_assert_direct_edge(
dag, "validate_dds.scd2_mutate", "validate_dds.scd2_run_student_load"
)
_assert_direct_edge(
dag, "validate_dds.scd2_run_student_load", "validate_dds.scd2_check"
)
_assert_direct_edge(dag, "validate_dds.scd2_check", "validate_dds.scd2_restore")
def test_no_gaps_after_restore(self):
"""no_gaps должен выполняться после restore (на чистых данных)."""
dag = _load_dag("airflow.dags.bookings_validate")
_assert_direct_edge(
dag, "validate_dds.scd2_restore", "validate_dds.check_dim_routes_no_gaps"
)
def test_restore_trigger_rule(self):
"""restore должен выполняться всегда (all_done), даже если check упал."""
dag = _load_dag("airflow.dags.bookings_validate")
restore_task = dag.task_dict["validate_dds.scd2_restore"]
assert restore_task.trigger_rule == "all_done"
def test_bookings_to_gp_dm_dag_structure():
"""Проверка структуры DAG bookings_to_gp_dm."""
dag = _load_dag("airflow.dags.bookings_to_gp_dm")
expected_tasks = {
"start_dm",
"load_dm_sales_report",
"dq_dm_sales_report",
"load_dm_route_performance",
"dq_dm_route_performance",
"load_dm_passenger_loyalty",
"dq_dm_passenger_loyalty",
"load_dm_airport_traffic",
"dq_dm_airport_traffic",
"load_dm_monthly_overview",
"dq_dm_monthly_overview",
"finish_dm_summary",
}
assert expected_tasks.issubset(dag.task_dict.keys())
marts = [
"sales_report",
"route_performance",
"passenger_loyalty",
"airport_traffic",
"monthly_overview",
]
for mart in marts:
# load -> dq
_assert_direct_edge(dag, f"load_dm_{mart}", f"dq_dm_{mart}")
# start -> load
_assert_reachable(dag, "start_dm", f"load_dm_{mart}")
# dq -> finish
_assert_reachable(dag, f"dq_dm_{mart}", "finish_dm_summary")