# План реализации: генератор долгоживущего demo-потока данных Цель: спроектировать генератор, который непрерывно производит правдоподобный поток событий для стенда и делает графики/аналитику визуально "живыми" без ручной подгрузки файлов. Документ фиксирует верхнеуровневую архитектуру, целевые параметры для наглядности и этапы внедрения с минимальным риском для текущего контура. --- ## 1) Цель и границы MVP ### Цель - Обеспечить поток данных "в долгую" (24/7 или во время демо-сессий). - Сохранить совместимость с текущим контрактом входных событий. - Получить в дашбордах заметные изменения метрик за 10-15 минут наблюдения. ### Границы MVP - Никаких изменений в текущих слоях `STG -> ODS -> DDS -> DM`. - Генератор работает как отдельный upstream к Kafka. - Формат output генератора идентичен текущему формату `data/*.jsonl` (1 JSON = 1 Kafka message). --- ## 2) Нефункциональные требования для "наглядного" демо - Запуск генерации: каждую минуту. - Время обработки одного запуска: до 40 секунд, чтобы не накапливались overlapping runs. - Контролируемая нагрузка для локального стенда: базово `~300 events/min`, пик `~900 events/min`. - Детерминизм при retry: повторный запуск того же `run_id` формирует тот же набор событий. - Устойчивость к "грязным" данным: ошибки должны фиксироваться downstream, но не валить пайплайн. --- ## 3) Целевая архитектура генератора ### 3.1 Позиция в текущем контуре `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. --- ## 4) Модель "наглядного" поведения (demo_visible_v1) ### 4.1 Базовые сценарии 1. `normal` (около 70% времени) - Стабильный фоновый трафик. 2. `promo` (около 20% времени) - Рост paid-трафика и CTR. - Умеренный рост конверсий. 3. `incident` (около 10% времени) - Рост доли ошибок и late events. - Просадка CR и качества данных. Переключение сценариев должно быть запланированным и видимым на горизонте 10-15 минут. ### 4.2 Целевые KPI-диапазоны - `CTR`: 3-8% - `CR`: 1-2.5% - `AOV` (средний чек): медиана 35-60, длинный хвост до 200+ - `DQ error rate`: 1-2%, во время `incident` до 4-5% ### 4.3 Распределения и зависимости - Нагрузка во времени: неравномерная (пики и просадки). - Сегменты: `new/returning`, `device`, `source`, `geo`. - Переходы внутри сессии: вероятностная цепочка `view -> click -> add_to_cart -> purchase`. - Поля не генерируются независимо: должны быть корреляции (пример: изменение source влияет на CTR/CR). --- ## 5) Контракт данных и идемпотентность ### 5.1 Контракт output - Формат сообщения: JSON-объект одной строкой. - Topic mapping: как в текущем процессе загрузки (`browser_events`, `location_events`, `device_events`, `geo_events`). - Совместимость с текущими парсингом и DDL обязательна. ### 5.2 Batch/run модель - Каждый минутный запуск имеет `batch_id` (`YYYYMMDDHHmm`) и `run_id`. - Seed вычисляется детерминированно от `batch_id` (+ version salt). - При retry того же `run_id` набор событий должен совпадать. --- ## 6) Оркестрация в Airflow (верхний уровень) Отдельный DAG `generator_minutely`: - `schedule`: каждую минуту. - `catchup=False` - `max_active_runs=1` - Короткий `execution_timeout`. Этапы DAG: 1. `prepare_context` - Рассчитать `batch_id`, `run_id`, seed, активный сценарий. 2. `generate_batch` - Сгенерировать события за минуту в памяти/временном буфере. 3. `publish_kafka` - Отправить события в Kafka. 4. `emit_metrics` - Зафиксировать метрики запуска (объем, ошибки, задержка). --- ## 7) Наблюдаемость и алерты ### 7.1 Метрики генератора (минимум) - `generator_events_total` - `generator_invalid_total` - `generator_duplicates_total` - `generator_late_events_total` - `generator_batch_duration_seconds` - `generator_publish_errors_total` ### 7.2 Минимальные алерты - Нет новых событий > 3 минут. - Длительность batch выше порога (например, > 45 сек). - Доля invalid выше ожидаемой (например, > 5% вне `incident` окна). --- ## 8) План внедрения по этапам ### Этап 1: Skeleton + совместимость - Реализовать каркас генератора без сложных распределений. - Включить публикацию в Kafka в текущем формате. - Проверить, что текущий `etl_pipeline` работает без изменений. Критерий готовности: - Поток стабильно идет 30+ минут. - STG/ODS/DDS/DM наполняются штатно. ### Этап 2: Сценарии и распределения - Добавить `normal/promo/incident`. - Включить корреляции полей и KPI-диапазоны. - Включить управляемую "грязь". Критерий готовности: - На дашбордах заметны смены сценариев и поведение KPI. ### Этап 3: Statefulness + reliability - Добавить хранилище state между минутами. - Доработать retry/идемпотентность и recovery. - Ввести базовые алерты на генератор. Критерий готовности: - При рестартах и ретраях поток остается предсказуемым. --- ## 9) Риски и меры снижения 1. Слишком "ровный" поток неинтересен для аналитики. - Мера: сценарные переключения и целевые KPI-паттерны. 2. Слишком тяжелая генерация перегружает стенд. - Мера: жесткие лимиты `events/min`, timeout и `max_active_runs=1`. 3. Ломается совместимость формата. - Мера: contract tests against current parser/ODS inserts. 4. Дубли при ретраях. - Мера: детерминированный seed + фиксированный `run_id`. --- ## 10) Acceptance criteria для demo-ready статуса Считаем задачу завершенной, если: 1. Генератор работает по расписанию 1 раз в минуту и стабильно публикует данные. 2. Текущий downstream пайплайн не требует изменений для приема потока. 3. За 10-15 минут на графиках видно: - смену структуры трафика, - заметное изменение CTR/CR, - реакцию DQ-метрик в `incident`. 4. Повторный retry run не генерирует новый "случайный" набор событий. 5. Есть базовые техметрики и алерты генератора. --- ## 11) Что намеренно не включаем в MVP - Продвинутое ML-моделирование поведения пользователей. - Сложные внешние reference-данные и enrichment в генераторе. - Полную эмуляцию всех edge-кейсов production-среды. Сначала приоритет: наглядность, повторяемость, совместимость со стендом.