"""Kafka Producer initialization.""" from typing import Dict import config from lambdacommon.common_config import logger from confluent_kafka.serialization import StringSerializer from kafka_utils.producer.event import EventProducer from kafka_utils.producer.serializer.simple_json import SimpleJSONSerializer TOPIC = config.KAFKA_TARGET_TOPIC HEADER_KEY = config.KAFKA_MESSAGE_HEADER_KEY 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 produce_for_id(message: Dict, id: str): """Produce a message to Kafka topic including id in header.""" get_producer().produce( topic=TOPIC, event_key=id, event_value=message, headers={HEADER_KEY: id.encode()}, 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)