From 6c4e0f4d1603b8cda1f197752116fa508cbdda34 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Thu, 11 Jun 2026 18:26:39 +0300 Subject: [PATCH] =?UTF-8?q?feat(generator):=20=D0=BF=D0=BE=D0=B4=D0=BA?= =?UTF-8?q?=D0=BB=D1=8E=D1=87=D0=B5=D0=BD=D0=B0=20=D0=BD=D0=BE=D0=B2=D0=B0?= =?UTF-8?q?=D1=8F=20=D0=BC=D0=BE=D0=B4=D0=B5=D0=BB=D1=8C=20=D0=BA=20steady?= =?UTF-8?q?-stream=20=D1=81=D0=B5=D1=80=D0=B2=D0=B8=D1=81=D1=83?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - после калибровки потока и state v2 генератор нужно принять как рабочий steady-stream источник, а не как исторически сломанный прототип. - Что: - добавлен сервисный тест multi-event визита с мок-публикацией во все четыре Kafka-топика. - compose позволяет переопределять демо-параметры генератора без правки файла, сохраняя внутренние контейнерные адреса. - README, OPERATIONS, KNOWN_ISSUES и карточка задачи синхронизированы с новой моделью и state v2. - Проверка: - uv run --with-requirements generator/requirements.txt pytest generator/tests -q. - git diff --check. - GEN_STATE_RESET=true GEN_POPULATION_MAX=123 docker compose config. --- .../issues/07-service-integration-and-docs.md | 27 +++-- docker-compose.yml | 26 +++-- docs/OPERATIONS.md | 19 ++++ generator/KNOWN_ISSUES.md | 99 ++++++++++--------- generator/README.md | 49 ++++++--- generator/tests/test_service.py | 79 +++++++++++++++ 6 files changed, 220 insertions(+), 79 deletions(-) diff --git a/.scratch/feature-data-generator/issues/07-service-integration-and-docs.md b/.scratch/feature-data-generator/issues/07-service-integration-and-docs.md index c722fd7..e979ed4 100644 --- a/.scratch/feature-data-generator/issues/07-service-integration-and-docs.md +++ b/.scratch/feature-data-generator/issues/07-service-integration-and-docs.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: ready-for-human # Подключение к сервису и документации @@ -18,17 +18,17 @@ Status: ready-for-agent ## Acceptance criteria -- [ ] Генератор в режиме steady-stream пишет связанные сообщения в +- [x] Генератор в режиме steady-stream пишет связанные сообщения в `browser_events`, `location_events`, `device_events`, `geo_events`. -- [ ] Существующие команды запуска и остановки генератора остаются рабочими или +- [x] Существующие команды запуска и остановки генератора остаются рабочими или документация явно описывает замену. -- [ ] Метрики Prometheus продолжают показывать успешные тики, ошибки публикации +- [x] Метрики Prometheus продолжают показывать успешные тики, ошибки публикации и объём отправленных событий. -- [ ] Интеграционный тест или проверка с мок-публикацией подтверждает, что +- [x] Интеграционный тест или проверка с мок-публикацией подтверждает, что сервисный контур использует новую модель, а не старую плоскую генерацию. -- [ ] `generator/README.md` описывает новые параметры и новую модель без старых +- [x] `generator/README.md` описывает новые параметры и новую модель без старых предупреждений о сломанном `click_id`. -- [ ] `generator/KNOWN_ISSUES.md` обновлён: старый дефект закрыт или перенесён в +- [x] `generator/KNOWN_ISSUES.md` обновлён: старый дефект закрыт или перенесён в исторический раздел, не как актуальный блокер. ## Blocked by @@ -37,3 +37,16 @@ Status: ready-for-agent ## Comments +- 2026-06-11: добавлена проверка сервисных тиков с мок-публикацией без Kafka: + steady-stream публикует несколько событий одного визита с общим `click_id`, + разными `event_id` и согласованными записями во всех четырёх топиках; история + тиков остаётся `success`. +- 2026-06-11: `docker-compose.yml` теперь пробрасывает переменные генератора + через `${VAR:-default}`, поэтому команды вида `GEN_STATE_RESET=true docker + compose up -d generator` реально меняют окружение контейнера. Внутренние + `KAFKA_BOOTSTRAP_SERVERS` и `GEN_DATA_DIR` оставлены безопасными значениями + для контейнера. +- 2026-06-11: `generator/README.md` и `generator/KNOWN_ISSUES.md` синхронизированы + с новой моделью, состоянием v2 и историческим статусом старого дефекта. +- Проверка: `uv run --with-requirements generator/requirements.txt pytest + generator/tests -q` — 113 passed; `git diff --check` — без замечаний. diff --git a/docker-compose.yml b/docker-compose.yml index 78fc632..21e220b 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -292,17 +292,23 @@ services: environment: KAFKA_BOOTSTRAP_SERVERS: kafka:29092 # Режим steady-stream: короткие тики 1-10 сек (по умолчанию 5) - GEN_TICK_SECONDS: "5" - GEN_LAMBDA_BASE_PER_MIN: "30" - GEN_JITTER_PCT: "20" - GEN_MIN_EVENTS_PER_TICK: "1" - GEN_MAX_EVENTS_PER_TICK: "50" + GEN_TICK_SECONDS: ${GEN_TICK_SECONDS:-5} + GEN_LAMBDA_BASE_PER_MIN: ${GEN_LAMBDA_BASE_PER_MIN:-30} + GEN_JITTER_PCT: ${GEN_JITTER_PCT:-20} + GEN_MIN_EVENTS_PER_TICK: ${GEN_MIN_EVENTS_PER_TICK:-1} + GEN_MAX_EVENTS_PER_TICK: ${GEN_MAX_EVENTS_PER_TICK:-50} + GEN_MAX_SESSION_EVENTS: ${GEN_MAX_SESSION_EVENTS:-30} + GEN_MAX_ACTIVE_SESSIONS: ${GEN_MAX_ACTIVE_SESSIONS:-200} + GEN_POPULATION_MAX: ${GEN_POPULATION_MAX:-300} + GEN_P_NEW_USER: ${GEN_P_NEW_USER:-0.15} + GEN_MIN_RETURN_MINUTES: ${GEN_MIN_RETURN_MINUTES:-30} GEN_DATA_DIR: /data - GEN_ENABLED: "true" - GEN_METRICS_PORT: "9109" - # State management (сохранение состояния между рестартами) - GEN_STATE_ENABLED: "true" - GEN_STATE_RESET: "false" + GEN_SEED: ${GEN_SEED:-} + GEN_ENABLED: ${GEN_ENABLED:-true} + GEN_METRICS_PORT: ${GEN_METRICS_PORT:-9109} + # Сохранение состояния между рестартами + GEN_STATE_ENABLED: ${GEN_STATE_ENABLED:-true} + GEN_STATE_RESET: ${GEN_STATE_RESET:-false} # PYTHONUNBUFFERED для сразу видеть логи PYTHONUNBUFFERED: "1" ports: diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index f50bab0..d79dffb 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -97,6 +97,23 @@ make generator-logs | `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_STATE_ENABLED` | Сохранять state v2 между рестартами | `true` | +| `GEN_STATE_RESET` | Сбросить state при старте | `false` | + +В `docker-compose.yml` через окружение переопределяются демо-параметры модели и +режима, например: + +```bash +GEN_STATE_RESET=true GEN_LAMBDA_BASE_PER_MIN=60 docker compose up -d generator +``` + +Контейнерные `KAFKA_BOOTSTRAP_SERVERS` и `GEN_DATA_DIR` в compose оставлены +внутренними значениями `kafka:29092` и `/data`. ### Топик истории @@ -118,6 +135,8 @@ Prometheus метрики доступны на `http://localhost:9109/metrics`: - `generator_events_total` — счётчик отправленных событий - `generator_publish_errors_total` — ошибки публикации - `generator_tick_duration_seconds` — длительность тика +- `generator_last_success_timestamp` — время последнего успешного или частично + успешного тика ### Мониторинг через Grafana diff --git a/generator/KNOWN_ISSUES.md b/generator/KNOWN_ISSUES.md index 652e1a7..31b2c64 100644 --- a/generator/KNOWN_ISSUES.md +++ b/generator/KNOWN_ISSUES.md @@ -1,35 +1,41 @@ # Генератор: известные проблемы и контекст для доработки -> **Статус (2026-06-06):** ветка `feature/data-generator` **не влита** в `main`. -> Генератор **не используется** как источник данных для витрин DM/дашборда. -> Витрины и Superset-дашборд строятся на **статическом сиде** (`data/*.jsonl`). -> Причина — ниже. Это не «сырой код по мелочи», а концептуальный дефект -> генеративной модели, который надо осознанно чинить перед использованием. +> **Статус (2026-06-11):** исторический дефект старой плоской генерации закрыт +> для режима `steady-stream`. Генератор строит визиты с общим `click_id`, +> монотонным временем событий, путём по страницам воронки, популяцией +> возвращающихся пользователей и состоянием v2 для активных визитов. > -> **Обновление (2026-06-11):** первый срез модели визита реализован в -> `generate_batch()`: один публичный вызов строит один `click_id` с несколькими -> событиями, марковским путём по страницам и запланированными строго растущими -> метками времени. Остальные пункты ниже остаются полезным историческим -> контекстом и списком следующих шагов: популяция возвращающихся пользователей, -> межсессионные паузы и полноценное состояние активных визитов ещё не закрыты. +> Эта заметка больше не является предупреждением «генератор концептуально +> сломан». Она оставлена как учебный разбор старого дефекта и как место для +> небольших остаточных ограничений. Заметка написана при дизайне Superset-дашборда (ветка `docs/advanced-clickstream-course`): разбирались, почему на дашборде -`Unique Users == Unique Sessions`, и по ходу вскрылось, что генератор -семантику сессии не чинит, а ломает сильнее. Чтобы при возвращении к -генератору не переоткрывать это заново — фиксирую понимание целиком. +`Unique Users == Unique Sessions`, и по ходу вскрылось, что старая реализация +генератора на тот момент семантику сессии не чинила, а ломала сильнее. Разбор +оставлен, чтобы не переоткрывать этот дефект заново и показать, почему новая +модель устроена иначе. ## TL;DR - **Модель интенсивности потока (сколько событий и когда) — нормальная.** - Poisson по тикам + дневной коэффициент + jitter. Её можно оставить. -- **Генеративная модель сущностей исправляется по шагам.** Срез одного визита - уже не штампует свежий `click_id` на каждое событие, но полная иерархия - пользователь → несколько визитов → события ещё требует популяции - возвращающихся пользователей. -- **Вывод:** прежде чем использовать генератор как полноценный источник, - доделать оставшиеся уровни модели сущностей. Математику интенсивности - трогать не обязательно. + Poisson по тикам + дневной коэффициент + jitter сохранены. +- **Генеративная модель сущностей переписана.** Поток больше не штампует свежий + `click_id` на каждое событие: визит живёт несколько событий, пользователь + может вернуться в новом визите после кулдауна. +- **Вывод:** старый дефект `Sessions == Events` не считается актуальным + блокером. Дальше генератор можно улучшать уже как работающую учебную модель, + а не как концептуально сломанный источник. + +## Текущие ограничения + +- Тип события остаётся `pageview` для всех событий. Это осознанное ограничение: + текущий дашборд и уроки строят воронку по `page_url_path`, а не по + `event_type`. +- Device/geo-профиль пользователя стабилен между визитами. Смену устройства + генератор пока не моделирует. +- При долгом простое больше 30 минут активный визит закрывается без досылки + остатка. Популяция пользователей при этом сохраняется. ## Доменная модель (как задумано в DDL) @@ -59,9 +65,9 @@ user_domain_id (постоянный пользователь, cookie) Это адекватная модель *интенсивности во времени*. Претензий к ней нет. -## Исторический корневой дефект: модель сущностей в `generate_batch()` +## Исторический корневой дефект: старая модель сущностей -До среза от 2026-06-11 прежняя реализация `generate_batch()` в монолитном +До переработки 2026-06-11 прежняя реализация в монолитном `generator/generator.py` на каждое событие в батче делала примерно следующее: ```python @@ -94,33 +100,34 @@ device_event = {**base_device, "click_id": new_click_id} # user_domain_i | `user_domain_id` | 1:1 с `click_id` | переиспользуется (потолок ~99) | | Семантика | `Sessions == Users` (вырождено по пользователю) | `Sessions == Events` (сессия = одно событие) | -Парадокс: **статический сид как учебная основа лучше**, потому что на нём -`click_id` несёт осмысленную семантику визита. Генератор её ломает. +Парадокс был таким: **статический сид как учебная основа был лучше**, потому что +на нём `click_id` нёс осмысленную семантику визита, а старый генератор её +ломал. -## Что перепроверить и переделать перед использованием +## Что сделано в новой модели -Чинить нужно **генеративную модель сущностей**, а не математику интенсивности: +Чинили именно **генеративную модель сущностей**, не переписывая математику +интенсивности: 1. **Иерархическая генерация вместо плоской выборки:** - - поддерживать популяцию пользователей с *постоянным* `user_domain_id`; - - пользователь со временем открывает 1..N **сессий** (новый `click_id` на - сессию, с межсессионными паузами — модель «вернувшегося пользователя»); - - сессия порождает последовательность из 1..M **событий**, разделяющих один - `click_id` и общий device/geo, упорядоченных по времени (правдоподобный путь - по страницам). -2. **Распределения, требующие проверки математики:** - - события на сессию (например, geometric/NB — длина визита); - - сессии на пользователя за период (возвраты); - - межсессионные интервалы (тайм-аут неактивности как граница сессии). -3. **Время событий** внутри сессии должно расти монотонно, а не быть `now()` для - всего батча. -4. **`event_id`/`click_id`** уже всегда новые (`uuid4`) — при иерархической - модели `click_id` должен переиспользоваться внутри сессии, а не на каждое - событие (см. оговорку в `README.md`, раздел State Recovery). + - есть ограниченная популяция пользователей с постоянным `user_domain_id`; + - пользователь со временем открывает 1..N визитов; + - визит порождает последовательность событий с одним `click_id`, общим + device/geo-контекстом и путём по страницам. +2. **Тиковый слой:** + - событийный бюджет тика превращается в рождения визитов через среднюю длину + визита; + - активные визиты живут между тиками; + - выпускаются только события, у которых наступила запланированная метка + времени. +3. **Время событий** внутри визита строго растёт и не прилипает к одному + `now()` для всего батча. +4. **Состояние v2** сохраняет популяцию, активные визиты, накопленный бюджет + рождения визитов, номер тика и состояние ГПСЧ. -После такой переделки на потоке естественно получится здоровая пирамида -`users < sessions < events`, и дашборд сможет показывать разницу -«пользователь vs сессия» честными числами. +После этой переделки поток на длинном окне и при штатных параметрах даёт +здоровую пирамиду `users < sessions < events`, и дашборд может показывать +разницу «пользователь vs сессия» честными числами. ## Ссылки diff --git a/generator/README.md b/generator/README.md index 0213449..ef0687c 100644 --- a/generator/README.md +++ b/generator/README.md @@ -1,13 +1,13 @@ # Генератор событий (MVP rev5) -> Перед использованием как источник витрин прочитать -> [KNOWN_ISSUES.md](./KNOWN_ISSUES.md): часть старого дефекта уже исправлена -> (один `click_id` на визит, путь по страницам, монотонное время, активные -> визиты между тиками), но популяция возвращающихся пользователей и -> восстановление активных визитов после рестарта ещё остаются следующими шагами. - Автономный генератор событий для Kafka с режимом `steady-stream`. +Генератор строит поток по иерархии `пользователь → визит → событие`: один +`click_id` живёт весь визит, события визита идут по страницам воронки с +монотонно растущим временем, а тиковый слой держит популяцию возвращающихся +пользователей и активные визиты между тиками. Исторический дефект старой +плоской генерации описан в [KNOWN_ISSUES.md](./KNOWN_ISSUES.md). + ## Архитектура ``` @@ -29,7 +29,7 @@ generator-service -> Kafka topics -> (потребители отдельно) | `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` | сериализуемое состояние генератора | +| `src/clickstream_generator/state.py` | сериализуемое состояние генератора v2 | | `src/clickstream_generator/metrics.py` | Prometheus-метрики | | `src/clickstream_generator/service.py` | основной цикл сервиса | | `generator.py` | запуск сервиса и совместимый фасад | @@ -79,6 +79,16 @@ generator-service -> Kafka topics -> (потребители отдельно) | `GEN_STATE_ENABLED` | Сохранять состояние между рестартами | `true` | | `GEN_STATE_RESET` | Сбросить состояние при старте | `false` | +В `docker-compose.yml` параметры генеративной модели и режима проброшены через +подстановку окружения, то есть их можно менять без правки файла: + +```bash +GEN_LAMBDA_BASE_PER_MIN=60 GEN_POPULATION_MAX=500 docker compose up -d generator +``` + +Контейнерные значения `KAFKA_BOOTSTRAP_SERVERS` и `GEN_DATA_DIR` в compose +оставлены безопасными внутренними значениями `kafka:29092` и `/data`. + ### Режим "раз в минуту" (для демо) Для контролируемых демо можно установить: @@ -197,22 +207,28 @@ docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \ --from-beginning ``` -## State Recovery (восстановление состояния) +## Восстановление состояния Генератор сохраняет своё состояние между перезапусками в Kafka-топик `generator_state` (compact topic). Это позволяет: - Продолжить нумерацию тиков с места остановки (continuity) - Сохранить последовательность случайных чисел (RNG state) -- Восстановить интенсивность генерации после рестарта - -**Важно:** восстанавливается continuity по номеру тика и интенсивности, но не гарантируется отсутствие дублирования событий — `event_id` и `click_id` всегда генерируются заново (`uuid4()`). +- Восстановить популяцию пользователей +- Продолжить активные визиты после короткого простоя ### Как работает -1. После каждого успешного тика состояние сохраняется в `generator_state` -2. При старте генератор читает последнее состояние из топика -3. Если состояние найдено - продолжает с сохранённого tick -4. Если нет - начинает с tick=1 +1. После каждого успешного тика состояние v2 сохраняется в `generator_state`. +2. При старте генератор читает последнее состояние из топика. +3. Если состояние найдено, сервис восстанавливает номер тика, состояние ГПСЧ, + популяцию пользователей, накопленный бюджет рождения визитов и активные + визиты. +4. Если состояния нет или оно невалидно, генератор начинает с чистого листа. + +Активный визит после простоя до 30 минут продолжается со своими исходными +запланированными метками времени. Если следующий шаг визита просрочен больше +чем на 30 минут, визит закрывается без досылки остатка: для демо это выглядит +как пользователь, который ушёл, пока стенд был остановлен. ### Топик `generator_state` @@ -250,7 +266,8 @@ docker compose exec kafka /opt/kafka/bin/kafka-topics.sh \ GEN_STATE_ENABLED=false docker compose up -d generator ``` -При отключенном state management генератор всегда начинает с tick=1, RNG инициализируется с GEN_SEED (или случайно). +При отключенном сохранении состояния генератор всегда начинает с `tick=1`, а +ГПСЧ инициализируется с `GEN_SEED` или случайно. ## Тестирование diff --git a/generator/tests/test_service.py b/generator/tests/test_service.py index 66a43a4..97ec95e 100644 --- a/generator/tests/test_service.py +++ b/generator/tests/test_service.py @@ -6,6 +6,7 @@ import logging import random from dataclasses import replace from datetime import datetime, timezone +from time import sleep as real_sleep from unittest.mock import MagicMock, patch import pytest @@ -13,6 +14,7 @@ from generator import ( Config, EventDictionary, EventGenerator, + EXPECTED_VISIT_EVENTS, GeneratorService, GeneratorState, KafkaBatchHistory, @@ -102,6 +104,83 @@ class TestGeneratorServiceDisabled: assert "disabled" in caplog.text.lower() or "GEN_ENABLED" in caplog.text +class TestGeneratorServiceSteadyStream: + """Проверки сервисного тика без настоящей Kafka.""" + + def test_service_ticks_publish_connected_multi_event_visit(self, base_config): + """Сервисные тики публикуют несколько связанных событий одного визита.""" + config = replace(base_config, tick_seconds=1, max_session_events=3) + service = GeneratorService(config) + service.publisher = MagicMock() + service.publisher.publish.side_effect = ( + lambda topic, events: (len(events), 0) + ) + service.history = MagicMock() + service._running = True + sleep_calls = 0 + + def stop_after_second_tick(sleep_seconds): + nonlocal sleep_calls + sleep_calls += 1 + if sleep_calls == 1: + real_sleep(sleep_seconds) + else: + service._running = False + + with patch.object( + service.generator, + "_calculate_events_count", + side_effect=[int(EXPECTED_VISIT_EVENTS), 0], + ), patch.object( + service.generator, + "_visit_pause_seconds", + return_value=0.05, + ), patch("clickstream_generator.service.time.sleep") as sleep_mock: + sleep_mock.side_effect = stop_after_second_tick + + service._main_loop() + + published = {} + for call in service.publisher.publish.call_args_list: + topic, events = call.args + published.setdefault(topic, []).extend(events) + + browser_events = published["browser_events"] + location_events = published["location_events"] + device_events = published["device_events"] + geo_events = published["geo_events"] + + assert set(published) == { + "browser_events", + "location_events", + "device_events", + "geo_events", + } + assert len(browser_events) >= 2 + assert len({event["click_id"] for event in browser_events}) == 1 + assert len({event["event_id"] for event in browser_events}) == len(browser_events) + assert {event["event_id"] for event in location_events} == { + event["event_id"] + for event in browser_events + } + assert {event["click_id"] for event in device_events} == { + browser_events[0]["click_id"] + } + assert {event["click_id"] for event in geo_events} == { + browser_events[0]["click_id"] + } + + history_records = [ + call.args[0] + for call in service.history.add.call_args_list + ] + assert [record.status for record in history_records] == ["success", "success"] + assert sum(record.sent_browser for record in history_records) == len(browser_events) + assert sum(record.sent_location for record in history_records) == len(location_events) + assert sum(record.sent_device for record in history_records) == len(device_events) + assert sum(record.sent_geo for record in history_records) == len(geo_events) + + class TestGeneratorServiceStateV2: """Тесты подключения state v2 к сервисному запуску."""