diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/NameplateKafkaConsumer.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/NameplateKafkaConsumer.java index b9c6261..9f50859 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/NameplateKafkaConsumer.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/NameplateKafkaConsumer.java @@ -19,7 +19,9 @@ import org.springframework.stereotype.Component; import org.springframework.web.bind.annotation.RestController; import javax.annotation.Resource; -import java.util.concurrent.*; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; /** * @ClassName NameplateKafkaConsumer @@ -49,10 +51,11 @@ public class NameplateKafkaConsumer { private int messageCount = 0; @KafkaListener(topics = "${analyticCenterKafka.consumer.topic}", groupId = "${analyticCenterKafka.consumer.group}" , - containerFactory = "analyticCenterConsumerFactory") + containerFactory = "analyticCenterConsumerFactory", + concurrency = "3") public void listen(String recordMessages, Acknowledgment ack, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) Integer partitionId, - @Header(KafkaHeaders.OFFSET) Long offset) throws InterruptedException { + @Header(KafkaHeaders.OFFSET) Long offset) { long startTime = System.currentTimeMillis(); messageCount++; @@ -67,7 +70,6 @@ public class NameplateKafkaConsumer { NameplateTableKafkaDTO tmNameplateCorpus = JSON.parseObject(recordMessages, NameplateTableKafkaDTO.class); log.info("nameplateKafkaConsumerParseType: {},size:{}", tmNameplateCorpus.getType(), tmNameplateCorpus.getData().size()); if(tmNameplateCorpus.getType().equals("INSERT")){ - CountDownLatch latch = new CountDownLatch(tmNameplateCorpus.getData().size()); tmNameplateCorpus.getData().forEach(nameplate -> CompletableFuture.runAsync(() -> { try { @@ -76,18 +78,14 @@ public class NameplateKafkaConsumer { } catch (Exception e) { log.error("nameplateKafkaConsumer铭牌解析处理出错 {}:{},ex:{}", partitionId, offset, e.getMessage()); - }finally { - log.info("nameplateKafkaConsumerack.acknowledge"); - latch.countDown(); } }, executor)); try { + log.info("nameplateKafkaConsumerack.acknowledge:{}"); ack.acknowledge(); // 确保无论成功与否都提交 offset } catch (IllegalStateException e) { log.warn("Offset 已提交,跳过重复提交"); } - boolean allDone = latch.await(2, TimeUnit.MINUTES); - log.info("nameplateKafkaConsumerack.acknowledge:{}",allDone); } }else { ack.acknowledge(); // 空消息直接跳过