From c74f26798fde238f7ba1c2028d21b7f35cab3489 Mon Sep 17 00:00:00 2001 From: zren25 Date: Tue, 11 Mar 2025 11:32:38 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=81=E5=BE=AE=E8=AF=AD=E6=96=99ai=E8=A7=A3?= =?UTF-8?q?=E6=9E=90=E5=BC=80=E5=8F=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../volvo/ai/analytic/center/dto/PageDto.java | 22 ++ .../analytic/center/job/QiWeiCorpusJob.java | 37 ++++ .../TmOdsVdqwMessagearchivingMapper.java | 4 +- .../TmOdsVdqwMessagearchivingService.java | 2 +- .../TmOdsVdqwMessagearchivingServiceImpl.java | 200 +++++++++++------- .../ai/analytic/center/utils/ConstantStr.java | 14 ++ .../TmOdsVdqwMessagearchivingMapper.xml | 21 +- .../mapper/TmTelephoneCorpusMapper.xml | 5 + 8 files changed, 227 insertions(+), 78 deletions(-) create mode 100644 ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/PageDto.java create mode 100644 ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/QiWeiCorpusJob.java create mode 100644 ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/utils/ConstantStr.java create mode 100644 ai-analytic-center-biz/src/main/resources/mapper/TmTelephoneCorpusMapper.xml diff --git a/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/PageDto.java b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/PageDto.java new file mode 100644 index 0000000..b7f8cf9 --- /dev/null +++ b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/PageDto.java @@ -0,0 +1,22 @@ +package com.volvo.ai.analytic.center.dto; +/** + * @ClassName PageDto + * @Description + * @Author renzhen + * @Date 2025-03-11 10:37 + * @Version 1.0 + **/ + +import lombok.Data; +import org.springframework.cloud.context.config.annotation.RefreshScope; + +@Data +public class PageDto { + + // 每页大小 + private static int currentPage = 1; // 当前页码 + public static int getTotalPages(int totalCount,int pageSize) { + return (int) Math.ceil((double) totalCount / pageSize); + } + +} 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 new file mode 100644 index 0000000..c84283c --- /dev/null +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/QiWeiCorpusJob.java @@ -0,0 +1,37 @@ +package com.volvo.ai.analytic.center.job; + +import com.volvo.ai.analytic.center.service.TmOdsVdqwMessagearchivingService; +import com.volvo.common.core.util.ResultMsg; +import com.xxl.job.core.handler.annotation.XxlJob; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Component; +import org.springframework.web.bind.annotation.PostMapping; +import org.springframework.web.bind.annotation.RequestBody; +import org.springframework.web.bind.annotation.RestController; + + +@Slf4j +@Component +@RestController +public class QiWeiCorpusJob { + + @Autowired + private TmOdsVdqwMessagearchivingService tmOdsVdqwMessagearchivingService; + + /** + * 企微语料处理 + */ + @XxlJob("QiWeiCorpusTask") + @PostMapping("qiWeiCorpusTask") + public ResultMsg qiWeiCorpusTask(@RequestBody String paramJson) { + try { + log.info("qiWeiCorpusTask 企微语料查询处理:{}",paramJson); + tmOdsVdqwMessagearchivingService.runQiWeiCorpusDify(paramJson); + } 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/mapper/TmOdsVdqwMessagearchivingMapper.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TmOdsVdqwMessagearchivingMapper.java index 7ef6384..30a6b6b 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TmOdsVdqwMessagearchivingMapper.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TmOdsVdqwMessagearchivingMapper.java @@ -5,6 +5,7 @@ import com.volvo.ai.analytic.center.dto.corpus.OdsVdqwMessageOTD; import com.volvo.ai.analytic.center.entity.TmOdsVdqwMessagearchiving; import org.apache.ibatis.annotations.Mapper; import org.apache.ibatis.annotations.Param; +import org.springframework.stereotype.Repository; import java.util.List; @@ -16,7 +17,8 @@ import java.util.List; @Mapper public interface TmOdsVdqwMessagearchivingMapper extends BaseMapper { - List queryOdsVdqwMessageByData(@Param("statTime") String statTime, @Param("endTime") String endTime); + int countOdsVdqwMessageByData(@Param("statTime") String statTime, @Param("endTime") String endTime); + List queryOdsVdqwMessageByData(@Param("statTime") String statTime, @Param("endTime") String endTime, @Param("offset") int offset, @Param("pageSize") int pageSize); List queryOdsVdqwMessageByFromUserIdAndAcceptUserId(@Param("statTime") String statTime, @Param("endTime") String endTime, @Param("userIdList") List userIdList); 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 a69ad81..3c84187 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,6 +11,6 @@ import java.util.*; */ public interface TmOdsVdqwMessagearchivingService extends IService { - void runQiWeiCorpusDify(); + void runQiWeiCorpusDify(String paramJson); } \ No newline at end of file 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 96e4cc8..38d1a07 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 @@ -7,6 +7,7 @@ import com.alibaba.fastjson.JSONObject; import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl; import com.volvo.ai.analytic.center.constant.Constant; +import com.volvo.ai.analytic.center.dto.PageDto; import com.volvo.ai.analytic.center.dto.corpus.OdsVdqwMessageOTD; import com.volvo.ai.analytic.center.dto.req.DiFyReq; import com.volvo.ai.analytic.center.dto.req.RunMaskingRuleInput; @@ -20,6 +21,7 @@ import com.volvo.ai.analytic.center.service.DataMaskingRuleService; import com.volvo.ai.analytic.center.service.DiFyService; import com.volvo.ai.analytic.center.service.TmOdsVdqwMessagearchivingService; import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService; +import com.volvo.ai.analytic.center.utils.ConstantStr; import com.volvo.ai.analytic.center.utils.FlowResultSplitUtil; import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections.CollectionUtils; @@ -29,14 +31,19 @@ import org.apache.rocketmq.client.producer.SendResult; 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.messaging.support.MessageBuilder; 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.concurrent.CompletableFuture; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; /** @@ -44,6 +51,7 @@ import java.util.Map; * @author rz * @date 2025-03-04 */ +@RefreshScope @Slf4j @Service public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl implements TmOdsVdqwMessagearchivingService { @@ -62,9 +70,13 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl messageList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByData(statTime,endTime); - messageList.stream().forEach(item->{ - log.info("企微语料内容:FromUserId:{}, AcceptUserId:{}",item.getFromUserId(), item.getAcceptUserId()); - // 1,vdqw_workuserinfo 这个表对应是 B端认证中心userId - // 2,vdqw_externalcontact 这个表对应是 企微客户unionId - String unionId = getUnionId(Arrays.asList(item.getFromUserId(), item.getAcceptUserId())); - String userId = getUserId(Arrays.asList(item.getFromUserId(), item.getAcceptUserId())); - log.info("企微查询信息为空:unionId:{}, userId:{}", unionId,userId); - if(StringUtils.isEmpty(unionId) || StringUtils.isNotBlank(userId)){ - log.info("企微查询信息为空 "); - return; + public void runQiWeiCorpusDify(String paramJson) { + log.info("runQiWeiCorpusDify paramJson {}", paramJson); + // 获取当前日期 + LocalDate today = LocalDate.now(); + // 获取前一天日期 + LocalDate yesterday = today.minusDays(1); + // 格式化输出 + String formattedDate = yesterday.toString(); // 默认格式为 yyyy-MM-dd + String statTime = formattedDate.concat(" 00:00:00"); + String endTime = formattedDate.concat(" 23:59:59"); + if (StringUtils.isNotBlank(paramJson)) { + JSONObject paramJsonObj = JSONObject.parseObject(paramJson); + if (null != paramJsonObj && paramJsonObj.containsKey("statTime") && paramJsonObj.containsKey("endTime")) { + statTime = paramJsonObj.getString("statTime"); + endTime = paramJsonObj.getString("endTime"); } + } + String finalStatTime = statTime; + String finalEndTime = endTime; - if(StringUtils.isNotBlank(item.getFromUserId()) && StringUtils.isNotBlank(item.getAcceptUserId())){ + Integer total = tmOdsVdqwMessagearchivingMapper.countOdsVdqwMessageByData(statTime, endTime); + int totalPages = PageDto.getTotalPages(total, pageSize); + + // 获取消息列表 + int optimalThreadPoolSize = Runtime.getRuntime().availableProcessors() + 1; + log.info("获取的线程数:{}",optimalThreadPoolSize); + // 创建线程池 + ExecutorService executor = Executors.newFixedThreadPool(optimalThreadPoolSize); // 根据需求调整线程池大小 - List contetnList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByFromUserIdAndAcceptUserId(statTime,endTime, Arrays.asList(item.getFromUserId(), item.getAcceptUserId())); - OdsVdqwMessageOTD maxMsgTimeItem = contetnList.stream() - .max((o1, o2) -> o1.getMsgTime().compareTo(o2.getMsgTime())) - .orElse(null); - contetnList.stream().forEach(contentItem->{ + for (int i = 1; i <= totalPages; i++) { + int offset = (i - 1) * pageSize; + List messageList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByData(statTime, endTime, offset, pageSize); + // 处理查询到的数据 + // 使用 CompletableFuture 并行处理 + CompletableFuture[] futures = messageList.stream() + .map(item -> CompletableFuture.runAsync(() -> { + try { + processItem(item, finalStatTime, finalEndTime); + } catch (Exception e) { + log.error("处理消息失败: FromUserId={}, AcceptUserId={}, 异常: {}", + item.getFromUserId(), item.getAcceptUserId(), e.getMessage(), e); + } + }, executor)) + .toArray(CompletableFuture[]::new); + // 等待所有任务完成 + CompletableFuture.allOf(futures).join(); + } + // 关闭线程池 + executor.shutdown(); - Map inputMap = new HashMap(); - DiFyReq diFyImageReq = new DiFyReq(); - diFyImageReq.setUser("11111"); - diFyImageReq.setFlowId("app-peJXSjHjVKdkYxjdOUuPnZ5b"); + } - JSONObject contentJson = JSONObject.parseObject(contentItem.getContent()); - String content = contentJson.getString("content"); - List maskingRuleItems = dataMaskingRuleService.getDataMaskingRuleListByApplicationChannel(Constant.CHANNEL_DCC); + private void processItem(OdsVdqwMessageOTD item, String statTime, String endTime) { + log.info("企微语料内容:FromUserId:{}, AcceptUserId:{}", item.getFromUserId(), item.getAcceptUserId()); + // 1,vdqw_workuserinfo 这个表对应是 B端认证中心userId + // 2,vdqw_externalcontact 这个表对应是 企微客户unionId + String unionId = getUnionId(Arrays.asList(item.getFromUserId(), item.getAcceptUserId())); + String userId = getUserId(Arrays.asList(item.getFromUserId(), item.getAcceptUserId())); + log.info("企微查询信息:unionId:{}, userId:{}", unionId, userId); + if (StringUtils.isEmpty(unionId) || StringUtils.isNotBlank(userId)) { + log.info("企微查询信息为空 "); + return; + } - RunMaskingRuleInput runMaskingRuleInput = new RunMaskingRuleInput(); - runMaskingRuleInput.setDataMaskingRules(maskingRuleItems); - // 拼接 role 和 text - String chat = contentItem.getFromUserId().concat(":").concat(content); - runMaskingRuleInput.setOldStr(chat); - String corpusChat = dataMaskingRuleService.runMaskingRule(runMaskingRuleInput); + if (StringUtils.isNotBlank(item.getFromUserId()) && StringUtils.isNotBlank(item.getAcceptUserId())) { + List contetnList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByFromUserIdAndAcceptUserId(statTime, endTime, Arrays.asList(item.getFromUserId(), item.getAcceptUserId())); + OdsVdqwMessageOTD maxMsgTimeItem = contetnList.stream() + .max((o1, o2) -> o1.getMsgTime().compareTo(o2.getMsgTime())) + .orElse(null); + contetnList.forEach(contentItem -> { + Map inputMap = new HashMap<>(); + DiFyReq diFyImageReq = new DiFyReq(); + diFyImageReq.setUser(ConstantStr.corpus_user); + diFyImageReq.setFlowId(qiweiToken); - inputMap.put("chat",corpusChat); + JSONObject contentJson = JSONObject.parseObject(contentItem.getContent()); + String content = contentJson.getString("content"); + List maskingRuleItems = dataMaskingRuleService.getDataMaskingRuleListByApplicationChannel(Constant.CHANNEL_DCC); - inputMap.put("model",tmTelephoneCorpusService.getCarModelList()); - diFyImageReq.setInputs(inputMap); + RunMaskingRuleInput runMaskingRuleInput = new RunMaskingRuleInput(); + runMaskingRuleInput.setDataMaskingRules(maskingRuleItems); + // 拼接 role 和 text + String chat = contentItem.getFromUserId().concat(":").concat(content); + runMaskingRuleInput.setOldStr(chat); + String corpusChat = dataMaskingRuleService.runMaskingRule(runMaskingRuleInput); - // 获取配置 - JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode()); - log.info("runDify execDifyFlow {}",execDifyFlow); - if(null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")){ + inputMap.put("chat", corpusChat); - String text = execDifyFlow.getJSONObject("outputs").getString("text"); - String resultStrOne = FlowResultSplitUtil.flowOutputTextSplit(text, "任务1", "任务2"); - String resultStrTwo =FlowResultSplitUtil.flowOutputTextSplit(text, "任务2", null); - Map ltoMap = new HashMap(); - ltoMap.put("analysisRecordId", execDifyFlow.getString("aiAnalysisRequestId")); - ltoMap.put("analysisScene", "1"); - ltoMap.put("unionId", unionId); - ltoMap.put("consultantId", userId); - ltoMap.put("communicateDate", DateUtil.format(maxMsgTimeItem.getMsgTime(), DatePattern.NORM_DATETIME_PATTERN)); - ltoMap.put("analysisResult", resultStrOne); - ltoMap.put("analysisDetail", resultStrTwo); - // 发送MQ + inputMap.put("model", tmTelephoneCorpusService.getCarModelList()); + diFyImageReq.setInputs(inputMap); - log.info("send mq {}",ltoMap); - rocketMqTemplate.asyncSend(topic, MessageBuilder.withPayload(JSON.toJSONString(ltoMap)).build(), - new SendCallback() { - @Override - public void onSuccess(SendResult sendResult) { - log.info("企微语料发送MQ成功 消息体:{}", JSON.toJSONString(ltoMap)); - } - @Override - public void onException(Throwable e) { - log.error("企微语料发送MQ异常 消息体:{}, 异常:", JSON.toJSONString(ltoMap), e); - } - }, 10000); - - } - - }); - } - - - - }); + // 获取配置 + JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode()); + log.info("runDify execDifyFlow {}", execDifyFlow); + if (null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")) { + String text = execDifyFlow.getJSONObject("outputs").getString("text"); + String resultStrOne = FlowResultSplitUtil.flowOutputTextSplit(text, "任务1", "任务2"); + String resultStrTwo = FlowResultSplitUtil.flowOutputTextSplit(text, "任务2", null); + Map ltoMap = new HashMap<>(); + ltoMap.put("analysisRecordId", execDifyFlow.getString("aiAnalysisRequestId")); + ltoMap.put("analysisScene", "1"); + ltoMap.put("unionId", unionId); + ltoMap.put("consultantId", userId); + ltoMap.put("communicateDate", DateUtil.format(maxMsgTimeItem.getMsgTime(), DatePattern.NORM_DATETIME_PATTERN)); + ltoMap.put("analysisResult", resultStrOne); + ltoMap.put("analysisDetail", resultStrTwo); + // 发送MQ + log.info("send mq {}", ltoMap); + rocketMqTemplate.asyncSend(topic, MessageBuilder.withPayload(JSON.toJSONString(ltoMap)).build(), + new SendCallback() { + @Override + public void onSuccess(SendResult sendResult) { + log.info("企微语料发送MQ成功 消息体:{}", JSON.toJSONString(ltoMap)); + } + @Override + public void onException(Throwable e) { + log.error("企微语料发送MQ异常 消息体:{}, 异常:", JSON.toJSONString(ltoMap), e); + } + }, 10000); + } + }); + } } private String getUnionId(List userIds){ diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/utils/ConstantStr.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/utils/ConstantStr.java new file mode 100644 index 0000000..6e7aec6 --- /dev/null +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/utils/ConstantStr.java @@ -0,0 +1,14 @@ +package com.volvo.ai.analytic.center.utils; +/** + * @ClassName ConstantStr + * @Description + * @Author renzhen + * @Date 2025-03-11 11:02 + * @Version 1.0 + **/ + +public class ConstantStr { + + public static final String CHANNEL_DCC = "Channel_Dcc"; + public static final String corpus_user = "corpush_user"; +} diff --git a/ai-analytic-center-biz/src/main/resources/mapper/TmOdsVdqwMessagearchivingMapper.xml b/ai-analytic-center-biz/src/main/resources/mapper/TmOdsVdqwMessagearchivingMapper.xml index babce66..4cdc87e 100644 --- a/ai-analytic-center-biz/src/main/resources/mapper/TmOdsVdqwMessagearchivingMapper.xml +++ b/ai-analytic-center-biz/src/main/resources/mapper/TmOdsVdqwMessagearchivingMapper.xml @@ -1,6 +1,24 @@ - + + +