"""Sends a message to a topic.""" from confluent_kafka import Producer from lambdacommon.common_config import logger import config p = Producer({'bootstrap.servers': config.KAFKA_BROKERS, 'security.protocol': 'SSL'}) def delivery_report(err, msg): """Log delivery results from callbacks for each message produced, triggered by poll() or flush().""" if err is not None: logger.exception('Message delivery failed: {}'.format(err)) else: print('Message delivered to {} [{}]'.format(msg.topic(), msg.partition())) def produce_message(data, topic): """Send message to topic.""" # Trigger any available delivery report callbacks from previous produce() calls p.poll(0) # 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. p.produce(topic, data.encode('utf-8'), callback=delivery_report) # Wait for any outstanding messages to be delivered and delivery report # callbacks to be triggered. p.flush()