diff --git a/plans/generator_demo_stream_plan.md b/plans/generator_demo_stream_plan.md index 2525bb4..ed547c6 100644 --- a/plans/generator_demo_stream_plan.md +++ b/plans/generator_demo_stream_plan.md @@ -1,234 +1,205 @@ -# План реализации: генератор долгоживущего demo-потока данных +# План реализации: простой автономный генератор (rev4) -Цель: спроектировать генератор, который непрерывно производит правдоподобный поток событий для стенда и делает графики/аналитику визуально "живыми" без ручной подгрузки файлов. - -Документ фиксирует верхнеуровневую архитектуру, целевые параметры для наглядности и этапы внедрения с минимальным риском для текущего контура. +Дата ревизии: 14 февраля 2026. --- -## 1) Цель и границы MVP +## 1) Цель -### Цель +Сделать **простой** генератор, который: -- Обеспечить поток данных "в долгую" (24/7 или во время демо-сессий). -- Сохранить совместимость с текущим контрактом входных событий. -- Получить в дашбордах заметные изменения метрик за 10-15 минут наблюдения. - -### Границы MVP - -- Никаких изменений в текущих слоях `STG -> ODS -> DDS -> DM`. -- Генератор работает как отдельный upstream к Kafka. -- Формат output генератора идентичен текущему формату `data/*.jsonl` (1 JSON = 1 Kafka message). +- работает автономно и не зависит от потребителей; +- может стабильно работать несколько часов; +- публикует данные в текущие Kafka-топики по существующему контракту; +- реально внедряется за короткий срок, без «проекта на недели». --- -## 2) Нефункциональные требования для "наглядного" демо +## 2) Что делаем в MVP (и что не делаем) -- Запуск генерации: каждую минуту. -- Время обработки одного запуска: до 40 секунд, чтобы не накапливались overlapping runs. -- Контролируемая нагрузка для локального стенда: базово `~300 events/min`, пик `~900 events/min`. -- Детерминизм при retry: повторный запуск того же `run_id` формирует тот же набор событий. -- Устойчивость к "грязным" данным: ошибки должны фиксироваться downstream, но не валить пайплайн. +### Делаем + +- один режим генерации; +- один автономный сервис (контейнер); +- простая конфигурация через env; +- базовые метрики и логи; +- минимальная история запусков. + +### Не делаем в MVP + +- много сценариев (`promo/incident/...`); +- сложные state machines; +- сложный replay; +- сложную оркестрацию с зависимостями от ETL; +- «идеальный прод» (exactly-once и т.д.). --- -## 3) Целевая архитектура генератора +## 3) Архитектура MVP -### 3.1 Позиция в текущем контуре +`generator-service -> Kafka topics -> (потребители отдельно)` -`Generator -> Kafka topics -> ClickHouse STG -> ODS -> DDS -> DM` - -Текущие DAG'и `ddl_init`, `kafka_load`, `etl_pipeline` остаются рабочими. Генератор добавляется как отдельный режим поставки данных. - -### 3.2 Логические компоненты - -1. `Scheduler` -- Триггер раз в минуту через отдельный Airflow DAG `generator_minutely`. - -2. `Scenario Engine` -- Вычисляет профиль минуты: сколько и каких событий генерировать. -- Управляет сценариями `normal`, `promo`, `incident`. - -3. `State Store` -- Хранит состояние между запусками: активные пользователи/сессии, текущий сценарий, служебные seed/run метаданные. -- MVP-вариант: таблица в ClickHouse (`ods.generator_state`) или отдельный компактный state-файл в volume. - -4. `Event Builder` -- Генерирует события строго по действующему контракту полей и типов. -- Добавляет служебные поля только при обратной совместимости (например, `generator_run_id`, `scenario_version`). - -5. `Corruption Injector` -- Добавляет управляемую долю неидеальных событий (malformed, missing fields, duplicate, late events). - -6. `Publisher` -- Публикует события в те же Kafka-топики, что использует текущий ingest. - -7. `Metrics Emitter` -- Пишет техметрики генератора для Prometheus/Grafana. +Генератор не вызывает Airflow DAG-и и не ждёт их. +Потребители запускаются по своему расписанию/логике. --- -## 4) Модель "наглядного" поведения (demo_visible_v1) +## 4) Единственный режим генерации: `steady` -### 4.1 Базовые сценарии +Логика режима: -1. `normal` (около 70% времени) -- Стабильный фоновый трафик. +- каждую минуту публикуем фиксированный объём событий; +- распределяем события по 4 топикам: + - `browser_events` + - `location_events` + - `device_events` + - `geo_events` +- слегка «оживляем» поток: + - варьируем объём в небольшом диапазоне; + - обновляем `event_timestamp`; + - сохраняем реалистичные связи `event_id <-> location`, `click_id <-> device/geo`. -2. `promo` (около 20% времени) -- Рост paid-трафика и CTR. -- Умеренный рост конверсий. +Этого достаточно, чтобы стенд жил часами и данные выглядели не статично. -3. `incident` (около 10% времени) -- Рост доли ошибок и late events. -- Просадка CR и качества данных. +### 4.1 Минимальная статистическая модель (Poisson) -Переключение сценариев должно быть запланированным и видимым на горизонте 10-15 минут. +Чтобы линия не была «ровной», используем простую интенсивность событий: -### 4.2 Целевые KPI-диапазоны +- число событий на тик: `N_t ~ Poisson(lambda_t)`; +- базовая интенсивность: `lambda_base` (например, 200 событий/мин); +- плавный профиль времени: `lambda_t = lambda_base * hour_factor(t)`. -- `CTR`: 3-8% -- `CR`: 1-2.5% -- `AOV` (средний чек): медиана 35-60, длинный хвост до 200+ -- `DQ error rate`: 1-2%, во время `incident` до 4-5% +Где `hour_factor(t)` можно сделать очень простым: -### 4.3 Распределения и зависимости +- дневное окно: `1.2` +- ночное окно: `0.7` +- остальное время: `1.0` -- Нагрузка во времени: неравномерная (пики и просадки). -- Сегменты: `new/returning`, `device`, `source`, `geo`. -- Переходы внутри сессии: вероятностная цепочка `view -> click -> add_to_cart -> purchase`. -- Поля не генерируются независимо: должны быть корреляции (пример: изменение source влияет на CTR/CR). +Плюс добавляем «защиту от шума»: + +- `N_min` и `N_max` (жёсткие границы); +- опционально короткое сглаживание по 3 последним тикам. + +Итог: поведение уже похоже на живой поток, но код остаётся компактным. --- -## 5) Контракт данных и идемпотентность +## 5) Как упростить реализацию -### 5.1 Контракт output +Чтобы не писать сложную генеративную модель: -- Формат сообщения: JSON-объект одной строкой. -- Topic mapping: как в текущем процессе загрузки (`browser_events`, `location_events`, `device_events`, `geo_events`). -- Совместимость с текущими парсингом и DDL обязательна. +1. Берём существующие JSONL как «базовый словарь» валидных событий. +2. На каждом тике семплируем записи из этого словаря. +3. Перегенерируем только необходимые поля (`event_id`, `click_id`, `event_timestamp`) с сохранением связности. +4. Публикуем в Kafka. -### 5.2 Batch/run модель - -- Каждый минутный запуск имеет `batch_id` (`YYYYMMDDHHmm`) и `run_id`. -- Seed вычисляется детерминированно от `batch_id` (+ version salt). -- При retry того же `run_id` набор событий должен совпадать. +Плюс: +- быстро; +- совместимо с текущим ODS-парсингом; +- минимум риска «сломать контракт». --- -## 6) Оркестрация в Airflow (верхний уровень) +## 6) Минимальная конфигурация сервиса -Отдельный DAG `generator_minutely`: +Через env: -- `schedule`: каждую минуту. -- `catchup=False` -- `max_active_runs=1` -- Короткий `execution_timeout`. +- `GEN_TICK_SECONDS` (по умолчанию `60`) +- `GEN_EVENTS_PER_TICK` (например, `200`) +- `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` (верхняя граница) -Этапы DAG: +Опционально: -1. `prepare_context` -- Рассчитать `batch_id`, `run_id`, seed, активный сценарий. - -2. `generate_batch` -- Сгенерировать события за минуту в памяти/временном буфере. - -3. `publish_kafka` -- Отправить события в Kafka. - -4. `emit_metrics` -- Зафиксировать метрики запуска (объем, ошибки, задержка). +- `GEN_ENABLED` (быстро включать/выключать цикл) +- `GEN_HOUR_PROFILE` (простая карта коэффициентов по часам) --- -## 7) Наблюдаемость и алерты +## 7) Минимальная наблюдаемость -### 7.1 Метрики генератора (минимум) +### Логи + +- старт/стоп сервиса; +- batch_id, объём отправки по топикам; +- длительность тика; +- ошибки публикации. + +### Метрики (минимум) - `generator_events_total` -- `generator_invalid_total` -- `generator_duplicates_total` -- `generator_late_events_total` -- `generator_batch_duration_seconds` - `generator_publish_errors_total` +- `generator_tick_duration_seconds` +- `generator_last_success_timestamp` -### 7.2 Минимальные алерты +### Минимальная история -- Нет новых событий > 3 минут. -- Длительность batch выше порога (например, > 45 сек). -- Доля invalid выше ожидаемой (например, > 5% вне `incident` окна). +Таблица `meta.generator_batches` (или файл/лог на первом шаге): + +- `batch_id` +- `started_at` +- `finished_at` +- `sent_total` +- `sent_browser/location/device/geo` +- `status` --- -## 8) План внедрения по этапам +## 8) План внедрения (короткий и реалистичный) -### Этап 1: Skeleton + совместимость +### Шаг 1. Skeleton сервиса -- Реализовать каркас генератора без сложных распределений. -- Включить публикацию в Kafka в текущем формате. -- Проверить, что текущий `etl_pipeline` работает без изменений. +- отдельная папка `generator/`; +- бесконечный цикл с тиком; +- подключение к Kafka; +- публикация в 4 топика. -Критерий готовности: -- Поток стабильно идет 30+ минут. -- STG/ODS/DDS/DM наполняются штатно. +Готово, если: +- сервис работает 1+ час без падений. -### Этап 2: Сценарии и распределения +### Шаг 2. Контрактная генерация из словаря -- Добавить `normal/promo/incident`. -- Включить корреляции полей и KPI-диапазоны. -- Включить управляемую "грязь". +- чтение базовых JSONL; +- семплирование + обновление ключевых полей; +- проверка, что downstream не ломается. -Критерий готовности: -- На дашбордах заметны смены сценариев и поведение KPI. +Готово, если: +- ODS/DDS наполняются штатно. -### Этап 3: Statefulness + reliability +### Шаг 3. Логи/метрики/история batch -- Добавить хранилище state между минутами. -- Доработать retry/идемпотентность и recovery. -- Ввести базовые алерты на генератор. +- добавить базовые метрики; +- писать историю batch; +- оформить runbook запуска/проверки. -Критерий готовности: -- При рестартах и ретраях поток остается предсказуемым. +Готово, если: +- можно показать историю работы стенда за несколько часов. --- -## 9) Риски и меры снижения +## 9) Критерии успеха MVP -1. Слишком "ровный" поток неинтересен для аналитики. -- Мера: сценарные переключения и целевые KPI-паттерны. +MVP успешен, если: -2. Слишком тяжелая генерация перегружает стенд. -- Мера: жесткие лимиты `events/min`, timeout и `max_active_runs=1`. - -3. Ломается совместимость формата. -- Мера: contract tests against current parser/ODS inserts. - -4. Дубли при ретраях. -- Мера: детерминированный seed + фиксированный `run_id`. +1. Генератор автономно работает 2-4 часа. +2. Потребители можно останавливать/запускать отдельно, генератор продолжает работу. +3. Данные остаются совместимыми с текущим пайплайном. +4. Есть базовые метрики и понятные логи. +5. Есть история batch-ов для учебного разбора. --- -## 10) Acceptance criteria для demo-ready статуса +## 10) Что делаем потом (после рабочего MVP) -Считаем задачу завершенной, если: +Когда простой генератор стабильно работает: -1. Генератор работает по расписанию 1 раз в минуту и стабильно публикует данные. -2. Текущий downstream пайплайн не требует изменений для приема потока. -3. За 10-15 минут на графиках видно: -- смену структуры трафика, -- заметное изменение CTR/CR, -- реакцию DQ-метрик в `incident`. -4. Повторный retry run не генерирует новый "случайный" набор событий. -5. Есть базовые техметрики и алерты генератора. - ---- - -## 11) Что намеренно не включаем в MVP - -- Продвинутое ML-моделирование поведения пользователей. -- Сложные внешние reference-данные и enrichment в генераторе. -- Полную эмуляцию всех edge-кейсов production-среды. - -Сначала приоритет: наглядность, повторяемость, совместимость со стендом. +1. Добавляем второй режим (например, «spike»). +2. Делаем инкрементальных потребителей и расписание ETL. +3. Расширяем учебные кейсы по observability и DQ. +4. Усиливаем статистику: например, Gamma-Poisson/Negative Binomial для более «рваного» трафика. +Принцип: сначала работающий простой baseline, потом расширение.