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.
This commit is contained in:
2026-07-04 21:06:30 +03:00
parent c672ed0329
commit dd4af822d2
12 changed files with 841 additions and 13 deletions
@@ -0,0 +1,250 @@
"""Чистая логика пульта Airflow для стартовой истории."""
from __future__ import annotations
import os
import urllib.error
import urllib.request
from contextlib import contextmanager
from pathlib import Path
from typing import Iterator
from clickstream_generator.config import Config
from clickstream_generator.kafka_io import (
KafkaStateManager,
KafkaStartupHistoryManifest,
ensure_topics,
)
from clickstream_generator.launch import build_launch_env
from clickstream_generator.service import GeneratorService
from clickstream_generator.startup_history_artifact import (
KafkaRawPublisher,
KafkaTopicInspector,
compare_clickhouse_stats_to_manifest,
import_startup_history_artifact,
load_startup_history_artifact,
validate_startup_history_artifact,
)
AIRFLOW_DATA_DIR = "/opt/airflow/data"
KAFKA_BOOTSTRAP_SERVERS = "kafka:29092"
GENERATOR_METRICS_URL = "http://generator:9109/metrics"
WORLD_OVERRIDE_KEYS = {
"GEN_SEED",
"GEN_MODEL_T0",
"GEN_MODEL_TIMEZONE",
"GEN_MODEL_TIME_SPEED",
"GEN_TICK_SECONDS",
"GEN_LAMBDA_BASE_PER_MIN",
"GEN_JITTER_PCT",
"GEN_MIN_EVENTS_PER_TICK",
"GEN_MAX_EVENTS_PER_TICK",
}
STG_EMPTY_SQL = """
SELECT
(SELECT count() FROM stg.browser_raw) AS browser_raw,
(SELECT count() FROM stg.location_raw) AS location_raw,
(SELECT count() FROM stg.device_raw) AS device_raw,
(SELECT count() FROM stg.geo_raw) AS geo_raw
"""
CLICKHOUSE_STATS_SQL = """
WITH
toDateTime64('{model_t0}', 6) AS t0,
toDateTime64('{model_t_end}', 6) AS t_end
SELECT
count() AS events,
uniqExact(click_id) AS visits,
uniqExact(user_domain_id) AS users,
toString(min(event_ts)) AS min_event_timestamp,
toString(max(event_ts)) AS max_event_timestamp
FROM dm.v_events_enriched
WHERE event_ts >= t0 AND event_ts < t_end
"""
def build_control_env(
operation: str,
*,
profile_name: str,
duration: str | None = None,
artifact_path: str | None = None,
overrides: dict[str, str] | None = None,
) -> dict[str, str]:
"""Готовит env для операции пульта без доступа к Docker."""
if operation not in {"backfill", "import"}:
raise ValueError("operation must be backfill or import")
if operation == "import" and not artifact_path:
raise ValueError("artifact_path is required for import")
selected_overrides = {
key: value
for key, value in (overrides or {}).items()
if key in WORLD_OVERRIDE_KEYS and value != ""
}
env = build_launch_env(
"backfill",
profile_name=profile_name,
duration=duration or None,
overrides=selected_overrides,
)
env["KAFKA_BOOTSTRAP_SERVERS"] = KAFKA_BOOTSTRAP_SERVERS
env["GEN_DATA_DIR"] = AIRFLOW_DATA_DIR
env["GEN_METRICS_ENABLED"] = "false"
if artifact_path:
env["GEN_STARTUP_HISTORY_ARTIFACT"] = artifact_path
return env
def assert_stand_clean(kafka_bootstrap_servers: str, clickhouse_hook) -> None:
"""Проверяет, что backfill/import не смешает миры."""
try:
KafkaTopicInspector(kafka_bootstrap_servers).assert_data_topics_empty()
except RuntimeError as exc:
raise RuntimeError(
"Стенд не чистый: в Kafka data-топиках уже есть сообщения. "
"Выполните make clean с консоли и повторите операцию."
) from exc
assert_stg_tables_empty(clickhouse_hook)
def assert_stg_tables_empty(clickhouse_hook) -> None:
"""Падает, если в STG уже есть строки."""
result = clickhouse_hook.execute(STG_EMPTY_SQL)
counts = result[0] if result else ()
names = ("stg.browser_raw", "stg.location_raw", "stg.device_raw", "stg.geo_raw")
dirty = [
f"{name}={int(count)}"
for name, count in zip(names, counts)
if int(count) > 0
]
if dirty:
raise RuntimeError(
"Стенд не чистый: в STG уже есть строки ("
+ ", ".join(dirty)
+ "). Выполните make clean с консоли и повторите операцию."
)
def assert_live_generator_not_running(url: str = GENERATOR_METRICS_URL) -> None:
"""Мягко предупреждает о live-сервисе по HTTP-метрикам."""
try:
urllib.request.urlopen(url, timeout=2).close()
except (urllib.error.URLError, TimeoutError, OSError):
return
raise RuntimeError(
"Live-генератор отвечает на metrics-порту. Остановите его с консоли "
"перед backfill/import, чтобы не смешать миры."
)
def run_backfill(env: dict[str, str]) -> None:
"""Запускает backfill в процессе Airflow worker."""
with _patched_environ(env):
config = Config()
GeneratorService(config).start()
def run_import(env: dict[str, str], artifact_path: str) -> dict:
"""Импортирует портативный артефакт в Kafka."""
with _patched_environ(env):
config = Config()
artifact = load_startup_history_artifact(artifact_path)
ensure_topics(config.kafka_bootstrap_servers)
publisher = KafkaRawPublisher(config.kafka_bootstrap_servers)
state_manager = KafkaStateManager(config.kafka_bootstrap_servers)
manifest_manager = KafkaStartupHistoryManifest(config.kafka_bootstrap_servers)
topic_inspector = KafkaTopicInspector(config.kafka_bootstrap_servers)
try:
return import_startup_history_artifact(
artifact,
publisher=publisher,
state_manager=state_manager,
manifest_manager=manifest_manager,
expected_config=config,
topic_inspector=topic_inspector,
)
finally:
publisher.close()
state_manager.close()
manifest_manager.close()
def validate_import_artifact(env: dict[str, str], artifact_path: str) -> None:
"""Проверяет артефакт и настройки import до записи в Kafka."""
with _patched_environ(env):
config = Config()
artifact = load_startup_history_artifact(artifact_path)
validate_startup_history_artifact(artifact, expected_config=config)
def load_manifest_from_kafka(
kafka_bootstrap_servers: str = KAFKA_BOOTSTRAP_SERVERS,
) -> dict:
"""Читает manifest стартовой истории из compact-топика."""
manager = KafkaStartupHistoryManifest(kafka_bootstrap_servers)
try:
manifest = manager.load()
finally:
manager.close()
if not manifest:
raise RuntimeError("Manifest стартовой истории не найден в Kafka.")
return manifest
def assert_clickhouse_matches_manifest(manifest: dict, clickhouse_hook) -> None:
"""Сверяет контрольные числа ClickHouse с manifest."""
model_t0 = _clickhouse_datetime_literal(str(manifest["model_t0"]))
model_t_end = _clickhouse_datetime_literal(str(manifest["model_t_end"]))
result = clickhouse_hook.execute(
CLICKHOUSE_STATS_SQL.format(model_t0=model_t0, model_t_end=model_t_end)
)
if not result:
raise RuntimeError("ClickHouse не вернул контрольные числа.")
row = result[0]
stats = {
"events": str(row[0]),
"visits": str(row[1]),
"users": str(row[2]),
"min_event_timestamp": str(row[3]),
"max_event_timestamp": str(row[4]),
}
mismatches = compare_clickhouse_stats_to_manifest(manifest, stats)
if mismatches:
raise RuntimeError(
"ClickHouse расходится с manifest: " + ", ".join(mismatches)
)
def default_artifact_path(operation: str) -> str:
"""Возвращает путь артефакта по умолчанию в общем томе data."""
filename = (
"startup-history-import.json"
if operation == "import"
else "startup-history.json"
)
return str(Path(AIRFLOW_DATA_DIR) / filename)
@contextmanager
def _patched_environ(env: dict[str, str]) -> Iterator[None]:
old_values = {key: os.environ.get(key) for key in env}
os.environ.update(env)
try:
yield
finally:
for key, value in old_values.items():
if value is None:
os.environ.pop(key, None)
else:
os.environ[key] = value
def _clickhouse_datetime_literal(value: str) -> str:
normalized = value.replace("T", " ").removesuffix("Z")
if len(normalized) >= 6 and normalized[-6] in "+-" and normalized[-3] == ":":
normalized = normalized[:-6]
return normalized
@@ -67,6 +67,10 @@ class Config:
metrics_port: int = field(
default_factory=lambda: int(os.getenv("GEN_METRICS_PORT", "9109"))
)
metrics_enabled: bool = field(
default_factory=lambda: os.getenv("GEN_METRICS_ENABLED", "true").lower()
== "true"
)
state_enabled: bool = field(
default_factory=lambda: os.getenv("GEN_STATE_ENABLED", "true").lower() == "true"
)
@@ -62,8 +62,11 @@ class GeneratorService:
logger.warning("Generator is disabled (GEN_ENABLED=false)")
return
logger.info(f"Starting metrics server on port {self.config.metrics_port}")
start_http_server(self.config.metrics_port)
if self.config.metrics_enabled:
logger.info(f"Starting metrics server on port {self.config.metrics_port}")
start_http_server(self.config.metrics_port)
else:
logger.info("Metrics HTTP server is disabled")
logger.info("Starting generator service...")
logger.info(