import logging from typing import Any from confluent_kafka import KafkaError, KafkaException, Message, Producer logger = logging.getLogger(__name__) class KafkaProducer: POLL_TIMEOUT = 0.1 def __init__(self, config: dict[str, Any]) -> None: self._config = config self._producer: Producer | None = None @property def producer(self) -> Producer: if self._producer is None: raise RuntimeError("Producer not started") return self._producer def start(self) -> None: self._producer = Producer(self._config) logger.info("Kafka producer started.") def close(self, timeout: float = 10.0) -> None: if self._producer is None: return remaining = self.producer.flush(timeout) if remaining > 0: logger.warning( "Kafka producer closed with undelivered messages.", extra={"kafka_remaining_messages": remaining}, ) self._producer = None logger.info("Kafka producer closed.") def publish(self, *, topic: str, value: str, key: str) -> None: try: self.producer.produce( topic, value, key=key, on_delivery=self._delivery_callback ) self.producer.poll(self.POLL_TIMEOUT) logger.info( "Message queued for delivery to `%s` topic with key `%s`", topic, key, extra={ "kafka_publish": True, "kafka_topic": topic, "kafka_key": key, }, ) except BufferError: logger.error( "Kafka producer queue is full. Message not sent to `%s` topic with key `%s`", topic, key, extra={ "kafka_topic": topic, "kafka_key": key, }, ) raise except KafkaException as e: logger.error( "Failed to publish message to `%s` topic with key `%s`", topic, key, extra={ "kafka_topic": topic, "kafka_key": key, "kafka_error": str(e), }, ) raise except Exception as e: logger.error( "Unexpected error while publishing message to `%s` topic with key `%s`", topic, key, extra={ "kafka_topic": topic, "kafka_key": key, "kafka_error": str(e), }, ) raise @staticmethod def _delivery_callback(err: KafkaError | None, msg: Message) -> None: topic = msg.topic() if err is not None: logger.error( "Message delivery failed to `%s` topic", topic, extra={ "kafka_topic": topic, "kafka_partition": msg.partition(), "kafka_error": str(err), }, ) else: logger.debug( "Message delivered successfully to `%s` topic", topic, extra={ "kafka_topic": topic, "kafka_partition": msg.partition(), "kafka_offset": msg.offset(), }, )