From 4d01cafb3d977b89e226eca312f8eb40dd14d23f Mon Sep 17 00:00:00 2001 From: jashan7167 Date: Thu, 2 Apr 2026 14:30:20 +0530 Subject: [PATCH] testing the dashboard-ingest-v2 --- .../java/com/ingestpipeline/consumer/EnrichmentConsumer.java | 2 +- .../java/com/ingestpipeline/consumer/TransformConsumer.java | 2 +- .../src/main/resources/application.properties | 2 ++ 3 files changed, 4 insertions(+), 2 deletions(-) diff --git a/business-services/dashboard-ingest-v2/src/main/java/com/ingestpipeline/consumer/EnrichmentConsumer.java b/business-services/dashboard-ingest-v2/src/main/java/com/ingestpipeline/consumer/EnrichmentConsumer.java index 603a7beb33f..b1f5440a080 100644 --- a/business-services/dashboard-ingest-v2/src/main/java/com/ingestpipeline/consumer/EnrichmentConsumer.java +++ b/business-services/dashboard-ingest-v2/src/main/java/com/ingestpipeline/consumer/EnrichmentConsumer.java @@ -41,7 +41,7 @@ public class EnrichmentConsumer implements KafkaConsumer { @Autowired private IESService elasticService; - @KafkaListener(id = INTENT, groupId = INTENT, topics = { Constants.KafkaTopics.TRANSFORMED_DATA}, containerFactory = Constants.BeanContainerFactory.INCOMING_KAFKA_LISTENER) + @KafkaListener(id = INTENT, groupId = "${kafka.consumer.group.enrichment:enrichment-v2}", topics = { Constants.KafkaTopics.TRANSFORMED_DATA}, containerFactory = Constants.BeanContainerFactory.INCOMING_KAFKA_LISTENER) public void processMessage(final Map incomingData, @Header(KafkaHeaders.RECEIVED_TOPIC) final String topic) { diff --git a/business-services/dashboard-ingest-v2/src/main/java/com/ingestpipeline/consumer/TransformConsumer.java b/business-services/dashboard-ingest-v2/src/main/java/com/ingestpipeline/consumer/TransformConsumer.java index 63e0ea5eb2c..c5d0f85de82 100644 --- a/business-services/dashboard-ingest-v2/src/main/java/com/ingestpipeline/consumer/TransformConsumer.java +++ b/business-services/dashboard-ingest-v2/src/main/java/com/ingestpipeline/consumer/TransformConsumer.java @@ -37,7 +37,7 @@ public class TransformConsumer implements KafkaConsumer { private ApplicationProperties applicationProperties; @Override - @KafkaListener(id = INTENT, groupId = INTENT, topics = { Constants.KafkaTopics.VALID_DATA }, containerFactory = Constants.BeanContainerFactory.INCOMING_KAFKA_LISTENER) + @KafkaListener(id = INTENT, groupId = "${kafka.consumer.group.transform:transform-v2}", topics = { Constants.KafkaTopics.VALID_DATA }, containerFactory = Constants.BeanContainerFactory.INCOMING_KAFKA_LISTENER) public void processMessage(Map incomingData, @Header(KafkaHeaders.RECEIVED_TOPIC) final String topic) { LOGGER.info("##KafkaMessageAlert## : key:" + topic + ":" + "value:" + incomingData.size()); diff --git a/business-services/dashboard-ingest-v2/src/main/resources/application.properties b/business-services/dashboard-ingest-v2/src/main/resources/application.properties index 10f41fc81ef..cd5332775fe 100644 --- a/business-services/dashboard-ingest-v2/src/main/resources/application.properties +++ b/business-services/dashboard-ingest-v2/src/main/resources/application.properties @@ -21,6 +21,8 @@ kafka.consumer.config.auto_commit_interval=100 kafka.consumer.config.session_timeout=15000 kafka.consumer.config.group_id=pipeline-group-v2 kafka.consumer.config.auto_offset_reset=earliest +kafka.consumer.group.transform=transform-v2 +kafka.consumer.group.enrichment=enrichment-v2 # KAFKA PRODUCER CONFIGURATIONS kafka.producer.config.retries_config=0