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 571caf7..1868e45 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 @@ -63,7 +63,6 @@ public class CorpusProcessKafkaProducer { @KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}") public void listen( @Payload List messageList, - @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) List keys, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) List partitions, @Header(KafkaHeaders.OFFSET) List offsets ) { @@ -77,11 +76,10 @@ public class CorpusProcessKafkaProducer { for (int i = 0; i < messageList.size(); i++) { String message = messageList.get(i); log.debug("corpusProcessKafkaProducer message:{}", message); - String key = keys.get(i); int partition = partitions.get(i); long offset = offsets.get(i); - log.debug("corpusProcessKafkaProducer message:key:{}, topic={}, partition={}, offset={}", "your-topic", key,partition, offset); + log.debug("corpusProcessKafkaProducer message: topic={}, partition={}, offset={}", "your-topic",partition, offset); // 提交异步任务 CompletableFuture.runAsync(() -> processSingleRecord(message), executor);