From d33a5cbb72d320ce11fa4f5d0cbedf74a2be3ad3 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 30 Jun 2025 14:12:02 -0300 Subject: [PATCH] SIENTIAPDE-1110 Refactor Kafka activity to use local kafka_connector variable for improved readability and ensure proper closure of the Kafka consumer after polling. --- scouter/activities/kafka.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/scouter/activities/kafka.py b/scouter/activities/kafka.py index a05b292..8fe1454 100644 --- a/scouter/activities/kafka.py +++ b/scouter/activities/kafka.py @@ -56,7 +56,7 @@ class Kafka(BaseActivity): topic = input_data["topic"] # Subscribe to the specified topic - self.kafka_connector = KafkaConsumer( + kafka_connector = KafkaConsumer( topic, bootstrap_servers=self.bootstrap_servers, auto_offset_reset="earliest", @@ -69,7 +69,7 @@ class Kafka(BaseActivity): message_values = [] # Poll for messages - records = self.kafka_connector.poll(timeout_ms=self.polling_time) + records = kafka_connector.poll(timeout_ms=self.polling_time) self.debug( f"Polled {len(records)} records from topic: {topic}", @@ -95,4 +95,6 @@ class Kafka(BaseActivity): metadata=metadata ) + kafka_connector.close() + return DataFrame(message_values).to_dict()