Files
ddadminandClaude Fable 5 271f1afdb3 test(generator): вычищены прозаические тесты из контракта DAG
Тесты на дословные фразы документов и комментариев удалены: прозу
сторожит ревью, а не pytest. Оставлены структурные инварианты (порядок
задач, монтирования, версии, живые перекрёстные ссылки). В шапке файла —
критерий «инвариант должен переживать честную переписку текста».

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-23 18:41:09 +03:00
..

Генератор событий

Автономный генератор событий для Kafka. В штатном стенде он создаёт стартовую историю через backfill, а затем может продолжить поток в режиме live.

Генератор строит поток по иерархии пользователь → визит → событие: один click_id живёт весь визит, события визита идут по страницам воронки с монотонно растущим временем, а тиковый слой держит популяцию возвращающихся пользователей и активные визиты между тиками. Исторический дефект старой плоской генерации описан в 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 Имя профиля запуска для логов daily-wave
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 параметры генеративной модели и режима проброшены через подстановку окружения, то есть их можно менять без правки файла:

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. Для одноразового backfill-запуска через compose используйте docker compose run --rm generator, а не docker compose up generator: у штатного сервиса включён restart policy.

Контейнерные значения KAFKA_BOOTSTRAP_SERVERS и GEN_DATA_DIR в compose оставлены безопасными внутренними значениями kafka:29092 и /data.

Штатный чистый путь всего стенда запускается из корня репозитория:

make generated-history-analytics

Эта команда очищает ClickHouse, Kafka-топики данных, state и manifest генератора, создаёт стартовую историю, прогоняет STG -> ODS -> DDS -> DM и проверяет Superset metadata. Файлы data/*.jsonl при этом не грузятся в Kafka: они пока используются только как фактура для генератора.

По умолчанию используется учебный профиль daily-wave: 3 суток с суточной волной. В live он идёт с GEN_MODEL_TIME_SPEED=60 и GEN_TICK_SECONDS=1: модельные сутки проходят примерно за 24 настенные минуты. Плоский профиль ci на 6 часов остаётся служебным для автоматических тестов.

Разовую длительность можно задать без ручного расчёта GEN_MODEL_T_END:

GEN_HISTORY_DURATION=2d make generated-history-analytics

Режим "раз в минуту" (для ручных экспериментов)

Для ручных экспериментов можно установить:

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

# Запустить только генератор
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 test
make lint

Основной ручной путь запуска — глаголы generator-backfill, generator-continue и generator-reset. Старые GEN_RUN_MODE, GEN_STATE_RESET и GEN_MODEL_T_END остаются низкоуровневым способом для отладки.

generator-backfill, startup-history-import и generator-reset перед записью проверяют, что 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 Время последнего успешного тика

Проверка метрик

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).

# Открыть дашборд
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

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

Просмотр текущего состояния

docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server kafka:29092 \
  --topic generator_state \
  --from-beginning \
  --property print.key=true

Сброс состояния (начать сначала)

# Вариант 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

Отключение сохранения состояния

GEN_STATE_ENABLED=false docker compose up -d generator

При отключенном сохранении состояния генератор всегда начинает с tick=1, а ГПСЧ инициализируется с GEN_SEED или случайно.

Тестирование

Тесты написаны на pytest.

Запуск тестов

# Через Makefile (рекомендуется)
make test
make lint

# Вручную через Docker
docker build -t generator:test .
docker run --rm -v $(PWD):/workspace -w /workspace/generator generator:test pytest tests/ -q

# Конкретный файл тестов
docker run --rm -v $(PWD):/workspace -w /workspace/generator generator:test pytest tests/test_generation.py -v

make test запускает тесты генератора, контрактные тесты верхнего уровня и проверку compose-конфигурации. make lint проверяет синтаксис Python и Bash, compose-конфиг и пробелы в diff. Долгие стендовые проверки вроде make generated-history-runtime-check запускаются отдельно.

Структура тестов

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

Интеграционный тест

# Запустить стек с генератором
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