"""Lambda send_failure function module.""" import json from config import app_logger as logger from confluent_kafka import Producer from config import BROKER_SERVERS from config import DEAD_LETTER_KAFKA_TOPIC def delivery_report(err, msg): """Log delivery results from callbacks for each message produced, triggered by flush().""" if err is not None: logger.exception('Message delivery failed: {}'.format(err)) else: logger.info('Message delivered to {} [{}]'.format(msg.topic(), msg.partition())) def handler(event, context): """Lambda entry point.""" try: logger.info('Starting to write failed event to Kafka topic.') producer = Producer({'bootstrap.servers': BROKER_SERVERS, 'security.protocol': 'SSL'}) # Asynchronously produce a message, the delivery report callback # will be triggered from poll() above, or flush() below, when the message has # been successfully delivered or failed permanently. producer.produce( DEAD_LETTER_KAFKA_TOPIC, json.dumps(event).encode('utf-8'), callback=delivery_report) # Wait for any outstanding messages to be delivered and delivery report callback. producer.flush() logger.info(f'Successfully wrote failed message to Kafka topic: {DEAD_LETTER_KAFKA_TOPIC}') return {'status': 'success'} except Exception as err: logger.exception(str(err)) raise err