Files
clickstream-ch-kafka-supers…/generator/src/clickstream_generator/kafka_io.py
T
ddadminandClaude Fable 5 103ac021c8 feat(generator): инкрементальные счётчики manifest без перечитки Kafka
- Зачем:
  - world_next_day перечитывал всю историю топиков Kafka ради
    накопительных счётчиков — время прогона росло с возрастом мира
    (issue #5, находка F9).
- Что:
  - счётчики засеваются при import из уже прочитанного артефакта и при
    backfill из потока; next-day продвигает их только событиями нового
    дня, полного чтения Kafka больше нет;
  - катящаяся контрольная сумма — сумма SHA-256 событий по модулю 2^256
    (инкремент равен полному пересчёту), старый формат артефакта
    принимается без изменений;
  - точные множества click_id/user_domain_id вынесены из manifest в
    цепочку контент-адресуемых фрагментов (<=10 000 ID, SHA-256-цепочка,
    отдельный топик counter_chunks) — потолок сообщения Kafka не грозит,
    предел 900 000 байт проверяется явно с понятной ошибкой;
  - порядок записи всюду: фрагменты -> manifest -> state; старое локальное
    состояние отклоняется с подсказкой перезапустить import;
  - документация manifest/state обновлена (ARCHITECTURE, OPERATIONS,
    runbook startup-history).
- Проверка:
  - make test (216+31) и make lint зелёные;
  - живая приёмка на чистом стенде: import 235 с; три прогона
    world_next_day — 716/718/716 с (плоское время, O(нового дня));
    мир 3->6 дней, 561 942 события; make generated-history-chain-check —
    все порции и стыки однородны;
  - тест равенства инкремента и полного пересчёта:
    test_incremental_counters_equal_full_recompute.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-22 23:47:02 +03:00

586 lines
21 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Kafka-интеграция генератора."""
import hashlib
import json
import logging
import sys
import time
from dataclasses import dataclass
from datetime import datetime
from clickstream_generator.metrics import METRICS_ERRORS_TOTAL, METRICS_EVENTS_TOTAL
from clickstream_generator.state import GeneratorState
from clickstream_generator.state import UnsupportedStateVersionError
logger = logging.getLogger("generator")
DATA_TOPICS = ("browser_events", "location_events", "device_events", "geo_events")
MAX_COMPACT_MESSAGE_BYTES = 900_000
_kafka_imported = False
KafkaProducer = None
KafkaError = None
def _import_kafka():
"""Лениво импортирует Kafka-клиент, чтобы тесты могли работать без Kafka."""
global _kafka_imported, KafkaProducer, KafkaError
if not _kafka_imported:
from kafka import KafkaProducer
from kafka.errors import KafkaError
_kafka_imported = True
return KafkaProducer, KafkaError
def _with_retry(operation, max_retries: int = 5, base_delay: float = 1.0, max_delay: float = 30.0):
"""Выполняет операцию с экспоненциальным backoff."""
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}): "
f"{e}. Retrying in {delay:.1f}s..."
)
time.sleep(delay)
else:
logger.error(f"Operation failed after {max_retries} attempts: {e}")
raise last_exception
def _facade_attr(name: str, fallback):
"""Берёт совместимый mock из фасада generator, если тест его подменил."""
facade = sys.modules.get("generator")
return getattr(facade, name, fallback) if facade else fallback
def _retry(operation, **kwargs):
return _facade_attr("_with_retry", _with_retry)(operation, **kwargs)
def _kafka_importer():
return _facade_attr("_import_kafka", _import_kafka)
def kafka_event_value_json(event: dict) -> str:
"""Возвращает JSON value ровно в формате Kafka producer генератора."""
return json.dumps(event)
def kafka_event_value_bytes(event: dict) -> bytes:
"""Возвращает Kafka value bytes для события генератора."""
return kafka_event_value_json(event).encode("utf-8")
@dataclass
class BatchRecord:
"""Запись об отправленном батче."""
batch_id: str
started_at: datetime
finished_at: datetime
sent_total: int
sent_browser: int
sent_location: int
sent_device: int
sent_geo: int
status: str
error_message: str | None = None
def to_dict(self) -> dict:
"""Конвертирует в словарь для сериализации."""
return {
"batch_id": self.batch_id,
"started_at": self.started_at.isoformat(),
"finished_at": self.finished_at.isoformat(),
"sent_total": self.sent_total,
"sent_browser": self.sent_browser,
"sent_location": self.sent_location,
"sent_device": self.sent_device,
"sent_geo": self.sent_geo,
"status": self.status,
"error_message": self.error_message,
}
class KafkaStateManager:
"""Управление состоянием генератора в Kafka compact topic."""
STATE_TOPIC = "generator_state"
STATE_KEY = "default"
def __init__(self, bootstrap_servers: str):
self.bootstrap_servers = bootstrap_servers
KafkaProducerCls, _ = _kafka_importer()()
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:
"""Сохраняет состояние в топик."""
value = state.to_dict()
def _do_send():
self.producer.send(self.STATE_TOPIC, key=self.STATE_KEY, value=value)
_retry(_do_send, max_retries=3, base_delay=0.5)
def flush(self) -> None:
"""Сбрасывает буфер с retry."""
def _do_flush():
self.producer.flush()
_retry(_do_flush, max_retries=3, base_delay=0.5)
def close(self) -> None:
"""Закрывает соединение."""
try:
self.producer.close()
except Exception as e:
logger.debug(f"Error closing producer (ignored): {e}")
def load(self) -> GeneratorState | None:
"""Загружает последнее состояние из топика."""
from kafka import KafkaConsumer
logger.info(f"Loading state from topic {self.STATE_TOPIC}")
def _do_load():
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()
return last_state
try:
last_state = _retry(_do_load, max_retries=3, base_delay=0.5)
if last_state:
logger.info(
f"Restored state: tick={last_state.get('tick')}, "
f"last_batch_id={last_state.get('last_batch_id')}"
)
try:
restored = GeneratorState.from_dict(last_state)
except UnsupportedStateVersionError:
logger.error(
"Unsupported generator state version %s",
last_state.get("version", "1.0"),
)
raise
except ValueError:
logger.warning("State data was invalid, starting fresh")
return None
return restored
logger.info("No previous state found, starting fresh")
return None
except UnsupportedStateVersionError:
raise
except Exception as e:
logger.warning(f"Failed to load state: {e}, starting fresh")
return None
class KafkaStartupHistoryManifest:
"""Хранение манифеста стартовой истории в Kafka compact topic."""
MANIFEST_TOPIC = "generator_startup_history_manifest"
COUNTER_TOPIC = "generator_startup_history_counter_chunks"
MANIFEST_KEY = "default"
COUNTER_CHUNK_KEY_PREFIX = "id-set:"
def __init__(self, bootstrap_servers: str):
self.bootstrap_servers = bootstrap_servers
KafkaProducerCls, _ = _kafka_importer()()
logger.info(
f"Connecting to Kafka for startup history manifest 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 startup history manifest successfully")
def save(self, manifest: dict) -> None:
"""Сохраняет манифест стартовой истории."""
self._assert_message_size(manifest, "manifest")
def _do_send():
self.producer.send(
self.MANIFEST_TOPIC,
key=self.MANIFEST_KEY,
value=manifest,
)
_retry(_do_send, max_retries=3, base_delay=0.5)
def save_counter_chunk(self, chunk_sha256: str, chunk: dict) -> None:
"""Сохраняет неизменяемый фрагмент точных множеств идентификаторов."""
if not isinstance(chunk_sha256, str) or len(chunk_sha256) != 64:
raise ValueError("хеш фрагмента множеств идентификаторов неверен")
actual_sha256 = hashlib.sha256(
json.dumps(
chunk,
sort_keys=True,
separators=(",", ":"),
ensure_ascii=True,
).encode("utf-8")
).hexdigest()
if actual_sha256 != chunk_sha256:
raise ValueError("хеш фрагмента множеств идентификаторов не совпадает")
self._assert_message_size(chunk, "фрагмент множеств идентификаторов")
def _do_send():
self.producer.send(
self.COUNTER_TOPIC,
key=self.COUNTER_CHUNK_KEY_PREFIX + chunk_sha256,
value=chunk,
)
_retry(_do_send, max_retries=3, base_delay=0.5)
@staticmethod
def _assert_message_size(value: dict, label: str) -> None:
size = len(json.dumps(value, ensure_ascii=True).encode("utf-8"))
if size > MAX_COMPACT_MESSAGE_BYTES:
raise ValueError(
f"{label} занимает {size} байт и превышает безопасный предел "
f"{MAX_COMPACT_MESSAGE_BYTES} байт одного сообщения Kafka"
)
def flush(self) -> None:
"""Сбрасывает буфер с retry."""
def _do_flush():
self.producer.flush()
_retry(_do_flush, max_retries=3, base_delay=0.5)
def close(self) -> None:
"""Закрывает соединение."""
try:
self.producer.close()
except Exception as e:
logger.debug(f"Error closing manifest producer (ignored): {e}")
def load(self) -> dict | None:
"""Загружает последний манифест стартовой истории."""
from kafka import KafkaConsumer
logger.info(f"Loading startup history manifest from topic {self.MANIFEST_TOPIC}")
def _do_load():
consumer = KafkaConsumer(
self.MANIFEST_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_manifest = None
for message in consumer:
if message.key and message.key.decode("utf-8") == self.MANIFEST_KEY:
last_manifest = message.value
consumer.close()
return last_manifest
try:
return _retry(_do_load, max_retries=3, base_delay=0.5)
except Exception as e:
logger.warning(f"Failed to load startup history manifest: {e}")
return None
class KafkaDataTopicReader:
"""Читает устойчивый снимок data-топиков для накопительного manifest."""
def __init__(self, bootstrap_servers: str):
self.bootstrap_servers = bootstrap_servers
def load(self) -> dict[str, list[dict]]:
"""Возвращает все события data-топиков после остановки live-потока."""
from kafka import KafkaConsumer, TopicPartition
consumer = KafkaConsumer(
bootstrap_servers=self.bootstrap_servers,
enable_auto_commit=False,
value_deserializer=lambda value: json.loads(value.decode("utf-8")),
)
topics = {topic: [] for topic in DATA_TOPICS}
try:
partitions = []
for topic in DATA_TOPICS:
topic_partitions = consumer.partitions_for_topic(topic) or set()
partitions.extend(
TopicPartition(topic, partition)
for partition in topic_partitions
)
if not partitions:
return topics
consumer.assign(partitions)
consumer.seek_to_beginning(*partitions)
end_offsets = consumer.end_offsets(partitions)
empty_polls = 0
while any(
consumer.position(partition) < end_offsets[partition]
for partition in partitions
):
records = consumer.poll(timeout_ms=1000)
if not records:
empty_polls += 1
if empty_polls >= 5:
raise RuntimeError(
"не удалось дочитать data-топики до зафиксированных "
"конечных смещений"
)
continue
empty_polls = 0
for partition, messages in records.items():
for message in messages:
if (
message.offset < end_offsets[partition]
and message.value is not None
):
topics[message.topic].append(message.value)
finally:
consumer.close()
return topics
def ensure_topics(bootstrap_servers: str) -> None:
"""Создаёт служебные топики, если их ещё нет."""
from kafka import KafkaAdminClient
from kafka.admin import NewTopic
from kafka.errors import TopicAlreadyExistsError
def _create_topics():
admin_client = KafkaAdminClient(bootstrap_servers=bootstrap_servers)
try:
history_topic = NewTopic(
name=KafkaBatchHistory.HISTORY_TOPIC,
num_partitions=1,
replication_factor=1,
)
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",
},
)
manifest_topic = NewTopic(
name=KafkaStartupHistoryManifest.MANIFEST_TOPIC,
num_partitions=1,
replication_factor=1,
topic_configs={
"cleanup.policy": "compact",
"min.cleanable.dirty.ratio": "0.1",
"delete.retention.ms": "100",
},
)
counter_topic = NewTopic(
name=KafkaStartupHistoryManifest.COUNTER_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,
manifest_topic,
counter_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()
_retry(_create_topics, max_retries=5, base_delay=1.0)
class KafkaBatchHistory:
"""Хранение истории batch в Kafka."""
HISTORY_TOPIC = "generator_batch_history"
def __init__(self, bootstrap_servers: str):
self.bootstrap_servers = bootstrap_servers
self.producer = None
self._connect()
def _connect(self):
"""Устанавливает соединение с Kafka с retry."""
KafkaProducerCls, _ = _kafka_importer()()
def _do_connect():
logger.info(f"Connecting to Kafka for history 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 history successfully")
_retry(_do_connect, max_retries=5, base_delay=1.0)
def add(self, record: BatchRecord):
"""Добавляет запись в историю."""
key = record.batch_id
value = record.to_dict()
def _do_send():
self.producer.send(self.HISTORY_TOPIC, key=key, value=value)
try:
_retry(_do_send, max_retries=3, base_delay=0.5)
except Exception as e:
logger.warning(f"Failed to send history record after retries: {e}, attempting reconnect")
self._connect()
_retry(_do_send, max_retries=2, base_delay=0.5)
def flush(self):
"""Сбрасывает буфер с retry."""
def _do_flush():
self.producer.flush()
_retry(_do_flush, max_retries=3, base_delay=0.5)
def close(self):
"""Закрывает соединение."""
try:
if self.producer:
self.producer.close()
except Exception as e:
logger.debug(f"Error closing history producer (ignored): {e}")
class KafkaPublisher:
"""Публикация событий в Kafka с retry и реконнектом."""
def __init__(self, bootstrap_servers: str):
self.bootstrap_servers = bootstrap_servers
self.producer = None
self._connect()
def _connect(self):
"""Устанавливает соединение с Kafka с retry."""
KafkaProducerCls, _ = _kafka_importer()()
def _do_connect():
logger.info(f"Connecting to Kafka at {self.bootstrap_servers}")
self.producer = KafkaProducerCls(
bootstrap_servers=self.bootstrap_servers,
value_serializer=kafka_event_value_bytes,
key_serializer=lambda k: k.encode("utf-8") if k else None,
batch_size=16384,
linger_ms=100,
retries=3,
retry_backoff_ms=1000,
)
logger.info("Connected to Kafka successfully")
_retry(_do_connect, max_retries=5, base_delay=1.0)
def _publish_with_retry(self, topic: str, events: list[dict]) -> tuple[int, int]:
"""Внутренняя функция публикации с retry на уровне batch."""
if not self.producer:
raise RuntimeError("Producer not connected")
sent = 0
errors = 0
futures = []
for event in events:
key = event.get("event_id") or event.get("click_id")
try:
future = self.producer.send(topic, key=key, value=event)
futures.append(future)
except Exception as e:
logger.error(f"Failed to send message to {topic}: {e}")
errors += 1
METRICS_ERRORS_TOTAL.labels(topic=topic).inc()
for future in futures:
try:
future.get(timeout=10)
sent += 1
METRICS_EVENTS_TOTAL.labels(topic=topic).inc()
except Exception as e:
logger.error(f"Failed to confirm message delivery: {e}")
errors += 1
METRICS_ERRORS_TOTAL.labels(topic=topic).inc()
return sent, errors
def publish(self, topic: str, events: list[dict]) -> tuple[int, int]:
"""Публикует события в топик с retry и автоматическим реконнектом."""
def _do_publish():
return self._publish_with_retry(topic, events)
try:
return _retry(_do_publish, max_retries=3, base_delay=0.5)
except Exception as e:
logger.warning(f"Publish failed after retries: {e}, attempting reconnect")
self._connect()
return _retry(_do_publish, max_retries=2, base_delay=0.5)
def flush(self):
"""Сбрасывает буфер с retry."""
def _do_flush():
if self.producer:
self.producer.flush()
_retry(_do_flush, max_retries=3, base_delay=0.5)
def close(self):
"""Закрывает соединение."""
try:
if self.producer:
self.producer.close()
except Exception as e:
logger.debug(f"Error closing producer (ignored): {e}")