"""Приёмники: куда уезжают канонические байты. Про содержимое они не знают. Приёмник глуп по замыслу (спека генератора, раздел 4): он берёт готовый байт события и доставляет его. Ни формы, ни темпа он не решает — форму задал сериализатор, темп задаёт проигрыватель. Отсюда и весь их интерфейс: принять событие, дождаться принятого, закрыться. Приёмников два, а режимов три: «Kafka пачкой» и «Kafka с темпом ×60» различаются не приёмником, а тем, кто его зовёт. Ускорять доставку, зная о времени события, значило бы вернуть приёмнику знание о содержимом — ровно то, чего правило не хочет. **Ключа у сообщения Kafka нет** — решение тикета #41, и вот довод. `WatchID` уникален у каждого события, поэтому ключом он не был бы ключом: обещание Kafka про ключ — «сообщения одного ключа лежат в одном разделе и сохраняют порядок», а у ряда, где ключи не повторяются, обещать нечего. Читатель же прочёл бы такой ключ как смысл, которого в нём нет. Настоящий ключ здесь — `ClientID` (события одной куки по порядку), и он тикетом не назначен: раскладку по разделам стенд намеренно оставляет транспорту (спека, раздел 2), а лаба «какая нода читала топик» живёт как раз тем, что раскладка не предрешена. Раздел для сообщения без ключа librdkafka выбирает случайно, но подряд идущие сообщения на десяток миллисекунд липнут к одному — измеренные умолчания и их следствия записаны в спеке (раздел 9). """ from pathlib import Path from typing import Any, Protocol # Сколько ждать разбора очереди отправителя, когда она заполнилась, и сколько — # доставки остатка при закрытии. Оба числа — потолок ожидания, а не пауза: # обычно ждать не приходится вовсе. QUEUE_WAIT_SECONDS = 1.0 FLUSH_WAIT_SECONDS = 60.0 class Sink(Protocol): """Приёмник: принимает байты события и доводит их до места.""" def send(self, payload: bytes) -> None: """Принять одно событие.""" def flush(self) -> None: """Дождаться, пока принятое дойдёт до места. Нужно не только при закрытии: доставка асинхронна, и без этого рубежа «доставка 0,1 с» в логе означала бы лишь то, что события успели лечь в очередь отправителя. Проигрыватель ставит рубеж в конце каждого дня. """ def close(self) -> None: """Закрыть приёмник, дождавшись всего принятого.""" class FileSink: """Файл: одно событие — одна строка. Построчность — контракт файла, а не удобство: файл читают построчно, и склейка двух событий в строку сломала бы разбор целиком. Файл — кэш чистой функции (спека, раздел 4): потерял — пересчитал, поэтому места в репозитории ему не отведено. """ def __init__(self, path: Path) -> None: path.parent.mkdir(parents=True, exist_ok=True) self._file = path.open("wb") def send(self, payload: bytes) -> None: self._file.write(payload + b"\n") def flush(self) -> None: self._file.flush() def close(self) -> None: self._file.close() class KafkaSink: """Kafka: одно событие — одно сообщение. «Пачкой» относится к темпу отправки, а не к упаковке: события идут подряд без пауз, но каждое отдельным сообщением (спека, раздел 4). Отправка асинхронная — отправитель копит сообщения в своей очереди и шлёт их пачками сам; отсюда `poll`, который отдаёт нам отчёты о доставке, и `flush`, без которого хвост очереди уехал бы в никуда вместе с процессом. """ def __init__(self, brokers: str, topic: str) -> None: # Импорт внутри: клиент — необязательная часть пакета, и без него # генератор пишет в файл (спека, раздел 9). from confluent_kafka import Producer self._topic = topic self._failure: str | None = None self._producer = Producer({"bootstrap.servers": brokers}) def send(self, payload: bytes) -> None: while True: try: self._producer.produce( self._topic, value=payload, on_delivery=self._report ) break except BufferError: # Очередь отправителя полна: ждём, пока брокер её разберёт. # `poll` здесь и работа, и пауза — он же отдаёт отчёты. self._producer.poll(QUEUE_WAIT_SECONDS) self._producer.poll(0) self._raise_failure() def flush(self) -> None: remaining = self._producer.flush(FLUSH_WAIT_SECONDS) if remaining: raise RuntimeError( f"Kafka не приняла {remaining} сообщений за" f" {FLUSH_WAIT_SECONDS:.0f} с: брокер недоступен или не успевает" ) self._raise_failure() def close(self) -> None: # Закрывать у отправителя нечего — важно лишь не бросить хвост # очереди: он уехал бы в никуда вместе с процессом. self.flush() def _report(self, error: Any, message: Any) -> None: """Отчёт о доставке: первую неудачу запоминаем, остальные не важны.""" if error is not None and self._failure is None: self._failure = str(error) def _raise_failure(self) -> None: """Неудачную доставку превращаем в остановку прогона. Молча потерянное сообщение — худший исход: счёт в хранилище разойдётся с числом отправленного, а причина будет забыта. """ if self._failure is not None: raise RuntimeError(f"Kafka не приняла сообщение: {self._failure}")