--- apiVersion: v1 kind: ConfigMap metadata: name: kafka-config namespace: message-hub data: kafka.py: | from ssl import SSLContext import ssl from aiokafka.helpers import create_ssl_context from faststream.kafka import KafkaBroker from faststream.security import BaseSecurity, SASLPlaintext, SASLScram512 from pydantic_settings import BaseSettings, SettingsConfigDict class KafkaSettings(BaseSettings): HOST: str = 'localhost' PORT: int = 9092 USERNAME: str | None = None PASSWORD: str | None = None SECURITY_PROTOCOL: str = 'PLAINTEXT' SASL_MECHANISM: str | None = None SSL_CAFILE: str | None = None model_config = SettingsConfigDict(env_prefix='KAFKA_', env_file='.env', extra='ignore') @property def bootstrap_servers(self) -> list[str]: return [f'{self.HOST}:{str(self.PORT)}'] def get_ssl_context(self) -> SSLContext: context = create_ssl_context(cafile=self.SSL_CAFILE) context.check_hostname = False context.verify_mode = ssl.CERT_NONE return context def get_security(self) -> BaseSecurity | None: use_ssl = self.SECURITY_PROTOCOL in ('SSL', 'SASL_SSL') ssl_context = self.get_ssl_context() if use_ssl else None if self.SASL_MECHANISM == 'PLAINTEXT': return SASLPlaintext( username=self.USERNAME or '', password=self.PASSWORD or '', ssl_context=ssl_context, use_ssl=use_ssl, ) if self.SASL_MECHANISM == 'SCRAM-SHA-512': return SASLScram512( username=self.USERNAME or '', password=self.PASSWORD or '', ssl_context=ssl_context, use_ssl=use_ssl, ) if use_ssl: return BaseSecurity(ssl_context=ssl_context) return None @property def broker(self) -> KafkaBroker: return KafkaBroker( bootstrap_servers=self.bootstrap_servers, security=self.get_security(), logger=None, )