import json import logging import uuid from typing import Any from confluent_kafka import KafkaError, Message, Producer logger = logging.getLogger(__name__) class KafkaClient: def __init__( self, producer: Producer, kafka_event_topic: str, kafka_status_update_topic: str ) -> None: self.producer = producer self.kafka_event_topic = kafka_event_topic self.kafka_status_update_topic = kafka_status_update_topic def send_event_message(self, data: dict[str, Any]) -> None: self.send_message( self.kafka_event_topic, key=str(uuid.uuid4()), value=json.dumps(data), ) def send_fan_state(self, data: dict[str, Any]) -> None: self.send_message( topic=self.kafka_status_update_topic, key=str(uuid.uuid4()), value=json.dumps(data), ) def send_message(self, topic: str, key: str, value: str) -> None: self.producer.produce( topic, key=key, value=value, callback=self.kafka_producer_callback, ) self.producer.flush() @staticmethod def kafka_producer_callback(error: KafkaError | None, message: Message) -> None: if error: logger.error("Kafka producer error %s", error) else: logger.info("Message %s delivered", message.key())