diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusProcessKafkaProducer.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusProcessKafkaProducer.java index c4018f4..213d3ab 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusProcessKafkaProducer.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusProcessKafkaProducer.java @@ -61,7 +61,7 @@ public class CorpusProcessKafkaProducer { @KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}") - public void listen(@Payload List messageList, + public void listen(@Payload List messageList, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) List partitions, @Header(KafkaHeaders.OFFSET) List offsets ) { @@ -73,7 +73,7 @@ public class CorpusProcessKafkaProducer { } try { for (int i = 0; i < messageList.size(); i++) { - String message =String.valueOf( messageList.get(i)); + String message = messageList.get(i); log.debug("corpusProcessKafkaProducer message:{}", message); int partition = partitions.get(i); long offset = offsets.get(i);