From 946d953598a65b41a4e871e1238c908559f61bc7 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 1 Jul 2025 10:57:12 -0300 Subject: [PATCH] SIENTIAPDE-1110 Refactor Kafka activity to improve message consumption by using getmany for batch retrieval, enhancing performance and simplifying message handling logic. --- scouter/activities/kafka.py | 23 ++++++++++------------- 1 file changed, 10 insertions(+), 13 deletions(-) diff --git a/scouter/activities/kafka.py b/scouter/activities/kafka.py index ec0aea4..6276455 100644 --- a/scouter/activities/kafka.py +++ b/scouter/activities/kafka.py @@ -71,22 +71,19 @@ class Kafka(BaseActivity): consumer = self.consumers[topic] - await consumer.subscribe([topic]) + consumer.subscribe(topics=[topic]) message_values = [] - end_time = asyncio.get_event_loop().time() + (self.polling_time / 1000) - while asyncio.get_event_loop().time() < end_time: - try: - msg = await asyncio.wait_for(consumer.getone(), timeout=(end_time - asyncio.get_event_loop().time())) - message_values.append(msg.value) - except asyncio.TimeoutError: - break - except Exception as e: - self.error(f"Error while consuming: {e}", metadata=metadata) - break + messages = await consumer.getmany(timeout_ms=self.polling_time) - self.debug( + for tp, msgs in messages.items(): + msg_topic = tp.topic + if msg_topic == topic: + for msg in msgs: + message_values.append(msg.value) + + self.info( f"Loaded {len(message_values)} messages from topic: {topic}", metadata=metadata ) @@ -96,7 +93,7 @@ class Kafka(BaseActivity): metadata=metadata ) - await consumer.unsubscribe() + consumer.unsubscribe() if not message_values: return {}