diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/QiWeiCorpusJob.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/QiWeiCorpusJob.java index b70c2a5..3b69246 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/QiWeiCorpusJob.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/QiWeiCorpusJob.java @@ -25,7 +25,6 @@ public class QiWeiCorpusJob { * 企微语料处理 */ @XxlJob("QiWeiCorpusTask") - @PostMapping("QiWeiCorpusTask") public ResultMsg qiWeiCorpusTask(@RequestBody String paramJson) { try { // 获取任务参数 @@ -43,4 +42,23 @@ public class QiWeiCorpusJob { } return ResultMsg.ok(); } + + @PostMapping("qiWeiCorpusTaskMock") + public ResultMsg qiWeiCorpusTaskMock(@RequestBody String paramJson) { + try { + // 获取任务参数 + String param = XxlJobHelper.getJobParam(); + if(StringUtils.isEmpty(param)){ + param = paramJson; + } + // 执行业务逻辑 + XxlJobHelper.log("任务参数: {}", param); + log.info("qiWeiCorpusTask 企微语料查询处理:{}",param); + tmOdsVdqwMessagearchivingService.runQiWeiCorpusDifyMock(param); + } catch (Exception e) { + log.error("processMessageByTask 定时任务补偿处理消息异常",e.getMessage()); + throw new RuntimeException(e); + } + return ResultMsg.ok(); + } } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/TmOdsVdqwMessagearchivingService.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/TmOdsVdqwMessagearchivingService.java index 748f88b..bab4a0e 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/TmOdsVdqwMessagearchivingService.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/TmOdsVdqwMessagearchivingService.java @@ -11,4 +11,7 @@ public interface TmOdsVdqwMessagearchivingService extends IService implements AiAnalysisRequestLogsService { @@ -22,8 +24,10 @@ public class AiAnalysisRequestLogsServiceImpl extends ServiceImpl 0; } else { + aiAnalysisRequestLogs.setUpdateTime(new Date()); return aiAnalysisRequestLogsMapper.update(aiAnalysisRequestLogs, queryWrapper) > 0; } } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmOdsVdqwMessagearchivingServiceImpl.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmOdsVdqwMessagearchivingServiceImpl.java index 8280860..c30ad3a 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmOdsVdqwMessagearchivingServiceImpl.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmOdsVdqwMessagearchivingServiceImpl.java @@ -33,10 +33,7 @@ import org.springframework.stereotype.Service; import javax.annotation.Resource; import java.time.LocalDate; -import java.util.Arrays; -import java.util.HashMap; -import java.util.List; -import java.util.Map; +import java.util.*; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -125,7 +122,6 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl messageList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByData(statTime, endTime, offset, pageSize, retry); @@ -259,4 +255,70 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl mockMessageList = new ArrayList<>(); + List messageList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByData(statTime, endTime, 0, pageSize, retry); + for(int i = 1; i <= 4000; i++){ + if(mockMessageList.size()<2000){ + mockMessageList.addAll(messageList); + } + } + log.info("runQiWeiCorpusDifyMock 数据量:",mockMessageList.size()); + // mock + for (int i = 1; i <= 2000; i++) { + log.info("runQiWeiCorpusDifyMock 数据量:{},第:{} 批次",mockMessageList.size(),i); + // 处理查询到的数据 + // 使用 CompletableFuture 并行处理 + final int[] count = {1}; + CompletableFuture[] futures = mockMessageList.stream() + .map(item -> CompletableFuture.runAsync(() -> { + log.info("runQiWeiCorpusDifyMock 异步执行{},第:{} 批次", count[0],JSONObject.toJSONString(item)); + try { + processItem(item, finalStatTime, finalEndTime); + } catch (Exception e) { + log.error("处理企微语料失败: FromUserId={}, AcceptUserId={}, 异常: {}", + item.getFromUserId(), item.getAcceptUserId(), e.getMessage(), e); + } + count[0] = count[0] +1; + }, executor)) + .toArray(CompletableFuture[]::new); + // 等待所有任务完成 + CompletableFuture.allOf(futures).join(); + } + // 关闭线程池 + executor.shutdown(); + + } } \ No newline at end of file