From 34f8f108b339f972d3ed187ee2c459a0e4d0ff30 Mon Sep 17 00:00:00 2001 From: ZLI263 Date: Fri, 19 Sep 2025 13:06:39 +0800 Subject: [PATCH] =?UTF-8?q?=E6=B7=BB=E5=8A=A0=E6=97=A5=E5=BF=97=EF=BC=8C?= =?UTF-8?q?=E4=BF=AE=E6=94=B9=E5=B9=B6=E5=8F=91=E9=80=BB=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../center/mq/NameplateKafkaConsumer.java | 123 +++++++++--------- 1 file changed, 65 insertions(+), 58 deletions(-) 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 83cb4cd..02661ab 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,11 +19,14 @@ import org.springframework.web.bind.annotation.RestController; import javax.annotation.Resource; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.Semaphore; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; /** * @ClassName NameplateKafkaConsumer - * @Description 消费tm_nameplate_corpus表binlog的Kafka消息 + * @Description 消费tm_nameplate_corpus表binlog的Kafka消息 * @Author renzhen * @Date 2025-03-04 10:18 * @Version 1.0 @@ -35,68 +38,72 @@ import java.util.concurrent.Semaphore; @RefreshScope public class NameplateKafkaConsumer { - @Autowired - private TmNameplateCorpusService tmNameplateCorpusService; + @Autowired + private TmNameplateCorpusService tmNameplateCorpusService; - @Autowired - @Resource(name = "threadPoolTaskExecutor") - private ThreadPoolTaskExecutor executor; - private Semaphore semaphore = new Semaphore(10); // 限制并发数 + @Autowired + @Resource(name = "threadPoolTaskExecutor") + private ThreadPoolTaskExecutor executor; + private Semaphore semaphore = new Semaphore(2); // 限制并发数 - @KafkaListener(topics = "${analyticCenterKafka.consumer.topic}", // = smart_assistant_nameplate_topic - groupId = "${analyticCenterKafka.consumer.group}" , //smart_assistant_nameplate_topic_group - containerFactory = "analyticCenterConsumerFactory", - concurrency = "3") - public void listen(String recordMessages, Acknowledgment ack, - @Header(KafkaHeaders.RECEIVED_PARTITION_ID) Integer partitionId, - @Header(KafkaHeaders.OFFSET) Long offset) { - long startTime = System.currentTimeMillis(); - log.info("nameplateKafkaConsumer 当前线程: {}, 线程ID: {},计数:{}", Thread.currentThread().getName(), Thread.currentThread().getId()); - log.info("nameplateKafkaConsumerMessage,消息:{}", recordMessages); - // 初始化绑定 Consumer - if (StringUtils.isEmpty(recordMessages)) { - ack.acknowledge(); - return; - } + @KafkaListener(topics = "${analyticCenterKafka.consumer.topic}", // = smart_assistant_nameplate_topic + groupId = "${analyticCenterKafka.consumer.group}", //smart_assistant_nameplate_topic_group + containerFactory = "analyticCenterConsumerFactory", + concurrency = "3") + public void listen(String recordMessages, Acknowledgment ack, + @Header(KafkaHeaders.RECEIVED_PARTITION_ID) Integer partitionId, + @Header(KafkaHeaders.OFFSET) Long offset) { - try { - NameplateTableKafkaDTO tmNameplateCorpus = JSON.parseObject(recordMessages, NameplateTableKafkaDTO.class); - if (!"INSERT".equals(tmNameplateCorpus.getType()) || CollectionUtils.isEmpty(tmNameplateCorpus.getData())) { - ack.acknowledge(); - } + if (StringUtils.isEmpty(recordMessages)) { + ack.acknowledge(); + return; + } - tmNameplateCorpus.getData().forEach(nameplate -> - CompletableFuture.runAsync(() -> { - try { - log.info(String.format("尝试获取许可,当前可用许可数: %d", - semaphore.availablePermits())); - if (semaphore.availablePermits() == 0) - log.info("获取许可 失败,❌ 任务被中断"); - semaphore.acquire(); // 获取许可 - log.info(String.format("获取许可✅ 成功后,当前可用许可数:: %d", - semaphore.availablePermits())); - tmNameplateCorpusService.processItem(nameplate); - } catch (Exception e) { - log.error("corpusPortrait画像铭牌异步任务执行失败", e); - } finally { - log.info("corpusPortrait 释放锁,availablePermits {}", semaphore.availablePermits()); - try { - ack.acknowledge(); // 提交 offset - } catch (IllegalStateException e) { - log.warn("Offset 已提交,跳过重复提交"); - } - semaphore.release(); // 释放许可 - log.error("释放许可,当前可用许可数:{}",semaphore.availablePermits()); - } - }, executor)); + try { + NameplateTableKafkaDTO tmNameplateCorpus = JSON.parseObject(recordMessages, NameplateTableKafkaDTO.class); + if (!"INSERT".equals(tmNameplateCorpus.getType()) || CollectionUtils.isEmpty(tmNameplateCorpus.getData())) { + ack.acknowledge(); + return; + } - log.info("nameplateKafkaConsumeracknowledge:{}",partitionId, offset); - // 手动提交 offset - } catch (Exception e) { - log.info("nameplateKafkaConsumerFailed to process message: {}", e); - } - log.info("nameplateKafkaConsumer消息处理开始,耗时:{}", System.currentTimeMillis() - startTime); - } + // 使用CountDownLatch等待所有异步任务完成 + CountDownLatch latch = new CountDownLatch(tmNameplateCorpus.getData().size()); + AtomicBoolean hasError = new AtomicBoolean(false); + + tmNameplateCorpus.getData().forEach(nameplate -> + CompletableFuture.runAsync(() -> { + try { + if (!semaphore.tryAcquire(5, TimeUnit.SECONDS)) { + log.warn("获取许可超时,跳过处理"); + hasError.set(true); + return; + } + + tmNameplateCorpusService.processItem(nameplate); + } catch (Exception e) { + log.error("处理铭牌数据失败", e); + hasError.set(true); + } finally { + semaphore.release(); + latch.countDown(); + } + }, executor)); + + // 等待所有任务完成 + if (latch.await(30, TimeUnit.SECONDS)) { + if (!hasError.get()) { + ack.acknowledge(); // 只有所有任务成功才提交offset + } else { + log.error("部分任务失败,不提交offset"); + } + } else { + log.error("任务执行超时,不提交offset"); + } + } catch (Exception e) { + log.error("处理消息失败", e); + } + + } }