Files
clickstream-ch-kafka-supers…/plans/generator_demo_stream_plan.md
T
ddadminandDmitry Dementiev 7c6cfd7d18 docs(docs): обновлён план генератора по интеграции Prometheus
- Зачем:
  - зафиксировать решение об observability генератора на уровне MVP-плана
- Что:
  - добавлены требования по /metrics, scrape_config и target generator:9109
  - добавлен env-параметр GEN_METRICS_PORT
  - обновлены шаг внедрения и критерии успеха
- Проверка:
  - проверен diff только для plans/generator_demo_stream_plan.md
2026-06-09 17:26:31 +03:00

9.0 KiB
Raw Blame History

План реализации: простой автономный генератор (rev5)

Дата ревизии: 14 февраля 2026.


1) Цель

Сделать простой генератор, который:

  • работает автономно и не зависит от потребителей;
  • может стабильно работать несколько часов;
  • публикует данные в текущие Kafka-топики по существующему контракту;
  • реально внедряется за короткий срок, без «проекта на недели».

2) Что делаем в MVP (и что не делаем)

Делаем

  • один режим генерации;
  • один автономный сервис (контейнер);
  • простая конфигурация через env;
  • базовые метрики и логи;
  • экспорт метрик в Prometheus;
  • минимальная история запусков.

Не делаем в MVP

  • много сценариев (promo/incident/...);
  • сложные state machines;
  • сложный replay;
  • сложную оркестрацию с зависимостями от ETL;
  • «идеальный прод» (exactly-once и т.д.).

3) Архитектура MVP

generator-service -> Kafka topics -> (потребители отдельно)

Генератор не вызывает Airflow DAG-и и не ждёт их.
Потребители запускаются по своему расписанию/логике.


4) Единственный режим генерации: steady-stream

Логика режима:

  • публикуем в Kafka постепенно, короткими тиками (например, каждые 1-10 секунд);
  • на каждом тике отправляем небольшую порцию сообщений;
  • держим целевую интенсивность в events/min, а не крупный минутный batch;
  • распределяем события по 4 топикам:
    • browser_events
    • location_events
    • device_events
    • geo_events
  • слегка «оживляем» поток:
    • варьируем объём в небольшом диапазоне;
    • обновляем event_timestamp;
    • сохраняем реалистичные связи event_id <-> location, click_id <-> device/geo.

Этого достаточно, чтобы стенд жил часами, выглядел как реальный streaming и не создавал искусственных «пакетов раз в минуту».

4.1 Минимальная статистическая модель (Poisson)

Чтобы линия не была «ровной», используем простую интенсивность:

  • базовая интенсивность в минуту: lambda_minute (например, 200 событий/мин);
  • для текущего тика:
    • lambda_tick = lambda_minute * hour_factor(t) * tick_seconds / 60;
    • N_t ~ Poisson(lambda_tick).

Где hour_factor(t) можно сделать очень простым:

  • дневное окно: 1.2
  • ночное окно: 0.7
  • остальное время: 1.0

Плюс добавляем «защиту от шума»:

  • N_min и N_max (жёсткие границы);
  • опционально короткое сглаживание по 3 последним тикам.

Итог: поведение уже похоже на живой поток, но код остаётся компактным.

Важно: режим «раз в минуту» не удаляем полностью, но рассматриваем только как опцию (GEN_TICK_SECONDS=60) для контролируемых демо.


5) Как упростить реализацию

Чтобы не писать сложную генеративную модель:

  1. Берём существующие JSONL как «базовый словарь» валидных событий.
  2. На каждом тике семплируем записи из этого словаря.
  3. Перегенерируем только необходимые поля (event_id, click_id, event_timestamp) с сохранением связности.
  4. Публикуем в Kafka.

Плюс:

  • быстро;
  • совместимо с текущим ODS-парсингом;
  • минимум риска «сломать контракт».

6) Минимальная конфигурация сервиса

Через env:

  • GEN_TICK_SECONDS (по умолчанию 5)
  • GEN_JITTER_PCT (например, 20)
  • GEN_SEED (для воспроизводимости)
  • KAFKA_BOOTSTRAP_SERVERS
  • GEN_LAMBDA_BASE_PER_MIN (базовый lambda для Poisson)
  • GEN_MIN_EVENTS_PER_TICK (нижняя граница)
  • GEN_MAX_EVENTS_PER_TICK (верхняя граница)
  • GEN_METRICS_PORT (порт HTTP-эндпоинта /metrics, по умолчанию 9109)

Опционально:

  • GEN_ENABLED (быстро включать/выключать цикл)
  • GEN_HOUR_PROFILE (простая карта коэффициентов по часам)
  • GEN_TICK_SECONDS=60 (опциональный режим «раз в минуту»)

7) Минимальная наблюдаемость

Логи

  • старт/стоп сервиса;
  • batch_id, объём отправки по топикам;
  • длительность тика;
  • ошибки публикации.

Метрики (минимум)

  • generator_events_total
  • generator_publish_errors_total
  • generator_tick_duration_seconds
  • generator_last_success_timestamp

Интеграция с Prometheus (MVP)

  • генератор поднимает HTTP-эндпоинт /metrics (например, на :9109);
  • в configs/prometheus.yml добавляется scrape_config для job generator;
  • target внутри docker-сети: generator:9109;
  • в runbook добавляется проверка, что метрики генератора видны в Prometheus.

Минимальная история

Таблица meta.generator_batches (или файл/лог на первом шаге):

  • batch_id
  • started_at
  • finished_at
  • sent_total
  • sent_browser/location/device/geo
  • status

8) План внедрения (короткий и реалистичный)

Шаг 1. Skeleton сервиса

  • отдельная папка generator/;
  • бесконечный цикл с тиком;
  • подключение к Kafka;
  • публикация в 4 топика.

Готово, если:

  • сервис работает 1+ час без падений.

Шаг 2. Контрактная генерация из словаря

  • чтение базовых JSONL;
  • семплирование + обновление ключевых полей;
  • проверка, что downstream не ломается.

Готово, если:

  • ODS/DDS наполняются штатно.

Шаг 3. Логи/метрики/история batch

  • добавить базовые метрики;
  • подключить экспорт метрик в Prometheus (/metrics + scrape_config);
  • писать историю batch;
  • оформить runbook запуска/проверки.

Готово, если:

  • можно показать историю работы стенда за несколько часов.

9) Критерии успеха MVP

MVP успешен, если:

  1. Генератор автономно работает 2-4 часа.
  2. Потребители можно останавливать/запускать отдельно, генератор продолжает работу.
  3. Данные остаются совместимыми с текущим пайплайном.
  4. Поток публикуется равномерно малыми порциями (без искусственного минутного burst).
  5. Есть история batch-ов для учебного разбора.
  6. Метрики генератора стабильно собираются Prometheus.

10) Что делаем потом (после рабочего MVP)

Когда простой генератор стабильно работает:

  1. Добавляем второй режим (например, «spike»).
  2. Делаем инкрементальных потребителей и расписание ETL.
  3. Расширяем учебные кейсы по observability и DQ.
  4. Усиливаем статистику: например, Gamma-Poisson/Negative Binomial для более «рваного» трафика.

Принцип: сначала работающий простой baseline, потом расширение.