import logging from types import TracebackType from typing import Any, Self from confluent_kafka import KafkaError, KafkaException, Message, Producer from fansifter_common.utils.functional import lazy_proxy from app.config import settings logger = logging.getLogger(__name__) class KafkaProducer: def __init__(self, config: dict[str, Any]) -> None: self._producer = Producer(config) def __enter__(self) -> Self: return self def __exit__( self, exc_type: type[BaseException] | None, exc_val: BaseException | None, exc_tb: TracebackType | None, ) -> None: self.close() @property def producer(self) -> Producer: if self._producer is None: raise RuntimeError("Producer not started") return self._producer 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 flush(self, timeout: float = 10.0) -> int: return self.producer.flush(timeout) 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(0) logger.debug( "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(), }, ) def get_kafka_producer() -> KafkaProducer: config = { "bootstrap.servers": settings.kafka_bootstrap_servers, "logger": logger, } if settings.kafka_use_ssl: config["security.protocol"] = "SSL" return KafkaProducer(config) kafka_producer = lazy_proxy(get_kafka_producer)