"""Kafka Producer.""" from typing import Dict from confluent_kafka.serialization import StringSerializer from kafka_utils.producer.event import EventProducer from kafka_utils.producer.serializer.simple_json import SimpleJSONSerializer from lambdacommon.common_config import logger import config BOOTSTRAP_SERVERS = config.KAFKA_BOOTSTRAP_SERVERS PROTOCOL = config.KAFKA_SECURITY_PROTOCOL def get_producer(): """Get kafka producer.""" return EventProducer( bootstrap_servers=BOOTSTRAP_SERVERS, key_serializer=StringSerializer(), value_serializer=SimpleJSONSerializer(), security_protocol=PROTOCOL ) def send_message_to_kaka_topic(message: Dict, key: str, topic: str): """Produce a message to Kafka topic.""" get_producer().produce( topic=topic, event_key=key, event_value=message, callback=msg_delivery_callback ) def msg_delivery_callback(err, msg): """Handle Kafka message delivery result.""" if err: logger.error(f'There was an error producing message for {msg.key()} to Kafka') raise Exception(err)