--- apiVersion: v1 kind: ConfigMap metadata: name: kafka-configmap namespace: flows data: kafka.py: | import ssl from functools import partial import anyio from faststream.kafka import KafkaBroker from faststream.security import BaseSecurity from flow.config import logger, settings from flow.models.events import BaseEvent def _build_broker() -> KafkaBroker: ssl_context = ssl.create_default_context() ssl_context.check_hostname = False ssl_context.verify_mode = ssl.CERT_NONE if settings.kafka.ssl_cert: ssl_context.load_verify_locations(cadata=settings.kafka.ssl_cert) security = BaseSecurity( ssl_context=ssl_context, use_ssl=True, ) return KafkaBroker( bootstrap_servers=[settings.kafka.host], security=security, ) broker = _build_broker() async def _publish(event: BaseEvent, key: str) -> None: await broker.publish( message=event, topic=event.metadata.event_type, key=key.encode(), ) def publish_event(event: BaseEvent, key: str) -> None: if not settings.enable_events: return try: anyio.from_thread.run(partial(_publish, event, key)) except Exception: logger.exception( "Failed to publish kafka event %s", event.metadata.event_type, ) __all__ = [ "broker", "publish_event", ]