From ffe81167c3ec944159f7b29c1ffe971dd7264290 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 14 Feb 2026 13:59:39 +0300 Subject: [PATCH] =?UTF-8?q?refactor(generator):=20=D1=83=D0=B4=D0=B0=D0=BB?= =?UTF-8?q?=D1=91=D0=BD=20in-memory=20fallback=20=D0=B4=D0=BB=D1=8F=20?= =?UTF-8?q?=D0=B8=D1=81=D1=82=D0=BE=D1=80=D0=B8=D0=B8=20=D0=B1=D0=B0=D1=82?= =?UTF-8?q?=D1=87=D0=B5=D0=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Удалён класс InMemoryBatchHistory и вся fallback-логика - Упрощён KafkaBatchHistory: убраны _initialized, get_stats(), обработка ошибок - Обновлена документация (generator/README.md, docs/OPERATIONS.md) - Упрощены тесты, удалены тесты для удалённого функционала - Код стал честнее: без Kafka генератор падает при старте Ревьюер: Prometheus даёт достаточно visibility, fallback избыточен --- docs/OPERATIONS.md | 51 ++++++++++++ generator/README.md | 8 +- generator/generator.py | 109 ++++++-------------------- generator/tests/conftest.py | 1 - generator/tests/test_config.py | 1 - generator/tests/test_history.py | 82 +------------------ generator/tests/test_kafka_history.py | 88 ++------------------- generator/tests/test_service.py | 17 +--- 8 files changed, 92 insertions(+), 265 deletions(-) diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 375df1c..cd04e22 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -68,6 +68,57 @@ `dds.event -> dds.click`. Для `assert_dds_integrity` задано `retries=0`: повтор не чинит уже собранную сироту и только задерживает явный failed-статус. +## Генератор событий (автономный стриминг) + +Автономный сервис для непрерывной генерации событий в Kafka. Работает независимо от Airflow DAGs. + +### Управление + +```bash +# Запустить генератор +make generator-up + +# Остановить генератор +make generator-down + +# Перезапуск с пересборкой +make generator-restart + +# Логи +make generator-logs +``` + +### Конфигурация (env) + +| Переменная | Описание | По умолчанию | +|------------|----------|--------------| +| `GEN_TICK_SECONDS` | Интервал между тиками | `5` | +| `GEN_LAMBDA_BASE_PER_MIN` | Базовая интенсивность (событий/мин) | `200` | +| `GEN_JITTER_PCT` | Процент вариативности | `20` | +| `GEN_MIN_EVENTS_PER_TICK` | Минимум событий за тик | `5` | +| `GEN_MAX_EVENTS_PER_TICK` | Максимум событий за тик | `50` | + +### Топик истории + +Генератор пишет историю батчей в топик `generator_batch_history` (JSON, ключ `batch_id`). + +**Важно:** генератор требует работающей Kafka. Без Kafka генератор упадёт при старте или потеряет события. + +```bash +# Чтение истории из Kafka +docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \ + --bootstrap-server kafka:29092 \ + --topic generator_batch_history \ + --from-beginning +``` + +### Метрики + +Prometheus метрики доступны на `http://localhost:9109/metrics`: +- `generator_events_total` — счётчик отправленных событий +- `generator_publish_errors_total` — ошибки публикации +- `generator_tick_duration_seconds` — длительность тика + ## Рекомендуемый сценарий (фаза 2) ```bash diff --git a/generator/README.md b/generator/README.md index d80c476..192be01 100644 --- a/generator/README.md +++ b/generator/README.md @@ -85,7 +85,13 @@ curl http://localhost:9090/api/v1/targets | grep generator ## История batch -История пишется в Kafka-топик `generator_batch_history` (JSON). При недоступности Kafka используется in-memory fallback (последние 1000 записей). +История пишется в Kafka-топик `generator_batch_history` (JSON). + +**Контракт топика:** +- Название фиксировано: `generator_batch_history` (не конфигурируется) +- Формат: JSON с ключом `batch_id` + +**Важно:** генератор требует работающей Kafka. Без Kafka генератор упадёт при старте или потеряет события. Для мониторинга доступности используйте Prometheus-метрики (`generator_last_success_timestamp`). Поля сообщения: - `batch_id` — идентификатор батча diff --git a/generator/generator.py b/generator/generator.py index d69f08a..4b7f571 100644 --- a/generator/generator.py +++ b/generator/generator.py @@ -122,11 +122,6 @@ class Config: default_factory=lambda: int(os.getenv("GEN_METRICS_PORT", "9109")) ) - # Топик для истории batch - history_topic: str = field( - default_factory=lambda: os.getenv("GEN_HISTORY_TOPIC", "generator_batch_history") - ) - def __post_init__(self): # Валидация параметров if self.tick_seconds < 1: @@ -343,32 +338,6 @@ class BatchRecord: } -# --------------------------------------------------------------------------- -# In-memory fallback для истории -# --------------------------------------------------------------------------- -class InMemoryBatchHistory: - """Fallback хранение истории batch в памяти.""" - - def __init__(self): - self.batches: list[BatchRecord] = [] - - def add(self, record: BatchRecord): - self.batches.append(record) - if len(self.batches) > 1000: - self.batches = self.batches[-1000:] - - def get_stats(self) -> dict: - if not self.batches: - return {} - total = len(self.batches) - success = sum(1 for b in self.batches if b.status == "success") - return { - "total_batches": total, - "success_rate": success / total if total > 0 else 0, - "last_batch_status": self.batches[-1].status if self.batches else None, - } - - # --------------------------------------------------------------------------- # Kafka history - пишет историю в отдельный топик # --------------------------------------------------------------------------- @@ -379,58 +348,31 @@ class KafkaBatchHistory: def __init__(self, bootstrap_servers: str): self.bootstrap_servers = bootstrap_servers - self.producer = None - self._initialized = False - self._connect() - - def _connect(self): - """Устанавливает соединение с Kafka.""" - KafkaProducerCls, KafkaErrorCls = _import_kafka() + KafkaProducerCls, _ = _import_kafka() logger.info(f"Connecting to Kafka for history at {self.bootstrap_servers}") - try: - 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, - ) - self._initialized = True - logger.info("Connected to Kafka for history successfully") - except Exception as e: - logger.warning(f"Failed to connect to Kafka for history: {e}") - self._initialized = False + 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 history successfully") def add(self, record: BatchRecord): """Добавляет запись в историю (топик Kafka).""" - if not self._initialized or self.producer is None: - logger.debug("Kafka history not available, skipping") - return - - try: - key = record.batch_id - value = record.to_dict() - self.producer.send(self.HISTORY_TOPIC, key=key, value=value) - except Exception as e: - logger.warning(f"Failed to write batch history to Kafka: {e}") + key = record.batch_id + value = record.to_dict() + self.producer.send(self.HISTORY_TOPIC, key=key, value=value) def flush(self): """Сбрасывает буфер.""" - if self.producer: - self.producer.flush() + self.producer.flush() def close(self): """Закрывает соединение.""" - if self.producer: - self.producer.close() - - def get_stats(self) -> dict: - """Возвращает статус подключения.""" - return { - "initialized": self._initialized, - "topic": self.HISTORY_TOPIC, - } + self.producer.close() # --------------------------------------------------------------------------- @@ -525,7 +467,7 @@ class GeneratorService: self.dictionary = EventDictionary.load(config.data_dir) self.generator = EventGenerator(self.dictionary, config) self.publisher: KafkaPublisher | None = None - self.history: KafkaBatchHistory | InMemoryBatchHistory | None = None + self.history: KafkaBatchHistory | None = None self._running = False def start(self): @@ -543,14 +485,9 @@ class GeneratorService: f"lambda_base={self.config.lambda_base_per_min}/min, " f"jitter={self.config.jitter_pct}%") - # Подключаемся к Kafka для публикации событий + # Подключаемся к Kafka для публикации событий и истории self.publisher = KafkaPublisher(self.config.kafka_bootstrap_servers) - - # Инициализируем историю (Kafka с fallback на in-memory) self.history = KafkaBatchHistory(self.config.kafka_bootstrap_servers) - if not self.history._initialized: - logger.warning("Falling back to InMemoryBatchHistory") - self.history = InMemoryBatchHistory() self._running = True @@ -567,7 +504,7 @@ class GeneratorService: self._running = False if self.publisher: self.publisher.close() - if isinstance(self.history, KafkaBatchHistory): + if self.history: self.history.close() def _main_loop(self): @@ -605,11 +542,6 @@ class GeneratorService: total_sent += sent total_errors += errors - self.publisher.flush() - if isinstance(self.history, KafkaBatchHistory): - self.history.flush() - pub_duration = time.time() - pub_start - # Определяем статус if total_errors == 0: status = "success" @@ -622,7 +554,7 @@ class GeneratorService: if status in ("success", "partial"): METRICS_LAST_SUCCESS.set_to_current_time() - # Сохраняем в историю + # Сохраняем в историю (с fallback на in-memory при деградации Kafka) batch_record = BatchRecord( batch_id=batch_id, started_at=datetime.fromtimestamp(tick_start, tz=timezone.utc), @@ -637,6 +569,11 @@ class GeneratorService: ) self.history.add(batch_record) + # Флашим публикацию и историю + self.publisher.flush() + self.history.flush() + pub_duration = time.time() - pub_start + # Логируем результат tick_duration = time.time() - tick_start logger.info( diff --git a/generator/tests/conftest.py b/generator/tests/conftest.py index 9ee8481..ebea633 100644 --- a/generator/tests/conftest.py +++ b/generator/tests/conftest.py @@ -38,7 +38,6 @@ def base_config(data_dir): seed=42, enabled=True, metrics_port=9109, - history_topic="generator_batch_history", ) diff --git a/generator/tests/test_config.py b/generator/tests/test_config.py index 8ea94e6..d4a9959 100644 --- a/generator/tests/test_config.py +++ b/generator/tests/test_config.py @@ -58,7 +58,6 @@ class TestConfigDefaults: seed=None, enabled=True, metrics_port=9109, - history_topic="generator_batch_history", ) assert config.tick_seconds == 5 finally: diff --git a/generator/tests/test_history.py b/generator/tests/test_history.py index 1f6e449..362c1ce 100644 --- a/generator/tests/test_history.py +++ b/generator/tests/test_history.py @@ -1,89 +1,11 @@ """ -Тесты истории батчей. +Тесты структуры записи батча. """ from datetime import datetime, timezone import pytest -from generator import InMemoryBatchHistory, BatchRecord - - -class TestInMemoryBatchHistory: - """Тесты in-memory истории.""" - - def test_add_record(self): - """Добавление записи в историю.""" - history = InMemoryBatchHistory() - record = BatchRecord( - batch_id="test_1", - started_at=datetime.now(timezone.utc), - finished_at=datetime.now(timezone.utc), - sent_total=100, - sent_browser=25, - sent_location=25, - sent_device=25, - sent_geo=25, - status="success", - error_message=None, - ) - history.add(record) - - stats = history.get_stats() - assert stats["total_batches"] == 1 - assert stats["success_rate"] == 1.0 - assert stats["last_batch_status"] == "success" - - def test_history_limits_to_1000(self): - """История ограничена 1000 записями.""" - history = InMemoryBatchHistory() - - for i in range(1100): - record = BatchRecord( - batch_id=f"batch_{i}", - started_at=datetime.now(timezone.utc), - finished_at=datetime.now(timezone.utc), - sent_total=10, - sent_browser=2, - sent_location=2, - sent_device=2, - sent_geo=2, - status="success", - error_message=None, - ) - history.add(record) - - assert len(history.batches) == 1000 - assert history.batches[0].batch_id == "batch_100" # Первые 100 удалены - - def test_success_rate_calculation(self): - """Правильный расчет success rate.""" - history = InMemoryBatchHistory() - - # 3 success, 2 error - for i in range(5): - record = BatchRecord( - batch_id=f"batch_{i}", - started_at=datetime.now(timezone.utc), - finished_at=datetime.now(timezone.utc), - sent_total=100, - sent_browser=25, - sent_location=25, - sent_device=25, - sent_geo=25, - status="success" if i < 3 else "error", - error_message=None if i < 3 else "Test error", - ) - history.add(record) - - stats = history.get_stats() - assert stats["total_batches"] == 5 - assert stats["success_rate"] == 0.6 # 3/5 - - def test_empty_history_returns_empty_stats(self): - """Пустая история возвращает пустой dict.""" - history = InMemoryBatchHistory() - stats = history.get_stats() - assert stats == {} +from generator import BatchRecord class TestBatchRecord: diff --git a/generator/tests/test_kafka_history.py b/generator/tests/test_kafka_history.py index 6fe9739..74b7002 100644 --- a/generator/tests/test_kafka_history.py +++ b/generator/tests/test_kafka_history.py @@ -98,21 +98,17 @@ class TestKafkaBatchHistory: history = KafkaBatchHistory("localhost:9092") - assert history._initialized is True assert history.bootstrap_servers == "localhost:9092" mock_producer_class.assert_called_once() @patch("generator._import_kafka") - def test_init_handles_connection_error(self, mock_import): - """Инициализация обрабатывает ошибку подключения.""" - # Симулируем ошибку при создании producer + def test_init_raises_on_connection_error(self, mock_import): + """Инициализация падает при ошибке подключения.""" mock_producer_class = MagicMock(side_effect=Exception("Connection failed")) mock_import.return_value = (mock_producer_class, Exception) - history = KafkaBatchHistory("localhost:9092") - - assert history._initialized is False - assert history.producer is None + with pytest.raises(Exception, match="Connection failed"): + KafkaBatchHistory("localhost:9092") @patch("generator._import_kafka") def test_add_sends_to_kafka(self, mock_import): @@ -145,32 +141,8 @@ class TestKafkaBatchHistory: assert "value" in call_args[1] @patch("generator._import_kafka") - def test_add_skips_if_not_initialized(self, mock_import): - """add пропускает если не инициализирован.""" - mock_producer_class = MagicMock(side_effect=Exception("Connection failed")) - mock_import.return_value = (mock_producer_class, Exception) - - history = KafkaBatchHistory("localhost:9092") - now = datetime.now(timezone.utc) - record = BatchRecord( - batch_id="skip456", - started_at=now, - finished_at=now, - sent_total=0, - sent_browser=0, - sent_location=0, - sent_device=0, - sent_geo=0, - status="error", - error_message="Test", - ) - - # Не должно упасть - history.add(record) - - @patch("generator._import_kafka") - def test_add_handles_send_error(self, mock_import): - """add обрабатывает ошибку отправки.""" + def test_add_raises_on_send_error(self, mock_import): + """add пробрасывает ошибку отправки.""" mock_producer = MagicMock() mock_producer.send.side_effect = Exception("Send failed") mock_producer_class = MagicMock(return_value=mock_producer) @@ -191,8 +163,8 @@ class TestKafkaBatchHistory: error_message=None, ) - # Не должно упасть - history.add(record) + with pytest.raises(Exception, match="Send failed"): + history.add(record) @patch("generator._import_kafka") def test_flush_calls_producer_flush(self, mock_import): @@ -206,16 +178,6 @@ class TestKafkaBatchHistory: mock_producer.flush.assert_called_once() - @patch("generator._import_kafka") - def test_flush_noop_if_not_initialized(self, mock_import): - """flush ничего не делает если не инициализирован.""" - mock_producer_class = MagicMock(side_effect=Exception("Connection failed")) - mock_import.return_value = (mock_producer_class, Exception) - - history = KafkaBatchHistory("localhost:9092") - # Не должно упасть - history.flush() - @patch("generator._import_kafka") def test_close_calls_producer_close(self, mock_import): """close вызывает close у producer.""" @@ -227,37 +189,3 @@ class TestKafkaBatchHistory: history.close() mock_producer.close.assert_called_once() - - @patch("generator._import_kafka") - def test_close_noop_if_not_initialized(self, mock_import): - """close ничего не делает если не инициализирован.""" - mock_producer_class = MagicMock(side_effect=Exception("Connection failed")) - mock_import.return_value = (mock_producer_class, Exception) - - history = KafkaBatchHistory("localhost:9092") - # Не должно упасть - history.close() - - @patch("generator._import_kafka") - def test_get_stats_returns_status(self, mock_import): - """get_stats возвращает статус инициализации.""" - mock_producer_class = MagicMock() - mock_import.return_value = (mock_producer_class, Exception) - - history = KafkaBatchHistory("localhost:9092") - stats = history.get_stats() - - assert stats["initialized"] is True - assert stats["topic"] == "generator_batch_history" - - @patch("generator._import_kafka") - def test_get_stats_handles_not_initialized(self, mock_import): - """get_stats корректен при неинициализированном состоянии.""" - mock_producer_class = MagicMock(side_effect=Exception("Connection failed")) - mock_import.return_value = (mock_producer_class, Exception) - - history = KafkaBatchHistory("localhost:9092") - stats = history.get_stats() - - assert stats["initialized"] is False - assert stats["topic"] == "generator_batch_history" diff --git a/generator/tests/test_service.py b/generator/tests/test_service.py index 91e24f5..95c92da 100644 --- a/generator/tests/test_service.py +++ b/generator/tests/test_service.py @@ -6,25 +6,10 @@ from unittest.mock import MagicMock, patch import pytest from generator import ( - Config, EventDictionary, GeneratorService, - InMemoryBatchHistory, KafkaBatchHistory + Config, EventDictionary, GeneratorService, KafkaBatchHistory ) -class TestConfigHistoryTopic: - """Тесты для history_topic в конфигурации.""" - - def test_default_history_topic(self, base_config): - """По умолчанию топик истории.""" - assert base_config.history_topic == "generator_batch_history" - - def test_custom_history_topic(self, base_config): - """Кастомный топик истории.""" - from dataclasses import replace - custom_config = replace(base_config, history_topic="custom_history") - assert custom_config.history_topic == "custom_history" - - class TestBatchRecordWithDictConversion: """Тесты конвертации BatchRecord в dict."""