"""Kafka executor class.""" from collections.abc import Callable from typing import Optional, Any from kafka import KafkaAdminClient from kafka import KafkaConsumer from kafka import KafkaProducer from kafka.errors import NoBrokersAvailable from kafka.errors import NodeNotReadyError from dbdeploy.dtos import KafkaCluster from dbdeploy.util.kafka.helpers import serialize_to_json class KafkaExecutor: """Factory class to abstract Kafka operations.""" def __init__(self, kafka_cluster: KafkaCluster): """Initialize executor. Args: kafka_cluster: DTO which holds kafka cluster params. """ self.basic_config = dict( security_protocol=kafka_cluster.security_protocol, bootstrap_servers=kafka_cluster.bootstrap_brokers) def admin_client(self, client_id: str, **kwargs: str) -> KafkaAdminClient: """Return admin client instance.""" try: return KafkaAdminClient( client_id=client_id, **kwargs, **self.basic_config), None except (NoBrokersAvailable, NodeNotReadyError) as err: return None, err def producer( self, client_id: str, key_serializer: Callable[[dict[str, str] | None], bytes | None] = serialize_to_json, value_serializer: Callable[[dict[str, str] | None], bytes | None] = serialize_to_json, **kwargs: str) -> tuple[KafkaProducer, None] | tuple[None, Any]: """Return producer instance.""" try: return KafkaProducer( client_id=client_id, key_serializer=key_serializer, value_serializer=value_serializer, **kwargs, **self.basic_config), None except (NoBrokersAvailable, NodeNotReadyError) as err: return None, err def consumer( self, client_id: str, auto_offset_reset: str, group_id: Optional[str] = None, topics: Optional[str] = None, enable_auto_commit: bool = True, **kwargs: str) -> KafkaConsumer: """Return consumer instance.""" try: if not topics: return KafkaConsumer( client_id=client_id, group_id=group_id, auto_offset_reset=auto_offset_reset, enable_auto_commit=enable_auto_commit, **self.basic_config), None return KafkaConsumer( topics, client_id=client_id, group_id=group_id, auto_offset_reset=auto_offset_reset, enable_auto_commit=enable_auto_commit, **kwargs, **self.basic_config), None except (NoBrokersAvailable, NodeNotReadyError) as err: return None, err