diff --git a/Makefile b/Makefile index 6b96431..e3b17f1 100644 --- a/Makefile +++ b/Makefile @@ -18,11 +18,11 @@ up: # Остановить и удалить контейнеры/сети текущего проекта down: - $(COMPOSE) down + $(COMPOSE) --profile live-generator down # Полная очистка окружения проекта (включая volumes) clean: - $(COMPOSE) down -v --remove-orphans + $(COMPOSE) --profile live-generator down -v --remove-orphans logs: $(COMPOSE) logs -f --tail=200 $(service) diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 3b1bddb..ffb0fbe 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -219,6 +219,9 @@ Kafka-топики данных, state и manifest генератора. `make g PROFILE=daily-wave make generator-backfill ``` +Перед записью команда проверяет с хоста, что Kafka data-топики и STG пустые. +Если там уже есть данные, она останавливается с подсказкой `make clean`. + Старый способ через `GEN_RUN_MODE`, `GEN_STATE_RESET` и `GEN_MODEL_T_END` остаётся низкоуровневым путём для отладки прямого `docker compose run`. В обычной проверке используйте `make generated-history-analytics`, чтобы не забыть @@ -237,6 +240,9 @@ docker compose exec -T kafka /opt/kafka/bin/kafka-console-consumer.sh \ Live-продолжение стартует с этого state. Используйте тот же профиль или ту же длительность, что были у backfill. Команда сама выставит `GEN_STATE_RESET=false`. +Если state записан старой версией генератора, запуск теперь падает громко: +сначала очистите стенд через `make clean` и пересоздайте историю либо явно +начните новый мир через `make generator-reset`. ```bash PROFILE=daily-wave make generator-continue diff --git a/docs/runbooks/startup-history.md b/docs/runbooks/startup-history.md index 4a2274b..bb041b4 100644 --- a/docs/runbooks/startup-history.md +++ b/docs/runbooks/startup-history.md @@ -59,8 +59,9 @@ Для `import` — что читать; пусто — `/opt/airflow/data/startup-history-import.json`. Backfill и import работают только на чистом стенде. Если Kafka-топики данных или -STG уже непустые, DAG упадёт до записи и подскажет `make clean`. Это защита от -смешивания разных миров. +STG уже непустые, DAG упадёт до записи и подскажет `make clean`. Консольные +команды `make generator-backfill` и `make startup-history-import` делают такую же +предпроверку с хоста. Это защита от смешивания разных миров. Границы пульта: @@ -145,6 +146,12 @@ PROFILE=daily-wave make generator-continue полей. Это защита от смешения разных миров. Для намеренного нового мира используйте `make generator-reset` или `make clean`. +После обновления кода старый state может оказаться в старом формате. При +`GEN_STATE_RESET=false` это теперь громкий отказ, а не тихий старт с нуля поверх +старой истории. Оператору нужно выбрать одно из двух: очистить стенд через +`make clean` и заново создать стартовую историю, либо осознанно начать новый мир +через `GEN_STATE_RESET=true` / `make generator-reset`. + После нестандартного мира `make generator-continue` нужно запускать с теми же настройками, что были у backfill/import. При расхождении генератор громко покажет поля, которые не совпали. diff --git a/generator/README.md b/generator/README.md index 388edda..078cc33 100644 --- a/generator/README.md +++ b/generator/README.md @@ -190,6 +190,11 @@ make generator-test и `generator-reset`. Старые `GEN_RUN_MODE`, `GEN_STATE_RESET` и `GEN_MODEL_T_END` остаются низкоуровневым способом для отладки. +`generator-backfill` и `startup-history-import` перед записью проверяют, что +Kafka data-топики и STG пустые. `generator-continue` продолжает только +совместимый state; старый формат state при `GEN_STATE_RESET=false` даёт отказ с +подсказкой очистить стенд или явно начать новый мир. + ## Метрики Prometheus Генератор экспортирует метрики на `:9109/metrics`: diff --git a/generator/src/clickstream_generator/airflow_control.py b/generator/src/clickstream_generator/airflow_control.py index c651246..c6e6fae 100644 --- a/generator/src/clickstream_generator/airflow_control.py +++ b/generator/src/clickstream_generator/airflow_control.py @@ -25,6 +25,11 @@ from clickstream_generator.startup_history_artifact import ( load_startup_history_artifact, validate_startup_history_artifact, ) +from clickstream_generator.stand_clean import ( + STG_EMPTY_SQL, + assert_kafka_data_topics_empty, + assert_stg_counts_empty, +) AIRFLOW_DATA_DIR = "/opt/airflow/data" @@ -43,14 +48,6 @@ WORLD_OVERRIDE_KEYS = { "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, @@ -101,13 +98,7 @@ def build_control_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_kafka_data_topics_empty(kafka_bootstrap_servers) assert_stg_tables_empty(clickhouse_hook) @@ -115,18 +106,7 @@ 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 с консоли и повторите операцию." - ) + assert_stg_counts_empty(counts) def assert_live_generator_not_running(url: str = GENERATOR_METRICS_URL) -> None: diff --git a/generator/src/clickstream_generator/kafka_io.py b/generator/src/clickstream_generator/kafka_io.py index 3bd80ba..4da2ca6 100644 --- a/generator/src/clickstream_generator/kafka_io.py +++ b/generator/src/clickstream_generator/kafka_io.py @@ -9,6 +9,7 @@ from datetime import datetime from clickstream_generator.metrics import METRICS_ERRORS_TOTAL, METRICS_EVENTS_TOTAL from clickstream_generator.state import GeneratorState +from clickstream_generator.state import UnsupportedStateVersionError logger = logging.getLogger("generator") @@ -178,14 +179,24 @@ class KafkaStateManager: f"Restored state: tick={last_state.get('tick')}, " f"last_batch_id={last_state.get('last_batch_id')}" ) - restored = GeneratorState.from_dict_safe(last_state) - if restored is None: + try: + restored = GeneratorState.from_dict(last_state) + except UnsupportedStateVersionError: + logger.error( + "Unsupported generator state version %s", + last_state.get("version", "1.0"), + ) + raise + except ValueError: logger.warning("State data was invalid, starting fresh") + return None return restored logger.info("No previous state found, starting fresh") return None + except UnsupportedStateVersionError: + raise except Exception as e: logger.warning(f"Failed to load state: {e}, starting fresh") return None diff --git a/generator/src/clickstream_generator/service.py b/generator/src/clickstream_generator/service.py index 3ac0f81..3d7a7f4 100644 --- a/generator/src/clickstream_generator/service.py +++ b/generator/src/clickstream_generator/service.py @@ -31,6 +31,7 @@ from clickstream_generator.startup_history_artifact import ( generation_settings_from_config, write_startup_history_artifact, ) +from clickstream_generator.state import UnsupportedStateVersionError logger = logging.getLogger("generator") @@ -98,7 +99,21 @@ class GeneratorService: return if not self.config.state_reset: - restored_state = self.state_manager.load() + try: + restored_state = self.state_manager.load() + except UnsupportedStateVersionError as exc: + logger.error( + "Readable generator state uses an old format: version %s. " + "Run make clean or set GEN_STATE_RESET=true only if you " + "intentionally start a new world.", + exc.found_version, + ) + raise IncompatibleStateError( + f"state version {exc.found_version} is not supported " + f"by this code version ({exc.expected_version}); " + "run make clean, or set GEN_STATE_RESET=true to start " + "a new world intentionally" + ) from exc if restored_state: try: self.manifest_manager = KafkaStartupHistoryManifest( diff --git a/generator/src/clickstream_generator/stand_clean.py b/generator/src/clickstream_generator/stand_clean.py new file mode 100644 index 0000000..ea5d6e9 --- /dev/null +++ b/generator/src/clickstream_generator/stand_clean.py @@ -0,0 +1,168 @@ +"""Общая предпроверка чистого стенда перед новым миром.""" + +from __future__ import annotations + +import argparse +import sys +import urllib.error +import urllib.request +from collections.abc import Iterable + +from clickstream_generator.config import Config +from clickstream_generator.kafka_io import ( + KafkaStateManager, + KafkaStartupHistoryManifest, +) +from clickstream_generator.service import GeneratorService, IncompatibleStateError +from clickstream_generator.state import UnsupportedStateVersionError +from clickstream_generator.startup_history_artifact import KafkaTopicInspector + + +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 +""" + +STG_TABLE_NAMES = ( + "stg.browser_raw", + "stg.location_raw", + "stg.device_raw", + "stg.geo_raw", +) + + +def assert_kafka_data_topics_empty(kafka_bootstrap_servers: str) -> None: + """Падает, если в Kafka data-топиках уже есть сообщения.""" + try: + KafkaTopicInspector(kafka_bootstrap_servers).assert_data_topics_empty() + except RuntimeError as exc: + raise RuntimeError( + "Стенд не чистый: в Kafka data-топиках уже есть сообщения. " + "Выполните make clean и повторите операцию." + ) from exc + + +def assert_stg_counts_empty(counts: Iterable[int]) -> None: + """Падает, если в STG уже есть строки.""" + dirty = [ + f"{name}={int(count)}" + for name, count in zip(STG_TABLE_NAMES, counts) + if int(count) > 0 + ] + if dirty: + raise RuntimeError( + "Стенд не чистый: в STG уже есть строки (" + + ", ".join(dirty) + + "). Выполните make clean и повторите операцию." + ) + + +def parse_stg_counts(value: str) -> tuple[int, int, int, int]: + """Разбирает TSV-строку ClickHouse с четырьмя счётчиками STG.""" + parts = value.strip().split() + if len(parts) != len(STG_TABLE_NAMES): + raise RuntimeError( + "Не удалось проверить чистоту STG: ClickHouse вернул " + f"{len(parts)} значений вместо {len(STG_TABLE_NAMES)}." + ) + return tuple(int(part) for part in parts) # type: ignore[return-value] + + +def assert_stand_clean(kafka_bootstrap_servers: str, stg_counts: Iterable[int]) -> None: + """Проверяет Kafka data-топики и STG перед backfill/import.""" + assert_kafka_data_topics_empty(kafka_bootstrap_servers) + assert_stg_counts_empty(stg_counts) + + +def assert_live_generator_not_running(live_metrics_url: str) -> None: + """Падает, если live-генератор отвечает на metrics-порту.""" + try: + urllib.request.urlopen(live_metrics_url, timeout=2).close() + except (urllib.error.URLError, TimeoutError, OSError): + return + raise RuntimeError( + "Live-генератор уже запущен. Остановите его через make generator-down " + "или make clean перед новым backfill/import." + ) + + +def assert_continue_has_state_or_clean_stand( + kafka_bootstrap_servers: str, + stg_counts: Iterable[int], + config: Config, +) -> None: + """Разрешает continue только от совместимого state или чистого стенда.""" + state_manager = KafkaStateManager(kafka_bootstrap_servers) + try: + state = state_manager.load() + except UnsupportedStateVersionError as exc: + raise RuntimeError( + f"State записан старой версией {exc.found_version}. " + "Выполните make clean или явно начните новый мир через make generator-reset." + ) from exc + finally: + state_manager.close() + + if state is None: + assert_stand_clean(kafka_bootstrap_servers, stg_counts) + return + + service = GeneratorService(config) + manifest_manager = KafkaStartupHistoryManifest(kafka_bootstrap_servers) + try: + manifest = manifest_manager.load() + finally: + manifest_manager.close() + + if service._is_startup_history_state(state, manifest): + return + if service._is_startup_history_marker(state): + mismatches = service._startup_history_mismatch_fields(state, manifest) + raise RuntimeError( + "State стартовой истории несовместим с текущими настройками: " + + ", ".join(mismatches) + + ". Выполните make clean или используйте тот же профиль." + ) + try: + service._validate_live_state_config(state) + except IncompatibleStateError as exc: + raise RuntimeError(str(exc)) from exc + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser( + description="Проверить, что стенд чист перед новым миром генератора." + ) + parser.add_argument("--mode", choices=("clean", "continue", "live"), default="clean") + parser.add_argument("--kafka-bootstrap-servers") + parser.add_argument("--stg-counts") + parser.add_argument("--live-metrics-url") + args = parser.parse_args(argv) + + try: + if args.live_metrics_url: + assert_live_generator_not_running(args.live_metrics_url) + if args.mode == "live": + return 0 + if not args.kafka_bootstrap_servers or not args.stg_counts: + raise RuntimeError("Нужны --kafka-bootstrap-servers и --stg-counts.") + stg_counts = parse_stg_counts(args.stg_counts) + if args.mode == "continue": + assert_continue_has_state_or_clean_stand( + args.kafka_bootstrap_servers, + stg_counts, + Config(), + ) + else: + assert_stand_clean(args.kafka_bootstrap_servers, stg_counts) + except RuntimeError as exc: + print(f"Ошибка: {exc}", file=sys.stderr) + return 1 + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/generator/src/clickstream_generator/state.py b/generator/src/clickstream_generator/state.py index 93029df..1780ba8 100644 --- a/generator/src/clickstream_generator/state.py +++ b/generator/src/clickstream_generator/state.py @@ -11,6 +11,19 @@ logger = logging.getLogger("generator") STATE_VERSION = "3.0" +KNOWN_OLD_STATE_VERSIONS = {"2.0"} + + +class UnsupportedStateVersionError(ValueError): + """State записан в старом формате и требует явного решения оператора.""" + + def __init__(self, found_version: str, expected_version: str = STATE_VERSION): + self.found_version = found_version + self.expected_version = expected_version + super().__init__( + f"unsupported state version {found_version}; " + f"current version is {expected_version}" + ) def _nested_list_to_tuple(obj): @@ -234,6 +247,8 @@ class GeneratorState: raise ValueError("state must be an object") version = data.get("version", "1.0") if version != STATE_VERSION: + if str(version) in KNOWN_OLD_STATE_VERSIONS: + raise UnsupportedStateVersionError(str(version)) logger.warning( "Unsupported generator state version %s, will start fresh", version, @@ -283,6 +298,8 @@ class GeneratorState: active_visits=data.get("active_visits", []), pending_visit_births=data.get("pending_visit_births", 0.0), ) + except UnsupportedStateVersionError: + raise except Exception as e: logger.warning(f"Invalid state format, will start fresh: {e}") raise ValueError(f"Invalid state: {e}") diff --git a/generator/tests/test_airflow_control.py b/generator/tests/test_airflow_control.py index 9627790..19f4ce3 100644 --- a/generator/tests/test_airflow_control.py +++ b/generator/tests/test_airflow_control.py @@ -68,7 +68,7 @@ def test_assert_stand_clean_rejects_non_empty_stg_before_writes(): def test_assert_stand_clean_rejects_non_empty_kafka_with_make_clean_hint(monkeypatch): """Непустые Kafka-топики дают ту же подсказку про make clean.""" - from clickstream_generator import airflow_control + from clickstream_generator import airflow_control, stand_clean class Inspector: def __init__(self, bootstrap_servers): @@ -81,12 +81,52 @@ def test_assert_stand_clean_rejects_non_empty_kafka_with_make_clean_hint(monkeyp def execute(self, sql): return [(0, 0, 0, 0)] - monkeypatch.setattr(airflow_control, "KafkaTopicInspector", Inspector) + monkeypatch.setattr(stand_clean, "KafkaTopicInspector", Inspector) with pytest.raises(RuntimeError, match="make clean"): airflow_control.assert_stand_clean("kafka:29092", Hook()) +def test_shared_stand_clean_rejects_host_stg_counts(): + """Хостовая предпроверка использует ту же STG-границу, что пульт.""" + from clickstream_generator.stand_clean import assert_stg_counts_empty + + with pytest.raises(RuntimeError, match="stg.location_raw=3.*make clean"): + assert_stg_counts_empty((0, 3, 0, 0)) + + +def test_continue_rejects_dirty_stand_without_state(base_config, monkeypatch): + """Continue без state не стартует новый мир поверх грязного STG.""" + from clickstream_generator import stand_clean + + class StateManager: + def __init__(self, bootstrap_servers): + self.bootstrap_servers = bootstrap_servers + + def load(self): + return None + + def close(self): + pass + + class Inspector: + def __init__(self, bootstrap_servers): + self.bootstrap_servers = bootstrap_servers + + def assert_data_topics_empty(self): + pass + + monkeypatch.setattr(stand_clean, "KafkaStateManager", StateManager) + monkeypatch.setattr(stand_clean, "KafkaTopicInspector", Inspector) + + with pytest.raises(RuntimeError, match="stg.browser_raw=1.*make clean"): + stand_clean.assert_continue_has_state_or_clean_stand( + "localhost:9092", + (1, 0, 0, 0), + base_config, + ) + + def test_check_manifest_compares_clickhouse_stats(base_config): """Check падает, когда контрольные числа ClickHouse расходятся с manifest.""" from clickstream_generator.airflow_control import assert_clickhouse_matches_manifest diff --git a/generator/tests/test_generator_control_dag_contract.py b/generator/tests/test_generator_control_dag_contract.py index 3a9e445..8cb350d 100644 --- a/generator/tests/test_generator_control_dag_contract.py +++ b/generator/tests/test_generator_control_dag_contract.py @@ -66,6 +66,66 @@ def test_compose_mounts_generator_code_without_socket_and_gates_live_service(): assert "profiles:\n - live-generator" in text +def test_down_and_clean_include_live_generator_profile(): + """Остановка стенда удаляет профильный live-генератор.""" + text = (REPO_ROOT / "Makefile").read_text(encoding="utf-8") + + assert "down:\n\t$(COMPOSE) --profile live-generator down" in text + assert "clean:\n\t$(COMPOSE) --profile live-generator down -v --remove-orphans" in text + + +def test_console_generator_paths_precheck_clean_stand_before_writes(): + """Backfill/import с хоста проверяют чистый стенд до записи.""" + run_generator = (REPO_ROOT / "scripts" / "run_generator.sh").read_text( + encoding="utf-8" + ) + import_artifact = ( + REPO_ROOT / "scripts" / "import_startup_history_artifact.sh" + ).read_text(encoding="utf-8") + + assert "bash \"${SCRIPT_DIR}/assert_stand_clean.sh\" clean" in run_generator + assert run_generator.index("assert_stand_clean.sh") < run_generator.index( + "run --rm --no-deps" + ) + assert "bash \"${SCRIPT_DIR}/assert_stand_clean.sh\" clean" in import_artifact + assert import_artifact.index("assert_stand_clean.sh") < import_artifact.index( + "startup_history_artifact import" + ) + + +def test_console_continue_checks_state_or_clean_stand_before_live_start(): + """Continue не стартует новый live-мир на грязном стенде без state.""" + run_generator = (REPO_ROOT / "scripts" / "run_generator.sh").read_text( + encoding="utf-8" + ) + + assert "elif [[ \"${VERB}\" == \"continue\" ]]" in run_generator + assert "bash \"${SCRIPT_DIR}/assert_stand_clean.sh\" continue" in run_generator + assert run_generator.index("assert_stand_clean.sh\" continue") < run_generator.index( + "up -d --build generator" + ) + + +def test_clean_start_paths_include_live_generator_profile(): + """Все чистые сбросы видят профильный live-генератор.""" + generated_history = ( + REPO_ROOT / "scripts" / "run_generated_history_analytics.sh" + ).read_text(encoding="utf-8") + + assert "--profile live-generator down -v --remove-orphans" in generated_history + + +def test_host_clean_precheck_checks_live_before_kafka_and_stg(): + """Host backfill/import сначала проверяет, что live-генератор не запущен.""" + text = (REPO_ROOT / "scripts" / "assert_stand_clean.sh").read_text( + encoding="utf-8" + ) + + assert "--mode live" in text + assert text.index("--mode live") < text.index("clickhouse-client") + assert text.index("clickhouse-client") < text.index("--mode \"${CHECK_MODE}\"") + + def test_airflow_and_generator_kafka_dependency_versions_match(): """Airflow и генератор используют одну версию kafka-python.""" airflow_req = (REPO_ROOT / "airflow" / "requirements.txt").read_text(encoding="utf-8") diff --git a/generator/tests/test_service.py b/generator/tests/test_service.py index f02b238..eccfdca 100644 --- a/generator/tests/test_service.py +++ b/generator/tests/test_service.py @@ -22,6 +22,7 @@ from generator import ( KafkaBatchHistory, TickStreamGenerator, ) +from clickstream_generator.state import UnsupportedStateVersionError class TestBatchRecordWithDictConversion: @@ -1162,6 +1163,31 @@ class TestGeneratorServiceState: state_manager.load.assert_not_called() assert service._tick == 0 + def test_known_old_state_version_fails_instead_of_starting_new_world( + self, base_config, caplog + ): + """GEN_STATE_RESET=false не скрывает старый state за чистым стартом.""" + state_manager = MagicMock() + state_manager.load.side_effect = UnsupportedStateVersionError("2.0", "3.0") + + with patch("clickstream_generator.service.start_http_server"), \ + patch("clickstream_generator.service.ensure_topics"), \ + patch("clickstream_generator.service.KafkaPublisher"), \ + patch("clickstream_generator.service.KafkaBatchHistory"), \ + patch( + "clickstream_generator.service.KafkaStateManager", + return_value=state_manager, + ), \ + patch.object(GeneratorService, "_main_loop", return_value=None), \ + caplog.at_level(logging.ERROR, logger="generator"): + + service = GeneratorService(base_config) + with pytest.raises(ValueError, match="state version 2.0.*GEN_STATE_RESET=true"): + service.start() + + assert service._tick == 0 + assert "Readable generator state uses an old format" in caplog.text + def test_invalid_restored_state_starts_fresh(self, base_config, caplog): """Сервис не падает, если state ссылается на неизвестный профиль.""" state_manager = MagicMock() diff --git a/generator/tests/test_state.py b/generator/tests/test_state.py index 1a0cebc..3dc6873 100644 --- a/generator/tests/test_state.py +++ b/generator/tests/test_state.py @@ -10,6 +10,7 @@ from unittest.mock import MagicMock, patch import pytest from generator import GeneratorState, KafkaStateManager +from clickstream_generator.state import UnsupportedStateVersionError def _make_valid_rng_state(seed: int = 42): @@ -423,6 +424,14 @@ class TestGeneratorStateValidation: with pytest.raises(ValueError, match="version"): GeneratorState.from_dict(data) + def test_from_dict_rejects_known_old_state_version_loudly(self): + """Известная старая версия state даёт отдельную ошибку миграции.""" + data = _make_valid_state_data() + data["version"] = "2.0" + + with pytest.raises(UnsupportedStateVersionError, match="2.0.*3.0"): + GeneratorState.from_dict(data) + class TestKafkaStateManager: """Тесты менеджера состояния.""" diff --git a/scripts/assert_stand_clean.sh b/scripts/assert_stand_clean.sh new file mode 100755 index 0000000..a3872ab --- /dev/null +++ b/scripts/assert_stand_clean.sh @@ -0,0 +1,60 @@ +#!/usr/bin/env bash +# +# Проверяет, что новый мир генератора не смешается со старым. + +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +REPO_ROOT="$(cd "${SCRIPT_DIR}/.." && pwd)" + +COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" +KAFKA_BOOTSTRAP_SERVERS_HOST="${KAFKA_BOOTSTRAP_SERVERS_HOST:-localhost:9092}" +GENERATOR_METRICS_URL_HOST="${GENERATOR_METRICS_URL_HOST:-http://localhost:9109/metrics}" +CLICKHOUSE_SERVICE="${CLICKHOUSE_SERVICE:-clickhouse}" +CLICKHOUSE_DB="${CLICKHOUSE_DB:-default}" +CLICKHOUSE_USER="${CLICKHOUSE_USER:-default}" +CLICKHOUSE_PASSWORD="${CLICKHOUSE_PASSWORD:-123456}" +CHECK_MODE="${1:-clean}" + +cd "${REPO_ROOT}" + +if [[ "${CHECK_MODE}" != "clean" && "${CHECK_MODE}" != "continue" ]]; then + echo "Ошибка: неизвестный режим предпроверки: ${CHECK_MODE}" >&2 + exit 2 +fi + +if [[ "${CHECK_MODE}" == "clean" ]]; then + echo "Предпроверка: live-генератор не должен быть запущен." + PYTHONPATH="${REPO_ROOT}/generator/src" \ + uv run python -m clickstream_generator.stand_clean \ + --mode live \ + --live-metrics-url "${GENERATOR_METRICS_URL_HOST}" +fi + +echo "Предпроверка: Kafka data-топики и STG должны быть пустыми или state совместим." + +if ! stg_counts="$(${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --database="${CLICKHOUSE_DB}" \ + --format=TSV \ + --query " +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 +")"; then + echo "Ошибка: не удалось проверить STG через ClickHouse." >&2 + echo "Проверьте, что ClickHouse запущен и DDL применён: make ddl" >&2 + echo "${stg_counts}" >&2 + exit 1 +fi + +KAFKA_BOOTSTRAP_SERVERS="${KAFKA_BOOTSTRAP_SERVERS_HOST}" \ +GEN_DATA_DIR="${REPO_ROOT}/data" \ +PYTHONPATH="${REPO_ROOT}/generator/src" \ + uv run python -m clickstream_generator.stand_clean \ + --mode "${CHECK_MODE}" \ + --kafka-bootstrap-servers "${KAFKA_BOOTSTRAP_SERVERS_HOST}" \ + --stg-counts "${stg_counts}" diff --git a/scripts/import_startup_history_artifact.sh b/scripts/import_startup_history_artifact.sh index bbd8945..04436d9 100755 --- a/scripts/import_startup_history_artifact.sh +++ b/scripts/import_startup_history_artifact.sh @@ -43,13 +43,16 @@ echo "GEN_MODEL_T0=${GEN_MODEL_T0}" echo "GEN_MODEL_T_END=${GEN_MODEL_T_END}" echo "" -echo "Шаг 1: запуск Kafka" -${COMPOSE_BIN} up -d kafka +echo "Шаг 1: запуск Kafka и ClickHouse" +${COMPOSE_BIN} up -d kafka clickhouse -echo "Шаг 2: сборка образа генератора" +echo "Шаг 2: предпроверка чистого стенда" +bash "${SCRIPT_DIR}/assert_stand_clean.sh" clean + +echo "Шаг 3: сборка образа генератора" ${COMPOSE_BIN} build generator -echo "Шаг 3: воспроизведение артефакта в Kafka" +echo "Шаг 4: воспроизведение артефакта в Kafka" ${COMPOSE_BIN} run --rm --no-deps \ -v "${ARTIFACT_DIR}:/startup-artifacts:ro" \ -e GEN_LAUNCH_PROFILE="${GEN_LAUNCH_PROFILE}" \ diff --git a/scripts/run_generated_history_analytics.sh b/scripts/run_generated_history_analytics.sh index 5c26ad3..6433b8c 100644 --- a/scripts/run_generated_history_analytics.sh +++ b/scripts/run_generated_history_analytics.sh @@ -59,7 +59,7 @@ echo "" if [[ "${CLEAN_START}" == "1" ]]; then echo "Шаг 0: очистка volumes ClickHouse/Kafka/state" - ${COMPOSE_BIN} down -v --remove-orphans + ${COMPOSE_BIN} --profile live-generator down -v --remove-orphans else echo "Шаг 0: CLEAN_START=0, очистка пропущена" fi diff --git a/scripts/run_generator.sh b/scripts/run_generator.sh index 4578861..66598ab 100644 --- a/scripts/run_generator.sh +++ b/scripts/run_generator.sh @@ -52,13 +52,16 @@ fi echo "" if [[ "${VERB}" == "backfill" ]]; then - echo "Шаг 1: запуск Kafka" - ${COMPOSE_BIN} up -d kafka + echo "Шаг 1: запуск Kafka и ClickHouse" + ${COMPOSE_BIN} up -d kafka clickhouse - echo "Шаг 2: сборка образа генератора" + echo "Шаг 2: предпроверка чистого стенда" + bash "${SCRIPT_DIR}/assert_stand_clean.sh" clean + + echo "Шаг 3: сборка образа генератора" ${COMPOSE_BIN} build generator - echo "Шаг 3: backfill стартовой истории" + echo "Шаг 4: backfill стартовой истории" ${COMPOSE_BIN} run --rm --no-deps \ -e GEN_RUN_MODE="${GEN_RUN_MODE}" \ -e GEN_STATE_RESET="${GEN_STATE_RESET}" \ @@ -74,6 +77,15 @@ if [[ "${VERB}" == "backfill" ]]; then -e GEN_MIN_EVENTS_PER_TICK="${GEN_MIN_EVENTS_PER_TICK}" \ -e GEN_MAX_EVENTS_PER_TICK="${GEN_MAX_EVENTS_PER_TICK}" \ generator +elif [[ "${VERB}" == "continue" ]]; then + echo "Шаг 1: запуск Kafka и ClickHouse" + ${COMPOSE_BIN} up -d kafka clickhouse + + echo "Шаг 2: проверка state или чистого стенда" + bash "${SCRIPT_DIR}/assert_stand_clean.sh" continue + + echo "Шаг 3: запуск live-генератора" + ${COMPOSE_BIN} up -d --build generator else echo "Шаг 1: запуск live-генератора" ${COMPOSE_BIN} up -d --build generator