From 2586bb5de1132da19f9846b9c5495fa769ff24e0 Mon Sep 17 00:00:00 2001 From: ZLI263 Date: Fri, 19 Sep 2025 10:59:32 +0800 Subject: [PATCH] =?UTF-8?q?=E9=93=AD=E7=89=8C=E7=9A=84=E5=A4=9A=E7=BA=BF?= =?UTF-8?q?=E7=A8=8B=E6=8E=A7=E5=88=B6=EF=BC=8C=20=E4=BF=AE=E6=94=B9?= =?UTF-8?q?=E4=BC=81=E5=BE=AE=E7=94=BB=E5=83=8F=E7=9A=84tag?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../center/config/ExecutorConfig.java | 14 +++++++------- .../center/mq/NameplateKafkaConsumer.java | 10 ++++++++-- .../impl/TmNameplateCorpusServiceImpl.java | 19 +++++++++++++------ .../TmOdsVdqwMessagearchivingServiceImpl.java | 7 +++---- 4 files changed, 31 insertions(+), 19 deletions(-) diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/config/ExecutorConfig.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/config/ExecutorConfig.java index c4eda2d..1cbf197 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/config/ExecutorConfig.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/config/ExecutorConfig.java @@ -20,14 +20,14 @@ public class ExecutorConfig { @Value("${task.pool.corePoolSize}") private int corePoolSize; - @Value("${task.pool.maxPoolSize}") - private int maxPoolSize; +// @Value("${task.pool.maxPoolSize}") + private int maxPoolSize=12; - @Value("${task.pool.keepAliveSeconds}") - private int keepAliveSeconds; +// @Value("${task.pool.keepAliveSeconds}") + private int keepAliveSeconds=12000; - @Value("${task.pool.queueCapacity}") - private int queueCapacity; +// @Value("${task.pool.queueCapacity}") + private int queueCapacity=20; private ThreadPoolTaskExecutor executor; private ThreadPoolExecutor blockingExecutor; @@ -38,7 +38,7 @@ public class ExecutorConfig { executor = new ThreadPoolTaskExecutor(); log.info("Thread-process-pool-Initializing: corePoolSize: {}, maxPoolSize: {}, queueCapacity: {}, keepAliveSeconds: {}",corePoolSize, maxPoolSize, queueCapacity, keepAliveSeconds); - // 设置核心线程数 + // 设置核心线程数 Thread-process-pool-Initializing: corePoolSize: 10, maxPoolSize: 300, queueCapacity: 1000, keepAliveSeconds: 60 executor.setCorePoolSize(corePoolSize); // 设置最大线程数为5 executor.setMaxPoolSize(maxPoolSize); 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 d267792..9156f7b 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,6 +19,7 @@ import org.springframework.web.bind.annotation.RestController; import javax.annotation.Resource; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.Semaphore; /** * @ClassName NameplateKafkaConsumer @@ -40,7 +41,7 @@ public class NameplateKafkaConsumer { @Autowired @Resource(name = "threadPoolTaskExecutor") private ThreadPoolTaskExecutor executor; - + private Semaphore semaphore = new Semaphore(10); // 限制并发数 @KafkaListener(topics = "${analyticCenterKafka.consumer.topic}", // = smart_assistant_nameplate_topic groupId = "${analyticCenterKafka.consumer.group}" , //smart_assistant_nameplate_topic_group @@ -51,12 +52,13 @@ public class NameplateKafkaConsumer { @Header(KafkaHeaders.OFFSET) Long offset) { long startTime = System.currentTimeMillis(); log.info("nameplateKafkaConsumer 当前线程: {}, 线程ID: {},计数:{}", Thread.currentThread().getName(), Thread.currentThread().getId()); - log.info("nameplateKafkaConsumerMessage总数:{},消息:{}", recordMessages); + log.info("nameplateKafkaConsumerMessage,消息:{}", recordMessages); // 初始化绑定 Consumer if (StringUtils.isEmpty(recordMessages)) { ack.acknowledge(); return; } + try { NameplateTableKafkaDTO tmNameplateCorpus = JSON.parseObject(recordMessages, NameplateTableKafkaDTO.class); if (!"INSERT".equals(tmNameplateCorpus.getType()) || CollectionUtils.isEmpty(tmNameplateCorpus.getData())) { @@ -66,9 +68,13 @@ public class NameplateKafkaConsumer { tmNameplateCorpus.getData().forEach(nameplate -> CompletableFuture.runAsync(() -> { try { + semaphore.acquire(); // 获取许可 tmNameplateCorpusService.processItem(nameplate); } catch (Exception e) { log.error("corpusPortrait画像铭牌异步任务执行失败", e); + } finally { + log.info("corpusPortrait 释放锁,availablePermits {}", semaphore.availablePermits()); + semaphore.release(); // 释放许可 } }, executor)); diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmNameplateCorpusServiceImpl.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmNameplateCorpusServiceImpl.java index a769038..fa11d00 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmNameplateCorpusServiceImpl.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmNameplateCorpusServiceImpl.java @@ -22,8 +22,10 @@ import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.cloud.context.config.annotation.RefreshScope; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import org.springframework.stereotype.Service; +import javax.annotation.Resource; import java.time.LocalDate; import java.util.*; import java.util.concurrent.CompletableFuture; @@ -73,7 +75,9 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl { + CompletableFuture summaryTask = CompletableFuture.runAsync(() -> { //第一个业务场景 开始: 总结和分类业务场景 - JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT_NAMEPLATE.getCode(), + JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT_NAMEPLATE.getCode(), JSONObject.toJSONString(corpusReportDTO), null); log.info("runDify execDifyFlow ,铭牌语料 ,总结和分类业务,返回: {}", execDifyFlow); - }, executorService); + }, executor); long endTime = System.currentTimeMillis(); log.info("第一个业务场景(总结和分类)执行时间: {} ms", (endTime - startTime)); //第一个业务场景, 结束 long startTime2 = System.currentTimeMillis(); - CompletableFuture.runAsync(() -> { + CompletableFuture portraitTask = CompletableFuture.runAsync(() -> { //"analysisScene": "分析类型", // 1、企微会话 2、AI通话录音 3、AI铭牌(客流) 4、AI铭牌(试驾) DiFyReq diFyImageReq2 = new DiFyReq(); diFyImageReq2.setUser(ConstantStr.corpus_user); @@ -199,9 +203,12 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl inputMap2 = new HashMap<>(); - inputMap2.put("businessId", unionId); + inputMap2.put("businessId", unionId+","+corpusReportDTO.getUserId()); + inputMap2.put("consultantId", corpusReportDTO.getUserId());//顾问id inputMap2.put("communicateDate", DateUtil.format(maxMsgTimeItem.getMsgTime(), DatePattern.NORM_DATETIME_PATTERN)); inputMap2.put("analysisScene", "1"); inputMap2.put("version", 2);