"""Kafka Queue Worker.""" import json import logging import sys from kafka import KafkaConsumer, KafkaProducer from switchboard_consumer.config import GROUP_ID, INBOUND, OUTBOUND from switchboard_consumer.logic import processing logger = logging.getLogger('kafka') logger.addHandler(logging.StreamHandler(sys.stdout)) logger.setLevel(logging.INFO) class KafkaWorker: """Kafka worker that dispatches the messages received.""" def __init__(self, configuration): """Create a KafkaWorker. Args: configuration: the Kafka bootstrap servers string """ self.consumer = KafkaConsumer( INBOUND, group_id=GROUP_ID, max_poll_records=1, **configuration ) self.producer = KafkaProducer( value_serializer=lambda v: json.dumps(v).encode('utf-8'), **configuration ) def run(self, switchboard_client, orchard_client): """Run the worker.""" for message in self.consumer: responses = processing.process(message.value, switchboard_client, orchard_client) if responses: for response in responses: self.producer.send(OUTBOUND, response).get()