From 1ce61d4bfc679bb35d576b82ea766efd616044c3 Mon Sep 17 00:00:00 2001 From: zren25 Date: Mon, 31 Mar 2025 14:04:42 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E6=94=B9kafka=E7=9B=91=E5=90=AC?= =?UTF-8?q?=E5=AF=B9=E8=B1=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../center/mq/CorpusProcessKafkaProducer.java | 15 ++++++--------- 1 file changed, 6 insertions(+), 9 deletions(-) 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 3cc40ec..c530bc6 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 @@ -43,22 +43,19 @@ public class CorpusProcessKafkaProducer { @PostMapping("corpusProcessKafkaConsumer") @KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}") - public void listen(List> recordMessage) { + public void listen(List recordMessages) { try { - log.info("CorpusProcessKafkaProducer Received message: {}", recordMessage); + log.info("CorpusProcessKafkaProducer Received message: {}", recordMessages); // 获取消息列表 int optimalThreadPoolSize = Runtime.getRuntime().availableProcessors() + 2; log.info("获取的线程数:{}", optimalThreadPoolSize); ExecutorService executor = Executors.newFixedThreadPool(optimalThreadPoolSize); - if(CollectionUtils.isNotEmpty(recordMessage)){ - log.info("CorpusProcessKafkaProducer List size: {}", recordMessage.size()); - for (ConsumerRecord record : recordMessage) { - log.info("CorpusProcessKafkaProducer List record.value: {}",record.value()); - log.info("CorpusProcessKafkaProducer List record.key: {}",record.key()); - + if(CollectionUtils.isNotEmpty(recordMessages)){ + log.info("CorpusProcessKafkaProducer List size: {}", recordMessages.size()); + for (String message : recordMessages) { + log.info("CorpusProcessKafkaProducer message: {}",message); executor.submit(() -> { try { - String message = (String) record.value(); AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(message, AicorpusTelephoneDTO.class); log.info("aicorpusTelephone categoryCode:{}, display: {}", aicorpusTelephone.getCategoryCode(), aicorpusTelephone.getDisplay()); DisplayDTO display = objectMapper.readValue(aicorpusTelephone.getDisplay(), DisplayDTO.class);