import json import ssl from typing import Any from uuid import uuid4 from aiokafka.producer import AIOKafkaProducer class KafkaClient: def __init__( self, bootstrap_servers: str, use_ssl: bool, topic: str, ) -> None: self.bootstrap_servers = bootstrap_servers self.use_ssl = use_ssl self.topic = topic self._producer: AIOKafkaProducer | None = None async def put(self, data: dict[str, Any], key: str | None = None) -> None: await self.producer.send( topic=self.topic, key=(key or str(uuid4())).encode(), value=json.dumps(data).encode(), ) @property def producer(self) -> AIOKafkaProducer: if self._producer is None: raise RuntimeError("Kafka producer is not started") return self._producer async def start(self) -> None: if self._producer is None: ctx = ssl.SSLContext(protocol=ssl.PROTOCOL_TLSv1_2) self._producer = AIOKafkaProducer( bootstrap_servers=self.bootstrap_servers, security_protocol="SSL" if self.use_ssl else "PLAINTEXT", ssl_context=ctx if self.use_ssl else None, ) await self._producer.start() async def stop(self) -> None: if self._producer is not None: await self._producer.stop() self._producer = None