From efb07b0283de3db1e9a5f75e9fb484d31407ba6e Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 14 Jun 2026 17:23:12 +0300 Subject: [PATCH] =?UTF-8?q?feat(generator):=20=D0=B4=D0=BE=D0=B1=D0=B0?= =?UTF-8?q?=D0=B2=D0=BB=D0=B5=D0=BD=D0=BE=20=D0=BC=D0=BE=D0=B4=D0=B5=D0=BB?= =?UTF-8?q?=D1=8C=D0=BD=D0=BE=D0=B5=20=D0=B2=D1=80=D0=B5=D0=BC=D1=8F=20liv?= =?UTF-8?q?e-=D0=BF=D0=BE=D1=82=D0=BE=D0=BA=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - генератор должен писать события от модельной точки T0 и проверяться повторяемо в ClickHouse. - Что: - добавлены настройки модельного времени и передача модельной точки в live-тик. - дневной коэффициент считается по модельному времени и часовому поясу. - обновлены проверки, compose, документация и статус issue 02. - Проверка: - make generator-test. - два чистых ClickHouse-прогона с GEN_STATE_RESET=true дали одинаковые контрольные числа. --- .../issues/02-model-time-to-clickhouse.md | 31 +++++--- docker-compose.yml | 4 ++ docs/OPERATIONS.md | 46 ++++++++++++ generator/README.md | 9 +++ generator/src/clickstream_generator/config.py | 40 +++++++++++ .../src/clickstream_generator/generation.py | 8 +-- .../src/clickstream_generator/intensity.py | 18 ++++- .../src/clickstream_generator/service.py | 20 +++++- generator/tests/test_config.py | 19 +++++ generator/tests/test_service.py | 71 +++++++++++++++++++ 10 files changed, 248 insertions(+), 18 deletions(-) diff --git a/.scratch/generator-model-time-startup-history/issues/02-model-time-to-clickhouse.md b/.scratch/generator-model-time-startup-history/issues/02-model-time-to-clickhouse.md index 0d684f0..4684dd6 100644 --- a/.scratch/generator-model-time-startup-history/issues/02-model-time-to-clickhouse.md +++ b/.scratch/generator-model-time-startup-history/issues/02-model-time-to-clickhouse.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: ready-for-human # Модельное время до ClickHouse @@ -16,21 +16,36 @@ Status: ready-for-agent ## Acceptance criteria -- [ ] При фиксированных `GEN_SEED` и `T0` первый короткий прогон пишет события с +- [x] При фиксированных `GEN_SEED` и `T0` первый короткий прогон пишет события с модельными `event_timestamp`, начинающимися около `T0`. -- [ ] Повторный чистый прогон с теми же настройками даёт те же контрольные +- [x] Повторный чистый прогон с теми же настройками даёт те же контрольные числа в ClickHouse по правилу повторяемости из контракта задачи 1. -- [ ] Повторный чистый прогон явно сбрасывает или обходит старое состояние +- [x] Повторный чистый прогон явно сбрасывает или обходит старое состояние генератора, чтобы сервис не продолжил прошлый запуск из `generator_state`. -- [ ] Расчёт дневного коэффициента в этом срезе больше не зависит от реального +- [x] Расчёт дневного коэффициента в этом срезе больше не зависит от реального часа запуска процесса. -- [ ] Локальные тесты проверяют поведение через публичный интерфейс генератора +- [x] Локальные тесты проверяют поведение через публичный интерфейс генератора или сервиса, без привязки к внутреннему устройству часов. -- [ ] Есть команда или короткая инструкция для координатора: поднять стенд, +- [x] Есть команда или короткая инструкция для координатора: поднять стенд, прогнать поток, выполнить SQL-проверку в ClickHouse. -- [ ] Документация запуска не утверждает, что `event_timestamp` равен +- [x] Документация запуска не утверждает, что `event_timestamp` равен настенному времени. +## Решение + +Живой сервис получает модельную точку из `GEN_MODEL_T0`, считает дневной +коэффициент по модельному времени и передаёт эту же точку в тиковый поток. После +успешного тика модельная точка сдвигается на +`GEN_TICK_SECONDS * GEN_MODEL_TIME_SPEED`. + +Проверка ClickHouse выполнена двумя чистыми прогонами с `GEN_STATE_RESET=true`. +Оба раза получены одинаковые контрольные числа: 5 событий, диапазон +`event_ts = 2026-01-01 10:00:00.000000`, 5 уникальных `event_id` и 5 уникальных +`click_id`. + +Риски по восстановлению state v2 оставлены для задачи 04: текущий срез +доказывает чистый live-старт, а не возобновление после сбоя. + ## Blocked by - `.scratch/generator-model-time-startup-history/issues/01-time-and-startup-history-contract.md` diff --git a/docker-compose.yml b/docker-compose.yml index 21e220b..9570ef7 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -302,6 +302,10 @@ services: 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_MODEL_T0: ${GEN_MODEL_T0:-2026-01-01T00:00:00+00:00} + GEN_MODEL_TIMEZONE: ${GEN_MODEL_TIMEZONE:-UTC} + GEN_MODEL_TIME_SPEED: ${GEN_MODEL_TIME_SPEED:-1} + GEN_RUN_MODE: ${GEN_RUN_MODE:-live} GEN_DATA_DIR: /data GEN_SEED: ${GEN_SEED:-} GEN_ENABLED: ${GEN_ENABLED:-true} diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index d79dffb..6dba349 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -102,6 +102,10 @@ make generator-logs | `GEN_POPULATION_MAX` | Потолок активной популяции пользователей | `300` | | `GEN_P_NEW_USER` | Доля визитов новых пользователей | `0.15` | | `GEN_MIN_RETURN_MINUTES` | Минимальная пауза перед возвратом пользователя | `30` | +| `GEN_MODEL_T0` | Стартовая модельная точка, ISO 8601 с часовым поясом | `2026-01-01T00:00:00+00:00` | +| `GEN_MODEL_TIMEZONE` | Часовой пояс модельных часов для дневного коэффициента | `UTC` | +| `GEN_MODEL_TIME_SPEED` | Сколько модельных секунд проходит за одну настенную секунду | `1` | +| `GEN_RUN_MODE` | Режим генератора | `live` | | `GEN_STATE_ENABLED` | Сохранять state v2 между рестартами | `true` | | `GEN_STATE_RESET` | Сбросить state при старте | `false` | @@ -115,6 +119,48 @@ 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`. +### Проверка модельного времени в ClickHouse + +Для повторяемой проверки используйте чистый стенд и явный сброс состояния +генератора. `event_timestamp` в событиях — модельное время от `GEN_MODEL_T0`, а +не настенное время запуска процесса. + +```bash +make clean +docker compose up -d clickhouse kafka +make ddl + +GEN_SEED=4242 \ +GEN_MODEL_T0=2026-01-01T10:00:00+00:00 \ +GEN_MODEL_TIMEZONE=UTC \ +GEN_STATE_RESET=true \ +GEN_TICK_SECONDS=60 \ +GEN_LAMBDA_BASE_PER_MIN=60 \ +GEN_JITTER_PCT=0 \ +docker compose up -d --build generator + +# Убедитесь по логам, что Tick 1 завершился, а Tick 2 ещё не стартовал. +sleep 9 +docker compose logs --tail=40 generator +docker compose stop generator +sleep 5 +bash scripts/run_batch.sh + +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query " +SELECT + count() AS events, + min(event_ts) AS min_event_ts, + max(event_ts) AS max_event_ts, + uniqExact(event_id) AS unique_events, + uniqExact(click_id) AS unique_clicks +FROM ods.browser_event +FORMAT Vertical" +``` + +Повторите блок с теми же значениями. При чистом стенде и `GEN_STATE_RESET=true` +контрольные числа должны совпасть, а `min_event_ts` должен начинаться от +`2026-01-01 10:00:00`. + ### Топик истории Генератор пишет историю батчей в топик `generator_batch_history` (JSON, ключ `batch_id`). diff --git a/generator/README.md b/generator/README.md index ef0687c..e7a9a19 100644 --- a/generator/README.md +++ b/generator/README.md @@ -72,6 +72,10 @@ generator-service -> Kafka topics -> (потребители отдельно) | `GEN_POPULATION_MAX` | Потолок активной популяции пользователей | `300` | | `GEN_P_NEW_USER` | Вероятность отдать новый визит новому пользователю | `0.15` | | `GEN_MIN_RETURN_MINUTES` | Минимальная пауза перед возвратом пользователя | `30` | +| `GEN_MODEL_T0` | Стартовая модельная точка, ISO 8601 с часовым поясом | `2026-01-01T00:00:00+00:00` | +| `GEN_MODEL_TIMEZONE` | Часовой пояс модельных часов для дневного коэффициента | `UTC` | +| `GEN_MODEL_TIME_SPEED` | Сколько модельных секунд проходит за одну настенную секунду | `1` | +| `GEN_RUN_MODE` | Режим генератора | `live` | | `GEN_DATA_DIR` | Путь к JSONL файлам | `/data` | | `GEN_SEED` | Сид для воспроизводимости | — | | `GEN_ENABLED` | Включить генерацию | `true` | @@ -86,6 +90,11 @@ generator-service -> Kafka topics -> (потребители отдельно) GEN_LAMBDA_BASE_PER_MIN=60 GEN_POPULATION_MAX=500 docker compose up -d generator ``` +В живом режиме `event_timestamp` берётся из модельного времени: первый чистый +тик стартует от `GEN_MODEL_T0`, дальше модельная точка сдвигается на +`GEN_TICK_SECONDS * GEN_MODEL_TIME_SPEED`. Дневной коэффициент считается по +`GEN_MODEL_TIMEZONE`, а не по реальному часу запуска процесса. + Контейнерные значения `KAFKA_BOOTSTRAP_SERVERS` и `GEN_DATA_DIR` в compose оставлены безопасными внутренними значениями `kafka:29092` и `/data`. diff --git a/generator/src/clickstream_generator/config.py b/generator/src/clickstream_generator/config.py index a260196..d53c983 100644 --- a/generator/src/clickstream_generator/config.py +++ b/generator/src/clickstream_generator/config.py @@ -2,7 +2,18 @@ import os from dataclasses import dataclass, field +from datetime import datetime, timezone from pathlib import Path +from zoneinfo import ZoneInfo, ZoneInfoNotFoundError + + +def _parse_model_timestamp(value: str) -> datetime: + """Разбирает ISO-метку модельного времени и нормализует её к UTC.""" + normalized = value.replace("Z", "+00:00") + timestamp = datetime.fromisoformat(normalized) + if timestamp.tzinfo is None: + raise ValueError("GEN_MODEL_T0 must include timezone") + return timestamp.astimezone(timezone.utc) @dataclass(frozen=True) @@ -62,6 +73,20 @@ class Config: state_reset: bool = field( default_factory=lambda: os.getenv("GEN_STATE_RESET", "false").lower() == "true" ) + model_t0: datetime = field( + default_factory=lambda: _parse_model_timestamp( + os.getenv("GEN_MODEL_T0", "2026-01-01T00:00:00+00:00") + ) + ) + model_timezone: str = field( + default_factory=lambda: os.getenv("GEN_MODEL_TIMEZONE", "UTC") + ) + model_time_speed: float = field( + default_factory=lambda: float(os.getenv("GEN_MODEL_TIME_SPEED", "1")) + ) + run_mode: str = field( + default_factory=lambda: os.getenv("GEN_RUN_MODE", "live") + ) def __post_init__(self): if self.tick_seconds < 1: @@ -82,3 +107,18 @@ class Config: 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}") + if self.model_t0.tzinfo is None: + raise ValueError("GEN_MODEL_T0 must include timezone") + object.__setattr__( + self, + "model_t0", + self.model_t0.astimezone(timezone.utc), + ) + try: + ZoneInfo(self.model_timezone) + except ZoneInfoNotFoundError as e: + raise ValueError(f"Unknown GEN_MODEL_TIMEZONE: {self.model_timezone}") from e + if self.model_time_speed <= 0: + raise ValueError("GEN_MODEL_TIME_SPEED must be > 0") + if self.run_mode not in {"live", "backfill"}: + raise ValueError("GEN_RUN_MODE must be live or backfill") diff --git a/generator/src/clickstream_generator/generation.py b/generator/src/clickstream_generator/generation.py index ec60b7a..0f4ab69 100644 --- a/generator/src/clickstream_generator/generation.py +++ b/generator/src/clickstream_generator/generation.py @@ -127,13 +127,13 @@ class EventGenerator: pause = self.rng.lognormvariate(math.log(20.0), 0.9) return max(1.0, min(pause, 29 * 60.0)) - def _hour_factor(self) -> float: + def _hour_factor(self, now: datetime | None = None) -> float: """Совместимый wrapper над расчётом часового коэффициента.""" - return hour_factor() + return hour_factor(now, self.config.model_timezone) - def _calculate_events_count(self) -> int: + def _calculate_events_count(self, now: datetime | None = None) -> int: """Совместимый wrapper над расчётом событийного бюджета.""" - return calculate_events_count(self.config, self.rng) + return calculate_events_count(self.config, self.rng, now=now) def generate_batch( self, diff --git a/generator/src/clickstream_generator/intensity.py b/generator/src/clickstream_generator/intensity.py index a09838b..9e3c921 100644 --- a/generator/src/clickstream_generator/intensity.py +++ b/generator/src/clickstream_generator/intensity.py @@ -3,13 +3,17 @@ import math import random from datetime import datetime, timezone +from zoneinfo import ZoneInfo from clickstream_generator.config import Config -def hour_factor(now: datetime | None = None) -> float: +def hour_factor(now: datetime | None = None, model_timezone: str = "UTC") -> float: """Возвращает коэффициент интенсивности в зависимости от часа дня.""" current = now or datetime.now(timezone.utc) + if current.tzinfo is None: + current = current.replace(tzinfo=timezone.utc) + current = current.astimezone(ZoneInfo(model_timezone)) hour = current.hour if 9 <= hour <= 18: return 1.2 @@ -18,9 +22,17 @@ def hour_factor(now: datetime | None = None) -> float: return 1.0 -def calculate_events_count(config: Config, rng: random.Random) -> int: +def calculate_events_count( + config: Config, + rng: random.Random, + now: datetime | None = None, +) -> int: """Вычисляет количество событий для текущего тика (Poisson + jitter).""" - lambda_minute = config.lambda_base_per_min * hour_factor() + 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 diff --git a/generator/src/clickstream_generator/service.py b/generator/src/clickstream_generator/service.py index 5946f77..0714ce1 100644 --- a/generator/src/clickstream_generator/service.py +++ b/generator/src/clickstream_generator/service.py @@ -4,7 +4,7 @@ import logging import sys import time import uuid -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from prometheus_client import start_http_server @@ -42,6 +42,7 @@ class GeneratorService: self.state_manager: KafkaStateManager | None = None self._running = False self._tick = 0 + self._model_time = config.model_t0 def start(self): """Запускает основной цикл.""" @@ -139,16 +140,22 @@ class GeneratorService: self._tick += 1 tick_start = time.time() batch_id = str(uuid.uuid4())[:8] + model_time = self._model_time with METRICS_TICK_DURATION.time(): logger.info(f"=== Tick {self._tick} (batch_id={batch_id}) ===") try: - events_count = self.generator._calculate_events_count() + events_count = self.generator._calculate_events_count( + now=model_time, + ) logger.info(f"Generating with event budget ~{events_count}") gen_start = time.time() - batch = self.stream.generate_tick(events_count) + batch = self.stream.generate_tick( + events_count, + tick_started_at=model_time, + ) gen_duration = time.time() - gen_start pub_start = time.time() @@ -173,6 +180,7 @@ class GeneratorService: if status in ("success", "partial"): METRICS_LAST_SUCCESS.set_to_current_time() self._save_state(batch_id) + self._advance_model_time() self.publisher.flush() pub_duration = time.time() - pub_start @@ -235,6 +243,12 @@ class GeneratorService: logger.debug(f"Sleeping for {sleep_time:.1f}s until next tick") time.sleep(sleep_time) + def _advance_model_time(self) -> None: + """Сдвигает модельное время после успешного live-тика.""" + self._model_time = self._model_time + timedelta( + seconds=self.config.tick_seconds * self.config.model_time_speed, + ) + def main(): """Точка входа сервиса.""" diff --git a/generator/tests/test_config.py b/generator/tests/test_config.py index da4721f..81a296a 100644 --- a/generator/tests/test_config.py +++ b/generator/tests/test_config.py @@ -3,6 +3,7 @@ """ from pathlib import Path +from datetime import datetime, timedelta, timezone import pytest from generator import Config @@ -60,6 +61,24 @@ class TestConfigValidation: assert base_config.jitter_pct == 20 assert base_config.max_active_sessions < base_config.population_max + def test_model_t0_is_normalized_to_utc(self, base_config): + """model_t0 нормализуется к UTC даже при прямом создании Config.""" + from dataclasses import replace + + config = replace( + base_config, + model_t0=datetime( + 2026, + 1, + 1, + 13, + 0, + tzinfo=timezone(timedelta(hours=3)), + ), + ) + + assert config.model_t0 == datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + class TestConfigDefaults: """Тесты значений по умолчанию.""" diff --git a/generator/tests/test_service.py b/generator/tests/test_service.py index 97ec95e..a3ee7ee 100644 --- a/generator/tests/test_service.py +++ b/generator/tests/test_service.py @@ -107,6 +107,77 @@ class TestGeneratorServiceDisabled: class TestGeneratorServiceSteadyStream: """Проверки сервисного тика без настоящей Kafka.""" + def test_service_live_tick_uses_model_time_for_events_and_day_factor(self, base_config): + """Живой тик пишет события от T0 и не зависит от реального часа запуска.""" + 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, + ) + + night_wall_run = self._run_service_ticks( + config, + wall_now=datetime(2026, 6, 14, 3, 0, tzinfo=timezone.utc), + ticks_count=2, + ) + day_wall_run = self._run_service_ticks( + config, + wall_now=datetime(2026, 6, 14, 11, 0, tzinfo=timezone.utc), + ticks_count=2, + ) + + assert len(night_wall_run["browser_events"]) == len(day_wall_run["browser_events"]) + assert night_wall_run["browser_events"] + timestamps = { + event["event_timestamp"] + for event in night_wall_run["browser_events"] + } + assert timestamps == { + "2026-01-01 10:00:00.000000", + "2026-01-01 10:01:00.000000", + } + + def _run_service_ticks(self, config, wall_now: datetime, ticks_count: int): + 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 + + class FrozenDateTime(datetime): + @classmethod + def now(cls, tz=None): + if tz is None: + return wall_now.replace(tzinfo=None) + return wall_now.astimezone(tz) + + def stop_after_tick(_sleep_seconds): + nonlocal sleep_calls + sleep_calls += 1 + if sleep_calls >= ticks_count: + service._running = False + + with patch("clickstream_generator.intensity.datetime", FrozenDateTime), \ + patch("clickstream_generator.service.time.sleep", side_effect=stop_after_tick): + service._main_loop() + + published = {} + for call in service.publisher.publish.call_args_list: + topic, events = call.args + published.setdefault(topic, []).extend(events) + return published + def test_service_ticks_publish_connected_multi_event_visit(self, base_config): """Сервисные тики публикуют несколько связанных событий одного визита.""" config = replace(base_config, tick_seconds=1, max_session_events=3)