diff --git a/.scratch/feature-data-generator/issues/03-tick-stream-with-active-visits.md b/.scratch/feature-data-generator/issues/03-tick-stream-with-active-visits.md index 5a50ca3..3f9de5b 100644 --- a/.scratch/feature-data-generator/issues/03-tick-stream-with-active-visits.md +++ b/.scratch/feature-data-generator/issues/03-tick-stream-with-active-visits.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: ready-for-human # Поток по тикам с активными визитами @@ -17,15 +17,15 @@ Status: ready-for-agent ## Acceptance criteria -- [ ] Тест через публичный интерфейс генератора показывает, что один `click_id` +- [x] Тест через публичный интерфейс генератора показывает, что один `click_id` может появляться в выходе нескольких последовательных тиков. -- [ ] На каждом тике выпускаются только созревшие события активных визитов. -- [ ] `event_timestamp` отражает запланированное время события, а не время +- [x] На каждом тике выпускаются только созревшие события активных визитов. +- [x] `event_timestamp` отражает запланированное время события, а не время фактической отправки тика. -- [ ] Завершённые визиты больше не выпускают события в следующих тиках. -- [ ] При достижении потолка активных визитов новые рождения в этот тик +- [x] Завершённые визиты больше не выпускают события в следующих тиках. +- [x] При достижении потолка активных визитов новые рождения в этот тик пропускаются, а бюджет не копится бесконечно. -- [ ] Конфигурация запрещает состояние, где потолок активных визитов не меньше +- [x] Конфигурация запрещает состояние, где потолок активных визитов не меньше потолка популяции пользователей. ## Blocked by @@ -34,3 +34,8 @@ Status: ready-for-agent - `.scratch/feature-data-generator/issues/02-5-generator-service-cleanup.md` ## Comments + +2026-06-11: реализовано через `TickStreamGenerator` в тиковом слое. Проверка: +`uv run --with-requirements generator/requirements.txt pytest generator/tests -q` +— 81 passed. После ревью исправлено: входной `event_budget` снова трактуется +как бюджет событий, а не как число рождений визитов. diff --git a/generator/README.md b/generator/README.md index b6362c2..7237644 100644 --- a/generator/README.md +++ b/generator/README.md @@ -2,9 +2,9 @@ > Перед использованием как источник витрин прочитать > [KNOWN_ISSUES.md](./KNOWN_ISSUES.md): часть старого дефекта уже исправлена -> (один `click_id` на визит, путь по страницам, монотонное время), но популяция -> возвращающихся пользователей и полноценное состояние активных визитов ещё -> остаются следующими шагами. +> (один `click_id` на визит, путь по страницам, монотонное время, активные +> визиты между тиками), но популяция возвращающихся пользователей и +> восстановление активных визитов после рестарта ещё остаются следующими шагами. Автономный генератор событий для Kafka с режимом `steady-stream`. @@ -27,7 +27,7 @@ generator-service -> Kafka topics -> (потребители отдельно) | `src/clickstream_generator/dictionary.py` | загрузка и индексы исходных JSONL | | `src/clickstream_generator/generation.py` | генерация одного связанного визита | | `src/clickstream_generator/intensity.py` | расчёт событийного бюджета тика | -| `src/clickstream_generator/runtime.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/metrics.py` | Prometheus-метрики | @@ -38,8 +38,10 @@ generator-service -> Kafka topics -> (потребители отдельно) - Публикуем постепенно, **короткими тиками** (по умолчанию каждые 5 секунд) - На каждом тике отправляем небольшую порцию сообщений -- Держим целевую интенсивность `events/min` без крупных минутных batch: тик - набирается из одного или нескольких визитов до рассчитанного бюджета событий +- Держим целевую интенсивность `events/min` без крупных минутных batch: + рассчитанный событийный бюджет тика планирует новые визиты и списывается по + фактической длине этих визитов, а события выходят позже по своим + запланированным меткам времени - Распределяем события по 4 топикам: - `browser_events` - `location_events` @@ -48,6 +50,8 @@ generator-service -> Kafka topics -> (потребители отдельно) - Публичный вызов генеративного ядра строит один визит: общий `click_id`, разные `event_id`, общий device/geo-контекст, путь по страницам воронки и строго растущие запланированные `event_timestamp`. +- Тиковый слой хранит активные визиты между вызовами и выпускает только события, + у которых наступил `event_timestamp`; завершённые визиты удаляются из памяти. - Сохраняем связи `event_id <-> location`, `click_id <-> device/geo`. ## Конфигурация (env) @@ -61,6 +65,8 @@ generator-service -> Kafka topics -> (потребители отдельно) | `GEN_MIN_EVENTS_PER_TICK` | Минимум событий за тик | `5` | | `GEN_MAX_EVENTS_PER_TICK` | Максимум событий за тик | `50` | | `GEN_MAX_SESSION_EVENTS` | Потолок длины одного визита, защита от петель | `30` | +| `GEN_MAX_ACTIVE_SESSIONS` | Потолок одновременных активных визитов | `200` | +| `GEN_POPULATION_MAX` | Будущий потолок популяции пользователей; сейчас нужен для валидации активных визитов | `300` | | `GEN_DATA_DIR` | Путь к JSONL файлам | `/data` | | `GEN_SEED` | Сид для воспроизводимости | — | | `GEN_ENABLED` | Включить генерацию | `true` | diff --git a/generator/generator.py b/generator/generator.py index ba5578d..38238e1 100644 --- a/generator/generator.py +++ b/generator/generator.py @@ -34,7 +34,7 @@ from clickstream_generator.metrics import ( METRICS_LAST_SUCCESS, METRICS_TICK_DURATION, ) -from clickstream_generator.runtime import generate_tick_batch +from clickstream_generator.runtime import TickStreamGenerator, generate_tick_batch from clickstream_generator.service import GeneratorService, main from clickstream_generator.state import GeneratorState, _nested_list_to_tuple @@ -59,6 +59,7 @@ __all__ = [ "METRICS_EVENTS_TOTAL", "METRICS_LAST_SUCCESS", "METRICS_TICK_DURATION", + "TickStreamGenerator", "_import_kafka", "_nested_list_to_tuple", "_with_retry", diff --git a/generator/src/clickstream_generator/config.py b/generator/src/clickstream_generator/config.py index bfba52c..675e76c 100644 --- a/generator/src/clickstream_generator/config.py +++ b/generator/src/clickstream_generator/config.py @@ -30,6 +30,12 @@ class Config: max_session_events: int = field( default_factory=lambda: int(os.getenv("GEN_MAX_SESSION_EVENTS", "30")) ) + max_active_sessions: int = field( + default_factory=lambda: int(os.getenv("GEN_MAX_ACTIVE_SESSIONS", "200")) + ) + population_max: int = field( + default_factory=lambda: int(os.getenv("GEN_POPULATION_MAX", "300")) + ) data_dir: Path = field( default_factory=lambda: Path(os.getenv("GEN_DATA_DIR", "/data")) ) @@ -58,5 +64,11 @@ class Config: raise ValueError("GEN_LAMBDA_BASE_PER_MIN must be >= 1") if self.max_session_events < 1: raise ValueError("GEN_MAX_SESSION_EVENTS must be >= 1") + if self.max_active_sessions < 1: + raise ValueError("GEN_MAX_ACTIVE_SESSIONS must be >= 1") + if self.population_max < 1: + raise ValueError("GEN_POPULATION_MAX must be >= 1") + if self.max_active_sessions >= self.population_max: + raise ValueError("GEN_MAX_ACTIVE_SESSIONS must be < GEN_POPULATION_MAX") if not self.data_dir.exists(): raise ValueError(f"Data directory does not exist: {self.data_dir}") diff --git a/generator/src/clickstream_generator/generation.py b/generator/src/clickstream_generator/generation.py index ce31802..50bc276 100644 --- a/generator/src/clickstream_generator/generation.py +++ b/generator/src/clickstream_generator/generation.py @@ -129,7 +129,11 @@ class EventGenerator: """Совместимый wrapper над расчётом событийного бюджета.""" return calculate_events_count(self.config, self.rng) - def generate_batch(self, batch_size: int) -> dict[str, list[dict]]: + def generate_batch( + self, + batch_size: int, + planned_start_at: datetime | None = None, + ) -> dict[str, list[dict]]: """Генерирует один визит с сохранением связей.""" if not self.dictionary.browser_events: return { @@ -176,7 +180,9 @@ class EventGenerator: base_device = self.dictionary.device_by_click_id.get(base_click_id) base_geo = self.dictionary.geo_by_click_id.get(base_click_id) new_click_id = self._new_uuid() - planned_timestamp = datetime.now(timezone.utc) + planned_timestamp = planned_start_at or datetime.now(timezone.utc) + if planned_timestamp.tzinfo is not None: + planned_timestamp = planned_timestamp.astimezone(timezone.utc).replace(tzinfo=None) for base_browser, page_url_path in zip(base_browser_events, visit_path): base_location = self.dictionary.location_by_event_id.get(base_browser["event_id"]) diff --git a/generator/src/clickstream_generator/runtime.py b/generator/src/clickstream_generator/runtime.py index 2cbb798..614138f 100644 --- a/generator/src/clickstream_generator/runtime.py +++ b/generator/src/clickstream_generator/runtime.py @@ -1,30 +1,133 @@ -"""Переходный тиковый слой генератора.""" +"""Тиковый слой генератора с активными визитами между вызовами.""" + +from dataclasses import dataclass +from datetime import datetime, timezone +from weakref import WeakKeyDictionary from clickstream_generator.generation import EventGenerator -def generate_tick_batch(generator: EventGenerator, event_budget: int) -> dict[str, list[dict]]: - """Набирает тиковый батч из одного или нескольких полных визитов. +TOPICS = ("browser_events", "location_events", "device_events", "geo_events") - Это временный механизм до задачи 03. В ней тиковый слой начнёт хранить - активные визиты и выпускать только созревшие события. + +@dataclass +class ActiveVisit: + """Запланированный визит, который выпускается по тикам.""" + + batch: dict[str, list[dict]] + timestamps: list[datetime] + next_index: int = 0 + + @property + def is_finished(self) -> bool: + return self.next_index >= len(self.timestamps) + + +def _empty_batch() -> dict[str, list[dict]]: + return {topic: [] for topic in TOPICS} + + +def _parse_timestamp(value: str) -> datetime: + return datetime.fromisoformat(value.replace(" ", "T")) + + +def _normalize_tick_time(tick_started_at: datetime | None) -> datetime: + tick_time = tick_started_at or datetime.now(timezone.utc) + if tick_time.tzinfo is not None: + return tick_time.astimezone(timezone.utc).replace(tzinfo=None) + return tick_time + + +class TickStreamGenerator: + """Раскладывает визиты по тикам и хранит активные визиты.""" + + def __init__(self, generator: EventGenerator): + self.generator = generator + self.active_visits: list[ActiveVisit] = [] + self._pending_event_budget = 0 + + def generate_tick( + self, + event_budget: int, + tick_started_at: datetime | None = None, + ) -> dict[str, list[dict]]: + """Возвращает события, созревшие к текущему тику.""" + tick_time = _normalize_tick_time(tick_started_at) + tick_batch = _empty_batch() + + self._release_due_events(tick_time, tick_batch) + self._drop_finished_visits() + self._birth_visits(event_budget, tick_time) + self._release_due_events(tick_time, tick_batch) + self._drop_finished_visits() + + return tick_batch + + def _birth_visits(self, event_budget: int, tick_time: datetime) -> None: + if len(self.active_visits) >= self.generator.config.max_active_sessions: + self._pending_event_budget = 0 + return + + self._pending_event_budget += max(0, event_budget) + + while ( + self._pending_event_budget > 0 + and len(self.active_visits) < self.generator.config.max_active_sessions + ): + visit_batch = self.generator.generate_batch( + self.generator.config.max_session_events, + planned_start_at=tick_time, + ) + timestamps = [ + _parse_timestamp(event["event_timestamp"]) + for event in visit_batch["browser_events"] + ] + if not timestamps: + break + + self.active_visits.append(ActiveVisit(batch=visit_batch, timestamps=timestamps)) + self._pending_event_budget -= len(timestamps) + + if len(self.active_visits) >= self.generator.config.max_active_sessions: + self._pending_event_budget = 0 + + def _release_due_events( + self, + tick_time: datetime, + tick_batch: dict[str, list[dict]], + ) -> None: + for visit in self.active_visits: + while not visit.is_finished and visit.timestamps[visit.next_index] <= tick_time: + event_index = visit.next_index + for topic in TOPICS: + tick_batch[topic].append(visit.batch[topic][event_index]) + visit.next_index += 1 + + def _drop_finished_visits(self) -> None: + self.active_visits = [ + visit + for visit in self.active_visits + if not visit.is_finished + ] + + +_STREAMS: WeakKeyDictionary[EventGenerator, TickStreamGenerator] = WeakKeyDictionary() + + +def generate_tick_batch( + generator: EventGenerator, + event_budget: int, + tick_started_at: datetime | None = None, +) -> dict[str, list[dict]]: + """Совместимый фасад тикового генератора. + + `event_budget` здесь трактуется как бюджет событий за тик. Если потолок + активных визитов достигнут, новые рождения пропускаются и бюджет не + переносится на следующие тики. """ - tick_batch = { - "browser_events": [], - "location_events": [], - "device_events": [], - "geo_events": [], - } + stream = _STREAMS.get(generator) + if stream is None: + stream = TickStreamGenerator(generator) + _STREAMS[generator] = stream - remaining_events = event_budget - while remaining_events > 0: - visit_batch = generator.generate_batch(remaining_events) - generated_events = len(visit_batch["browser_events"]) - if generated_events == 0: - break - - for topic, events in visit_batch.items(): - tick_batch[topic].extend(events) - remaining_events -= generated_events - - return tick_batch + return stream.generate_tick(event_budget, tick_started_at=tick_started_at) diff --git a/generator/src/clickstream_generator/service.py b/generator/src/clickstream_generator/service.py index a6f3f3f..2029160 100644 --- a/generator/src/clickstream_generator/service.py +++ b/generator/src/clickstream_generator/service.py @@ -133,7 +133,7 @@ class GeneratorService: try: events_count = self.generator._calculate_events_count() - logger.info(f"Generating ~{events_count} base events") + logger.info(f"Generating with event budget ~{events_count}") gen_start = time.time() batch = generate_tick_batch(self.generator, events_count) diff --git a/generator/tests/test_config.py b/generator/tests/test_config.py index 41080d7..2201bd5 100644 --- a/generator/tests/test_config.py +++ b/generator/tests/test_config.py @@ -29,6 +29,12 @@ class TestConfigValidation: with pytest.raises(ValueError, match="GEN_MAX_SESSION_EVENTS"): replace(base_config, max_session_events=0) + def test_max_active_sessions_must_be_less_than_population(self, base_config): + """Потолок активных визитов должен быть меньше потолка популяции.""" + from dataclasses import replace + with pytest.raises(ValueError, match="GEN_MAX_ACTIVE_SESSIONS"): + replace(base_config, max_active_sessions=10, population_max=10) + def test_data_dir_must_exist(self, base_config): """data_dir должен существовать.""" from dataclasses import replace @@ -40,6 +46,7 @@ class TestConfigValidation: assert base_config.tick_seconds == 5 assert base_config.lambda_base_per_min == 200 assert base_config.jitter_pct == 20 + assert base_config.max_active_sessions < base_config.population_max class TestConfigDefaults: diff --git a/generator/tests/test_generation.py b/generator/tests/test_generation.py index afca78f..ea07356 100644 --- a/generator/tests/test_generation.py +++ b/generator/tests/test_generation.py @@ -4,10 +4,10 @@ import uuid from dataclasses import replace -from datetime import datetime +from datetime import datetime, timedelta import pytest -from generator import EventGenerator, EventDictionary, generate_tick_batch +from generator import EventGenerator, EventDictionary, TickStreamGenerator, generate_tick_batch ALLOWED_PAGE_PATHS = { @@ -91,16 +91,144 @@ class TestEventGeneration: assert len(batch["device_events"]) == browser_count assert len(batch["geo_events"]) == browser_count - def test_generate_tick_batch_fills_requested_event_budget(self, event_dictionary, base_config): - """Тиковый батч набирает целевой бюджет из одного или нескольких визитов.""" + def test_generate_tick_batch_uses_requested_event_budget(self, event_dictionary, base_config): + """Тиковый батч трактует входное число как событийный бюджет.""" generator = EventGenerator(event_dictionary, base_config) batch = generate_tick_batch(generator, 20) - assert len(batch["browser_events"]) == 20 - assert len(batch["location_events"]) == 20 - assert len(batch["device_events"]) == 20 - assert len(batch["geo_events"]) == 20 + browser_count = len(batch["browser_events"]) + + assert 1 <= browser_count <= 20 + assert len(batch["location_events"]) == browser_count + assert len(batch["device_events"]) == browser_count + assert len(batch["geo_events"]) == browser_count + assert len({event["click_id"] for event in batch["browser_events"]}) == browser_count + + def test_tick_stream_keeps_long_window_intensity_near_event_budget( + self, event_dictionary, base_config + ): + """Длинное окно не разгоняется от событийного бюджета к бюджету визитов.""" + config = replace( + base_config, + tick_seconds=5, + jitter_pct=0, + min_events_per_tick=17, + max_events_per_tick=17, + max_active_sessions=10_000, + population_max=10_001, + ) + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + tick_at = datetime(2026, 6, 11, 12, 0) + ticks_count = 12 * 60 + total_events = 0 + + for tick_index in range(ticks_count): + batch = stream.generate_tick( + event_budget=17, + tick_started_at=tick_at + timedelta(seconds=5 * tick_index), + ) + total_events += len(batch["browser_events"]) + + events_per_minute = total_events / (ticks_count * config.tick_seconds / 60) + + assert events_per_minute < 300 + + def test_tick_stream_keeps_active_visit_between_ticks(self, event_dictionary, base_config): + """Один визит может выпускать события в нескольких последовательных тиках.""" + config = replace(base_config, tick_seconds=60, max_session_events=5) + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + first_tick_at = datetime(2026, 6, 11, 12, 0) + + first_tick = stream.generate_tick(event_budget=1, tick_started_at=first_tick_at) + second_tick = stream.generate_tick( + event_budget=0, + tick_started_at=first_tick_at + timedelta(seconds=config.tick_seconds), + ) + + first_click_id = first_tick["browser_events"][0]["click_id"] + second_click_ids = { + event["click_id"] + for event in second_tick["browser_events"] + } + second_timestamps = _parse_event_timestamps(second_tick) + + assert first_click_id in second_click_ids + assert all( + timestamp < first_tick_at + timedelta(seconds=config.tick_seconds) + for timestamp in second_timestamps + ) + + def test_tick_stream_releases_only_matured_events(self, event_dictionary, base_config): + """Тик не выпускает будущие события активного визита.""" + generator = EventGenerator(event_dictionary, base_config) + stream = TickStreamGenerator(generator) + tick_at = datetime(2026, 6, 11, 12, 0) + + first_tick = stream.generate_tick(event_budget=1, tick_started_at=tick_at) + same_time_tick = stream.generate_tick(event_budget=0, tick_started_at=tick_at) + + assert len(first_tick["browser_events"]) == 1 + assert same_time_tick["browser_events"] == [] + assert same_time_tick["location_events"] == [] + assert same_time_tick["device_events"] == [] + assert same_time_tick["geo_events"] == [] + + def test_tick_stream_does_not_reemit_finished_visit(self, event_dictionary, base_config): + """Завершённый визит больше не выпускает события в следующих тиках.""" + config = replace(base_config, max_session_events=2) + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + tick_at = datetime(2026, 6, 11, 12, 0) + + first_tick = stream.generate_tick(event_budget=1, tick_started_at=tick_at) + final_tick = stream.generate_tick( + event_budget=0, + tick_started_at=tick_at + timedelta(hours=1), + ) + later_tick = stream.generate_tick( + event_budget=0, + tick_started_at=tick_at + timedelta(hours=2), + ) + + click_id = first_tick["browser_events"][0]["click_id"] + + assert {event["click_id"] for event in final_tick["browser_events"]} == {click_id} + assert later_tick["browser_events"] == [] + + def test_tick_stream_drops_births_when_active_limit_is_reached( + self, event_dictionary, base_config + ): + """При заполненном потолке новые рождения пропускаются без накопления бюджета.""" + config = replace( + base_config, + max_session_events=2, + max_active_sessions=1, + population_max=2, + ) + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + tick_at = datetime(2026, 6, 11, 12, 0) + + first_tick = stream.generate_tick(event_budget=5, tick_started_at=tick_at) + blocked_tick = stream.generate_tick(event_budget=5, tick_started_at=tick_at) + final_tick = stream.generate_tick( + event_budget=0, + tick_started_at=tick_at + timedelta(hours=1), + ) + later_tick = stream.generate_tick( + event_budget=0, + tick_started_at=tick_at + timedelta(hours=2), + ) + + click_id = first_tick["browser_events"][0]["click_id"] + + assert len(first_tick["browser_events"]) == 1 + assert blocked_tick["browser_events"] == [] + assert {event["click_id"] for event in final_tick["browser_events"]} == {click_id} + assert later_tick["browser_events"] == [] def test_generate_batch_creates_one_connected_visit(self, event_dictionary, base_config): """Публичный вызов генератора создаёт один связанный визит."""