From d7965e45b9df56109d6ec732c53892cf73d9bacc Mon Sep 17 00:00:00 2001 From: zren25 Date: Wed, 28 May 2025 15:32:57 +0800 Subject: [PATCH] =?UTF-8?q?=E5=88=A0=E9=99=A4=E6=97=A0=E7=94=A8=E4=BB=A3?= =?UTF-8?q?=E7=A0=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../center/job/NameplateCorpusJob.java | 23 ------- .../service/TmNameplateCorpusService.java | 2 - .../impl/TmNameplateCorpusServiceImpl.java | 68 ++----------------- .../TmOdsVdqwMessagearchivingServiceImpl.java | 20 ++---- 4 files changed, 10 insertions(+), 103 deletions(-) diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/NameplateCorpusJob.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/NameplateCorpusJob.java index 4f0955d..7c8d14d 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/NameplateCorpusJob.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/NameplateCorpusJob.java @@ -36,29 +36,6 @@ public class NameplateCorpusJob { private TmTelephoneCorpusMapper tmTelephoneCorpusMapper; @Value("${dify.corpus.nameplate.isUse}") private boolean isUse; - /** - * 铭牌语料处理 - */ - @XxlJob("nameplateCorpusTask") - @PostMapping("nameplateCorpusTask") - public ResultMsg nameplateCorpusTask(@RequestBody String paramJson) { - try { - // 获取任务参数 - String param = XxlJobHelper.getJobParam(); - if(StringUtils.isEmpty(param)){ - param = paramJson; - } - // 执行业务逻辑 - XxlJobHelper.log("任务参数: {}", param); - log.info("nameplateCorpusTask 铭牌语料查询处理:{}",param); - tmNameplateCorpusService.runNameplateCorpusDify(param); - } catch (Exception e) { - log.error("nameplateCorpusTask 定时任务补偿处理消息异常",e.getMessage()); - } - return ResultMsg.ok(); - } - - /** * 铭牌语料处理失败重试 * @param paramJson diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/TmNameplateCorpusService.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/TmNameplateCorpusService.java index 2bc25bc..2bd98cb 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/TmNameplateCorpusService.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/TmNameplateCorpusService.java @@ -13,8 +13,6 @@ import java.util.List; */ public interface TmNameplateCorpusService extends IService { - void runNameplateCorpusDify(String paramJson); - void runNameplateCorpusDifyRetry(String paramJson); void processItem(TmNameplateCorpus item); 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 f017646..d3a8504 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 @@ -25,14 +25,13 @@ import org.apache.rocketmq.spring.core.RocketMQTemplate; 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; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; /** @@ -77,58 +76,10 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl messageList = tmNameplateCorpusMapper.queryTmNameplateCorpusByData(statTime, endTime, offset, pageSize); - // 处理查询到的数据 - // 使用 CompletableFuture 并行处理 - CompletableFuture[] futures = messageList.stream() - .map(item -> CompletableFuture.runAsync(() -> { - try { - processItem(item); - } catch (Exception e) { - log.error("铭牌语料失败: customerFlowId={}, AcceptUserId={}, 异常: {}", - item.getCustomerFlowId(), e.getMessage(), e); - } - }, executor)) - .toArray(CompletableFuture[]::new); - // 等待所有任务完成 - CompletableFuture.allOf(futures).join(); - } - // 关闭线程池 - executor.shutdown(); - log.info("企微数据跑批结束 耗时:{}",System.currentTimeMillis()-startTime); - } @Override public void runNameplateCorpusDifyRetry(String paramJson) { long startTime = System.currentTimeMillis(); @@ -156,13 +107,6 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl messageList = tmNameplateCorpusMapper.queryTmNameplateCorpusRetry(statTime, endTime, offset, pageSize, customerFlowIds, retry); @@ -178,11 +122,7 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl