feat(generator): добавлено ускорение модельного времени

- Зачем:
  - ускоренный стенд должен проходить больше модельного времени и давать соответствующий событийный бюджет.
- Что:
  - расчёт интенсивности переведён на модельную длительность тика.
  - добавлена устойчивая выборка бюджета при больших λ.
  - усилены тесты ×K, дневного коэффициента и независимости от настенного часа.
- Проверка:
  - make generator-test.
  - review gate после issue 03 пройден после исправления underflow Poisson.
This commit is contained in:
2026-06-14 17:48:51 +03:00
parent efb07b0283
commit 589e2321b3
7 changed files with 237 additions and 47 deletions
@@ -1,4 +1,4 @@
Status: ready-for-agent Status: ready-for-human
# ×K и дневной коэффициент по модельному времени # ×K и дневной коэффициент по модельному времени
@@ -19,21 +19,68 @@ Status: ready-for-agent
## Acceptance criteria ## Acceptance criteria
- [ ] На ×1 поведение остаётся совместимым с обычным живым режимом. - [x] На ×1 поведение остаётся совместимым с обычным живым режимом.
- [ ] На ×K модельные `event_timestamp` за короткий реальный прогон покрывают - [x] На ×K модельные `event_timestamp` за короткий реальный прогон покрывают
примерно `K` раз большую модельную длительность. примерно `K` раз большую модельную длительность.
- [ ] Событийный бюджет считается по модельной длительности тика, а не по - [x] Событийный бюджет считается по модельной длительности тика, а не по
реальному времени сна процесса. реальному времени сна процесса.
- [ ] Дневной коэффициент меняется при переходе модельного времени через - [x] Дневной коэффициент меняется при переходе модельного времени через
дневные/ночные часы в зафиксированном часовом поясе, независимо от реального дневные/ночные часы в зафиксированном часовом поясе, независимо от реального
часа запуска. часа запуска.
- [ ] ClickHouse-проверка показывает повторяемые контрольные числа при тех же - [x] ClickHouse-проверка показывает повторяемые контрольные числа при тех же
`GEN_SEED`, `T0`, скорости и настройках. `GEN_SEED`, `T0`, скорости и настройках.
- [ ] Worker даёт промежуточный статус, если статистический прогон или стендовая - [x] Worker даёт промежуточный статус, если статистический прогон или стендовая
проверка занимает заметное время. проверка занимает заметное время.
- [ ] В `generator/README.md` или `docs/OPERATIONS.md` кратко описано, что ×K - [x] В `generator/README.md` или `docs/OPERATIONS.md` кратко описано, что ×K
ускоряет именно модельное время стенда. ускоряет именно модельное время стенда.
## Решение
`calculate_events_count()` теперь считает λ тика по модельной длительности:
`GEN_TICK_SECONDS * GEN_MODEL_TIME_SPEED`. При `GEN_MODEL_TIME_SPEED=1`
формула остаётся прежней.
Для больших λ старый Knuth-алгоритм заменён нормальным приближением, чтобы
ускорение ×K не упиралось в underflow `exp(-λ)` и событийный бюджет продолжал
расти вместе с модельной длительностью.
Добавлены проверки:
- средний бюджет ×10 примерно в 10 раз больше бюджета ×1;
- большой λ не залипает на старом underflow-ограничении;
- `hour_factor()` берёт час из заданного `GEN_MODEL_TIMEZONE`;
- сервисный live-тик при ×10 сдвигает `event_timestamp` с `10:00` на `10:10`;
- старый live-сценарий ×1 остаётся зелёным.
Документация в `generator/README.md` и `docs/OPERATIONS.md` уточняет, что ×K
ускоряет модельную длительность тика и событийный бюджет.
## Проверка
- `make generator-test``120 passed`.
- ClickHouse, быстрый прогон ×60 с `GEN_STATE_RESET=true`,
`GEN_MODEL_T0=2026-01-01T05:58:00+00:00`, `GEN_TICK_SECONDS=1`: за короткий
реальный прогон `event_ts` покрыл `2026-01-01 05:58:00`
`2026-01-01 06:09:00`; после 06:00 UTC поминутные числа выросли с `40/42` до
примерно `58..63`, что показывает смену ночного коэффициента.
- ClickHouse, повторяемость: два чистых запуска с одинаковыми
`GEN_SEED=4242`, `GEN_MODEL_T0=2026-01-01T10:00:00+00:00`,
`GEN_MODEL_TIME_SPEED=60`, `GEN_TICK_SECONDS=60`,
`GEN_STATE_RESET=true` дали одинаковые контрольные числа:
`events=73`, `unique_events=73`, `unique_clicks=73`,
`min_event_ts=max_event_ts=2026-01-01 10:00:00`.
## Риски и границы
- State v2 resume, backfill, manifest и startup history не менялись. Известные
риски по ним остаются задачами 04/05.
- Сервис всё ещё пишет операционные метки state и истории пачек по настенному
времени; это соответствует текущему контракту issue 03.
- Нормальное приближение для больших λ не является точным Poisson, но для
событийного бюджета учебного демо достаточно сохраняет масштаб, среднее и
дисперсию. Точную статистическую выборку можно выделить в отдельную задачу,
если она станет учебной целью.
## Blocked by ## Blocked by
- `.scratch/generator-model-time-startup-history/issues/02-model-time-to-clickhouse.md` - `.scratch/generator-model-time-startup-history/issues/02-model-time-to-clickhouse.md`
+3 -1
View File
@@ -123,7 +123,9 @@ GEN_STATE_RESET=true GEN_LAMBDA_BASE_PER_MIN=60 docker compose up -d generator
Для повторяемой проверки используйте чистый стенд и явный сброс состояния Для повторяемой проверки используйте чистый стенд и явный сброс состояния
генератора. `event_timestamp` в событиях — модельное время от `GEN_MODEL_T0`, а генератора. `event_timestamp` в событиях — модельное время от `GEN_MODEL_T0`, а
не настенное время запуска процесса. не настенное время запуска процесса. При `GEN_MODEL_TIME_SPEED=K` один тик
покрывает `GEN_TICK_SECONDS * K` модельных секунд, и событийный бюджет считается
по этой модельной длительности.
```bash ```bash
make clean make clean
+4 -2
View File
@@ -92,8 +92,10 @@ GEN_LAMBDA_BASE_PER_MIN=60 GEN_POPULATION_MAX=500 docker compose up -d generator
В живом режиме `event_timestamp` берётся из модельного времени: первый чистый В живом режиме `event_timestamp` берётся из модельного времени: первый чистый
тик стартует от `GEN_MODEL_T0`, дальше модельная точка сдвигается на тик стартует от `GEN_MODEL_T0`, дальше модельная точка сдвигается на
`GEN_TICK_SECONDS * GEN_MODEL_TIME_SPEED`. Дневной коэффициент считается по `GEN_TICK_SECONDS * GEN_MODEL_TIME_SPEED`. Событийный бюджет тика считается по
`GEN_MODEL_TIMEZONE`, а не по реальному часу запуска процесса. этой же модельной длительности, поэтому ×K даёт больше событий за короткий
реальный прогон. Дневной коэффициент считается по `GEN_MODEL_TIMEZONE`, а не по
реальному часу запуска процесса.
Контейнерные значения `KAFKA_BOOTSTRAP_SERVERS` и `GEN_DATA_DIR` в compose Контейнерные значения `KAFKA_BOOTSTRAP_SERVERS` и `GEN_DATA_DIR` в compose
оставлены безопасными внутренними значениями `kafka:29092` и `/data`. оставлены безопасными внутренними значениями `kafka:29092` и `/data`.
@@ -129,7 +129,7 @@ class EventGenerator:
def _hour_factor(self, now: datetime | None = None) -> float: def _hour_factor(self, now: datetime | None = None) -> float:
"""Совместимый wrapper над расчётом часового коэффициента.""" """Совместимый wrapper над расчётом часового коэффициента."""
return hour_factor(now, self.config.model_timezone) return hour_factor(now or self.config.model_t0, self.config.model_timezone)
def _calculate_events_count(self, now: datetime | None = None) -> int: def _calculate_events_count(self, now: datetime | None = None) -> int:
"""Совместимый wrapper над расчётом событийного бюджета.""" """Совместимый wrapper над расчётом событийного бюджета."""
@@ -2,17 +2,22 @@
import math import math
import random import random
from datetime import datetime, timezone from datetime import datetime
from zoneinfo import ZoneInfo from zoneinfo import ZoneInfo
from clickstream_generator.config import Config from clickstream_generator.config import Config
POISSON_KNUTH_MAX_LAMBDA = 100.0
def hour_factor(now: datetime | None = None, model_timezone: str = "UTC") -> float: def hour_factor(now: datetime | None = None, model_timezone: str = "UTC") -> float:
"""Возвращает коэффициент интенсивности в зависимости от часа дня.""" """Возвращает коэффициент интенсивности в зависимости от часа дня."""
current = now or datetime.now(timezone.utc) current = now
if current is None:
raise ValueError("now is required for hour_factor")
if current.tzinfo is None: if current.tzinfo is None:
current = current.replace(tzinfo=timezone.utc) current = current.replace(tzinfo=ZoneInfo("UTC"))
current = current.astimezone(ZoneInfo(model_timezone)) current = current.astimezone(ZoneInfo(model_timezone))
hour = current.hour hour = current.hour
if 9 <= hour <= 18: if 9 <= hour <= 18:
@@ -22,18 +27,10 @@ def hour_factor(now: datetime | None = None, model_timezone: str = "UTC") -> flo
return 1.0 return 1.0
def calculate_events_count( def _sample_poisson(lambda_tick: float, rng: random.Random) -> int:
config: Config, """Разыгрывает Poisson без underflow на больших λ."""
rng: random.Random, if lambda_tick >= POISSON_KNUTH_MAX_LAMBDA:
now: datetime | None = None, return max(0, round(rng.gauss(lambda_tick, math.sqrt(lambda_tick))))
) -> int:
"""Вычисляет количество событий для текущего тика (Poisson + jitter)."""
if now is None:
factor = hour_factor()
else:
factor = hour_factor(now, config.model_timezone)
lambda_minute = config.lambda_base_per_min * factor
lambda_tick = lambda_minute * (config.tick_seconds / 60.0)
count = 0 count = 0
threshold = math.exp(-lambda_tick) threshold = math.exp(-lambda_tick)
@@ -41,7 +38,22 @@ def calculate_events_count(
while product > threshold: while product > threshold:
product *= rng.random() product *= rng.random()
count += 1 count += 1
count -= 1 return count - 1
def calculate_events_count(
config: Config,
rng: random.Random,
now: datetime | None = None,
) -> int:
"""Вычисляет количество событий для текущего тика (Poisson + jitter)."""
model_time = now or config.model_t0
factor = hour_factor(model_time, config.model_timezone)
lambda_minute = config.lambda_base_per_min * factor
model_tick_seconds = config.tick_seconds * config.model_time_speed
lambda_tick = lambda_minute * (model_tick_seconds / 60.0)
count = _sample_poisson(lambda_tick, rng)
if config.jitter_pct > 0: if config.jitter_pct > 0:
jitter_factor = 1.0 + rng.uniform( jitter_factor = 1.0 + rng.uniform(
+96 -18
View File
@@ -6,7 +6,7 @@ import json
import random import random
import uuid import uuid
from dataclasses import replace from dataclasses import replace
from datetime import datetime, timedelta from datetime import datetime, timedelta, timezone
import pytest import pytest
from generator import ( from generator import (
@@ -16,6 +16,7 @@ from generator import (
TickStreamGenerator, TickStreamGenerator,
calculate_events_count, calculate_events_count,
generate_tick_batch, generate_tick_batch,
hour_factor,
) )
@@ -1030,14 +1031,19 @@ class TestEventGeneration:
class TestPoissonDistribution: class TestPoissonDistribution:
"""Тесты статистической модели.""" """Тесты статистической модели."""
def test_event_budget_mean_follows_lambda_and_hour_factor( def test_hour_factor_uses_model_timezone(self):
self, base_config, monkeypatch """Дневной коэффициент считается по заданному часовому поясу модели."""
): assert hour_factor(
datetime(2026, 1, 1, 2, 30, tzinfo=timezone.utc),
"Europe/Moscow",
) == 0.7
assert hour_factor(
datetime(2026, 1, 1, 6, 30, tzinfo=timezone.utc),
"Europe/Moscow",
) == 1.2
def test_event_budget_mean_follows_lambda_and_hour_factor(self, base_config):
"""Средний событийный бюджет следует λ и часовому коэффициенту.""" """Средний событийный бюджет следует λ и часовому коэффициенту."""
monkeypatch.setattr(
"clickstream_generator.intensity.hour_factor",
lambda: 1.2,
)
config = replace( config = replace(
base_config, base_config,
tick_seconds=60, tick_seconds=60,
@@ -1045,6 +1051,8 @@ class TestPoissonDistribution:
jitter_pct=0, jitter_pct=0,
min_events_per_tick=1, min_events_per_tick=1,
max_events_per_tick=100, max_events_per_tick=100,
model_t0=datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc),
model_timezone="UTC",
) )
rng = random.Random(config.seed) rng = random.Random(config.seed)
@@ -1053,20 +1061,16 @@ class TestPoissonDistribution:
assert mean_budget == pytest.approx(30 * 1.2, rel=0.15) assert mean_budget == pytest.approx(30 * 1.2, rel=0.15)
def test_default_tick_budget_floor_does_not_outgrow_target_lambda( def test_default_tick_budget_floor_does_not_outgrow_target_lambda(self, base_config):
self, base_config, monkeypatch
):
"""Дефолтная нижняя граница бюджета не разгоняет lambda=30 на тике 5 секунд.""" """Дефолтная нижняя граница бюджета не разгоняет lambda=30 на тике 5 секунд."""
monkeypatch.setattr(
"clickstream_generator.intensity.hour_factor",
lambda: 1.0,
)
config = replace( config = replace(
base_config, base_config,
tick_seconds=5, tick_seconds=5,
lambda_base_per_min=30, lambda_base_per_min=30,
jitter_pct=0, jitter_pct=0,
max_events_per_tick=50, max_events_per_tick=50,
model_t0=datetime(2026, 1, 1, 7, 0, tzinfo=timezone.utc),
model_timezone="UTC",
) )
rng = random.Random(config.seed) rng = random.Random(config.seed)
@@ -1075,6 +1079,69 @@ class TestPoissonDistribution:
assert events_per_minute == pytest.approx(config.lambda_base_per_min, rel=0.20) assert events_per_minute == pytest.approx(config.lambda_base_per_min, rel=0.20)
def test_event_budget_uses_model_tick_duration(self, base_config):
"""При ускорении событийный бюджет растёт по модельной длительности тика."""
base = replace(
base_config,
tick_seconds=60,
lambda_base_per_min=30,
jitter_pct=0,
min_events_per_tick=1,
max_events_per_tick=10_000,
model_time_speed=1,
)
accelerated = replace(base, model_time_speed=10)
model_tick_at = datetime(2026, 1, 1, 7, 0)
normal_rng = random.Random(base.seed)
accelerated_rng = random.Random(accelerated.seed)
normal_samples = [
calculate_events_count(base, normal_rng, now=model_tick_at)
for _ in range(300)
]
accelerated_samples = [
calculate_events_count(accelerated, accelerated_rng, now=model_tick_at)
for _ in range(300)
]
normal_mean = sum(normal_samples) / len(normal_samples)
accelerated_mean = sum(accelerated_samples) / len(accelerated_samples)
assert accelerated_mean / normal_mean == pytest.approx(10, rel=0.15)
def test_large_event_budget_does_not_stick_on_knuth_underflow(self, base_config):
"""Для λ > 1000 средний бюджет растёт вместе с целевой интенсивностью."""
config = replace(
base_config,
tick_seconds=60,
lambda_base_per_min=1200,
jitter_pct=0,
min_events_per_tick=0,
max_events_per_tick=10_000,
model_t0=datetime(2026, 1, 1, 7, 0, tzinfo=timezone.utc),
model_timezone="UTC",
)
rng = random.Random(config.seed)
samples = [calculate_events_count(config, rng) for _ in range(500)]
mean_budget = sum(samples) / len(samples)
assert mean_budget == pytest.approx(1200, rel=0.05)
assert mean_budget > 1000
def test_event_generator_hour_factor_defaults_to_model_t0(
self, event_dictionary, base_config
):
"""Wrapper без аргумента берёт модельную точку, а не настенный час."""
config = replace(
base_config,
model_t0=datetime(2026, 1, 1, 6, 30, tzinfo=timezone.utc),
model_timezone="Europe/Moscow",
)
generator = EventGenerator(event_dictionary, config)
assert generator._hour_factor() == 1.2
def test_calculate_events_respects_bounds(self, event_dictionary, base_config): def test_calculate_events_respects_bounds(self, event_dictionary, base_config):
"""Расчет количества событий уважает границы.""" """Расчет количества событий уважает границы."""
generator = EventGenerator(event_dictionary, base_config) generator = EventGenerator(event_dictionary, base_config)
@@ -1084,10 +1151,21 @@ class TestPoissonDistribution:
assert all(s >= base_config.min_events_per_tick for s in samples) assert all(s >= base_config.min_events_per_tick for s in samples)
assert all(s <= base_config.max_events_per_tick for s in samples) assert all(s <= base_config.max_events_per_tick for s in samples)
def test_jitter_increases_variance(self, event_dictionary, base_config, config_no_jitter): def test_jitter_increases_variance(self, event_dictionary, base_config):
"""Jitter увеличивает дисперсию.""" """Jitter увеличивает дисперсию."""
gen_with = EventGenerator(event_dictionary, base_config) config_with_jitter = replace(
gen_without = EventGenerator(event_dictionary, config_no_jitter) base_config,
tick_seconds=60,
lambda_base_per_min=30,
jitter_pct=50,
min_events_per_tick=1,
max_events_per_tick=100,
model_t0=datetime(2026, 1, 1, 7, 0, tzinfo=timezone.utc),
model_timezone="UTC",
)
config_without_jitter = replace(config_with_jitter, jitter_pct=0)
gen_with = EventGenerator(event_dictionary, config_with_jitter)
gen_without = EventGenerator(event_dictionary, config_without_jitter)
samples_with = [gen_with._calculate_events_count() for _ in range(200)] samples_with = [gen_with._calculate_events_count() for _ in range(200)]
samples_without = [gen_without._calculate_events_count() for _ in range(200)] samples_without = [gen_without._calculate_events_count() for _ in range(200)]
+50 -1
View File
@@ -135,6 +135,11 @@ class TestGeneratorServiceSteadyStream:
) )
assert len(night_wall_run["browser_events"]) == len(day_wall_run["browser_events"]) assert len(night_wall_run["browser_events"]) == len(day_wall_run["browser_events"])
assert night_wall_run["budget_model_times"] == [
datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc),
datetime(2026, 1, 1, 10, 1, tzinfo=timezone.utc),
]
assert day_wall_run["budget_model_times"] == night_wall_run["budget_model_times"]
assert night_wall_run["browser_events"] assert night_wall_run["browser_events"]
timestamps = { timestamps = {
event["event_timestamp"] event["event_timestamp"]
@@ -145,6 +150,38 @@ class TestGeneratorServiceSteadyStream:
"2026-01-01 10:01:00.000000", "2026-01-01 10:01:00.000000",
} }
def test_service_live_tick_advances_event_timestamps_by_model_speed(self, base_config):
"""При ×K сервис сдвигает события на ускоренный модельный шаг."""
model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc)
config = replace(
base_config,
tick_seconds=60,
lambda_base_per_min=600,
jitter_pct=0,
min_events_per_tick=1,
max_events_per_tick=1000,
max_session_events=1,
max_active_sessions=250,
population_max=251,
model_t0=model_t0,
model_time_speed=10,
)
published = self._run_service_ticks(
config,
wall_now=datetime(2026, 6, 14, 11, 0, tzinfo=timezone.utc),
ticks_count=2,
)
timestamps = {
event["event_timestamp"]
for event in published["browser_events"]
}
assert timestamps == {
"2026-01-01 10:00:00.000000",
"2026-01-01 10:10:00.000000",
}
def _run_service_ticks(self, config, wall_now: datetime, ticks_count: int): def _run_service_ticks(self, config, wall_now: datetime, ticks_count: int):
service = GeneratorService(config) service = GeneratorService(config)
service.publisher = MagicMock() service.publisher = MagicMock()
@@ -154,6 +191,8 @@ class TestGeneratorServiceSteadyStream:
service.history = MagicMock() service.history = MagicMock()
service._running = True service._running = True
sleep_calls = 0 sleep_calls = 0
budget_model_times = []
original_calculate_events_count = service.generator._calculate_events_count
class FrozenDateTime(datetime): class FrozenDateTime(datetime):
@classmethod @classmethod
@@ -168,7 +207,16 @@ class TestGeneratorServiceSteadyStream:
if sleep_calls >= ticks_count: if sleep_calls >= ticks_count:
service._running = False service._running = False
with patch("clickstream_generator.intensity.datetime", FrozenDateTime), \ def calculate_events_count(now=None):
budget_model_times.append(now)
return original_calculate_events_count(now=now)
with patch("clickstream_generator.service.datetime", FrozenDateTime), \
patch.object(
service.generator,
"_calculate_events_count",
side_effect=calculate_events_count,
), \
patch("clickstream_generator.service.time.sleep", side_effect=stop_after_tick): patch("clickstream_generator.service.time.sleep", side_effect=stop_after_tick):
service._main_loop() service._main_loop()
@@ -176,6 +224,7 @@ class TestGeneratorServiceSteadyStream:
for call in service.publisher.publish.call_args_list: for call in service.publisher.publish.call_args_list:
topic, events = call.args topic, events = call.args
published.setdefault(topic, []).extend(events) published.setdefault(topic, []).extend(events)
published["budget_model_times"] = budget_model_times
return published return published
def test_service_ticks_publish_connected_multi_event_visit(self, base_config): def test_service_ticks_publish_connected_multi_event_visit(self, base_config):