- Зачем: - зафиксировать решение об observability генератора на уровне MVP-плана - Что: - добавлены требования по /metrics, scrape_config и target generator:9109 - добавлен env-параметр GEN_METRICS_PORT - обновлены шаг внедрения и критерии успеха - Проверка: - проверен diff только для plans/generator_demo_stream_plan.md
222 lines
9.0 KiB
Markdown
222 lines
9.0 KiB
Markdown
# План реализации: простой автономный генератор (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, потом расширение.
|