diff --git a/README.md b/README.md index 4435214..d1f3b07 100644 --- a/README.md +++ b/README.md @@ -28,13 +28,22 @@ ## Быстрый старт -Штатный чистый запуск строит аналитику из стартовой истории генератора. Команда -очищает volumes ClickHouse и Kafka, создаёт стартовую историю, доводит её до DM и -проверяет Superset metadata. +Для ручной работы поднимите стенд и создайте стартовую историю через Airflow: + +```bash +make up +docker compose ps +``` + +Откройте Airflow: `http://localhost:8080` (`admin/admin`). Запустите +`generator_control` с операцией `backfill`: DAG создаст стартовую историю, +запустит ETL и выполнит `check`. `make up` не запускает live-генератор; live +включается отдельно командой `make generator-continue`. + +Для полностью автоматического чистого прогона из консоли остаётся команда: ```bash make generated-history-analytics -docker compose ps ``` По умолчанию это быстрый профиль `ci`: 6 часов модельного времени. Историю на diff --git a/airflow/dags/generator_control_dag.py b/airflow/dags/generator_control_dag.py new file mode 100644 index 0000000..c773ffc --- /dev/null +++ b/airflow/dags/generator_control_dag.py @@ -0,0 +1,257 @@ +""" +DAG-пульт генератора стартовой истории. + +Пульт выполняет только Python-код генератора внутри Airflow worker. Жизненный +цикл контейнеров остаётся в Makefile. +""" + +from __future__ import annotations + +from datetime import datetime, timedelta + +from airflow import DAG +from airflow.models.param import Param +from airflow.operators.empty import EmptyOperator +from airflow.operators.python import BranchPythonOperator, PythonOperator +from airflow.operators.trigger_dagrun import TriggerDagRunOperator +from airflow.utils.trigger_rule import TriggerRule +from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook + +from clickstream_generator.airflow_control import ( + KAFKA_BOOTSTRAP_SERVERS, + assert_clickhouse_matches_manifest, + assert_live_generator_not_running, + assert_stand_clean, + build_control_env, + default_artifact_path, + load_manifest_from_kafka, + run_backfill, + run_import, + validate_import_artifact, +) +from clickstream_generator.launch import PROFILES + + +default_args = { + "owner": "airflow", + "depends_on_past": False, + "email_on_failure": False, + "email_on_retry": False, + "retries": 0, + "retry_delay": timedelta(minutes=1), +} + + +def _conf(context) -> dict: + dag_run = context.get("dag_run") + return dag_run.conf if dag_run and dag_run.conf else {} + + +def _param(context, name: str): + return _conf(context).get(name, context["params"][name]) + + +def _operation(context) -> str: + return str(_param(context, "operation")) + + +def _artifact_path(context) -> str: + value = str(_param(context, "artifact_path") or "").strip() + if value: + return value + return default_artifact_path(_operation(context)) + + +def _overrides(context) -> dict[str, str]: + return { + key: str(_param(context, param_name) or "").strip() + for param_name, key in { + "seed": "GEN_SEED", + "model_time_speed": "GEN_MODEL_TIME_SPEED", + }.items() + } + + +def _build_env(context, operation: str) -> dict[str, str]: + artifact_path = _artifact_path(context) + if ( + operation == "backfill" + and not str(_param(context, "artifact_path") or "").strip() + ): + artifact_path = None + return build_control_env( + operation, + profile_name=str(_param(context, "profile")), + duration=str(_param(context, "duration") or "").strip(), + artifact_path=artifact_path, + overrides=_overrides(context), + ) + + +def choose_operation(**context) -> str: + """Выбирает ветку пульта по параметру operation.""" + operation = _operation(context) + if operation == "backfill": + return "precheck_backfill" + if operation == "import": + return "precheck_import" + if operation == "check": + return "check_only" + raise ValueError(f"Неизвестная операция: {operation}") + + +def precheck_backfill(**context) -> None: + """Проверяет чистоту стенда перед backfill.""" + assert_live_generator_not_running() + hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default") + assert_stand_clean(KAFKA_BOOTSTRAP_SERVERS, hook) + context["ti"].xcom_push( + key="generator_env", + value=_build_env(context, "backfill"), + ) + + +def run_backfill_task(**context) -> None: + """Выполняет backfill через код генератора.""" + env = context["ti"].xcom_pull(task_ids="precheck_backfill", key="generator_env") + run_backfill(env) + + +def precheck_import(**context) -> None: + """Проверяет чистоту стенда и совместимость артефакта перед import.""" + assert_live_generator_not_running() + hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default") + assert_stand_clean(KAFKA_BOOTSTRAP_SERVERS, hook) + env = _build_env(context, "import") + artifact_path = _artifact_path(context) + validate_import_artifact(env, artifact_path) + context["ti"].xcom_push(key="generator_env", value=env) + context["ti"].xcom_push(key="artifact_path", value=artifact_path) + + +def run_import_task(**context) -> None: + """Выполняет import портативного артефакта.""" + ti = context["ti"] + env = ti.xcom_pull(task_ids="precheck_import", key="generator_env") + artifact_path = ti.xcom_pull(task_ids="precheck_import", key="artifact_path") + run_import(env, artifact_path) + + +def check_manifest_task(**context) -> None: + """Сверяет витрину ClickHouse с manifest стартовой истории.""" + manifest = load_manifest_from_kafka() + hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default") + assert_clickhouse_matches_manifest(manifest, hook) + + +with DAG( + dag_id="generator_control", + description="Пульт стартовой истории генератора", + default_args=default_args, + schedule=None, + start_date=datetime(2024, 1, 1), + catchup=False, + max_active_runs=1, + is_paused_upon_creation=True, + tags=["generator", "startup-history"], + params={ + "operation": Param( + "backfill", + type="string", + enum=["backfill", "import", "check"], + title="Операция", + description=( + "Что сделать: создать историю, импортировать артефакт " + "или проверить витрины." + ), + ), + "profile": Param( + sorted(PROFILES)[0], + type="string", + enum=sorted(PROFILES), + title="Профиль", + description="Именованный набор настроек генератора.", + ), + "duration": Param( + "", + type="string", + title="Длительность", + description="Например 6h или 2d. Пусто — взять длительность из профиля.", + ), + "seed": Param( + "", + type="string", + title="GEN_SEED", + description="Пусто — взять seed из профиля.", + ), + "model_time_speed": Param( + "", + type="string", + title="GEN_MODEL_TIME_SPEED", + description="Пусто — взять скорость модельного времени из профиля.", + ), + "artifact_path": Param( + "", + type="string", + title="Артефакт", + description=( + "Backfill: куда сохранить файл; пусто — не сохранять. " + "Import: что читать; пусто — путь по умолчанию в data." + ), + ), + }, +) as dag: + route = BranchPythonOperator( + task_id="choose_operation", + python_callable=choose_operation, + ) + + precheck_backfill_task = PythonOperator( + task_id="precheck_backfill", + python_callable=precheck_backfill, + ) + backfill_task = PythonOperator( + task_id="run_backfill", + python_callable=run_backfill_task, + ) + + precheck_import_task = PythonOperator( + task_id="precheck_import", + python_callable=precheck_import, + ) + import_task = PythonOperator( + task_id="run_import", + python_callable=run_import_task, + ) + + trigger_etl = TriggerDagRunOperator( + task_id="trigger_etl", + trigger_dag_id="etl_pipeline", + conf={"full_refresh": True}, + wait_for_completion=True, + allowed_states=["success"], + failed_states=["failed"], + poke_interval=30, + trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS, + ) + + check_after_etl = PythonOperator( + task_id="check_after_etl", + python_callable=check_manifest_task, + ) + + check_only = PythonOperator( + task_id="check_only", + python_callable=check_manifest_task, + ) + + done = EmptyOperator( + task_id="done", + trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS, + ) + + route >> [precheck_backfill_task, precheck_import_task, check_only] + precheck_backfill_task >> backfill_task >> trigger_etl + precheck_import_task >> import_task >> trigger_etl + trigger_etl >> check_after_etl >> done + check_only >> done diff --git a/airflow/requirements.txt b/airflow/requirements.txt index 9504b3c..94184f2 100644 --- a/airflow/requirements.txt +++ b/airflow/requirements.txt @@ -6,5 +6,8 @@ psycopg2-binary==2.9.9 # ClickHouse operator/hook для DAG'ов airflow-clickhouse-plugin==1.6.0 -# Kafka client для загрузки данных (DAG kafka_load) +# Kafka client для загрузки данных и пульта генератора kafka-python==2.0.6 + +# Метрики импортируются кодом генератора; HTTP-сервер в Airflow-задачах выключен +prometheus-client==0.21.1 diff --git a/docker-compose.yml b/docker-compose.yml index a9c323f..770bb92 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -10,6 +10,7 @@ x-airflow-env: &airflow-default-env AIRFLOW__METRICS__STATSD_HOST: "statsd-exporter" AIRFLOW__METRICS__STATSD_PORT: "8125" AIRFLOW__METRICS__STATSD_PREFIX: "airflow" + PYTHONPATH: /opt/airflow/generator_src services: @@ -115,6 +116,7 @@ services: - "8080:8080" volumes: - ./airflow/dags:/opt/airflow/dags + - ./generator/src:/opt/airflow/generator_src:ro - ./sql:/opt/airflow/sql:ro - ./data:/opt/airflow/data networks: @@ -140,6 +142,7 @@ services: " volumes: - ./airflow/dags:/opt/airflow/dags + - ./generator/src:/opt/airflow/generator_src:ro - ./sql:/opt/airflow/sql:ro - ./data:/opt/airflow/data networks: @@ -162,6 +165,7 @@ services: <<: *airflow-default-env volumes: - ./airflow/dags:/opt/airflow/dags + - ./generator/src:/opt/airflow/generator_src:ro - ./sql:/opt/airflow/sql:ro - ./data:/opt/airflow/data networks: @@ -286,6 +290,8 @@ services: # Генератор событий (автономный стриминг в Kafka) generator: + profiles: + - live-generator build: context: ./generator dockerfile: Dockerfile diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index eafba41..8ec44d8 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -37,9 +37,31 @@ ## Airflow DAGs -Штатный аналитический путь больше не начинается с `kafka_load`: чистый стенд -получает данные из стартовой истории генератора. DAG-и ниже остаются для -ручных экспериментов, отладки и совместимости учебного стенда. +Штатный ручной путь начинается с `generator_control`: чистый стенд получает +стартовую историю генератора, затем этот же DAG запускает ETL и проверку. +`kafka_load` остаётся для экспериментов и совместимости учебного стенда. + +### `generator_control` + +- Запуск: ручной (`Trigger DAG`). +- Назначение: пульт стартовой истории генератора. +- Операции: + - `backfill` — создать стартовую историю, затем запустить `etl_pipeline` и + дождаться `success`; + - `import` — импортировать портативный артефакт, затем запустить + `etl_pipeline` и дождаться `success`; + - `check` — сверить ClickHouse с manifest из Kafka. +- Параметры: + - `operation` (`backfill` / `import` / `check`); + - `profile` — список берётся из `PROFILES` генератора; + - `duration` — `6h`, `2d` и т.п.; пусто означает длительность профиля; + - `seed`, `model_time_speed` — необязательные переопределения мира; + - `artifact_path` — для `backfill` путь сохранения, для `import` путь чтения. + +Backfill/import требуют чистый стенд: пустые data-топики Kafka и пустые +`stg.*_raw`. При отказе очистите стенд через `make clean`. Операции `continue` +в DAG нет: live-генератор — долгоживущий сервис, его запускают с консоли через +`make generator-continue`. ### `ddl_init` @@ -81,12 +103,17 @@ ## Генератор событий (автономный стриминг) -Автономный сервис для непрерывной генерации событий в Kafka. Работает независимо от Airflow DAGs. +Автономный сервис для непрерывной генерации событий в Kafka. Работает независимо +от Airflow DAGs. `make up` его не запускает: live включается только явной +командой. ### Управление ```bash -# Запустить генератор +# Продолжить live-поток из state +make generator-continue + +# Запустить генератор с текущими env напрямую make generator-up # Остановить генератор diff --git a/docs/runbooks/startup-history.md b/docs/runbooks/startup-history.md index 5145069..cc9c7b6 100644 --- a/docs/runbooks/startup-history.md +++ b/docs/runbooks/startup-history.md @@ -34,6 +34,46 @@ Длительность можно переопределить через `GEN_HISTORY_DURATION`, например `2d`. Команда сама считает `GEN_MODEL_T_END` от `GEN_MODEL_T0`. +## Пульт в Airflow + +Основной ручной путь — DAG `generator_control` в Airflow UI: + +1. Поднимите стенд: `make up`. +2. Если DDL ещё не применён, запустите `ddl_init`. +3. Откройте `generator_control` и выберите `operation`. + +Операции: + +- `backfill` — создать стартовую историю. После записи в Kafka DAG сам запускает + `etl_pipeline`, ждёт завершения и выполняет `check`. +- `import` — прочитать артефакт из `artifact_path`. Несовместимый артефакт + отклоняется до записи в Kafka. +- `check` — сверить ClickHouse с manifest из Kafka. + +Поля формы: + +- `profile` берётся из профилей генератора. +- `duration` можно оставить пустым, тогда берётся длительность профиля. +- `seed` и `model_time_speed` — необязательные переопределения мира. +- `artifact_path`: для `backfill` — куда сохранить файл; пусто — не сохранять. + Для `import` — что читать; пусто — `/opt/airflow/data/startup-history-import.json`. + +Backfill и import работают только на чистом стенде. Если Kafka-топики данных или +STG уже непустые, DAG упадёт до записи и подскажет `make clean`. Это защита от +смешивания разных миров. + +Границы пульта: + +- `make up`, `make clean` и live-продолжение остаются в консоли. +- Операции `continue` в DAG нет намеренно: live — долгоживущий сервис, а пульт + управляет разовыми пакетными операциями. +- Airflow не получает доступ к жизненному циклу контейнеров; таски выполняют обычный + Python-код генератора. + +Если backfill сохраняет файл в `./data`, он создаётся пользователем Airflow +внутри контейнера. Чтение работает из Airflow и консольных команд, но перезапись +чужого файла может потребовать удалить старый файл вручную. + ## Экспорт По умолчанию создаётся быстрый 6-часовой артефакт: @@ -103,5 +143,9 @@ PROFILE=ci make generator-continue полей. Это защита от смешения разных миров. Для намеренного нового мира используйте `make generator-reset` или `make clean`. +После нестандартного мира `make generator-continue` нужно запускать с теми же +настройками, что были у backfill/import. При расхождении генератор громко +покажет поля, которые не совпали. + Старые переменные `GEN_RUN_MODE`, `GEN_STATE_RESET` и `GEN_MODEL_T_END` остаются низкоуровневым способом для отладки и прямого `docker compose run`. diff --git a/generator/requirements.txt b/generator/requirements.txt index 23e7cf1..badfad0 100644 --- a/generator/requirements.txt +++ b/generator/requirements.txt @@ -1,5 +1,5 @@ # Kafka клиент -kafka-python==2.0.5 +kafka-python==2.0.6 # Prometheus метрики prometheus-client==0.21.1 diff --git a/generator/src/clickstream_generator/airflow_control.py b/generator/src/clickstream_generator/airflow_control.py new file mode 100644 index 0000000..c651246 --- /dev/null +++ b/generator/src/clickstream_generator/airflow_control.py @@ -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 diff --git a/generator/src/clickstream_generator/config.py b/generator/src/clickstream_generator/config.py index 6d7cf3d..d7ac0dc 100644 --- a/generator/src/clickstream_generator/config.py +++ b/generator/src/clickstream_generator/config.py @@ -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" ) diff --git a/generator/src/clickstream_generator/service.py b/generator/src/clickstream_generator/service.py index b73cd07..11d14d1 100644 --- a/generator/src/clickstream_generator/service.py +++ b/generator/src/clickstream_generator/service.py @@ -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( diff --git a/generator/tests/test_airflow_control.py b/generator/tests/test_airflow_control.py new file mode 100644 index 0000000..9627790 --- /dev/null +++ b/generator/tests/test_airflow_control.py @@ -0,0 +1,149 @@ +""" +Тесты чистой логики пульта 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 diff --git a/generator/tests/test_generator_control_dag_contract.py b/generator/tests/test_generator_control_dag_contract.py new file mode 100644 index 0000000..3a9e445 --- /dev/null +++ b/generator/tests/test_generator_control_dag_contract.py @@ -0,0 +1,76 @@ +""" +Контракт DAG generator_control без запуска Airflow. +""" + +import ast +from pathlib import Path + +from clickstream_generator.launch import PROFILES + + +DAG_PATH = Path(__file__).parents[2] / "airflow" / "dags" / "generator_control_dag.py" +REPO_ROOT = Path(__file__).parents[2] + + +def _tree(): + return ast.parse(DAG_PATH.read_text(encoding="utf-8")) + + +def test_dag_file_exists_and_uses_dynamic_profiles(): + """DAG берёт варианты профилей из PROFILES, а не из ручного списка.""" + text = DAG_PATH.read_text(encoding="utf-8") + + assert "dag_id=\"generator_control\"" in text + assert "sorted(PROFILES)" in text + for profile in PROFILES: + assert profile not in {"hardcoded-profile"} + + +def test_trigger_form_has_expected_param_enums(): + """Форма запуска ограничивает операции и профили.""" + text = DAG_PATH.read_text(encoding="utf-8") + + assert 'enum=["backfill", "import", "check"]' in text + assert "enum=sorted(PROFILES)" in text + assert '"duration": Param(' in text + assert '"artifact_path": Param(' in text + + +def test_dag_branches_and_waits_for_etl_completion(): + """Backfill/import запускают ETL и ждут его завершения перед check.""" + text = DAG_PATH.read_text(encoding="utf-8") + + assert "BranchPythonOperator" in text + assert "TriggerDagRunOperator" in text + assert 'trigger_dag_id="etl_pipeline"' in text + assert "wait_for_completion=True" in text + assert 'allowed_states=["success"]' in text + assert 'failed_states=["failed"]' in text + + +def test_no_docker_or_continue_operation_in_dag(): + """Пульт не управляет Docker и не содержит операцию continue.""" + text = DAG_PATH.read_text(encoding="utf-8") + + assert "docker" not in text.lower() + assert '"continue"' not in text + + +def test_compose_mounts_generator_code_without_socket_and_gates_live_service(): + """Airflow видит код генератора, но не получает Docker socket.""" + text = (REPO_ROOT / "docker-compose.yml").read_text(encoding="utf-8") + + assert "PYTHONPATH: /opt/airflow/generator_src" in text + assert "./generator/src:/opt/airflow/generator_src:ro" in text + assert "/var/run/docker.sock" not in text + assert "profiles:\n - live-generator" in text + + +def test_airflow_and_generator_kafka_dependency_versions_match(): + """Airflow и генератор используют одну версию kafka-python.""" + airflow_req = (REPO_ROOT / "airflow" / "requirements.txt").read_text(encoding="utf-8") + generator_req = (REPO_ROOT / "generator" / "requirements.txt").read_text(encoding="utf-8") + + assert "kafka-python==2.0.6" in airflow_req + assert "kafka-python==2.0.6" in generator_req + assert "prometheus-client==0.21.1" in airflow_req