From 642789d4844ca34554595f114927fb473209a0cb Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Thu, 11 Jun 2026 15:00:50 +0300 Subject: [PATCH] =?UTF-8?q?refactor(generator):=20=D1=80=D0=B0=D0=B7=D0=BD?= =?UTF-8?q?=D0=B5=D1=81=D1=91=D0=BD=20=D1=81=D0=B5=D1=80=D0=B2=D0=B8=D1=81?= =?UTF-8?q?=20=D0=B3=D0=B5=D0=BD=D0=B5=D1=80=D0=B0=D1=82=D0=BE=D1=80=D0=B0?= =?UTF-8?q?=20=D0=BF=D0=BE=20src-=D0=BF=D0=B0=D0=BA=D0=B5=D1=82=D1=83?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - перед активными визитами нужно отделить генеративную модель от Kafka, состояния и сервисного цикла. - Что: - перенесены модули генератора в пакет `src/clickstream_generator`. - `generator.py` оставлен фасадом и точкой входа с совместимыми импортами. - обновлены Dockerfile, тесты, README, спека и issue 02.5. - Проверка: - `docker build -t generator:test generator`. - `docker run --rm -v /home/dmitry/sources/clickstream-ch-kafka-superset-demo:/workspace -w /workspace/generator generator:test pytest tests/ -q`. - `python -m py_compile generator.py src/clickstream_generator/*.py` в Docker. --- .../issues/02-5-generator-service-cleanup.md | 24 +- .../2026-06-11-generator-service-cleanup.md | 46 +- generator/Dockerfile | 4 +- generator/KNOWN_ISSUES.md | 12 +- generator/README.md | 19 + generator/generator.py | 1203 +---------------- .../src/clickstream_generator/__init__.py | 2 + generator/src/clickstream_generator/config.py | 62 + .../src/clickstream_generator/dictionary.py | 65 + .../src/clickstream_generator/generation.py | 211 +++ .../src/clickstream_generator/intensity.py | 41 + .../src/clickstream_generator/kafka_io.py | 363 +++++ .../src/clickstream_generator/metrics.py | 23 + .../src/clickstream_generator/runtime.py | 30 + .../src/clickstream_generator/service.py | 235 ++++ generator/src/clickstream_generator/state.py | 80 ++ generator/tests/conftest.py | 7 +- generator/tests/test_generation.py | 4 +- generator/tests/test_service_cleanup.py | 55 + 19 files changed, 1297 insertions(+), 1189 deletions(-) create mode 100644 generator/src/clickstream_generator/__init__.py create mode 100644 generator/src/clickstream_generator/config.py create mode 100644 generator/src/clickstream_generator/dictionary.py create mode 100644 generator/src/clickstream_generator/generation.py create mode 100644 generator/src/clickstream_generator/intensity.py create mode 100644 generator/src/clickstream_generator/kafka_io.py create mode 100644 generator/src/clickstream_generator/metrics.py create mode 100644 generator/src/clickstream_generator/runtime.py create mode 100644 generator/src/clickstream_generator/service.py create mode 100644 generator/src/clickstream_generator/state.py create mode 100644 generator/tests/test_service_cleanup.py diff --git a/.scratch/feature-data-generator/issues/02-5-generator-service-cleanup.md b/.scratch/feature-data-generator/issues/02-5-generator-service-cleanup.md index 1cd0022..9edff19 100644 --- a/.scratch/feature-data-generator/issues/02-5-generator-service-cleanup.md +++ b/.scratch/feature-data-generator/issues/02-5-generator-service-cleanup.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: ready-for-human # Уборка сервиса генератора перед активными визитами @@ -13,20 +13,20 @@ Status: ready-for-agent ## Acceptance criteria -- [ ] Генеративная модель визита отделена от Kafka, состояния, метрик и +- [x] Генеративная модель визита отделена от Kafka, состояния, метрик и сервисного цикла. -- [ ] Код загрузки и индексации JSONL-сида вынесен из сервисного слоя. -- [ ] Расчёт интенсивности отделён от генерации одного визита. -- [ ] Kafka publisher, Kafka-state и создание служебных топиков не смешаны с +- [x] Код загрузки и индексации JSONL-сида вынесен из сервисного слоя. +- [x] Расчёт интенсивности отделён от генерации одного визита. +- [x] Kafka publisher, Kafka-state и создание служебных топиков не смешаны с генеративной моделью. -- [ ] `generator/generator.py` остаётся тонкой точкой входа или совместимым +- [x] `generator/generator.py` остаётся тонкой точкой входа или совместимым фасадом, чтобы не ломать запуск и тесты без причины. -- [ ] `generate_tick_batch()` не остаётся в чистой генеративной модели визита: +- [x] `generate_tick_batch()` не остаётся в чистой генеративной модели визита: он вынесен в переходный тиковый слой и подготовлен к замене активными визитами в задаче 03. -- [ ] Dockerfile и запуск контейнера обновлены под новую структуру файлов. -- [ ] Все существующие тесты генератора проходят. -- [ ] README генератора кратко отражает новую структуру файлов. +- [x] Dockerfile и запуск контейнера обновлены под новую структуру файлов. +- [x] Все существующие тесты генератора проходят. +- [x] README генератора кратко отражает новую структуру файлов. ## Non-goals @@ -50,3 +50,7 @@ Status: ready-for-agent Задача появилась после разбора распухания сервиса: первые два среза уже внесли новую модель визита, а третий срез добавит состояние активных визитов. Перед ним нужно отделить модель, тиковый слой и Kafka-интеграцию. + +2026-06-11: исходники перенесены не россыпью в `generator/`, а в пакет +`generator/src/clickstream_generator/`. `generator/generator.py` оставлен тонким +фасадом и точкой входа для совместимости. diff --git a/docs/specs/2026-06-11-generator-service-cleanup.md b/docs/specs/2026-06-11-generator-service-cleanup.md index 5a9cc03..6e917a1 100644 --- a/docs/specs/2026-06-11-generator-service-cleanup.md +++ b/docs/specs/2026-06-11-generator-service-cleanup.md @@ -71,22 +71,29 @@ Целевая форма: -- `generator/config.py` — конфигурация и её валидация; -- `generator/dictionary.py` — загрузка и индексы исходных JSONL; -- `generator/generation.py` — чистая генеративная модель визита, страниц и пауз; -- `generator/intensity.py` — расчёт бюджета активности: Пуассон, часовой - коэффициент, jitter; -- `generator/kafka_io.py` — Kafka publisher, создание топиков, Kafka-state; -- `generator/state.py` — сериализуемое состояние генератора; -- `generator/service.py` — основной сервисный цикл и метрики; -- `generator/generator.py` — тонкая точка входа или совместимый фасад для - старых импортов тестов. +- `generator/generator.py` — тонкая точка входа и совместимый фасад для старых + импортов тестов; +- `generator/src/clickstream_generator/config.py` — конфигурация и её валидация; +- `generator/src/clickstream_generator/dictionary.py` — загрузка и индексы + исходных JSONL; +- `generator/src/clickstream_generator/generation.py` — чистая генеративная + модель визита, страниц и пауз; +- `generator/src/clickstream_generator/intensity.py` — расчёт бюджета активности: + Пуассон, часовой коэффициент, jitter; +- `generator/src/clickstream_generator/runtime.py` — переходный тиковый слой до + активных визитов; +- `generator/src/clickstream_generator/kafka_io.py` — Kafka publisher, создание + топиков, Kafka-state и история batch; +- `generator/src/clickstream_generator/state.py` — сериализуемое состояние + генератора; +- `generator/src/clickstream_generator/metrics.py` — Prometheus-метрики; +- `generator/src/clickstream_generator/service.py` — основной сервисный цикл. -Сейчас `generator/` не оформлен как Python-пакет: тесты импортируют -`generator/generator.py` как модуль `generator`, а Dockerfile копирует только -`generator.py`. В рамках этой уборки нужно сохранить совместимость импортов и -обновить контейнерный запуск под новую структуру файлов; превращать каталог в -пакет через `__init__.py` не требуется. +`generator.py` остаётся фасадом, потому что тесты и часть документации уже +используют импорт `from generator import ...`. Рабочий пакет называется +`clickstream_generator`, чтобы не путать его с каталогом `generator/` и файлом +`generator.py`. Dockerfile должен копировать `src/` и выставлять +`PYTHONPATH=/app/src`. Точное разбиение может быть чуть проще, если это уменьшит шум, но граница между моделью генерации, тиковым runtime и Kafka-интеграцией должна быть явной. @@ -97,6 +104,9 @@ новый слой состояния, и его нельзя удобно встроить в текущий монолит. - **Принято:** внешнее поведение генератора не меняется. Цель задачи — форма кода, а не новая модель данных. +- **Принято:** исходники живут в `generator/src/clickstream_generator/`, а не + россыпью в `generator/`. Причина: так структура ближе к обычному Python-пакету + и проще для учебного чтения. - **Принято:** `generate_tick_batch()` считать временным переходным механизмом. После задачи 03 тик должен выпускать созревшие события активных визитов. - **Отклонено:** ограничиться перестановкой функций внутри `generator.py`. @@ -118,13 +128,15 @@ - Сервисный контур продолжает публиковать в те же четыре Kafka-топика: `browser_events`, `location_events`, `device_events`, `geo_events`. - Dockerfile и запуск контейнера учитывают новую структуру файлов. +- Рабочий пакет импортируется как `clickstream_generator`, при этом старые + импорты через фасад `generator.py` продолжают работать. - В коде есть ясное место, куда задача 03 добавит активные визиты без расширения Kafka-слоя и без переписывания генерации одного визита. ## Documentation impact - После реализации уборки обновить `generator/README.md`: описать новую - структуру файлов и убрать формулировки, которые привязаны к одному - `generator.py`. + структуру `src/clickstream_generator` и убрать формулировки, которые привязаны + к одному `generator.py`. - Если по ходу будет принято решение удалить или заменить `generator_batch_history`, это должно быть отдельно отражено в README и в задаче реализации. diff --git a/generator/Dockerfile b/generator/Dockerfile index 1cb0d88..3e85845 100644 --- a/generator/Dockerfile +++ b/generator/Dockerfile @@ -12,8 +12,9 @@ COPY requirements.txt . RUN uv venv /app/.venv && \ uv pip install --python /app/.venv/bin/python -r requirements.txt -# Копируем код +# Копируем точку входа и пакет генератора COPY generator.py . +COPY src ./src # Не-root пользователь для безопасности RUN useradd -m -u 1000 generator && chown -R generator:generator /app @@ -21,6 +22,7 @@ USER generator # Активируем venv через PATH ENV PATH="/app/.venv/bin:$PATH" +ENV PYTHONPATH="/app/src" # Запуск генератора CMD ["python", "generator.py"] diff --git a/generator/KNOWN_ISSUES.md b/generator/KNOWN_ISSUES.md index 7cf3c5b..652e1a7 100644 --- a/generator/KNOWN_ISSUES.md +++ b/generator/KNOWN_ISSUES.md @@ -50,7 +50,8 @@ user_domain_id (постоянный пользователь, cookie) ## Что генератор делает правильно -Файл `generator/generator.py`, `_calculate_events_count()` (стр. ~232–260): +Файл `generator/src/clickstream_generator/intensity.py`, +`calculate_events_count()`: - **Poisson-процесс** прихода событий: λ на минуту → λ на тик → розыгрыш Пуассона. - **Дневной коэффициент** `_hour_factor()`: день (9–18) ×1.2, ночь (0–5) ×0.7. @@ -60,8 +61,8 @@ user_domain_id (постоянный пользователь, cookie) ## Исторический корневой дефект: модель сущностей в `generate_batch()` -До среза от 2026-06-11 файл `generator/generator.py`, `generate_batch()` на -каждое событие в батче делал примерно следующее: +До среза от 2026-06-11 прежняя реализация `generate_batch()` в монолитном +`generator/generator.py` на каждое событие в батче делала примерно следующее: ```python base_browser = self.rng.choice(self.dictionary.browser_events) # случайная сид-строка @@ -123,8 +124,9 @@ device_event = {**base_device, "click_id": new_click_id} # user_domain_i ## Ссылки -- Код: `generator/generator.py` — `_calculate_events_count()` (стр. ~232–260), - `generate_batch()` (стр. ~262–330). +- Код: `generator/src/clickstream_generator/intensity.py` — + `calculate_events_count()`; `generator/src/clickstream_generator/generation.py` + — `EventGenerator.generate_batch()`. - Доменная модель: `sql/ddl/dds/30_dds.sql`, `sql/ddl/dm/40_dm.sql`. - Контекст обсуждения: ветка `docs/advanced-clickstream-course`, дизайн Superset-дашборда (урок 6 курса). diff --git a/generator/README.md b/generator/README.md index e7feb20..b6362c2 100644 --- a/generator/README.md +++ b/generator/README.md @@ -16,6 +16,24 @@ generator-service -> Kafka topics -> (потребители отдельно) Генератор работает автономно и не зависит от потребителей (Airflow, ClickHouse). +### Структура кода + +Код разнесён в пакет `src/clickstream_generator/`. Сам `generator.py` остаётся +точкой входа и совместимым фасадом для старых импортов из тестов. + +| Файл | Назначение | +|------|------------| +| `src/clickstream_generator/config.py` | переменные окружения и валидация настроек | +| `src/clickstream_generator/dictionary.py` | загрузка и индексы исходных JSONL | +| `src/clickstream_generator/generation.py` | генерация одного связанного визита | +| `src/clickstream_generator/intensity.py` | расчёт событийного бюджета тика | +| `src/clickstream_generator/runtime.py` | временный тиковый слой до активных визитов | +| `src/clickstream_generator/kafka_io.py` | Kafka publisher, история batch, Kafka-state и служебные топики | +| `src/clickstream_generator/state.py` | сериализуемое состояние генератора | +| `src/clickstream_generator/metrics.py` | Prometheus-метрики | +| `src/clickstream_generator/service.py` | основной цикл сервиса | +| `generator.py` | запуск сервиса и совместимый фасад | + ## Режим работы: `steady-stream` - Публикуем постепенно, **короткими тиками** (по умолчанию каждые 5 секунд) @@ -246,6 +264,7 @@ generator/tests/ ├── test_history.py # Тесты структуры BatchRecord ├── test_kafka_history.py # Тесты KafkaBatchHistory ├── test_service.py # Тесты GeneratorService +├── test_service_cleanup.py # Контракт разбиения сервиса на модули └── test_state.py # Тесты GeneratorState и KafkaStateManager ``` diff --git a/generator/generator.py b/generator/generator.py index 1f88420..ba5578d 100644 --- a/generator/generator.py +++ b/generator/generator.py @@ -1,1175 +1,74 @@ #!/usr/bin/env python3 """ -Автономный генератор событий для Kafka (MVP rev5). +Совместимый фасад и точка входа генератора. -Режим 'steady-stream': публикуем постепенно, короткими тиками (1-10 сек), -держим целевую интенсивность events/min без крупных минутных batch. +Основной код разнесён по модулям рядом с этим файлом. Старые импорты вида +`from generator import Config` сохраняются для тестов и внешних запусков. """ -import json import logging -import math -import os -import random import sys -import time -import uuid -from dataclasses import dataclass, field -from datetime import datetime, timedelta, timezone from pathlib import Path -from typing import Any - -# Prometheus метрики -from prometheus_client import Counter, Gauge, Histogram, start_http_server - -# Kafka импортируем lazy для возможности тестирования без Kafka -_kafka_imported = False -KafkaProducer = None -KafkaError = None -def _import_kafka(): - global _kafka_imported, KafkaProducer, KafkaError - if not _kafka_imported: - from kafka import KafkaProducer - from kafka.errors import KafkaError - _kafka_imported = True - return KafkaProducer, KafkaError +SRC_DIR = Path(__file__).resolve().parent / "src" +if str(SRC_DIR) not in sys.path: + sys.path.insert(0, str(SRC_DIR)) + +from clickstream_generator.config import Config +from clickstream_generator.dictionary import EventDictionary +from clickstream_generator.generation import EventGenerator +from clickstream_generator.intensity import calculate_events_count, hour_factor +from clickstream_generator.kafka_io import ( + BatchRecord, + KafkaBatchHistory, + KafkaPublisher, + KafkaStateManager, + _import_kafka, + _with_retry, + ensure_topics, +) +from clickstream_generator.metrics import ( + METRICS_ERRORS_TOTAL, + METRICS_EVENTS_TOTAL, + METRICS_LAST_SUCCESS, + METRICS_TICK_DURATION, +) +from clickstream_generator.runtime import generate_tick_batch +from clickstream_generator.service import GeneratorService, main +from clickstream_generator.state import GeneratorState, _nested_list_to_tuple -# --------------------------------------------------------------------------- -# Настройка логирования -# --------------------------------------------------------------------------- logging.basicConfig( level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s", datefmt="%Y-%m-%d %H:%M:%S", ) -logger = logging.getLogger("generator") - -# --------------------------------------------------------------------------- -# Prometheus метрики -# --------------------------------------------------------------------------- -METRICS_EVENTS_TOTAL = Counter( - "generator_events_total", - "Total number of events sent to Kafka", - ["topic"] -) -METRICS_ERRORS_TOTAL = Counter( - "generator_publish_errors_total", - "Total number of publish errors", - ["topic"] -) -METRICS_TICK_DURATION = Histogram( - "generator_tick_duration_seconds", - "Duration of generator tick in seconds" -) -METRICS_LAST_SUCCESS = Gauge( - "generator_last_success_timestamp", - "Unix timestamp of last successful tick" -) - - -PAGE_START_DISTRIBUTION = [ - ("/home", 0.59), - ("/product_a", 0.20), - ("/product_b", 0.14), - ("/cart", 0.04), - ("/payment", 0.02), - ("/confirmation", 0.01), +__all__ = [ + "BatchRecord", + "Config", + "EventDictionary", + "EventGenerator", + "GeneratorService", + "GeneratorState", + "KafkaBatchHistory", + "KafkaPublisher", + "KafkaStateManager", + "METRICS_ERRORS_TOTAL", + "METRICS_EVENTS_TOTAL", + "METRICS_LAST_SUCCESS", + "METRICS_TICK_DURATION", + "_import_kafka", + "_nested_list_to_tuple", + "_with_retry", + "calculate_events_count", + "ensure_topics", + "generate_tick_batch", + "hour_factor", + "main", ] -PAGE_TRANSITIONS = { - "/home": [ - ("/home", 0.40), - ("/product_a", 0.28), - ("/product_b", 0.18), - ("/cart", 0.04), - (None, 0.10), - ], - "/product_a": [ - ("/home", 0.30), - ("/product_a", 0.16), - ("/product_b", 0.18), - ("/cart", 0.27), - (None, 0.09), - ], - "/product_b": [ - ("/home", 0.32), - ("/product_a", 0.15), - ("/product_b", 0.18), - ("/cart", 0.27), - (None, 0.08), - ], - "/cart": [ - ("/home", 0.20), - ("/product_a", 0.12), - ("/product_b", 0.10), - ("/cart", 0.10), - ("/payment", 0.42), - (None, 0.06), - ], - "/payment": [ - ("/home", 0.12), - ("/cart", 0.24), - ("/payment", 0.12), - ("/confirmation", 0.38), - (None, 0.14), - ], - "/confirmation": [ - ("/home", 0.35), - ("/product_a", 0.15), - ("/product_b", 0.10), - ("/confirmation", 0.05), - (None, 0.35), - ], -} - - -# --------------------------------------------------------------------------- -# Конфигурация через env -# --------------------------------------------------------------------------- -@dataclass(frozen=True) -class Config: - """Конфигурация генератора из переменных окружения.""" - - # Подключение к Kafka - kafka_bootstrap_servers: str = field( - default_factory=lambda: os.getenv("KAFKA_BOOTSTRAP_SERVERS", "kafka:29092") - ) - - # Параметры генерации - tick_seconds: int = field( - default_factory=lambda: int(os.getenv("GEN_TICK_SECONDS", "5")) - ) - lambda_base_per_min: int = field( - default_factory=lambda: int(os.getenv("GEN_LAMBDA_BASE_PER_MIN", "200")) - ) - jitter_pct: int = field( - default_factory=lambda: int(os.getenv("GEN_JITTER_PCT", "20")) - ) - min_events_per_tick: int = field( - default_factory=lambda: int(os.getenv("GEN_MIN_EVENTS_PER_TICK", "5")) - ) - max_events_per_tick: int = field( - default_factory=lambda: int(os.getenv("GEN_MAX_EVENTS_PER_TICK", "50")) - ) - max_session_events: int = field( - default_factory=lambda: int(os.getenv("GEN_MAX_SESSION_EVENTS", "30")) - ) - - # Пути к данным - data_dir: Path = field( - default_factory=lambda: Path(os.getenv("GEN_DATA_DIR", "/data")) - ) - - # Сид для воспроизводимости - seed: int | None = field( - default_factory=lambda: int(os.getenv("GEN_SEED")) - if os.getenv("GEN_SEED") - else None - ) - - # Включение/выключение генерации - enabled: bool = field( - default_factory=lambda: os.getenv("GEN_ENABLED", "true").lower() == "true" - ) - - # Порт для Prometheus метрик - metrics_port: int = field( - default_factory=lambda: int(os.getenv("GEN_METRICS_PORT", "9109")) - ) - - # Управление сохранением состояния - state_enabled: bool = field( - default_factory=lambda: os.getenv("GEN_STATE_ENABLED", "true").lower() == "true" - ) - state_reset: bool = field( - default_factory=lambda: os.getenv("GEN_STATE_RESET", "false").lower() == "true" - ) - - def __post_init__(self): - # Валидация параметров - if self.tick_seconds < 1: - raise ValueError("GEN_TICK_SECONDS must be >= 1") - if self.lambda_base_per_min < 1: - raise ValueError("GEN_LAMBDA_BASE_PER_MIN must be >= 1") - if self.max_session_events < 1: - raise ValueError("GEN_MAX_SESSION_EVENTS must be >= 1") - if not self.data_dir.exists(): - raise ValueError(f"Data directory does not exist: {self.data_dir}") - - -# --------------------------------------------------------------------------- -# Загрузка базового словаря событий -# --------------------------------------------------------------------------- -@dataclass -class EventDictionary: - """Базовый словарь событий из JSONL файлов.""" - - browser_events: list[dict[str, Any]] - location_events: list[dict[str, Any]] - device_events: list[dict[str, Any]] - geo_events: list[dict[str, Any]] - - # Индексы для быстрого поиска - browser_by_click_id: dict[str, list[dict]] = field(default_factory=dict) - location_by_event_id: dict[str, dict] = field(default_factory=dict) - device_by_click_id: dict[str, dict] = field(default_factory=dict) - geo_by_click_id: dict[str, dict] = field(default_factory=dict) - - def __post_init__(self): - # Строим индексы для связности - for browser in self.browser_events: - self.browser_by_click_id.setdefault(browser["click_id"], []).append(browser) - for loc in self.location_events: - self.location_by_event_id[loc["event_id"]] = loc - for dev in self.device_events: - self.device_by_click_id[dev["click_id"]] = dev - for geo in self.geo_events: - self.geo_by_click_id[geo["click_id"]] = geo - - @classmethod - def load(cls, data_dir: Path) -> "EventDictionary": - """Загружает события из JSONL файлов.""" - logger.info(f"Loading event dictionary from {data_dir}") - - def load_jsonl(filename: str) -> list[dict]: - path = data_dir / filename - events = [] - with open(path, "r", encoding="utf-8") as f: - for line in f: - line = line.strip() - if line: - events.append(json.loads(line)) - logger.info(f" Loaded {len(events)} events from {filename}") - return events - - browser_events = load_jsonl("browser_events.jsonl") - location_events = load_jsonl("location_events.jsonl") - device_events = load_jsonl("device_events.jsonl") - geo_events = load_jsonl("geo_events.jsonl") - - if not browser_events: - raise ValueError("browser_events.jsonl is empty or missing") - - return cls( - browser_events=browser_events, - location_events=location_events, - device_events=device_events, - geo_events=geo_events, - ) - - -# --------------------------------------------------------------------------- -# Генерация событий -# --------------------------------------------------------------------------- -class EventGenerator: - """Генератор событий с сохранением связности.""" - - def __init__(self, dictionary: EventDictionary, config: Config): - self.dictionary = dictionary - self.config = config - self.rng = random.Random(config.seed) - - def _new_uuid(self) -> str: - """Генерирует новый UUID.""" - return str(uuid.uuid4()) - - def _current_timestamp(self) -> str: - """Возвращает текущую метку времени в формате JSONL.""" - return datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S.%f") - - def _format_timestamp(self, timestamp: datetime) -> str: - """Форматирует запланированную метку времени для JSONL.""" - return timestamp.strftime("%Y-%m-%d %H:%M:%S.%f") - - def _weighted_choice(self, choices: list[tuple[Any, float]]) -> Any: - """Разыгрывает значение по списку весов.""" - point = self.rng.random() - cumulative = 0.0 - for value, weight in choices: - cumulative += weight - if point < cumulative: - return value - return choices[-1][0] - - def _generate_visit_path(self, max_events: int, min_events: int = 1) -> list[str]: - """Генерирует путь визита по страницам с защитой от бесконечных петель.""" - page = self._weighted_choice(PAGE_START_DISTRIBUTION) - path = [] - - while page is not None and len(path) < max_events: - path.append(page) - page = self._weighted_choice(PAGE_TRANSITIONS[page]) - - while len(path) < min_events and len(path) < max_events: - transitions = [ - (next_page, weight) - for next_page, weight in PAGE_TRANSITIONS[path[-1]] - if next_page is not None - ] - path.append(self._weighted_choice(transitions)) - - return path - - def _visit_pause_seconds(self) -> float: - """Разыгрывает паузу между событиями визита.""" - pause = self.rng.lognormvariate(math.log(45.0), 1.0) - return max(1.0, min(pause, 29 * 60.0)) - - def _hour_factor(self) -> float: - """Возвращает коэффициент интенсивности в зависимости от часа дня.""" - hour = datetime.now(timezone.utc).hour - # Дневное окно (9-18): 1.2 - # Ночное окно (0-5): 0.7 - # Остальное время: 1.0 - if 9 <= hour <= 18: - return 1.2 - elif 0 <= hour <= 5: - return 0.7 - return 1.0 - - def _calculate_events_count(self) -> int: - """Вычисляет количество событий для текущего тика (Poisson + jitter).""" - # Базовая интенсивность с учётом часа - lambda_minute = self.config.lambda_base_per_min * self._hour_factor() - - # Масштабируем на длительность тика - lambda_tick = lambda_minute * (self.config.tick_seconds / 60.0) - - # Генерируем Poisson - count = 0 - L = math.exp(-lambda_tick) - p = 1.0 - while p > L: - p *= self.rng.random() - count += 1 - count -= 1 - - # Применяем jitter (вариативность) - if self.config.jitter_pct > 0: - jitter_factor = 1.0 + self.rng.uniform( - -self.config.jitter_pct / 100.0, - self.config.jitter_pct / 100.0 - ) - count = int(count * jitter_factor) - - # Применяем границы - count = max(self.config.min_events_per_tick, min(count, self.config.max_events_per_tick)) - - return count - - def generate_batch(self, batch_size: int) -> dict[str, list[dict]]: - """ - Генерирует один визит с сохранением связей. - - Возвращает словарь {topic: [events]} - """ - if not self.dictionary.browser_events: - logger.warning("Event dictionary is empty, skipping batch generation") - return { - "browser_events": [], - "location_events": [], - "device_events": [], - "geo_events": [], - } - - batch = { - "browser_events": [], - "location_events": [], - "device_events": [], - "geo_events": [], - } - - if batch_size <= 0: - return batch - - max_visit_events = min(batch_size, self.config.max_session_events) - min_visit_events = min(2, max_visit_events) - visit_path = self._generate_visit_path(max_visit_events, min_visit_events) - - visit_candidates = [ - click_id for click_id, browser_events in self.dictionary.browser_by_click_id.items() - if ( - len(browser_events) >= len(visit_path) - and click_id in self.dictionary.device_by_click_id - and click_id in self.dictionary.geo_by_click_id - and all( - event["event_id"] in self.dictionary.location_by_event_id - for event in browser_events[:len(visit_path)] - ) - ) - ] - if visit_candidates: - base_click_id = self.rng.choice(visit_candidates) - base_browser_events = self.dictionary.browser_by_click_id[base_click_id][:len(visit_path)] - else: - # Крайний случай для очень малого сида: сохраняем форму визита, - # даже если приходится брать события с повторением. - base_browser = self.rng.choice(self.dictionary.browser_events) - base_click_id = base_browser["click_id"] - base_browser_events = [base_browser for _ in range(len(visit_path))] - - base_device = self.dictionary.device_by_click_id.get(base_click_id) - base_geo = self.dictionary.geo_by_click_id.get(base_click_id) - new_click_id = self._new_uuid() - - planned_timestamp = datetime.now(timezone.utc) - - for base_browser, page_url_path in zip(base_browser_events, visit_path): - base_location = self.dictionary.location_by_event_id.get(base_browser["event_id"]) - - # Генерируем новые ID - new_event_id = self._new_uuid() - new_timestamp = self._format_timestamp(planned_timestamp) - - # Создаём новое браузерное событие - browser_event = { - **base_browser, - "event_id": new_event_id, - "click_id": new_click_id, - "event_timestamp": new_timestamp, - } - batch["browser_events"].append(browser_event) - - # Связанное location событие - if base_location: - location_event = { - **base_location, - "event_id": new_event_id, - "page_url": f"http://www.dummywebsite.com{page_url_path}", - "page_url_path": page_url_path, - } - batch["location_events"].append(location_event) - - # Связанное device событие - if base_device: - device_event = { - **base_device, - "click_id": new_click_id, - } - batch["device_events"].append(device_event) - - # Связанное geo событие - if base_geo: - geo_event = { - **base_geo, - "click_id": new_click_id, - } - batch["geo_events"].append(geo_event) - - planned_timestamp += timedelta(seconds=self._visit_pause_seconds()) - - return batch - - def generate_tick_batch(self, event_budget: int) -> dict[str, list[dict]]: - """ - Генерирует батч тика из одного или нескольких визитов. - - `generate_batch()` остаётся публичным срезом одного визита. Для тика - нужно набрать рассчитанный бюджет событий, поэтому здесь несколько - визитов объединяются в один набор записей для публикации. - """ - tick_batch = { - "browser_events": [], - "location_events": [], - "device_events": [], - "geo_events": [], - } - - remaining_events = event_budget - while remaining_events > 0: - visit_batch = self.generate_batch(remaining_events) - generated_events = len(visit_batch["browser_events"]) - if generated_events == 0: - break - - for topic, events in visit_batch.items(): - tick_batch[topic].extend(events) - remaining_events -= generated_events - - return tick_batch - - -# --------------------------------------------------------------------------- -# Batch record для истории -# --------------------------------------------------------------------------- -@dataclass -class BatchRecord: - """Запись об отправленном батче.""" - - batch_id: str - started_at: datetime - finished_at: datetime - sent_total: int - sent_browser: int - sent_location: int - sent_device: int - sent_geo: int - status: str # 'success', 'partial', 'error' - error_message: str | None = None - - def to_dict(self) -> dict: - """Конвертирует в словарь для сериализации.""" - return { - "batch_id": self.batch_id, - "started_at": self.started_at.isoformat(), - "finished_at": self.finished_at.isoformat(), - "sent_total": self.sent_total, - "sent_browser": self.sent_browser, - "sent_location": self.sent_location, - "sent_device": self.sent_device, - "sent_geo": self.sent_geo, - "status": self.status, - "error_message": self.error_message, - } - - -# --------------------------------------------------------------------------- -# State record для сохранения состояния генератора -# --------------------------------------------------------------------------- -def _nested_list_to_tuple(obj): - """Рекурсивно преобразует list в tuple (для восстановления RNG state после JSON).""" - if isinstance(obj, list): - return tuple(_nested_list_to_tuple(x) for x in obj) - return obj - - -@dataclass -class GeneratorState: - """Состояние генератора для восстановления после рестарта. - - Используем JSON-safe сериализацию: - - rng_state от random.getstate() - кортеж из простых типов (int, tuple), - безопасно сериализуется в JSON напрямую без pickle - """ - - tick: int - rng_state: tuple # результат random.getstate() - JSON-serializable - last_batch_id: str - last_timestamp: datetime - version: str = "1.0" - - def to_dict(self) -> dict: - """Конвертирует в словарь для JSON-сериализации.""" - return { - "tick": self.tick, - "rng_state": self.rng_state, # tuple из int - JSON-serializable - "last_batch_id": self.last_batch_id, - "last_timestamp": self.last_timestamp.isoformat(), - "version": self.version, - } - - @classmethod - def from_dict(cls, data: dict) -> "GeneratorState": - """Создаёт состояние из словаря (JSON-only, без pickle). - - При невалидном rng_state логирует предупреждение и возвращает None - (вызывающий код должен обработать как "начать с чистого листа"). - """ - try: - rng_state_raw = data.get("rng_state") - if not rng_state_raw: - logger.warning("State missing rng_state field") - raise ValueError("rng_state is missing") - - rng_state = _nested_list_to_tuple(rng_state_raw) - - # Валидация: rng_state должен быть tuple и иметь минимальную структуру - if not isinstance(rng_state, tuple): - logger.warning(f"rng_state is not tuple: {type(rng_state)}") - raise ValueError("rng_state must be tuple") - - if len(rng_state) < 2: - logger.warning(f"rng_state has insufficient length: {len(rng_state)}") - raise ValueError("rng_state has insufficient length") - - # Проверка что можем создать RNG и вызвать setstate (тестовая валидация) - test_rng = random.Random() - test_rng.setstate(rng_state) - - return cls( - tick=data.get("tick", 0), - rng_state=rng_state, - last_batch_id=data.get("last_batch_id", ""), - last_timestamp=datetime.fromisoformat(data.get("last_timestamp", "1970-01-01T00:00:00+00:00")), - version=data.get("version", "1.0"), - ) - except Exception as e: - logger.warning(f"Invalid state format, will start fresh: {e}") - raise ValueError(f"Invalid state: {e}") - - @classmethod - def from_dict_safe(cls, data: dict) -> "GeneratorState | None": - """Безопасная загрузка state с graceful degradation. - - Returns: - GeneratorState если данные валидны, иначе None (начать с чистого листа). - """ - try: - return cls.from_dict(data) - except Exception: - # Уже залогировано в from_dict - return None - - -# --------------------------------------------------------------------------- -# Kafka state manager - сохраняет/восстанавливает состояние генератора -# --------------------------------------------------------------------------- -class KafkaStateManager: - """Управление состоянием генератора в Kafka (compact topic).""" - - STATE_TOPIC = "generator_state" - STATE_KEY = "default" # Для возможности нескольких генераторов в будущем - - def __init__(self, bootstrap_servers: str): - self.bootstrap_servers = bootstrap_servers - KafkaProducerCls, _ = _import_kafka() - - logger.info(f"Connecting to Kafka for state management at {self.bootstrap_servers}") - self.producer = KafkaProducerCls( - bootstrap_servers=self.bootstrap_servers, - value_serializer=lambda v: json.dumps(v).encode("utf-8"), - key_serializer=lambda k: k.encode("utf-8") if k else None, - retries=3, - retry_backoff_ms=1000, - ) - logger.info("Connected to Kafka for state management successfully") - - def save(self, state: GeneratorState) -> None: - """Сохраняет состояние в топик (compact topic - только последнее значение). - - Использует retry при сбоях подключения к Kafka. - """ - value = state.to_dict() - - def _do_send(): - self.producer.send(self.STATE_TOPIC, key=self.STATE_KEY, value=value) - - _with_retry(_do_send, max_retries=3, base_delay=0.5) - - def flush(self) -> None: - """Сбрасывает буфер с retry.""" - def _do_flush(): - self.producer.flush() - - _with_retry(_do_flush, max_retries=3, base_delay=0.5) - - def close(self) -> None: - """Закрывает соединение.""" - try: - self.producer.close() - except Exception as e: - logger.debug(f"Error closing producer (ignored): {e}") - - def load(self) -> GeneratorState | None: - """Загружает последнее состояние из топика. - - Для compact topic хранится только последнее значение для ключа, - поэтому читаем все сообщения и берём последнее с нужным ключом. - Использует retry при сбоях подключения к Kafka. - """ - from kafka import KafkaConsumer - - logger.info(f"Loading state from topic {self.STATE_TOPIC}") - - def _do_load(): - consumer = KafkaConsumer( - self.STATE_TOPIC, - bootstrap_servers=self.bootstrap_servers, - auto_offset_reset="earliest", - enable_auto_commit=False, - consumer_timeout_ms=5000, - value_deserializer=lambda v: json.loads(v.decode("utf-8")), - ) - - last_state = None - for message in consumer: - if message.key and message.key.decode("utf-8") == self.STATE_KEY: - last_state = message.value - - consumer.close() - return last_state - - try: - last_state = _with_retry(_do_load, max_retries=3, base_delay=0.5) - - if last_state: - logger.info(f"Restored state: tick={last_state.get('tick')}, " - f"last_batch_id={last_state.get('last_batch_id')}") - # Используем from_dict_safe для graceful degradation при битом state - restored = GeneratorState.from_dict_safe(last_state) - if restored is None: - logger.warning("State data was invalid, starting fresh") - return restored - else: - logger.info("No previous state found, starting fresh") - return None - - except Exception as e: - logger.warning(f"Failed to load state: {e}, starting fresh") - return None - - -# --------------------------------------------------------------------------- -# Утилиты для работы с Kafka с retry/backoff -# --------------------------------------------------------------------------- -def _with_retry(operation, max_retries: int = 5, base_delay: float = 1.0, max_delay: float = 30.0): - """Выполняет операцию с экспоненциальным backoff и ограниченным числом попыток. - - Args: - operation: функция для выполнения - max_retries: максимальное число попыток - base_delay: начальная задержка между попытками (сек) - max_delay: максимальная задержка между попытками (сек) - - Returns: - результат операции - - Raises: - последнее исключение после исчерпания попыток - """ - import time - - last_exception = None - for attempt in range(max_retries): - try: - return operation() - except Exception as e: - last_exception = e - if attempt < max_retries - 1: - delay = min(base_delay * (2 ** attempt), max_delay) - logger.warning(f"Operation failed (attempt {attempt + 1}/{max_retries}): {e}. Retrying in {delay:.1f}s...") - time.sleep(delay) - else: - logger.error(f"Operation failed after {max_retries} attempts: {e}") - raise last_exception - - -# --------------------------------------------------------------------------- -# Создание топиков для истории и состояния -# --------------------------------------------------------------------------- -def ensure_topics(bootstrap_servers: str) -> None: - """Создаёт необходимые топики если они не существуют (с retry на подключение).""" - from kafka import KafkaAdminClient - from kafka.admin import NewTopic - from kafka.errors import TopicAlreadyExistsError - - def _create_topics(): - admin_client = KafkaAdminClient(bootstrap_servers=bootstrap_servers) - try: - # Топик для истории батчей (обычный, с retention) - history_topic = NewTopic( - name=KafkaBatchHistory.HISTORY_TOPIC, - num_partitions=1, - replication_factor=1, - ) - - # Топик для состояния (compact - храним только последнее значение) - state_topic = NewTopic( - name=KafkaStateManager.STATE_TOPIC, - num_partitions=1, - replication_factor=1, - topic_configs={ - "cleanup.policy": "compact", - "min.cleanable.dirty.ratio": "0.1", - "delete.retention.ms": "100", - }, - ) - - for topic in [history_topic, state_topic]: - try: - admin_client.create_topics([topic]) - logger.info(f"Created topic: {topic.name}") - except TopicAlreadyExistsError: - logger.debug(f"Topic already exists: {topic.name}") - finally: - admin_client.close() - - _with_retry(_create_topics, max_retries=5, base_delay=1.0) - - -# --------------------------------------------------------------------------- -# Kafka history - пишет историю в отдельный топик -# --------------------------------------------------------------------------- -class KafkaBatchHistory: - """Хранение истории batch в Kafka (отдельный топик) с retry.""" - - HISTORY_TOPIC = "generator_batch_history" - - def __init__(self, bootstrap_servers: str): - self.bootstrap_servers = bootstrap_servers - self.producer = None - self._connect() - - def _connect(self): - """Устанавливает соединение с Kafka с retry.""" - KafkaProducerCls, KafkaErrorCls = _import_kafka() - - def _do_connect(): - logger.info(f"Connecting to Kafka for history at {self.bootstrap_servers}") - self.producer = KafkaProducerCls( - bootstrap_servers=self.bootstrap_servers, - value_serializer=lambda v: json.dumps(v).encode("utf-8"), - key_serializer=lambda k: k.encode("utf-8") if k else None, - retries=3, - retry_backoff_ms=1000, - ) - logger.info("Connected to Kafka for history successfully") - - _with_retry(_do_connect, max_retries=5, base_delay=1.0) - - def add(self, record: BatchRecord): - """Добавляет запись в историю (топик Kafka) с retry и реконнектом.""" - key = record.batch_id - value = record.to_dict() - - def _do_send(): - self.producer.send(self.HISTORY_TOPIC, key=key, value=value) - - try: - _with_retry(_do_send, max_retries=3, base_delay=0.5) - except Exception as e: - logger.warning(f"Failed to send history record after retries: {e}, attempting reconnect") - self._connect() - # Повторная попытка после реконнекта - _with_retry(_do_send, max_retries=2, base_delay=0.5) - - def flush(self): - """Сбрасывает буфер с retry.""" - def _do_flush(): - self.producer.flush() - - _with_retry(_do_flush, max_retries=3, base_delay=0.5) - - def close(self): - """Закрывает соединение.""" - try: - if self.producer: - self.producer.close() - except Exception as e: - logger.debug(f"Error closing history producer (ignored): {e}") - - -# --------------------------------------------------------------------------- -# Kafka publisher -# --------------------------------------------------------------------------- -class KafkaPublisher: - """Публикация событий в Kafka с retry и реконнектом.""" - - def __init__(self, bootstrap_servers: str): - self.bootstrap_servers = bootstrap_servers - self.producer = None - self._connect() - - def _connect(self): - """Устанавливает соединение с Kafka с retry.""" - KafkaProducerCls, KafkaErrorCls = _import_kafka() - - def _do_connect(): - logger.info(f"Connecting to Kafka at {self.bootstrap_servers}") - self.producer = KafkaProducerCls( - bootstrap_servers=self.bootstrap_servers, - value_serializer=lambda v: json.dumps(v).encode("utf-8"), - key_serializer=lambda k: k.encode("utf-8") if k else None, - batch_size=16384, - linger_ms=100, - retries=3, - retry_backoff_ms=1000, - ) - logger.info("Connected to Kafka successfully") - - _with_retry(_do_connect, max_retries=5, base_delay=1.0) - - def _publish_with_retry(self, topic: str, events: list[dict]) -> tuple[int, int]: - """Внутренняя функция публикации с retry на уровне batch.""" - if not self.producer: - raise RuntimeError("Producer not connected") - - sent = 0 - errors = 0 - futures = [] - - for event in events: - key = event.get("event_id") or event.get("click_id") - try: - future = self.producer.send(topic, key=key, value=event) - futures.append(future) - except Exception as e: - logger.error(f"Failed to send message to {topic}: {e}") - errors += 1 - METRICS_ERRORS_TOTAL.labels(topic=topic).inc() - - # Ждём подтверждений - for future in futures: - try: - future.get(timeout=10) - sent += 1 - METRICS_EVENTS_TOTAL.labels(topic=topic).inc() - except Exception as e: - logger.error(f"Failed to confirm message delivery: {e}") - errors += 1 - METRICS_ERRORS_TOTAL.labels(topic=topic).inc() - - return sent, errors - - def publish(self, topic: str, events: list[dict]) -> tuple[int, int]: - """ - Публикует события в топик с retry и автоматическим реконнектом. - - Returns: - (sent_count, error_count) - """ - def _do_publish(): - return self._publish_with_retry(topic, events) - - try: - return _with_retry(_do_publish, max_retries=3, base_delay=0.5) - except Exception as e: - logger.warning(f"Publish failed after retries: {e}, attempting reconnect") - self._connect() - # Повторная попытка после реконнекта - return _with_retry(_do_publish, max_retries=2, base_delay=0.5) - - def flush(self): - """Сбрасывает буфер с retry.""" - def _do_flush(): - if self.producer: - self.producer.flush() - - _with_retry(_do_flush, max_retries=3, base_delay=0.5) - - def close(self): - """Закрывает соединение.""" - try: - if self.producer: - self.producer.close() - except Exception as e: - logger.debug(f"Error closing producer (ignored): {e}") - - -# --------------------------------------------------------------------------- -# Основной цикл генератора -# --------------------------------------------------------------------------- -class GeneratorService: - """Основной сервис генератора.""" - - def __init__(self, config: Config): - self.config = config - self.dictionary = EventDictionary.load(config.data_dir) - self.generator = EventGenerator(self.dictionary, config) - self.publisher: KafkaPublisher | None = None - self.history: KafkaBatchHistory | None = None - self.state_manager: KafkaStateManager | None = None - self._running = False - self._tick = 0 # Текущий номер тика (восстанавливается из стейта) - - def start(self): - """Запускает основной цикл.""" - if not self.config.enabled: - logger.warning("Generator is disabled (GEN_ENABLED=false)") - return - - # Запускаем HTTP-сервер для Prometheus метрик - logger.info(f"Starting metrics server on port {self.config.metrics_port}") - start_http_server(self.config.metrics_port) - - logger.info("Starting generator service...") - logger.info(f"Configuration: tick={self.config.tick_seconds}s, " - f"lambda_base={self.config.lambda_base_per_min}/min, " - f"jitter={self.config.jitter_pct}%, " - f"state_enabled={self.config.state_enabled}, " - f"state_reset={self.config.state_reset}") - - # Подключаемся к Kafka для публикации событий и истории - # Сначала создаём топики если нужно - ensure_topics(self.config.kafka_bootstrap_servers) - self.publisher = KafkaPublisher(self.config.kafka_bootstrap_servers) - self.history = KafkaBatchHistory(self.config.kafka_bootstrap_servers) - - # Инициализируем state manager если включено - if self.config.state_enabled: - self.state_manager = KafkaStateManager(self.config.kafka_bootstrap_servers) - - # Восстанавливаем стейт если не требуется сброс - if not self.config.state_reset: - restored_state = self.state_manager.load() - if restored_state: - self._tick = restored_state.tick - self.generator.rng.setstate(restored_state.rng_state) - logger.info(f"Restored state: continuing from tick {self._tick}, " - f"last_batch_id={restored_state.last_batch_id}") - else: - logger.info("State reset requested, starting fresh") - else: - logger.info("State management disabled") - - self._running = True - - try: - self._main_loop() - except KeyboardInterrupt: - logger.info("Received shutdown signal") - finally: - self.stop() - - def stop(self): - """Останавливает сервис.""" - logger.info("Stopping generator service...") - self._running = False - if self.publisher: - self.publisher.close() - if self.history: - self.history.close() - if self.state_manager: - self.state_manager.close() - - def _save_state(self, batch_id: str) -> None: - """Сохраняет текущее состояние генератора.""" - if not self.state_manager or not self.config.state_enabled: - return - - try: - state = GeneratorState( - tick=self._tick, - rng_state=self.generator.rng.getstate(), - last_batch_id=batch_id, - last_timestamp=datetime.now(timezone.utc), - ) - self.state_manager.save(state) - self.state_manager.flush() - logger.debug(f"Saved state: tick={self._tick}, batch_id={batch_id}") - except Exception as e: - logger.warning(f"Failed to save state: {e}") - METRICS_ERRORS_TOTAL.labels(topic="state").inc() - - def _main_loop(self): - """Основной цикл тиков.""" - while self._running: - self._tick += 1 - tick_start = time.time() - batch_id = str(uuid.uuid4())[:8] - - with METRICS_TICK_DURATION.time(): - logger.info(f"=== Tick {self._tick} (batch_id={batch_id}) ===") - - try: - # Вычисляем количество событий - events_count = self.generator._calculate_events_count() - logger.info(f"Generating ~{events_count} base events") - - # Генерируем тиковый батч из одного или нескольких визитов - gen_start = time.time() - batch = self.generator.generate_tick_batch(events_count) - gen_duration = time.time() - gen_start - - # Публикуем в Kafka - pub_start = time.time() - total_sent = 0 - total_errors = 0 - - sent_counts = {} - for topic, events in batch.items(): - if events: - sent, errors = self.publisher.publish(topic, events) - sent_counts[topic] = {"sent": sent, "errors": errors} - total_sent += sent - total_errors += errors - - # Определяем статус - if total_errors == 0: - status = "success" - elif total_sent > 0: - status = "partial" - else: - status = "error" - - # Обновляем метрику последнего успешного тика и сохраняем стейт - if status in ("success", "partial"): - METRICS_LAST_SUCCESS.set_to_current_time() - self._save_state(batch_id) - - # Флашим публикацию событий - self.publisher.flush() - pub_duration = time.time() - pub_start - - # Сохраняем в историю (best-effort: ошибки не валят тик) - try: - batch_record = BatchRecord( - batch_id=batch_id, - started_at=datetime.fromtimestamp(tick_start, tz=timezone.utc), - finished_at=datetime.now(timezone.utc), - sent_total=total_sent, - sent_browser=sent_counts.get("browser_events", {}).get("sent", 0), - sent_location=sent_counts.get("location_events", {}).get("sent", 0), - sent_device=sent_counts.get("device_events", {}).get("sent", 0), - sent_geo=sent_counts.get("geo_events", {}).get("sent", 0), - status=status, - error_message=None if status == "success" else f"Errors: {total_errors}", - ) - self.history.add(batch_record) - self.history.flush() - except Exception as hist_err: - logger.warning(f"Failed to write batch history: {hist_err}") - METRICS_ERRORS_TOTAL.labels(topic="history").inc() - - # Логируем результат - tick_duration = time.time() - tick_start - logger.info( - f"Batch {batch_id} completed: " - f"sent={total_sent}, errors={total_errors}, " - f"gen_time={gen_duration:.3f}s, pub_time={pub_duration:.3f}s, " - f"total_time={tick_duration:.3f}s" - ) - - for topic, counts in sent_counts.items(): - if counts["sent"] > 0: - logger.info(f" {topic}: {counts['sent']} sent") - - except Exception as e: - logger.exception(f"Error in tick {self._tick}: {e}") - # Пытаемся записать ошибку в историю (best-effort) - try: - self.history.add( - BatchRecord( - batch_id=batch_id, - started_at=datetime.fromtimestamp(tick_start, tz=timezone.utc), - finished_at=datetime.now(timezone.utc), - sent_total=0, - sent_browser=0, - sent_location=0, - sent_device=0, - sent_geo=0, - status="error", - error_message=str(e), - ) - ) - self.history.flush() - except Exception as hist_err: - logger.warning(f"Failed to write error to history: {hist_err}") - - # Ждём до следующего тика - elapsed = time.time() - tick_start - sleep_time = max(0, self.config.tick_seconds - elapsed) - if sleep_time > 0: - logger.debug(f"Sleeping for {sleep_time:.1f}s until next tick") - time.sleep(sleep_time) - - -# --------------------------------------------------------------------------- -# Точка входа -# --------------------------------------------------------------------------- -def main(): - try: - config = Config() - service = GeneratorService(config) - service.start() - except Exception as e: - logger.exception(f"Fatal error: {e}") - sys.exit(1) - if __name__ == "__main__": main() diff --git a/generator/src/clickstream_generator/__init__.py b/generator/src/clickstream_generator/__init__.py new file mode 100644 index 0000000..7015c64 --- /dev/null +++ b/generator/src/clickstream_generator/__init__.py @@ -0,0 +1,2 @@ +"""Пакет генератора кликстрима.""" + diff --git a/generator/src/clickstream_generator/config.py b/generator/src/clickstream_generator/config.py new file mode 100644 index 0000000..bfba52c --- /dev/null +++ b/generator/src/clickstream_generator/config.py @@ -0,0 +1,62 @@ +"""Конфигурация генератора из переменных окружения.""" + +import os +from dataclasses import dataclass, field +from pathlib import Path + + +@dataclass(frozen=True) +class Config: + """Конфигурация генератора из переменных окружения.""" + + kafka_bootstrap_servers: str = field( + default_factory=lambda: os.getenv("KAFKA_BOOTSTRAP_SERVERS", "kafka:29092") + ) + tick_seconds: int = field( + default_factory=lambda: int(os.getenv("GEN_TICK_SECONDS", "5")) + ) + lambda_base_per_min: int = field( + default_factory=lambda: int(os.getenv("GEN_LAMBDA_BASE_PER_MIN", "200")) + ) + jitter_pct: int = field( + default_factory=lambda: int(os.getenv("GEN_JITTER_PCT", "20")) + ) + min_events_per_tick: int = field( + default_factory=lambda: int(os.getenv("GEN_MIN_EVENTS_PER_TICK", "5")) + ) + max_events_per_tick: int = field( + default_factory=lambda: int(os.getenv("GEN_MAX_EVENTS_PER_TICK", "50")) + ) + max_session_events: int = field( + default_factory=lambda: int(os.getenv("GEN_MAX_SESSION_EVENTS", "30")) + ) + data_dir: Path = field( + default_factory=lambda: Path(os.getenv("GEN_DATA_DIR", "/data")) + ) + seed: int | None = field( + default_factory=lambda: int(os.getenv("GEN_SEED")) + if os.getenv("GEN_SEED") + else None + ) + enabled: bool = field( + default_factory=lambda: os.getenv("GEN_ENABLED", "true").lower() == "true" + ) + metrics_port: int = field( + default_factory=lambda: int(os.getenv("GEN_METRICS_PORT", "9109")) + ) + state_enabled: bool = field( + default_factory=lambda: os.getenv("GEN_STATE_ENABLED", "true").lower() == "true" + ) + state_reset: bool = field( + default_factory=lambda: os.getenv("GEN_STATE_RESET", "false").lower() == "true" + ) + + def __post_init__(self): + if self.tick_seconds < 1: + raise ValueError("GEN_TICK_SECONDS must be >= 1") + if self.lambda_base_per_min < 1: + raise ValueError("GEN_LAMBDA_BASE_PER_MIN must be >= 1") + if self.max_session_events < 1: + raise ValueError("GEN_MAX_SESSION_EVENTS must be >= 1") + if not self.data_dir.exists(): + raise ValueError(f"Data directory does not exist: {self.data_dir}") diff --git a/generator/src/clickstream_generator/dictionary.py b/generator/src/clickstream_generator/dictionary.py new file mode 100644 index 0000000..d8d27d1 --- /dev/null +++ b/generator/src/clickstream_generator/dictionary.py @@ -0,0 +1,65 @@ +"""Загрузка и индексация исходных JSONL-событий.""" + +import json +import logging +from dataclasses import dataclass, field +from pathlib import Path +from typing import Any + + +logger = logging.getLogger("generator") + + +@dataclass +class EventDictionary: + """Базовый словарь событий из JSONL файлов.""" + + browser_events: list[dict[str, Any]] + location_events: list[dict[str, Any]] + device_events: list[dict[str, Any]] + geo_events: list[dict[str, Any]] + browser_by_click_id: dict[str, list[dict]] = field(default_factory=dict) + location_by_event_id: dict[str, dict] = field(default_factory=dict) + device_by_click_id: dict[str, dict] = field(default_factory=dict) + geo_by_click_id: dict[str, dict] = field(default_factory=dict) + + def __post_init__(self): + for browser in self.browser_events: + self.browser_by_click_id.setdefault(browser["click_id"], []).append(browser) + for loc in self.location_events: + self.location_by_event_id[loc["event_id"]] = loc + for dev in self.device_events: + self.device_by_click_id[dev["click_id"]] = dev + for geo in self.geo_events: + self.geo_by_click_id[geo["click_id"]] = geo + + @classmethod + def load(cls, data_dir: Path) -> "EventDictionary": + """Загружает события из JSONL файлов.""" + logger.info(f"Loading event dictionary from {data_dir}") + + def load_jsonl(filename: str) -> list[dict]: + path = data_dir / filename + events = [] + with open(path, "r", encoding="utf-8") as f: + for line in f: + line = line.strip() + if line: + events.append(json.loads(line)) + logger.info(f" Loaded {len(events)} events from {filename}") + return events + + browser_events = load_jsonl("browser_events.jsonl") + location_events = load_jsonl("location_events.jsonl") + device_events = load_jsonl("device_events.jsonl") + geo_events = load_jsonl("geo_events.jsonl") + + if not browser_events: + raise ValueError("browser_events.jsonl is empty or missing") + + return cls( + browser_events=browser_events, + location_events=location_events, + device_events=device_events, + geo_events=geo_events, + ) diff --git a/generator/src/clickstream_generator/generation.py b/generator/src/clickstream_generator/generation.py new file mode 100644 index 0000000..ce31802 --- /dev/null +++ b/generator/src/clickstream_generator/generation.py @@ -0,0 +1,211 @@ +"""Чистая генеративная модель визита.""" + +import math +import random +import uuid +from datetime import datetime, timedelta, timezone +from typing import Any + +from clickstream_generator.config import Config +from clickstream_generator.dictionary import EventDictionary +from clickstream_generator.intensity import calculate_events_count, hour_factor + + +PAGE_START_DISTRIBUTION = [ + ("/home", 0.59), + ("/product_a", 0.20), + ("/product_b", 0.14), + ("/cart", 0.04), + ("/payment", 0.02), + ("/confirmation", 0.01), +] + +PAGE_TRANSITIONS = { + "/home": [ + ("/home", 0.40), + ("/product_a", 0.28), + ("/product_b", 0.18), + ("/cart", 0.04), + (None, 0.10), + ], + "/product_a": [ + ("/home", 0.30), + ("/product_a", 0.16), + ("/product_b", 0.18), + ("/cart", 0.27), + (None, 0.09), + ], + "/product_b": [ + ("/home", 0.32), + ("/product_a", 0.15), + ("/product_b", 0.18), + ("/cart", 0.27), + (None, 0.08), + ], + "/cart": [ + ("/home", 0.20), + ("/product_a", 0.12), + ("/product_b", 0.10), + ("/cart", 0.10), + ("/payment", 0.42), + (None, 0.06), + ], + "/payment": [ + ("/home", 0.12), + ("/cart", 0.24), + ("/payment", 0.12), + ("/confirmation", 0.38), + (None, 0.14), + ], + "/confirmation": [ + ("/home", 0.35), + ("/product_a", 0.15), + ("/product_b", 0.10), + ("/confirmation", 0.05), + (None, 0.35), + ], +} + + +class EventGenerator: + """Генератор одного связанного визита.""" + + def __init__(self, dictionary: EventDictionary, config: Config): + self.dictionary = dictionary + self.config = config + self.rng = random.Random(config.seed) + + def _new_uuid(self) -> str: + """Генерирует новый UUID.""" + return str(uuid.uuid4()) + + def _current_timestamp(self) -> str: + """Возвращает текущую метку времени в формате JSONL.""" + return datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S.%f") + + def _format_timestamp(self, timestamp: datetime) -> str: + """Форматирует запланированную метку времени для JSONL.""" + return timestamp.strftime("%Y-%m-%d %H:%M:%S.%f") + + def _weighted_choice(self, choices: list[tuple[Any, float]]) -> Any: + """Разыгрывает значение по списку весов.""" + point = self.rng.random() + cumulative = 0.0 + for value, weight in choices: + cumulative += weight + if point < cumulative: + return value + return choices[-1][0] + + def _generate_visit_path(self, max_events: int, min_events: int = 1) -> list[str]: + """Генерирует путь визита по страницам с защитой от бесконечных петель.""" + page = self._weighted_choice(PAGE_START_DISTRIBUTION) + path = [] + + while page is not None and len(path) < max_events: + path.append(page) + page = self._weighted_choice(PAGE_TRANSITIONS[page]) + + while len(path) < min_events and len(path) < max_events: + transitions = [ + (next_page, weight) + for next_page, weight in PAGE_TRANSITIONS[path[-1]] + if next_page is not None + ] + path.append(self._weighted_choice(transitions)) + + return path + + def _visit_pause_seconds(self) -> float: + """Разыгрывает паузу между событиями визита.""" + pause = self.rng.lognormvariate(math.log(45.0), 1.0) + return max(1.0, min(pause, 29 * 60.0)) + + def _hour_factor(self) -> float: + """Совместимый wrapper над расчётом часового коэффициента.""" + return hour_factor() + + def _calculate_events_count(self) -> int: + """Совместимый wrapper над расчётом событийного бюджета.""" + return calculate_events_count(self.config, self.rng) + + def generate_batch(self, batch_size: int) -> dict[str, list[dict]]: + """Генерирует один визит с сохранением связей.""" + if not self.dictionary.browser_events: + return { + "browser_events": [], + "location_events": [], + "device_events": [], + "geo_events": [], + } + + batch = { + "browser_events": [], + "location_events": [], + "device_events": [], + "geo_events": [], + } + + if batch_size <= 0: + return batch + + max_visit_events = min(batch_size, self.config.max_session_events) + min_visit_events = min(2, max_visit_events) + visit_path = self._generate_visit_path(max_visit_events, min_visit_events) + + visit_candidates = [ + click_id for click_id, browser_events in self.dictionary.browser_by_click_id.items() + if ( + len(browser_events) >= len(visit_path) + and click_id in self.dictionary.device_by_click_id + and click_id in self.dictionary.geo_by_click_id + and all( + event["event_id"] in self.dictionary.location_by_event_id + for event in browser_events[:len(visit_path)] + ) + ) + ] + if visit_candidates: + base_click_id = self.rng.choice(visit_candidates) + base_browser_events = self.dictionary.browser_by_click_id[base_click_id][:len(visit_path)] + else: + base_browser = self.rng.choice(self.dictionary.browser_events) + base_click_id = base_browser["click_id"] + base_browser_events = [base_browser for _ in range(len(visit_path))] + + base_device = self.dictionary.device_by_click_id.get(base_click_id) + base_geo = self.dictionary.geo_by_click_id.get(base_click_id) + new_click_id = self._new_uuid() + planned_timestamp = datetime.now(timezone.utc) + + for base_browser, page_url_path in zip(base_browser_events, visit_path): + base_location = self.dictionary.location_by_event_id.get(base_browser["event_id"]) + new_event_id = self._new_uuid() + new_timestamp = self._format_timestamp(planned_timestamp) + + browser_event = { + **base_browser, + "event_id": new_event_id, + "click_id": new_click_id, + "event_timestamp": new_timestamp, + } + batch["browser_events"].append(browser_event) + + if base_location: + location_event = { + **base_location, + "event_id": new_event_id, + "page_url": f"http://www.dummywebsite.com{page_url_path}", + "page_url_path": page_url_path, + } + batch["location_events"].append(location_event) + + if base_device: + batch["device_events"].append({**base_device, "click_id": new_click_id}) + + if base_geo: + batch["geo_events"].append({**base_geo, "click_id": new_click_id}) + + planned_timestamp += timedelta(seconds=self._visit_pause_seconds()) + + return batch diff --git a/generator/src/clickstream_generator/intensity.py b/generator/src/clickstream_generator/intensity.py new file mode 100644 index 0000000..a09838b --- /dev/null +++ b/generator/src/clickstream_generator/intensity.py @@ -0,0 +1,41 @@ +"""Расчёт событийного бюджета тика.""" + +import math +import random +from datetime import datetime, timezone + +from clickstream_generator.config import Config + + +def hour_factor(now: datetime | None = None) -> float: + """Возвращает коэффициент интенсивности в зависимости от часа дня.""" + current = now or datetime.now(timezone.utc) + hour = current.hour + if 9 <= hour <= 18: + return 1.2 + if 0 <= hour <= 5: + return 0.7 + return 1.0 + + +def calculate_events_count(config: Config, rng: random.Random) -> int: + """Вычисляет количество событий для текущего тика (Poisson + jitter).""" + lambda_minute = config.lambda_base_per_min * hour_factor() + lambda_tick = lambda_minute * (config.tick_seconds / 60.0) + + count = 0 + threshold = math.exp(-lambda_tick) + product = 1.0 + while product > threshold: + product *= rng.random() + count += 1 + count -= 1 + + if config.jitter_pct > 0: + jitter_factor = 1.0 + rng.uniform( + -config.jitter_pct / 100.0, + config.jitter_pct / 100.0, + ) + count = int(count * jitter_factor) + + return max(config.min_events_per_tick, min(count, config.max_events_per_tick)) diff --git a/generator/src/clickstream_generator/kafka_io.py b/generator/src/clickstream_generator/kafka_io.py new file mode 100644 index 0000000..103cb20 --- /dev/null +++ b/generator/src/clickstream_generator/kafka_io.py @@ -0,0 +1,363 @@ +"""Kafka-интеграция генератора.""" + +import json +import logging +import sys +import time +from dataclasses import dataclass +from datetime import datetime + +from clickstream_generator.metrics import METRICS_ERRORS_TOTAL, METRICS_EVENTS_TOTAL +from clickstream_generator.state import GeneratorState + + +logger = logging.getLogger("generator") + +_kafka_imported = False +KafkaProducer = None +KafkaError = None + + +def _import_kafka(): + """Лениво импортирует Kafka-клиент, чтобы тесты могли работать без Kafka.""" + global _kafka_imported, KafkaProducer, KafkaError + if not _kafka_imported: + from kafka import KafkaProducer + from kafka.errors import KafkaError + _kafka_imported = True + return KafkaProducer, KafkaError + + +def _with_retry(operation, max_retries: int = 5, base_delay: float = 1.0, max_delay: float = 30.0): + """Выполняет операцию с экспоненциальным backoff.""" + last_exception = None + for attempt in range(max_retries): + try: + return operation() + except Exception as e: + last_exception = e + if attempt < max_retries - 1: + delay = min(base_delay * (2 ** attempt), max_delay) + logger.warning( + f"Operation failed (attempt {attempt + 1}/{max_retries}): " + f"{e}. Retrying in {delay:.1f}s..." + ) + time.sleep(delay) + else: + logger.error(f"Operation failed after {max_retries} attempts: {e}") + raise last_exception + + +def _facade_attr(name: str, fallback): + """Берёт совместимый mock из фасада generator, если тест его подменил.""" + facade = sys.modules.get("generator") + return getattr(facade, name, fallback) if facade else fallback + + +def _retry(operation, **kwargs): + return _facade_attr("_with_retry", _with_retry)(operation, **kwargs) + + +def _kafka_importer(): + return _facade_attr("_import_kafka", _import_kafka) + + +@dataclass +class BatchRecord: + """Запись об отправленном батче.""" + + batch_id: str + started_at: datetime + finished_at: datetime + sent_total: int + sent_browser: int + sent_location: int + sent_device: int + sent_geo: int + status: str + error_message: str | None = None + + def to_dict(self) -> dict: + """Конвертирует в словарь для сериализации.""" + return { + "batch_id": self.batch_id, + "started_at": self.started_at.isoformat(), + "finished_at": self.finished_at.isoformat(), + "sent_total": self.sent_total, + "sent_browser": self.sent_browser, + "sent_location": self.sent_location, + "sent_device": self.sent_device, + "sent_geo": self.sent_geo, + "status": self.status, + "error_message": self.error_message, + } + + +class KafkaStateManager: + """Управление состоянием генератора в Kafka compact topic.""" + + STATE_TOPIC = "generator_state" + STATE_KEY = "default" + + def __init__(self, bootstrap_servers: str): + self.bootstrap_servers = bootstrap_servers + KafkaProducerCls, _ = _kafka_importer()() + + logger.info(f"Connecting to Kafka for state management at {self.bootstrap_servers}") + self.producer = KafkaProducerCls( + bootstrap_servers=self.bootstrap_servers, + value_serializer=lambda v: json.dumps(v).encode("utf-8"), + key_serializer=lambda k: k.encode("utf-8") if k else None, + retries=3, + retry_backoff_ms=1000, + ) + logger.info("Connected to Kafka for state management successfully") + + def save(self, state: GeneratorState) -> None: + """Сохраняет состояние в топик.""" + value = state.to_dict() + + def _do_send(): + self.producer.send(self.STATE_TOPIC, key=self.STATE_KEY, value=value) + + _retry(_do_send, max_retries=3, base_delay=0.5) + + def flush(self) -> None: + """Сбрасывает буфер с retry.""" + def _do_flush(): + self.producer.flush() + + _retry(_do_flush, max_retries=3, base_delay=0.5) + + def close(self) -> None: + """Закрывает соединение.""" + try: + self.producer.close() + except Exception as e: + logger.debug(f"Error closing producer (ignored): {e}") + + def load(self) -> GeneratorState | None: + """Загружает последнее состояние из топика.""" + from kafka import KafkaConsumer + + logger.info(f"Loading state from topic {self.STATE_TOPIC}") + + def _do_load(): + consumer = KafkaConsumer( + self.STATE_TOPIC, + bootstrap_servers=self.bootstrap_servers, + auto_offset_reset="earliest", + enable_auto_commit=False, + consumer_timeout_ms=5000, + value_deserializer=lambda v: json.loads(v.decode("utf-8")), + ) + + last_state = None + for message in consumer: + if message.key and message.key.decode("utf-8") == self.STATE_KEY: + last_state = message.value + + consumer.close() + return last_state + + try: + last_state = _retry(_do_load, max_retries=3, base_delay=0.5) + + if last_state: + logger.info( + 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: + logger.warning("State data was invalid, starting fresh") + return restored + + logger.info("No previous state found, starting fresh") + return None + + except Exception as e: + logger.warning(f"Failed to load state: {e}, starting fresh") + return None + + +def ensure_topics(bootstrap_servers: str) -> None: + """Создаёт служебные топики, если их ещё нет.""" + from kafka import KafkaAdminClient + from kafka.admin import NewTopic + from kafka.errors import TopicAlreadyExistsError + + def _create_topics(): + admin_client = KafkaAdminClient(bootstrap_servers=bootstrap_servers) + try: + history_topic = NewTopic( + name=KafkaBatchHistory.HISTORY_TOPIC, + num_partitions=1, + replication_factor=1, + ) + state_topic = NewTopic( + name=KafkaStateManager.STATE_TOPIC, + num_partitions=1, + replication_factor=1, + topic_configs={ + "cleanup.policy": "compact", + "min.cleanable.dirty.ratio": "0.1", + "delete.retention.ms": "100", + }, + ) + + for topic in [history_topic, state_topic]: + try: + admin_client.create_topics([topic]) + logger.info(f"Created topic: {topic.name}") + except TopicAlreadyExistsError: + logger.debug(f"Topic already exists: {topic.name}") + finally: + admin_client.close() + + _retry(_create_topics, max_retries=5, base_delay=1.0) + + +class KafkaBatchHistory: + """Хранение истории batch в Kafka.""" + + HISTORY_TOPIC = "generator_batch_history" + + def __init__(self, bootstrap_servers: str): + self.bootstrap_servers = bootstrap_servers + self.producer = None + self._connect() + + def _connect(self): + """Устанавливает соединение с Kafka с retry.""" + KafkaProducerCls, _ = _kafka_importer()() + + def _do_connect(): + logger.info(f"Connecting to Kafka for history at {self.bootstrap_servers}") + self.producer = KafkaProducerCls( + bootstrap_servers=self.bootstrap_servers, + value_serializer=lambda v: json.dumps(v).encode("utf-8"), + key_serializer=lambda k: k.encode("utf-8") if k else None, + retries=3, + retry_backoff_ms=1000, + ) + logger.info("Connected to Kafka for history successfully") + + _retry(_do_connect, max_retries=5, base_delay=1.0) + + def add(self, record: BatchRecord): + """Добавляет запись в историю.""" + key = record.batch_id + value = record.to_dict() + + def _do_send(): + self.producer.send(self.HISTORY_TOPIC, key=key, value=value) + + try: + _retry(_do_send, max_retries=3, base_delay=0.5) + except Exception as e: + logger.warning(f"Failed to send history record after retries: {e}, attempting reconnect") + self._connect() + _retry(_do_send, max_retries=2, base_delay=0.5) + + def flush(self): + """Сбрасывает буфер с retry.""" + def _do_flush(): + self.producer.flush() + + _retry(_do_flush, max_retries=3, base_delay=0.5) + + def close(self): + """Закрывает соединение.""" + try: + if self.producer: + self.producer.close() + except Exception as e: + logger.debug(f"Error closing history producer (ignored): {e}") + + +class KafkaPublisher: + """Публикация событий в Kafka с retry и реконнектом.""" + + def __init__(self, bootstrap_servers: str): + self.bootstrap_servers = bootstrap_servers + self.producer = None + self._connect() + + def _connect(self): + """Устанавливает соединение с Kafka с retry.""" + KafkaProducerCls, _ = _kafka_importer()() + + def _do_connect(): + logger.info(f"Connecting to Kafka at {self.bootstrap_servers}") + self.producer = KafkaProducerCls( + bootstrap_servers=self.bootstrap_servers, + value_serializer=lambda v: json.dumps(v).encode("utf-8"), + key_serializer=lambda k: k.encode("utf-8") if k else None, + batch_size=16384, + linger_ms=100, + retries=3, + retry_backoff_ms=1000, + ) + logger.info("Connected to Kafka successfully") + + _retry(_do_connect, max_retries=5, base_delay=1.0) + + def _publish_with_retry(self, topic: str, events: list[dict]) -> tuple[int, int]: + """Внутренняя функция публикации с retry на уровне batch.""" + if not self.producer: + raise RuntimeError("Producer not connected") + + sent = 0 + errors = 0 + futures = [] + + for event in events: + key = event.get("event_id") or event.get("click_id") + try: + future = self.producer.send(topic, key=key, value=event) + futures.append(future) + except Exception as e: + logger.error(f"Failed to send message to {topic}: {e}") + errors += 1 + METRICS_ERRORS_TOTAL.labels(topic=topic).inc() + + for future in futures: + try: + future.get(timeout=10) + sent += 1 + METRICS_EVENTS_TOTAL.labels(topic=topic).inc() + except Exception as e: + logger.error(f"Failed to confirm message delivery: {e}") + errors += 1 + METRICS_ERRORS_TOTAL.labels(topic=topic).inc() + + return sent, errors + + def publish(self, topic: str, events: list[dict]) -> tuple[int, int]: + """Публикует события в топик с retry и автоматическим реконнектом.""" + def _do_publish(): + return self._publish_with_retry(topic, events) + + try: + return _retry(_do_publish, max_retries=3, base_delay=0.5) + except Exception as e: + logger.warning(f"Publish failed after retries: {e}, attempting reconnect") + self._connect() + return _retry(_do_publish, max_retries=2, base_delay=0.5) + + def flush(self): + """Сбрасывает буфер с retry.""" + def _do_flush(): + if self.producer: + self.producer.flush() + + _retry(_do_flush, max_retries=3, base_delay=0.5) + + def close(self): + """Закрывает соединение.""" + try: + if self.producer: + self.producer.close() + except Exception as e: + logger.debug(f"Error closing producer (ignored): {e}") diff --git a/generator/src/clickstream_generator/metrics.py b/generator/src/clickstream_generator/metrics.py new file mode 100644 index 0000000..cb9b6fb --- /dev/null +++ b/generator/src/clickstream_generator/metrics.py @@ -0,0 +1,23 @@ +"""Prometheus-метрики генератора.""" + +from prometheus_client import Counter, Gauge, Histogram + + +METRICS_EVENTS_TOTAL = Counter( + "generator_events_total", + "Total number of events sent to Kafka", + ["topic"], +) +METRICS_ERRORS_TOTAL = Counter( + "generator_publish_errors_total", + "Total number of publish errors", + ["topic"], +) +METRICS_TICK_DURATION = Histogram( + "generator_tick_duration_seconds", + "Duration of generator tick in seconds", +) +METRICS_LAST_SUCCESS = Gauge( + "generator_last_success_timestamp", + "Unix timestamp of last successful tick", +) diff --git a/generator/src/clickstream_generator/runtime.py b/generator/src/clickstream_generator/runtime.py new file mode 100644 index 0000000..2cbb798 --- /dev/null +++ b/generator/src/clickstream_generator/runtime.py @@ -0,0 +1,30 @@ +"""Переходный тиковый слой генератора.""" + +from clickstream_generator.generation import EventGenerator + + +def generate_tick_batch(generator: EventGenerator, event_budget: int) -> dict[str, list[dict]]: + """Набирает тиковый батч из одного или нескольких полных визитов. + + Это временный механизм до задачи 03. В ней тиковый слой начнёт хранить + активные визиты и выпускать только созревшие события. + """ + tick_batch = { + "browser_events": [], + "location_events": [], + "device_events": [], + "geo_events": [], + } + + remaining_events = event_budget + while remaining_events > 0: + visit_batch = generator.generate_batch(remaining_events) + generated_events = len(visit_batch["browser_events"]) + if generated_events == 0: + break + + for topic, events in visit_batch.items(): + tick_batch[topic].extend(events) + remaining_events -= generated_events + + return tick_batch diff --git a/generator/src/clickstream_generator/service.py b/generator/src/clickstream_generator/service.py new file mode 100644 index 0000000..a6f3f3f --- /dev/null +++ b/generator/src/clickstream_generator/service.py @@ -0,0 +1,235 @@ +"""Основной сервисный цикл генератора.""" + +import logging +import sys +import time +import uuid +from datetime import datetime, timezone + +from prometheus_client import start_http_server + +from clickstream_generator.config import Config +from clickstream_generator.dictionary import EventDictionary +from clickstream_generator.generation import EventGenerator +from clickstream_generator.kafka_io import ( + BatchRecord, + KafkaBatchHistory, + KafkaPublisher, + KafkaStateManager, + ensure_topics, +) +from clickstream_generator.metrics import ( + METRICS_ERRORS_TOTAL, + METRICS_LAST_SUCCESS, + METRICS_TICK_DURATION, +) +from clickstream_generator.runtime import generate_tick_batch +from clickstream_generator.state import GeneratorState + + +logger = logging.getLogger("generator") + + +class GeneratorService: + """Основной сервис генератора.""" + + def __init__(self, config: Config): + self.config = config + self.dictionary = EventDictionary.load(config.data_dir) + self.generator = EventGenerator(self.dictionary, config) + self.publisher: KafkaPublisher | None = None + self.history: KafkaBatchHistory | None = None + self.state_manager: KafkaStateManager | None = None + self._running = False + self._tick = 0 + + def start(self): + """Запускает основной цикл.""" + if not self.config.enabled: + 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) + + logger.info("Starting generator service...") + logger.info( + f"Configuration: tick={self.config.tick_seconds}s, " + f"lambda_base={self.config.lambda_base_per_min}/min, " + f"jitter={self.config.jitter_pct}%, " + f"state_enabled={self.config.state_enabled}, " + f"state_reset={self.config.state_reset}" + ) + + ensure_topics(self.config.kafka_bootstrap_servers) + self.publisher = KafkaPublisher(self.config.kafka_bootstrap_servers) + self.history = KafkaBatchHistory(self.config.kafka_bootstrap_servers) + + if self.config.state_enabled: + self.state_manager = KafkaStateManager(self.config.kafka_bootstrap_servers) + + if not self.config.state_reset: + restored_state = self.state_manager.load() + if restored_state: + self._tick = restored_state.tick + self.generator.rng.setstate(restored_state.rng_state) + logger.info( + f"Restored state: continuing from tick {self._tick}, " + f"last_batch_id={restored_state.last_batch_id}" + ) + else: + logger.info("State reset requested, starting fresh") + else: + logger.info("State management disabled") + + self._running = True + + try: + self._main_loop() + except KeyboardInterrupt: + logger.info("Received shutdown signal") + finally: + self.stop() + + def stop(self): + """Останавливает сервис.""" + logger.info("Stopping generator service...") + self._running = False + if self.publisher: + self.publisher.close() + if self.history: + self.history.close() + if self.state_manager: + self.state_manager.close() + + def _save_state(self, batch_id: str) -> None: + """Сохраняет текущее состояние генератора.""" + if not self.state_manager or not self.config.state_enabled: + return + + try: + state = GeneratorState( + tick=self._tick, + rng_state=self.generator.rng.getstate(), + last_batch_id=batch_id, + last_timestamp=datetime.now(timezone.utc), + ) + self.state_manager.save(state) + self.state_manager.flush() + logger.debug(f"Saved state: tick={self._tick}, batch_id={batch_id}") + except Exception as e: + logger.warning(f"Failed to save state: {e}") + METRICS_ERRORS_TOTAL.labels(topic="state").inc() + + def _main_loop(self): + """Основной цикл тиков.""" + while self._running: + self._tick += 1 + tick_start = time.time() + batch_id = str(uuid.uuid4())[:8] + + with METRICS_TICK_DURATION.time(): + logger.info(f"=== Tick {self._tick} (batch_id={batch_id}) ===") + + try: + events_count = self.generator._calculate_events_count() + logger.info(f"Generating ~{events_count} base events") + + gen_start = time.time() + batch = generate_tick_batch(self.generator, events_count) + gen_duration = time.time() - gen_start + + pub_start = time.time() + total_sent = 0 + total_errors = 0 + + sent_counts = {} + for topic, events in batch.items(): + if events: + sent, errors = self.publisher.publish(topic, events) + sent_counts[topic] = {"sent": sent, "errors": errors} + total_sent += sent + total_errors += errors + + if total_errors == 0: + status = "success" + elif total_sent > 0: + status = "partial" + else: + status = "error" + + if status in ("success", "partial"): + METRICS_LAST_SUCCESS.set_to_current_time() + self._save_state(batch_id) + + self.publisher.flush() + pub_duration = time.time() - pub_start + + try: + batch_record = BatchRecord( + batch_id=batch_id, + started_at=datetime.fromtimestamp(tick_start, tz=timezone.utc), + finished_at=datetime.now(timezone.utc), + sent_total=total_sent, + sent_browser=sent_counts.get("browser_events", {}).get("sent", 0), + sent_location=sent_counts.get("location_events", {}).get("sent", 0), + sent_device=sent_counts.get("device_events", {}).get("sent", 0), + sent_geo=sent_counts.get("geo_events", {}).get("sent", 0), + status=status, + error_message=None if status == "success" else f"Errors: {total_errors}", + ) + self.history.add(batch_record) + self.history.flush() + except Exception as hist_err: + logger.warning(f"Failed to write batch history: {hist_err}") + METRICS_ERRORS_TOTAL.labels(topic="history").inc() + + tick_duration = time.time() - tick_start + logger.info( + f"Batch {batch_id} completed: " + f"sent={total_sent}, errors={total_errors}, " + f"gen_time={gen_duration:.3f}s, pub_time={pub_duration:.3f}s, " + f"total_time={tick_duration:.3f}s" + ) + + for topic, counts in sent_counts.items(): + if counts["sent"] > 0: + logger.info(f" {topic}: {counts['sent']} sent") + + except Exception as e: + logger.exception(f"Error in tick {self._tick}: {e}") + try: + self.history.add( + BatchRecord( + batch_id=batch_id, + started_at=datetime.fromtimestamp(tick_start, tz=timezone.utc), + finished_at=datetime.now(timezone.utc), + sent_total=0, + sent_browser=0, + sent_location=0, + sent_device=0, + sent_geo=0, + status="error", + error_message=str(e), + ) + ) + self.history.flush() + except Exception as hist_err: + logger.warning(f"Failed to write error to history: {hist_err}") + + elapsed = time.time() - tick_start + sleep_time = max(0, self.config.tick_seconds - elapsed) + if sleep_time > 0: + logger.debug(f"Sleeping for {sleep_time:.1f}s until next tick") + time.sleep(sleep_time) + + +def main(): + """Точка входа сервиса.""" + try: + config = Config() + service = GeneratorService(config) + service.start() + except Exception as e: + logger.exception(f"Fatal error: {e}") + sys.exit(1) diff --git a/generator/src/clickstream_generator/state.py b/generator/src/clickstream_generator/state.py new file mode 100644 index 0000000..e25ebe3 --- /dev/null +++ b/generator/src/clickstream_generator/state.py @@ -0,0 +1,80 @@ +"""Сериализуемое состояние генератора.""" + +import logging +import random +from dataclasses import dataclass +from datetime import datetime + + +logger = logging.getLogger("generator") + + +def _nested_list_to_tuple(obj): + """Рекурсивно преобразует list в tuple для восстановления RNG state.""" + if isinstance(obj, list): + return tuple(_nested_list_to_tuple(x) for x in obj) + return obj + + +@dataclass +class GeneratorState: + """Состояние генератора для восстановления после рестарта.""" + + tick: int + rng_state: tuple + last_batch_id: str + last_timestamp: datetime + version: str = "1.0" + + def to_dict(self) -> dict: + """Конвертирует в словарь для JSON-сериализации.""" + return { + "tick": self.tick, + "rng_state": self.rng_state, + "last_batch_id": self.last_batch_id, + "last_timestamp": self.last_timestamp.isoformat(), + "version": self.version, + } + + @classmethod + def from_dict(cls, data: dict) -> "GeneratorState": + """Создаёт состояние из словаря.""" + try: + rng_state_raw = data.get("rng_state") + if not rng_state_raw: + logger.warning("State missing rng_state field") + raise ValueError("rng_state is missing") + + rng_state = _nested_list_to_tuple(rng_state_raw) + + if not isinstance(rng_state, tuple): + logger.warning(f"rng_state is not tuple: {type(rng_state)}") + raise ValueError("rng_state must be tuple") + + if len(rng_state) < 2: + logger.warning(f"rng_state has insufficient length: {len(rng_state)}") + raise ValueError("rng_state has insufficient length") + + test_rng = random.Random() + test_rng.setstate(rng_state) + + return cls( + tick=data.get("tick", 0), + rng_state=rng_state, + last_batch_id=data.get("last_batch_id", ""), + last_timestamp=datetime.fromisoformat( + data.get("last_timestamp", "1970-01-01T00:00:00+00:00") + ), + version=data.get("version", "1.0"), + ) + except Exception as e: + logger.warning(f"Invalid state format, will start fresh: {e}") + raise ValueError(f"Invalid state: {e}") + + @classmethod + def from_dict_safe(cls, data: dict) -> "GeneratorState | None": + """Безопасная загрузка state с graceful degradation.""" + try: + return cls.from_dict(data) + except Exception: + return None diff --git a/generator/tests/conftest.py b/generator/tests/conftest.py index 8a75e0a..250c8e4 100644 --- a/generator/tests/conftest.py +++ b/generator/tests/conftest.py @@ -5,8 +5,11 @@ Pytest fixtures для тестирования генератора. import sys from pathlib import Path -# Добавляем родительскую директорию в путь -sys.path.insert(0, str(Path(__file__).parent.parent)) +GENERATOR_DIR = Path(__file__).parent.parent + +# Добавляем фасад generator.py и src-пакет в путь +sys.path.insert(0, str(GENERATOR_DIR)) +sys.path.insert(0, str(GENERATOR_DIR / "src")) import pytest from generator import Config, EventDictionary diff --git a/generator/tests/test_generation.py b/generator/tests/test_generation.py index 3801ef1..afca78f 100644 --- a/generator/tests/test_generation.py +++ b/generator/tests/test_generation.py @@ -7,7 +7,7 @@ from dataclasses import replace from datetime import datetime import pytest -from generator import EventGenerator, EventDictionary +from generator import EventGenerator, EventDictionary, generate_tick_batch ALLOWED_PAGE_PATHS = { @@ -95,7 +95,7 @@ class TestEventGeneration: """Тиковый батч набирает целевой бюджет из одного или нескольких визитов.""" generator = EventGenerator(event_dictionary, base_config) - batch = generator.generate_tick_batch(20) + batch = generate_tick_batch(generator, 20) assert len(batch["browser_events"]) == 20 assert len(batch["location_events"]) == 20 diff --git a/generator/tests/test_service_cleanup.py b/generator/tests/test_service_cleanup.py new file mode 100644 index 0000000..090389f --- /dev/null +++ b/generator/tests/test_service_cleanup.py @@ -0,0 +1,55 @@ +""" +Контракт уборки сервиса генератора. +""" + +from pathlib import Path +import subprocess +import sys + +from generator import ( + Config, + EventDictionary, + EventGenerator, + GeneratorService, + KafkaPublisher, + GeneratorState, +) + + +def test_dockerfile_copies_split_python_modules(): + """Контейнерный запуск видит все модули генератора после разбиения.""" + dockerfile = (Path(__file__).parent.parent / "Dockerfile").read_text(encoding="utf-8") + + assert "COPY src ./src" in dockerfile + assert 'ENV PYTHONPATH="/app/src"' in dockerfile + + +def test_generator_facade_imports_without_pythonpath(): + """Локальный фасад сам находит src-пакет без внешнего PYTHONPATH.""" + generator_dir = Path(__file__).parent.parent + + result = subprocess.run( + [sys.executable, "-c", "import generator"], + cwd=generator_dir, + env={"PATH": ""}, + text=True, + capture_output=True, + check=False, + ) + + assert result.returncode == 0, result.stderr + + +def test_generator_facade_preserves_public_imports_after_split(): + """Старый импорт из generator работает, но классы живут в отдельных модулях.""" + assert Config.__module__ == "clickstream_generator.config" + assert EventDictionary.__module__ == "clickstream_generator.dictionary" + assert EventGenerator.__module__ == "clickstream_generator.generation" + assert KafkaPublisher.__module__ == "clickstream_generator.kafka_io" + assert GeneratorState.__module__ == "clickstream_generator.state" + assert GeneratorService.__module__ == "clickstream_generator.service" + + +def test_tick_batch_is_not_part_of_pure_visit_generation(): + """Временная сборка тика вынесена из чистой модели одного визита.""" + assert not hasattr(EventGenerator, "generate_tick_batch")