diff --git a/generator/README.md b/generator/README.md index e703f1e..f962d96 100644 --- a/generator/README.md +++ b/generator/README.md @@ -116,9 +116,11 @@ docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \ Генератор сохраняет своё состояние между перезапусками в Kafka-топик `generator_state` (compact topic). Это позволяет: -- Продолжить нумерацию тиков с места остановки +- Продолжить нумерацию тиков с места остановки (continuity) - Сохранить последовательность случайных чисел (RNG state) -- Избежать дублирования при рестарте +- Восстановить интенсивность генерации после рестарта + +**Важно:** восстанавливается continuity по номеру тика и интенсивности, но не гарантируется отсутствие дублирования событий — `event_id` и `click_id` всегда генерируются заново (`uuid4()`). ### Как работает diff --git a/generator/generator.py b/generator/generator.py index fd315df..d5db77b 100644 --- a/generator/generator.py +++ b/generator/generator.py @@ -6,12 +6,10 @@ держим целевую интенсивность events/min без крупных минутных batch. """ -import base64 import json import logging import math import os -import pickle import random import sys import time @@ -368,25 +366,33 @@ class BatchRecord: # --------------------------------------------------------------------------- # State record для сохранения состояния генератора # --------------------------------------------------------------------------- +def _nested_list_to_tuple(obj): + """Рекурсивно преобразует list в tuple (для восстановления RNG state после JSON).""" + if isinstance(obj, list): + return tuple(_nested_list_to_tuple(x) for x in obj) + return obj + + @dataclass class GeneratorState: - """Состояние генератора для восстановления после рестарта.""" + """Состояние генератора для восстановления после рестарта. + + Используем JSON-safe сериализацию: + - rng_state от random.getstate() - кортеж из простых типов (int, tuple), + безопасно сериализуется в JSON напрямую без pickle + """ tick: int - rng_state: tuple # результат random.getstate() + rng_state: tuple # результат random.getstate() - JSON-serializable 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") - + """Конвертирует в словарь для JSON-сериализации.""" return { "tick": self.tick, - "rng_state": rng_state_b64, + "rng_state": self.rng_state, # tuple из int - JSON-serializable "last_batch_id": self.last_batch_id, "last_timestamp": self.last_timestamp.isoformat(), "version": self.version, @@ -394,14 +400,10 @@ class GeneratorState: @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) - + """Создаёт состояние из словаря (JSON-only, без pickle).""" return cls( tick=data["tick"], - rng_state=rng_state, + rng_state=_nested_list_to_tuple(data["rng_state"]), # рекурсивно list->tuple last_batch_id=data["last_batch_id"], last_timestamp=datetime.fromisoformat(data["last_timestamp"]), version=data.get("version", "1.0"), @@ -445,7 +447,11 @@ class KafkaStateManager: self.producer.close() def load(self) -> GeneratorState | None: - """Загружает последнее состояние из топика.""" + """Загружает последнее состояние из топика. + + Для compact topic хранится только последнее значение для ключа, + поэтому читаем все сообщения и берём последнее с нужным ключом. + """ from kafka import KafkaConsumer logger.info(f"Loading state from topic {self.STATE_TOPIC}") @@ -479,44 +485,82 @@ class KafkaStateManager: return None +# --------------------------------------------------------------------------- +# Утилиты для работы с Kafka с retry/backoff +# --------------------------------------------------------------------------- +def _with_retry(operation, max_retries: int = 5, base_delay: float = 1.0, max_delay: float = 30.0): + """Выполняет операцию с экспоненциальным backoff и ограниченным числом попыток. + + Args: + operation: функция для выполнения + max_retries: максимальное число попыток + base_delay: начальная задержка между попытками (сек) + max_delay: максимальная задержка между попытками (сек) + + Returns: + результат операции + + Raises: + последнее исключение после исчерпания попыток + """ + import time + + last_exception = None + for attempt in range(max_retries): + try: + return operation() + except Exception as e: + last_exception = e + if attempt < max_retries - 1: + delay = min(base_delay * (2 ** attempt), max_delay) + logger.warning(f"Operation failed (attempt {attempt + 1}/{max_retries}): {e}. Retrying in {delay:.1f}s...") + time.sleep(delay) + else: + logger.error(f"Operation failed after {max_retries} attempts: {e}") + raise last_exception + + # --------------------------------------------------------------------------- # Создание топиков для истории и состояния # --------------------------------------------------------------------------- def ensure_topics(bootstrap_servers: str) -> None: - """Создаёт необходимые топики если они не существуют.""" + """Создаёт необходимые топики если они не существуют (с retry на подключение).""" from kafka import KafkaAdminClient from kafka.admin import NewTopic from kafka.errors import TopicAlreadyExistsError - admin_client = KafkaAdminClient(bootstrap_servers=bootstrap_servers) - try: - # Топик для истории батчей (обычный, с retention) - history_topic = NewTopic( - name=KafkaBatchHistory.HISTORY_TOPIC, - num_partitions=1, - replication_factor=1, - ) + def _create_topics(): + admin_client = KafkaAdminClient(bootstrap_servers=bootstrap_servers) + try: + # Топик для истории батчей (обычный, с retention) + history_topic = NewTopic( + name=KafkaBatchHistory.HISTORY_TOPIC, + num_partitions=1, + replication_factor=1, + ) - # Топик для состояния (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", - }, - ) + # Топик для состояния (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() + 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() + + _with_retry(_create_topics, max_retries=5, base_delay=1.0) # --------------------------------------------------------------------------- diff --git a/generator/tests/test_state.py b/generator/tests/test_state.py index 072c424..1272154 100644 --- a/generator/tests/test_state.py +++ b/generator/tests/test_state.py @@ -1,8 +1,7 @@ """ Тесты сохранения и восстановления состояния генератора. """ -import base64 -import pickle +import json from datetime import datetime, timezone from unittest.mock import MagicMock, patch @@ -48,7 +47,7 @@ class TestGeneratorState: assert state.version == "1.0" def test_to_dict_serialization(self): - """Сериализация в словарь.""" + """Сериализация в словарь (JSON-safe, без pickle).""" now = datetime.now(timezone.utc) rng_state = (3, (1, 2, 3), None) @@ -66,13 +65,14 @@ class TestGeneratorState: assert data["last_timestamp"] == now.isoformat() assert data["version"] == "1.0" - # Проверяем что rng_state сериализован через pickle+base64 + # Проверяем что rng_state сериализован как tuple (JSON-safe, без pickle) 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 + assert data["rng_state"] == rng_state + # Проверяем что можно сериализовать в JSON и восстановить + json_str = json.dumps(data) + restored_data = json.loads(json_str) + restored_state = GeneratorState.from_dict(restored_data) + assert restored_state.rng_state == rng_state def test_from_dict_deserialization(self): """Десериализация из словаря.""" @@ -289,3 +289,93 @@ class TestKafkaStateManager: manager.close() mock_producer.close.assert_called_once() + + +class TestJsonSafeState: + """Тесты JSON-safe сериализации state (без pickle).""" + + def test_json_roundtrip_with_nested_tuples(self): + """Проверка что nested tuple корректно восстанавливается после JSON.""" + from generator import _nested_list_to_tuple + + # Симулируем что получаем после json.loads() - все tuple становятся list + json_loaded = [3, [1, 2, 3], None] + + result = _nested_list_to_tuple(json_loaded) + + assert result == (3, (1, 2, 3), None) + assert isinstance(result, tuple) + assert isinstance(result[1], tuple) + + def test_json_roundtrip_rng_state(self): + """Полный цикл: rng.getstate() -> JSON -> from_dict -> setstate.""" + import random + + rng = random.Random(42) + # Делаем несколько вызовов + values_before = [rng.random() for _ in range(10)] + + # Сохраняем state + state = GeneratorState( + tick=100, + rng_state=rng.getstate(), + last_batch_id="test123", + last_timestamp=datetime.now(timezone.utc), + ) + + # Сериализуем через JSON (как в Kafka) + data = state.to_dict() + json_str = json.dumps(data) + restored_data = json.loads(json_str) + + # Восстанавливаем + restored_state = GeneratorState.from_dict(restored_data) + + # Проверяем что RNG state восстановлен корректно + rng2 = random.Random() + rng2.setstate(restored_state.rng_state) + + # Проверяем что следующие значения совпадают + values_after = [rng2.random() for _ in range(5)] + + # Оригинальный RNG должен дать те же значения + values_expected = [rng.random() for _ in range(5)] + + assert values_after == values_expected + + +class TestWithRetry: + """Тесты функции _with_retry.""" + + def test_success_on_first_attempt(self): + """Успех с первой попытки.""" + from generator import _with_retry + + operation = MagicMock(return_value="success") + + result = _with_retry(operation, max_retries=3, base_delay=0.01) + + assert result == "success" + assert operation.call_count == 1 + + def test_success_after_retries(self): + """Успех после нескольких попыток.""" + from generator import _with_retry + + operation = MagicMock(side_effect=[Exception("fail1"), Exception("fail2"), "success"]) + + result = _with_retry(operation, max_retries=3, base_delay=0.01) + + assert result == "success" + assert operation.call_count == 3 + + def test_failure_after_all_retries(self): + """Исчерпание всех попыток.""" + from generator import _with_retry + + operation = MagicMock(side_effect=Exception("always fails")) + + with pytest.raises(Exception, match="always fails"): + _with_retry(operation, max_retries=3, base_delay=0.01) + + assert operation.call_count == 3