"""Sends a message to a topic.""" from src.common.exceptions import exceptions from src.common import logger from confluent_kafka import Producer import config p = Producer( { 'bootstrap.servers': config.BROKER_SERVERS, 'message.timeout.ms': config.KAFKA_TIMEOUT, 'security.protocol': 'SSL', } ) def new_version_logger(err, msg): """Log newVersion results from callbacks for each message produced.""" if err is not None: logger.exception('Message newVersion failed: {}'.format(err)) raise exceptions.RetryableException(f'Kafka error: {err}') else: logger.info( 'Message newVersion to {} [{}]'.format(msg.topic(), msg.partition()) ) def produce_message(data, topic, key): """Send message to topic.""" # Asynchronously produce a message, the new_version 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'), key=key, callback=new_version_logger) # Wait for any outstanding messages to be delivered and new_version # callbacks to be triggered. p.flush()