From 6144bcf8ffd3ea8e4e273b3a8e98bc4ee3edc843 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 30 Jun 2025 14:11:30 -0300 Subject: [PATCH] SIENTIAPDE-1110 Update Kafka polling time in values.yaml and refactor Kafka subscription in kafka.py to use KafkaConsumer for improved message handling. --- scouter/activities/kafka.py | 9 ++++++++- values.yaml | 2 +- 2 files changed, 9 insertions(+), 2 deletions(-) diff --git a/scouter/activities/kafka.py b/scouter/activities/kafka.py index 089b28b..a05b292 100644 --- a/scouter/activities/kafka.py +++ b/scouter/activities/kafka.py @@ -56,7 +56,14 @@ class Kafka(BaseActivity): topic = input_data["topic"] # Subscribe to the specified topic - self.kafka_connector.subscribe([topic]) + self.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")) + ) # List to store message values message_values = [] diff --git a/values.yaml b/values.yaml index 48db1f1..8006b67 100644 --- a/values.yaml +++ b/values.yaml @@ -146,7 +146,7 @@ env: - name: KAFKA_BOOTSTRAP_SERVERS value: "kafka.kafka.svc.cluster.local:9092" - name: KAFKA_POLLING_TIME - value: "1000" + value: "10000" - name: REDIS_HOST value: "redis-master.redis.svc.cluster.local"