Files
clickstream-ch-kafka-supers…/generator/tests/test_airflow_control.py
T
ddadmin dd4af822d2 feat(airflow): добавлен пульт управления генератором
- Зачем:
  - нужен основной ручной интерфейс стенда для backfill/import/check без консольной матрицы переменных.
- Что:
  - добавлен DAG generator_control с параметрами Airflow, ветвлением операций и ожиданием ETL.
  - вынесена общая логика запуска и предпроверок генератора для Airflow.
  - обновлены compose-настройки, зависимости, тесты и документация по пульту.
- Проверка:
  - uv run --with pytest --with-requirements generator/requirements.txt pytest generator/tests -q.
  - docker compose config --quiet.
2026-07-04 21:06:30 +03:00

150 lines
4.9 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.
"""
Тесты чистой логики пульта Airflow для генератора.
"""
from dataclasses import replace
from datetime import datetime, timezone
from pathlib import Path
import pytest
def test_build_control_env_uses_profile_duration_and_airflow_data_dir():
"""Пульт строит env для backfill из профиля и каталога Airflow."""
from clickstream_generator.airflow_control import build_control_env
env = build_control_env(
"backfill",
profile_name="ci",
duration="",
artifact_path="/opt/airflow/data/startup-history.json",
overrides={"GEN_LAMBDA_BASE_PER_MIN": "120", "GEN_SEED": "99"},
)
assert env["GEN_RUN_MODE"] == "backfill"
assert env["GEN_HISTORY_DURATION"] == "6h"
assert env["GEN_MODEL_T_END"] == "2026-01-01T06:00:00+00:00"
assert env["GEN_DATA_DIR"] == "/opt/airflow/data"
assert env["GEN_METRICS_ENABLED"] == "false"
assert env["GEN_STARTUP_HISTORY_ARTIFACT"] == "/opt/airflow/data/startup-history.json"
assert env["GEN_LAMBDA_BASE_PER_MIN"] == "120"
assert env["GEN_SEED"] == "99"
def test_import_env_requires_artifact_path_and_uses_backfill_contract():
"""Import валидируется как тот же мир, что backfill."""
from clickstream_generator.airflow_control import build_control_env
with pytest.raises(ValueError, match="artifact_path"):
build_control_env("import", profile_name="ci", artifact_path="")
env = build_control_env(
"import",
profile_name="daily-wave",
artifact_path="/opt/airflow/data/history.json",
)
assert env["GEN_RUN_MODE"] == "backfill"
assert env["GEN_STATE_RESET"] == "true"
assert env["GEN_HISTORY_DURATION"] == "2d"
def test_assert_stand_clean_rejects_non_empty_stg_before_writes():
"""Backfill/import не стартуют на непустом STG."""
from clickstream_generator.airflow_control import assert_stg_tables_empty
class Hook:
def execute(self, sql):
self.sql = sql
return [(0, 2, 0, 0)]
hook = Hook()
with pytest.raises(RuntimeError, match="make clean"):
assert_stg_tables_empty(hook)
assert "stg.browser_raw" in hook.sql
assert "stg.location_raw" in hook.sql
def test_assert_stand_clean_rejects_non_empty_kafka_with_make_clean_hint(monkeypatch):
"""Непустые Kafka-топики дают ту же подсказку про make clean."""
from clickstream_generator import airflow_control
class Inspector:
def __init__(self, bootstrap_servers):
self.bootstrap_servers = bootstrap_servers
def assert_data_topics_empty(self):
raise RuntimeError("Kafka data topics are not empty: browser_events")
class Hook:
def execute(self, sql):
return [(0, 0, 0, 0)]
monkeypatch.setattr(airflow_control, "KafkaTopicInspector", Inspector)
with pytest.raises(RuntimeError, match="make clean"):
airflow_control.assert_stand_clean("kafka:29092", Hook())
def test_check_manifest_compares_clickhouse_stats(base_config):
"""Check падает, когда контрольные числа ClickHouse расходятся с manifest."""
from clickstream_generator.airflow_control import assert_clickhouse_matches_manifest
from clickstream_generator.startup_history_artifact import (
StartupHistoryArtifactBuilder,
build_manifest,
)
from test_startup_history_artifact import _batch, _state
state = _state()
builder = StartupHistoryArtifactBuilder()
builder.add_batch(_batch())
manifest = build_manifest(
config=replace(
base_config,
model_t0=state.model_t0,
model_t_end=state.model_timestamp,
model_time_speed=1,
model_timezone="UTC",
seed=42,
),
counters=builder.counters,
state=state,
)
class Hook:
def execute(self, sql):
self.sql = sql
return [(
0,
1,
1,
"2026-01-01 00:00:00.000000",
"2026-01-01 00:00:00.000000",
)]
with pytest.raises(RuntimeError, match="events"):
assert_clickhouse_matches_manifest(manifest, Hook())
def test_generator_metrics_server_can_be_disabled(base_config, monkeypatch):
"""Airflow-задача может запускать генератор без HTTP-сервера метрик."""
from clickstream_generator.service import GeneratorService
config = replace(base_config, enabled=False, metrics_enabled=False)
service = GeneratorService(config)
called = False
def start_http_server(_port):
nonlocal called
called = True
monkeypatch.setattr(
"clickstream_generator.service.start_http_server",
start_http_server,
)
service.start()
assert called is False