From 66b568838e3f112e464f41ff3de4695547eb724d Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 1 Jul 2025 08:21:00 -0300 Subject: [PATCH] SIENTIAPDE-1110 Refactor Kafka activity to use instance-level kafka_connector for improved message polling and management, replacing local variable usage with class attribute methods. --- scouter/activities/kafka.py | 13 +++---------- 1 file changed, 3 insertions(+), 10 deletions(-) diff --git a/scouter/activities/kafka.py b/scouter/activities/kafka.py index 64ad62e..ccd434f 100644 --- a/scouter/activities/kafka.py +++ b/scouter/activities/kafka.py @@ -58,20 +58,13 @@ class Kafka(BaseActivity): topic = input_data["topic"] # Subscribe to the specified topic - kafka_connector = KafkaConsumer( - topic, - bootstrap_servers=self.bootstrap_servers, - auto_offset_reset="earliest", - enable_auto_commit=True, - group_id=self.group_id, - value_deserializer=lambda x: json.loads(x.decode("utf-8")) - ) + self.kafka_connector.subscribe([topic]) # List to store message values message_values = [] # Poll for messages - records = kafka_connector.poll(timeout_ms=self.polling_time) + records = self.kafka_connector.poll(timeout_ms=self.polling_time) self.debug( f"Polled {len(records)} records from topic: {topic}", @@ -97,6 +90,6 @@ class Kafka(BaseActivity): metadata=metadata ) - kafka_connector.close() + self.kafka_connector.unsubscribe() return DataFrame(message_values).to_dict()