# Генератор событий Автономный генератор событий для Kafka. В штатном стенде он создаёт стартовую историю через `backfill`, а затем может продолжить поток в режиме `live`. Генератор строит поток по иерархии `пользователь → визит → событие`: один `click_id` живёт весь визит, события визита идут по страницам воронки с монотонно растущим временем, а тиковый слой держит популяцию возвращающихся пользователей и активные визиты между тиками. Исторический дефект старой плоской генерации описан в [KNOWN_ISSUES.md](./KNOWN_ISSUES.md). ## Архитектура ``` 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` | сериализуемое состояние генератора v3 | | `src/clickstream_generator/metrics.py` | Prometheus-метрики | | `src/clickstream_generator/service.py` | основной цикл сервиса | | `generator.py` | запуск сервиса и совместимый фасад | ## Режим работы: `live` - Публикуем постепенно, **короткими тиками** (по умолчанию каждые 5 секунд) - На каждом тике отправляем небольшую порцию сообщений - Держим целевую интенсивность `events/min` без крупных минутных batch: рассчитанный событийный бюджет тика копится как бюджет рождения визитов через ожидаемую среднюю длину визита, а события выходят позже по своим запланированным меткам времени - Распределяем события по 4 топикам: - `browser_events` - `location_events` - `device_events` - `geo_events` - Публичный вызов генеративного ядра строит один визит: общий `click_id`, разные `event_id`, общий device/geo-контекст, путь по страницам воронки и строго растущие запланированные `event_timestamp`. - Тиковый слой хранит активные визиты между вызовами и выпускает только события, у которых наступил `event_timestamp`; завершённые визиты удаляются из памяти. - Тиковый слой ведёт ограниченную популяцию пользователей: один `user_domain_id` может вернуться в новом `click_id` после кулдауна, а при переполнении вытесняется давно неактивный пользователь. - Сохраняем связи `event_id <-> location`, `click_id <-> device/geo`. ## Конфигурация (env) | Переменная | Описание | По умолчанию | |------------|----------|--------------| | `KAFKA_BOOTSTRAP_SERVERS` | Адрес Kafka | `kafka:29092` | | `GEN_TICK_SECONDS` | Интервал между тиками | `5` (1-10 сек рекомендуется) | | `GEN_LAMBDA_BASE_PER_MIN` | Базовая интенсивность (событий/мин) | `30` | | `GEN_JITTER_PCT` | Процент вариативности | `20` | | `GEN_MIN_EVENTS_PER_TICK` | Минимальный событийный бюджет тика | `1` | | `GEN_MAX_EVENTS_PER_TICK` | Максимальный событийный бюджет тика | `50` | | `GEN_MAX_SESSION_EVENTS` | Потолок длины одного визита, защита от петель | `30` | | `GEN_MAX_ACTIVE_SESSIONS` | Потолок одновременных активных визитов | `200` | | `GEN_POPULATION_MAX` | Потолок активной популяции пользователей | `300` | | `GEN_P_NEW_USER` | Вероятность отдать новый визит новому пользователю | `0.15` | | `GEN_MIN_RETURN_MINUTES` | Минимальная пауза перед возвратом пользователя | `30` | | `GEN_MODEL_T0` | Стартовая модельная точка, ISO 8601 с часовым поясом | `2026-01-01T00:00:00+00:00` | | `GEN_MODEL_T_END` | Правая граница стартовой истории для `backfill` | — | | `GEN_MODEL_TIMEZONE` | Часовой пояс модельных часов для дневного коэффициента | `UTC` | | `GEN_MODEL_TIME_SPEED` | Сколько модельных секунд проходит за одну настенную секунду | `1` | | `GEN_RUN_MODE` | Режим генератора | `live` | | `GEN_LAUNCH_PROFILE` | Имя профиля запуска для логов | `ci` | | `GEN_STARTUP_HISTORY_ARTIFACT` | JSON-файл для экспорта стартовой истории в режиме `backfill` | — | | `GEN_DATA_DIR` | Путь к JSONL файлам | `/data` | | `GEN_SEED` | Сид для воспроизводимости | — | | `GEN_ENABLED` | Включить генерацию | `true` | | `GEN_METRICS_PORT` | Порт для Prometheus | `9109` | | `GEN_STATE_ENABLED` | Сохранять состояние между рестартами | `true` | | `GEN_STATE_RESET` | Сбросить состояние при старте | `false` | В `docker-compose.yml` параметры генеративной модели и режима проброшены через подстановку окружения, то есть их можно менять без правки файла: ```bash GEN_LAMBDA_BASE_PER_MIN=60 GEN_POPULATION_MAX=500 docker compose up -d generator ``` В живом режиме `event_timestamp` берётся из модельного времени: первый чистый тик стартует от `GEN_MODEL_T0`, дальше модельная точка сдвигается на `GEN_TICK_SECONDS * GEN_MODEL_TIME_SPEED`. Событийный бюджет тика считается по этой же модельной длительности, поэтому ×K даёт больше событий за короткий реальный прогон. Дневной коэффициент считается по `GEN_MODEL_TIMEZONE`, а не по реальному часу запуска процесса. В режиме `backfill` генератор без сна проходит от `GEN_MODEL_T0` до `GEN_MODEL_T_END`, публикует события только за `[T0, T_end)`, сохраняет state v3 на `T_end` в `generator_state` и пишет manifest в compact-topic `generator_startup_history_manifest`. Live-запуск с теми же настройками использует этот manifest, чтобы продолжить ровно с `T_end` без настенной дельты. Экспорт и импорт портативного файла описаны в [`docs/runbooks/startup-history.md`](../docs/runbooks/startup-history.md). Для одноразового backfill-запуска через compose используйте `docker compose run --rm generator`, а не `docker compose up generator`: у штатного сервиса включён restart policy. Контейнерные значения `KAFKA_BOOTSTRAP_SERVERS` и `GEN_DATA_DIR` в compose оставлены безопасными внутренними значениями `kafka:29092` и `/data`. Штатный чистый путь всего стенда запускается из корня репозитория: ```bash make generated-history-analytics ``` Эта команда очищает ClickHouse, Kafka-топики данных, state и manifest генератора, создаёт стартовую историю, прогоняет STG -> ODS -> DDS -> DM и проверяет Superset metadata. Файлы `data/*.jsonl` при этом не грузятся в Kafka: они пока используются только как фактура для генератора. По умолчанию используется быстрый профиль `ci`: 6 часов модельного времени. Профиль `daily-wave` даёт 2 суток, чтобы была видна суточная волна. В live он идёт с `GEN_MODEL_TIME_SPEED=60` и `GEN_TICK_SECONDS=1`: модельные сутки проходят примерно за 24 настенные минуты. ```bash PROFILE=daily-wave make generated-history-analytics ``` Разовую длительность можно задать без ручного расчёта `GEN_MODEL_T_END`: ```bash GEN_HISTORY_DURATION=2d make generated-history-analytics ``` ### Режим "раз в минуту" (для ручных экспериментов) Для ручных экспериментов можно установить: ```bash GEN_TICK_SECONDS=60 GEN_MIN_EVENTS_PER_TICK=50 GEN_MAX_EVENTS_PER_TICK=500 ``` `GEN_MIN_EVENTS_PER_TICK` и `GEN_MAX_EVENTS_PER_TICK` задают нижнюю и верхнюю границы событийного бюджета тика. Этот бюджет сначала превращается в рождения визитов через среднюю длину визита, поэтому фактическое число отправленных событий в конкретном тике может отличаться. ## Управление через Makefile ```bash # Запустить только генератор make generator-up # Остановить генератор make generator-down # Смотреть логи make generator-logs # Перезапуск с пересборкой make generator-restart # Промотать стартовую историю make generator-backfill # Продолжить live-поток из state make generator-continue # Начать live-поток как новый мир make generator-reset # Чистый аналитический прогон всего стенда make generated-history-analytics # Запуск тестов make generator-test ``` Основной ручной путь запуска — глаголы `generator-backfill`, `generator-continue` и `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`: | Метрика | Тип | Описание | |---------|-----|----------| | `generator_events_total` | Counter | Всего отправлено событий (по топикам) | | `generator_publish_errors_total` | Counter | Ошибки публикации (по топикам) | | `generator_tick_duration_seconds` | Histogram | Длительность тика | | `generator_last_success_timestamp` | Gauge | Время последнего успешного тика | ### Проверка метрик ```bash curl http://localhost:9109/metrics curl http://localhost:9090/api/v1/targets | grep generator ``` ## Мониторинг в Grafana **Dashboard URL:** `http://localhost:3000/d/generator-overview` Дашборд "Generator Overview" предоставляет полную визуализацию работы генератора: ### Ключевые панели | Панель | Метрика | Описание | |--------|---------|----------| | **Events/min** | `rate(generator_events_total[1m]) * 60` | Текущая скорость генерации | | **Tick Duration** | `generator_tick_duration_seconds` | p50 и p99 длительности тика | | **Last Successful Tick** | `generator_last_success_timestamp` | Время последнего успешного тика | | **Events per Hour** | `increase(generator_events_total[1h])` | 24-часовое распределение по топикам (bar chart) | | **Errors** | `generator_publish_errors_total` | Общее число и rate ошибок | | **Generator Status** | derived | Активен ли генератор | ### Структура дашборда Дашборд разделён на 6 секций: 1. **Overview** — ключевые метрики (events/min, tick duration, last success) 2. **Events by Topic** — bar chart Events per Hour, rate by topic, total counters 3. **Errors** — total errors, error rate, errors by topic 4. **Tick Statistics** — duration distribution (p50/p95/p99), events per tick, hour factor 5. **Status** — generator status, generator health (heartbeat), time since last tick 6. **Info** — полезные команды и параметры конфигурации ### Доступ к дашборду Дашборд автоматически загружается в Grafana при старте контейнера (provisioning). ```bash # Открыть дашборд open http://localhost:3000/d/generator-overview # Перезагрузить provisioning (если дашборд не появился) curl -s -u admin:admin -X POST http://localhost:3000/api/admin/provisioning/dashboards/reload ``` ## История batch История пишется в Kafka-топик `generator_batch_history` (JSON). **Контракт топика:** - Название фиксировано: `generator_batch_history` (не конфигурируется) - Формат: JSON с ключом `batch_id` **Важно:** генератор требует работающей Kafka. Без Kafka генератор упадёт при старте или потеряет события. Для мониторинга доступности используйте Prometheus-метрики (`generator_last_success_timestamp`). Поля сообщения: - `batch_id` — идентификатор батча - `started_at` / `finished_at` — время начала/окончания (ISO format) - `sent_total` — всего отправлено - `sent_browser/location/device/geo` — по топикам - `status` — success/partial/error - `error_message` — описание ошибки (если есть) ### Чтение истории из Kafka ```bash docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 \ --topic generator_batch_history \ --from-beginning ``` ## Восстановление состояния Генератор сохраняет своё состояние между перезапусками в Kafka-топик `generator_state` (compact topic). Это позволяет: - Продолжить нумерацию тиков с места остановки (continuity) - Сохранить последовательность случайных чисел (RNG state) - Восстановить популяцию пользователей - Продолжить активные визиты после короткого простоя ### Как работает 1. После каждого успешного тика состояние v3 сохраняется в `generator_state`. 2. При старте генератор читает последнее состояние из топика. 3. Если состояние найдено, сервис восстанавливает номер тика, состояние ГПСЧ, популяцию пользователей, накопленный бюджет рождения визитов и активные визиты. Для активного визита state хранит `base_click_id` — донора браузерной и source-фактуры из статического сида. 4. Если состояния нет или оно невалидно, генератор начинает с чистого листа. Активный визит после простоя до 30 минут продолжается со своими исходными запланированными метками времени. Если следующий шаг визита просрочен больше чем на 30 минут, визит закрывается без досылки остатка: для демо это выглядит как пользователь, который ушёл, пока стенд был остановлен. ### Топик `generator_state` - **Название**: фиксировано `generator_state` - **Тип**: compact topic (хранится только последнее значение для каждого ключа) - **Ключ**: `default` (для возможности нескольких генераторов в будущем) - **Конфигурация**: `cleanup.policy=compact`, минимальный retention ### Просмотр текущего состояния ```bash docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 \ --topic generator_state \ --from-beginning \ --property print.key=true ``` ### Сброс состояния (начать сначала) ```bash # Вариант 1: через env (рекомендуется) GEN_STATE_RESET=true docker compose up -d generator # Вариант 2: удалить топик полностью docker compose exec kafka /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka:29092 \ --delete \ --topic generator_state ``` ### Отключение сохранения состояния ```bash GEN_STATE_ENABLED=false docker compose up -d generator ``` При отключенном сохранении состояния генератор всегда начинает с `tick=1`, а ГПСЧ инициализируется с `GEN_SEED` или случайно. ## Тестирование Тесты написаны на **pytest**. ### Запуск тестов ```bash # Через Makefile (рекомендуется) make generator-test # Вручную через Docker docker build -t generator:test . docker run --rm -v $(PWD):/workspace -w /workspace/generator generator:test pytest tests/ -v # Конкретный файл тестов docker run --rm -v $(PWD):/workspace -w /workspace/generator generator:test pytest tests/test_generation.py -v ``` ### Структура тестов ``` generator/tests/ ├── conftest.py # Fixtures pytest ├── test_config.py # Тесты конфигурации ├── test_generation.py # Тесты генерации событий ├── test_history.py # Тесты структуры BatchRecord ├── test_kafka_history.py # Тесты KafkaBatchHistory ├── test_service.py # Тесты GeneratorService ├── test_service_cleanup.py # Контракт разбиения сервиса на модули └── test_state.py # Тесты GeneratorState и KafkaStateManager ``` ### Интеграционный тест ```bash # Запустить стек с генератором make generator-up # Проверить логи make generator-logs # Проверить метрики curl http://localhost:9109/metrics # Проверить сообщения в Kafka docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:29092 --topic browser_events --from-beginning ```