diff --git a/docker-compose.yml b/docker-compose.yml index 714d429..9943fe9 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -300,6 +300,9 @@ services: GEN_DATA_DIR: /data GEN_ENABLED: "true" GEN_METRICS_PORT: "9109" + # State management (сохранение состояния между рестартами) + GEN_STATE_ENABLED: "true" + GEN_STATE_RESET: "false" # PYTHONUNBUFFERED для сразу видеть логи PYTHONUNBUFFERED: "1" ports: diff --git a/generator/README.md b/generator/README.md index a10b0fb..e703f1e 100644 --- a/generator/README.md +++ b/generator/README.md @@ -36,6 +36,8 @@ generator-service -> Kafka topics -> (потребители отдельно) | `GEN_SEED` | Сид для воспроизводимости | — | | `GEN_ENABLED` | Включить генерацию | `true` | | `GEN_METRICS_PORT` | Порт для Prometheus | `9109` | +| `GEN_STATE_ENABLED` | Сохранять состояние между рестартами | `true` | +| `GEN_STATE_RESET` | Сбросить состояние при старте | `false` | ### Режим "раз в минуту" (для демо) @@ -110,6 +112,59 @@ docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \ --from-beginning ``` +## State Recovery (восстановление состояния) + +Генератор сохраняет своё состояние между перезапусками в Kafka-топик `generator_state` (compact topic). Это позволяет: + +- Продолжить нумерацию тиков с места остановки +- Сохранить последовательность случайных чисел (RNG state) +- Избежать дублирования при рестарте + +### Как работает + +1. После каждого успешного тика состояние сохраняется в `generator_state` +2. При старте генератор читает последнее состояние из топика +3. Если состояние найдено - продолжает с сохранённого tick +4. Если нет - начинает с tick=1 + +### Топик `generator_state` + +- **Название**: фиксировано `generator_state` +- **Тип**: compact topic (хранится только последнее значение для каждого ключа) +- **Ключ**: `default` (для возможности нескольких генераторов в будущем) +- **Конфигурация**: `cleanup.policy=compact`, минимальный retention + +### Просмотр текущего состояния + +```bash +docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \ + --bootstrap-server kafka:29092 \ + --topic generator_state \ + --from-beginning \ + --property print.key=true +``` + +### Сброс состояния (начать сначала) + +```bash +# Вариант 1: через env (рекомендуется) +GEN_STATE_RESET=true docker compose up -d generator + +# Вариант 2: удалить топик полностью +docker compose exec kafka /opt/kafka/bin/kafka-topics.sh \ + --bootstrap-server kafka:29092 \ + --delete \ + --topic generator_state +``` + +### Отключение сохранения состояния + +```bash +GEN_STATE_ENABLED=false docker compose up -d generator +``` + +При отключенном state management генератор всегда начинает с tick=1, RNG инициализируется с GEN_SEED (или случайно). + ## Тестирование Тесты написаны на **pytest**. @@ -137,7 +192,8 @@ generator/tests/ ├── test_generation.py # Тесты генерации событий ├── test_history.py # Тесты структуры BatchRecord ├── test_kafka_history.py # Тесты KafkaBatchHistory -└── test_service.py # Тесты GeneratorService +├── test_service.py # Тесты GeneratorService +└── test_state.py # Тесты GeneratorState и KafkaStateManager ``` ### Интеграционный тест diff --git a/generator/generator.py b/generator/generator.py index 0c630f6..fd315df 100644 --- a/generator/generator.py +++ b/generator/generator.py @@ -6,10 +6,12 @@ держим целевую интенсивность events/min без крупных минутных batch. """ +import base64 import json import logging import math import os +import pickle import random import sys import time @@ -122,6 +124,14 @@ class Config: default_factory=lambda: int(os.getenv("GEN_METRICS_PORT", "9109")) ) + # Управление сохранением состояния + state_enabled: bool = field( + default_factory=lambda: os.getenv("GEN_STATE_ENABLED", "true").lower() == "true" + ) + state_reset: bool = field( + default_factory=lambda: os.getenv("GEN_STATE_RESET", "false").lower() == "true" + ) + def __post_init__(self): # Валидация параметров if self.tick_seconds < 1: @@ -356,25 +366,155 @@ class BatchRecord: # --------------------------------------------------------------------------- -# Создание топика для истории +# State record для сохранения состояния генератора # --------------------------------------------------------------------------- -def ensure_history_topic(bootstrap_servers: str) -> None: - """Создаёт топик для истории батчей если он не существует.""" +@dataclass +class GeneratorState: + """Состояние генератора для восстановления после рестарта.""" + + tick: int + rng_state: tuple # результат random.getstate() + last_batch_id: str + last_timestamp: datetime + version: str = "1.0" + + def to_dict(self) -> dict: + """Конвертирует в словарь для сериализации.""" + # Сериализуем rng_state через pickle + base64 + rng_state_bytes = pickle.dumps(self.rng_state) + rng_state_b64 = base64.b64encode(rng_state_bytes).decode("utf-8") + + return { + "tick": self.tick, + "rng_state": rng_state_b64, + "last_batch_id": self.last_batch_id, + "last_timestamp": self.last_timestamp.isoformat(), + "version": self.version, + } + + @classmethod + def from_dict(cls, data: dict) -> "GeneratorState": + """Создаёт состояние из словаря.""" + # Десериализуем rng_state + rng_state_bytes = base64.b64decode(data["rng_state"]) + rng_state = pickle.loads(rng_state_bytes) + + return cls( + tick=data["tick"], + rng_state=rng_state, + last_batch_id=data["last_batch_id"], + last_timestamp=datetime.fromisoformat(data["last_timestamp"]), + version=data.get("version", "1.0"), + ) + + +# --------------------------------------------------------------------------- +# Kafka state manager - сохраняет/восстанавливает состояние генератора +# --------------------------------------------------------------------------- +class KafkaStateManager: + """Управление состоянием генератора в Kafka (compact topic).""" + + STATE_TOPIC = "generator_state" + STATE_KEY = "default" # Для возможности нескольких генераторов в будущем + + def __init__(self, bootstrap_servers: str): + self.bootstrap_servers = bootstrap_servers + KafkaProducerCls, _ = _import_kafka() + + logger.info(f"Connecting to Kafka for state management at {self.bootstrap_servers}") + self.producer = KafkaProducerCls( + bootstrap_servers=self.bootstrap_servers, + value_serializer=lambda v: json.dumps(v).encode("utf-8"), + key_serializer=lambda k: k.encode("utf-8") if k else None, + retries=3, + retry_backoff_ms=1000, + ) + logger.info("Connected to Kafka for state management successfully") + + def save(self, state: GeneratorState) -> None: + """Сохраняет состояние в топик (compact topic - только последнее значение).""" + value = state.to_dict() + self.producer.send(self.STATE_TOPIC, key=self.STATE_KEY, value=value) + + def flush(self) -> None: + """Сбрасывает буфер.""" + self.producer.flush() + + def close(self) -> None: + """Закрывает соединение.""" + self.producer.close() + + def load(self) -> GeneratorState | None: + """Загружает последнее состояние из топика.""" + from kafka import KafkaConsumer + + logger.info(f"Loading state from topic {self.STATE_TOPIC}") + try: + consumer = KafkaConsumer( + self.STATE_TOPIC, + bootstrap_servers=self.bootstrap_servers, + auto_offset_reset="earliest", + enable_auto_commit=False, + consumer_timeout_ms=5000, + value_deserializer=lambda v: json.loads(v.decode("utf-8")), + ) + + last_state = None + for message in consumer: + if message.key and message.key.decode("utf-8") == self.STATE_KEY: + last_state = message.value + + consumer.close() + + if last_state: + logger.info(f"Restored state: tick={last_state.get('tick')}, " + f"last_batch_id={last_state.get('last_batch_id')}") + return GeneratorState.from_dict(last_state) + else: + logger.info("No previous state found, starting fresh") + return None + + except Exception as e: + logger.warning(f"Failed to load state: {e}, starting fresh") + return None + + +# --------------------------------------------------------------------------- +# Создание топиков для истории и состояния +# --------------------------------------------------------------------------- +def ensure_topics(bootstrap_servers: str) -> None: + """Создаёт необходимые топики если они не существуют.""" from kafka import KafkaAdminClient from kafka.admin import NewTopic from kafka.errors import TopicAlreadyExistsError admin_client = KafkaAdminClient(bootstrap_servers=bootstrap_servers) try: - new_topic = NewTopic( + # Топик для истории батчей (обычный, с retention) + history_topic = NewTopic( name=KafkaBatchHistory.HISTORY_TOPIC, num_partitions=1, replication_factor=1, ) - admin_client.create_topics([new_topic]) - logger.info(f"Created topic: {KafkaBatchHistory.HISTORY_TOPIC}") - except TopicAlreadyExistsError: - logger.debug(f"Topic already exists: {KafkaBatchHistory.HISTORY_TOPIC}") + + # Топик для состояния (compact - храним только последнее значение) + state_topic = NewTopic( + name=KafkaStateManager.STATE_TOPIC, + num_partitions=1, + replication_factor=1, + topic_configs={ + "cleanup.policy": "compact", + "min.cleanable.dirty.ratio": "0.1", + "delete.retention.ms": "100", + }, + ) + + for topic in [history_topic, state_topic]: + try: + admin_client.create_topics([topic]) + logger.info(f"Created topic: {topic.name}") + except TopicAlreadyExistsError: + logger.debug(f"Topic already exists: {topic.name}") finally: admin_client.close() @@ -509,7 +649,9 @@ class GeneratorService: self.generator = EventGenerator(self.dictionary, config) self.publisher: KafkaPublisher | None = None self.history: KafkaBatchHistory | None = None + self.state_manager: KafkaStateManager | None = None self._running = False + self._tick = 0 # Текущий номер тика (восстанавливается из стейта) def start(self): """Запускает основной цикл.""" @@ -524,14 +666,33 @@ class GeneratorService: logger.info("Starting generator service...") logger.info(f"Configuration: tick={self.config.tick_seconds}s, " f"lambda_base={self.config.lambda_base_per_min}/min, " - f"jitter={self.config.jitter_pct}%") + f"jitter={self.config.jitter_pct}%, " + f"state_enabled={self.config.state_enabled}, " + f"state_reset={self.config.state_reset}") # Подключаемся к Kafka для публикации событий и истории - # Сначала создаём топик для истории если нужно - ensure_history_topic(self.config.kafka_bootstrap_servers) + # Сначала создаём топики если нужно + ensure_topics(self.config.kafka_bootstrap_servers) self.publisher = KafkaPublisher(self.config.kafka_bootstrap_servers) self.history = KafkaBatchHistory(self.config.kafka_bootstrap_servers) + # Инициализируем state manager если включено + if self.config.state_enabled: + self.state_manager = KafkaStateManager(self.config.kafka_bootstrap_servers) + + # Восстанавливаем стейт если не требуется сброс + if not self.config.state_reset: + restored_state = self.state_manager.load() + if restored_state: + self._tick = restored_state.tick + self.generator.rng.setstate(restored_state.rng_state) + logger.info(f"Restored state: continuing from tick {self._tick}, " + f"last_batch_id={restored_state.last_batch_id}") + else: + logger.info("State reset requested, starting fresh") + else: + logger.info("State management disabled") + self._running = True try: @@ -549,18 +710,37 @@ class GeneratorService: self.publisher.close() if self.history: self.history.close() + if self.state_manager: + self.state_manager.close() + + def _save_state(self, batch_id: str) -> None: + """Сохраняет текущее состояние генератора.""" + if not self.state_manager or not self.config.state_enabled: + return + + try: + state = GeneratorState( + tick=self._tick, + rng_state=self.generator.rng.getstate(), + last_batch_id=batch_id, + last_timestamp=datetime.now(timezone.utc), + ) + self.state_manager.save(state) + self.state_manager.flush() + logger.debug(f"Saved state: tick={self._tick}, batch_id={batch_id}") + except Exception as e: + logger.warning(f"Failed to save state: {e}") + METRICS_ERRORS_TOTAL.labels(topic="state").inc() def _main_loop(self): """Основной цикл тиков.""" - tick = 0 - while self._running: - tick += 1 + self._tick += 1 tick_start = time.time() batch_id = str(uuid.uuid4())[:8] with METRICS_TICK_DURATION.time(): - logger.info(f"=== Tick {tick} (batch_id={batch_id}) ===") + logger.info(f"=== Tick {self._tick} (batch_id={batch_id}) ===") try: # Вычисляем количество событий @@ -593,9 +773,10 @@ class GeneratorService: else: status = "error" - # Обновляем метрику последнего успешного тика + # Обновляем метрику последнего успешного тика и сохраняем стейт if status in ("success", "partial"): METRICS_LAST_SUCCESS.set_to_current_time() + self._save_state(batch_id) # Флашим публикацию событий self.publisher.flush() @@ -635,7 +816,7 @@ class GeneratorService: logger.info(f" {topic}: {counts['sent']} sent") except Exception as e: - logger.exception(f"Error in tick {tick}: {e}") + logger.exception(f"Error in tick {self._tick}: {e}") # Пытаемся записать ошибку в историю (best-effort) try: self.history.add( diff --git a/generator/tests/conftest.py b/generator/tests/conftest.py index ebea633..8a75e0a 100644 --- a/generator/tests/conftest.py +++ b/generator/tests/conftest.py @@ -38,6 +38,8 @@ def base_config(data_dir): seed=42, enabled=True, metrics_port=9109, + state_enabled=True, + state_reset=False, ) diff --git a/generator/tests/test_state.py b/generator/tests/test_state.py new file mode 100644 index 0000000..072c424 --- /dev/null +++ b/generator/tests/test_state.py @@ -0,0 +1,291 @@ +""" +Тесты сохранения и восстановления состояния генератора. +""" +import base64 +import pickle +from datetime import datetime, timezone +from unittest.mock import MagicMock, patch + +import pytest + +from generator import GeneratorState, KafkaStateManager + + +class TestGeneratorState: + """Тесты структуры состояния генератора.""" + + def test_state_creation(self): + """Создание состояния с всеми полями.""" + now = datetime.now(timezone.utc) + rng_state = (3, (1, 2, 3), None) # Минимальный валидный state для random + + state = GeneratorState( + tick=42, + rng_state=rng_state, + last_batch_id="abc123", + last_timestamp=now, + version="1.0", + ) + + assert state.tick == 42 + assert state.rng_state == rng_state + assert state.last_batch_id == "abc123" + assert state.last_timestamp == now + assert state.version == "1.0" + + def test_default_version(self): + """Версия по умолчанию.""" + now = datetime.now(timezone.utc) + rng_state = (3, (1, 2, 3), None) + + state = GeneratorState( + tick=1, + rng_state=rng_state, + last_batch_id="test", + last_timestamp=now, + ) + + assert state.version == "1.0" + + def test_to_dict_serialization(self): + """Сериализация в словарь.""" + now = datetime.now(timezone.utc) + rng_state = (3, (1, 2, 3), None) + + state = GeneratorState( + tick=42, + rng_state=rng_state, + last_batch_id="abc123", + last_timestamp=now, + ) + + data = state.to_dict() + + assert data["tick"] == 42 + assert data["last_batch_id"] == "abc123" + assert data["last_timestamp"] == now.isoformat() + assert data["version"] == "1.0" + + # Проверяем что rng_state сериализован через pickle+base64 + assert "rng_state" in data + assert isinstance(data["rng_state"], str) + # Проверяем что можно десериализовать + decoded = base64.b64decode(data["rng_state"]) + restored_rng = pickle.loads(decoded) + assert restored_rng == rng_state + + def test_from_dict_deserialization(self): + """Десериализация из словаря.""" + now = datetime.now(timezone.utc) + rng_state = (3, (1, 2, 3), None) + + # Создаём исходное состояние + original = GeneratorState( + tick=42, + rng_state=rng_state, + last_batch_id="abc123", + last_timestamp=now, + ) + + # Сериализуем и десериализуем + data = original.to_dict() + restored = GeneratorState.from_dict(data) + + assert restored.tick == original.tick + assert restored.rng_state == original.rng_state + assert restored.last_batch_id == original.last_batch_id + assert restored.last_timestamp == original.last_timestamp + assert restored.version == original.version + + def test_roundtrip_with_real_random(self): + """Проверка что RNG state действительно восстанавливает последовательность.""" + import random + + # Создаём генератор и делаем несколько вызовов + rng = random.Random(12345) + values_before = [rng.random() for _ in range(5)] + + # Сохраняем состояние + state = GeneratorState( + tick=10, + rng_state=rng.getstate(), + last_batch_id="test", + last_timestamp=datetime.now(timezone.utc), + ) + + # Десериализуем + data = state.to_dict() + restored_state = GeneratorState.from_dict(data) + + # Создаём новый генератор с восстановленным состоянием + new_rng = random.Random() + new_rng.setstate(restored_state.rng_state) + + # Проверяем что следующие значения совпадают + values_after = [new_rng.random() for _ in range(5)] + + # Если state восстановлен корректно, values должны совпадать + # Но т.к. rng уже "прокручен" на 5 значений, берём следующие + rng.setstate(state.rng_state) + next_values = [rng.random() for _ in range(5)] + + assert next_values == values_after + + +class TestKafkaStateManager: + """Тесты менеджера состояния.""" + + def test_init(self): + """Инициализация менеджера.""" + with patch("generator._import_kafka") as mock_import: + mock_producer_class = MagicMock() + mock_import.return_value = (mock_producer_class, None) + + manager = KafkaStateManager("kafka:29092") + + assert manager.bootstrap_servers == "kafka:29092" + assert manager.STATE_TOPIC == "generator_state" + assert manager.STATE_KEY == "default" + mock_producer_class.assert_called_once() + + def test_save(self): + """Сохранение состояния.""" + with patch("generator._import_kafka") as mock_import: + mock_producer = MagicMock() + mock_producer_class = MagicMock(return_value=mock_producer) + mock_import.return_value = (mock_producer_class, None) + + manager = KafkaStateManager("kafka:29092") + + now = datetime.now(timezone.utc) + rng_state = (3, (1, 2, 3), None) + state = GeneratorState( + tick=42, + rng_state=rng_state, + last_batch_id="abc123", + last_timestamp=now, + ) + + manager.save(state) + + # Проверяем что producer.send был вызван + mock_producer.send.assert_called_once() + call_args = mock_producer.send.call_args + assert call_args[0][0] == "generator_state" # topic + assert call_args[1]["key"] == "default" # key + assert call_args[1]["value"]["tick"] == 42 + assert call_args[1]["value"]["last_batch_id"] == "abc123" + + def test_load_no_messages(self): + """Загрузка при отсутствии состояния.""" + with patch("generator._import_kafka") as mock_import, \ + patch("kafka.KafkaConsumer") as mock_consumer_class: + + mock_producer_class = MagicMock() + mock_import.return_value = (mock_producer_class, None) + + # Мокаем пустой consumer (нет сообщений) + mock_consumer = MagicMock() + mock_consumer.__iter__ = MagicMock(return_value=iter([])) + mock_consumer_class.return_value = mock_consumer + + manager = KafkaStateManager("kafka:29092") + result = manager.load() + + assert result is None + + def test_load_with_messages(self): + """Загрузка существующего состояния.""" + with patch("generator._import_kafka") as mock_import, \ + patch("kafka.KafkaConsumer") as mock_consumer_class: + + mock_producer_class = MagicMock() + mock_import.return_value = (mock_producer_class, None) + + now = datetime.now(timezone.utc) + rng_state = (3, (1, 2, 3), None) + state = GeneratorState( + tick=100, + rng_state=rng_state, + last_batch_id="xyz789", + last_timestamp=now, + ) + + # Мокаем consumer с сообщением + mock_message = MagicMock() + mock_message.key = b"default" + mock_message.value = state.to_dict() + + mock_consumer = MagicMock() + mock_consumer.__iter__ = MagicMock(return_value=iter([mock_message])) + mock_consumer_class.return_value = mock_consumer + + manager = KafkaStateManager("kafka:29092") + result = manager.load() + + assert result is not None + assert result.tick == 100 + assert result.last_batch_id == "xyz789" + + def test_load_ignores_wrong_key(self): + """Загрузка игнорирует сообщения с другим ключом.""" + with patch("generator._import_kafka") as mock_import, \ + patch("kafka.KafkaConsumer") as mock_consumer_class: + + mock_producer_class = MagicMock() + mock_import.return_value = (mock_producer_class, None) + + # Мокаем consumer с сообщением с неправильным ключом + mock_message = MagicMock() + mock_message.key = b"other_generator" # Другой ключ + mock_message.value = {"tick": 999} + + mock_consumer = MagicMock() + mock_consumer.__iter__ = MagicMock(return_value=iter([mock_message])) + mock_consumer_class.return_value = mock_consumer + + manager = KafkaStateManager("kafka:29092") + result = manager.load() + + # Не должно быть загружено, т.к. ключ не совпадает + assert result is None + + def test_load_exception_returns_none(self): + """При ошибке загрузки возвращается None.""" + with patch("generator._import_kafka") as mock_import, \ + patch("kafka.KafkaConsumer") as mock_consumer_class: + + mock_producer_class = MagicMock() + mock_import.return_value = (mock_producer_class, None) + + # Симулируем ошибку при создании consumer + mock_consumer_class.side_effect = Exception("Kafka unavailable") + + manager = KafkaStateManager("kafka:29092") + result = manager.load() + + assert result is None + + def test_flush(self): + """Сброс буфера.""" + with patch("generator._import_kafka") as mock_import: + mock_producer = MagicMock() + mock_producer_class = MagicMock(return_value=mock_producer) + mock_import.return_value = (mock_producer_class, None) + + manager = KafkaStateManager("kafka:29092") + manager.flush() + + mock_producer.flush.assert_called_once() + + def test_close(self): + """Закрытие соединения.""" + with patch("generator._import_kafka") as mock_import: + mock_producer = MagicMock() + mock_producer_class = MagicMock(return_value=mock_producer) + mock_import.return_value = (mock_producer_class, None) + + manager = KafkaStateManager("kafka:29092") + manager.close() + + mock_producer.close.assert_called_once()