"""Kafka Queue Worker.""" import itertools from kafka import KafkaConsumer from neo4j_sync_dlq import config from neo4j_sync_dlq.logic import processing from neo4j_sync_dlq.logic.kafka_message import KafkaMessage class KafkaWorker: """Kafka worker that dispatches the messages received.""" def __init__(self, dlq_topic, group_id, kafka_configuration): """Create KafkaWorker instance. Args: dlq_topic (str): DLQ topic name group_id (str): Kafka consumer group id kafka_configuration (dict): kafka connection options """ self.do_work = True self.consumer = KafkaConsumer( dlq_topic, group_id=group_id, **kafka_configuration ) # TODO: implement the message retry logic # self.producer = KafkaProducer( # value_serializer=lambda v: json.dumps(v).encode('utf-8'), # **configuration # ) def run(self): """Run the Kafka worker.""" while self.do_work: records = self.consumer.poll( config.KAFKA_POLL_TIMEOUT_MS, config.KAFKA_POLL_NUM_RECORDS) if records: for record in list(itertools.chain(*records.values())): message = KafkaMessage(record) processing.process(message)