From f72537a6c109e707f74ac0210c44f8bb5ca1721d Mon Sep 17 00:00:00 2001 From: zren25 Date: Fri, 14 Mar 2025 17:02:21 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=E5=A4=B1=E8=B4=A5=E8=A1=A5?= =?UTF-8?q?=E5=81=BF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../ai/analytic/center/job/CorpusFailJob.java | 132 +++++++++++++++--- .../center/mapper/AiAnalysisErrorsMapper.java | 7 + .../analytic/center/service/DiFyService.java | 2 +- .../service/TmTelephoneCorpusService.java | 2 +- .../impl/AiAnalysisErrorsServiceImpl.java | 2 +- .../center/service/impl/DiFyServiceImpl.java | 5 +- .../TmOdsVdqwMessagearchivingServiceImpl.java | 22 +-- .../impl/TmTelephoneCorpusServiceImpl.java | 30 ++-- 8 files changed, 146 insertions(+), 56 deletions(-) diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java index 7b4a461..310cb07 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java @@ -1,51 +1,147 @@ package com.volvo.ai.analytic.center.job; +import cn.hutool.core.date.DateUtil; +import com.alibaba.fastjson.JSON; +import com.alibaba.fastjson.JSONObject; +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import com.volvo.ai.analytic.center.dto.corpus.CorpusReportDTO; +import com.volvo.ai.analytic.center.dto.req.DiFyReq; import com.volvo.ai.analytic.center.entity.AiAnalysisErrors; +import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs; +import com.volvo.ai.analytic.center.entity.TmCorpusReport; import com.volvo.ai.analytic.center.enums.BusinessTypeEnum; -import com.volvo.ai.analytic.center.service.AiAnalysisErrorsService; -import com.volvo.common.core.util.ResultMsg; +import com.volvo.ai.analytic.center.enums.CategoryEnum; +import com.volvo.ai.analytic.center.mapper.AiAnalysisRequestLogsMapper; +import com.volvo.ai.analytic.center.service.*; +import com.volvo.ai.analytic.center.utils.FlowResultSplitUtil; import com.xxl.job.core.handler.annotation.XxlJob; import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections.CollectionUtils; +import org.apache.commons.lang3.StringUtils; +import org.apache.rocketmq.spring.core.RocketMQTemplate; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; -import org.springframework.web.bind.annotation.RestController; +import javax.annotation.Resource; +import java.util.Date; +import java.util.HashMap; import java.util.List; +import java.util.Map; @Slf4j @Component -@RestController public class CorpusFailJob { @Autowired private AiAnalysisErrorsService aiAnalysisErrorsService; + @Autowired + private AiAnalysisRequestLogsMapper aiAnalysisRequestLogsMapper; + @Autowired + private AiAnalysisRequestLogsService aiAnalysisRequestLogsService; + + @Autowired + private DiFyService diFyService; + + @Value("${rocketmq.producer.corpus.topic}") + private String topic; + + @Resource + private RocketMQTemplate rocketMqTemplate; + + @Autowired + private TmCorpusReportService tmCorpusReportService; + + @Autowired + private TmTelephoneCorpusService tmTelephoneCorpusService; + /** * 企微语料处理 */ @XxlJob("corpusFailTask") - public ResultMsg corpusFailTask() { - try { - log.info("语料解析失败重试处理"); - List aiAnalysisErrorsListlist = aiAnalysisErrorsService.queryAnalysisErrorList(BusinessTypeEnum.SMART_ASSISTANT.getCode()); - if(CollectionUtils.isNotEmpty(aiAnalysisErrorsListlist)){ + public void corpusFailTask() { - aiAnalysisErrorsListlist.stream().forEach(aiAnalysisErrors -> { + log.info("语料解析失败重试处理"); + List aiAnalysisErrorsListlist = aiAnalysisErrorsService.queryAnalysisErrorList(BusinessTypeEnum.SMART_ASSISTANT.getCode()); + if(CollectionUtils.isNotEmpty(aiAnalysisErrorsListlist)) { + log.info("语料解析失败重试处理 size:{}", aiAnalysisErrorsListlist.size()); + aiAnalysisErrorsListlist.stream().forEach(aiAnalysisErrors -> { - if(aiAnalysisErrors.getMaxRetryCount()>= aiAnalysisErrors.getRetryCount()){ - log.info("已超过最大重试次数!"); - return ; + LambdaQueryWrapper queryWrapper = new LambdaQueryWrapper<>(); + queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisErrors.getAiAnalysisRequestId()); + AiAnalysisRequestLogs oldAiAnalysisRequestLogs = aiAnalysisRequestLogsMapper.selectOne(queryWrapper); + + if (null != oldAiAnalysisRequestLogs) { + DiFyReq diFyReq = JSONObject.parseObject(oldAiAnalysisRequestLogs.getDifyRequest(), DiFyReq.class); + CorpusReportDTO corpusReportDTO = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class); + + Object difyoutResult = diFyService.getDiFyObject(diFyReq); + JSONObject execDifyFlow = JSONObject.parseObject(JSON.toJSONString(difyoutResult)).getJSONObject("data"); + 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); + if (StringUtils.isBlank(resultStrOne) || StringUtils.isBlank(resultStrTwo)) { + log.info("企微语料解析为空,text:{}", text); + return; + } + Map ltoMap = new HashMap<>(); + ltoMap.put("analysisRecordId", execDifyFlow.getString("aiAnalysisRequestId")); + ltoMap.put("analysisScene", corpusReportDTO.getAnalysisScene() + ""); + String tag = ""; + if (corpusReportDTO.getAnalysisScene() == 1) { //企微 + ltoMap.put("unionId", corpusReportDTO.getUnionId()); + ltoMap.put("consultantId", corpusReportDTO.getUserId()); + tag = CategoryEnum.ENTERPRISE_WECHAT.getCode(); + } else { + ltoMap.put("recordId", corpusReportDTO.getUserId()); + tag = CategoryEnum.PHONE_VOICE.getCode(); + } + + ltoMap.put("communicateDate", corpusReportDTO.getCorpusTime()); + ltoMap.put("analysisResult", resultStrOne.replace("#", "")); + ltoMap.put("analysisDetail", resultStrTwo); + // 发送MQ + log.info("send mq {}", ltoMap); + tmTelephoneCorpusService.sendMq(tag, JSONObject.toJSONString(ltoMap)); + try { + // 保存报告 + tmCorpusReportService.saveTmCorpusReport(TmCorpusReport.builder() + .aiAnalysisRequestId(oldAiAnalysisRequestLogs.getAiAnalysisRequestId()) + .corpusType(corpusReportDTO.getAnalysisScene()) + .corpusTime(DateUtil.parseTime(corpusReportDTO.getCorpusTime())) + .userId(corpusReportDTO.getUserId()) + .unionId(corpusReportDTO.getUnionId()) + .reportTitle(resultStrOne.replace("#", "")) + .reportInfo(resultStrTwo) + .isLike(0) + .isDeleted(0) + .version(0) + .createBy("system") + .updateBy("") + .createSqlby("") + .updateSqlby("") + .createTime(new Date()) + .build()); + } catch (Exception e) { + log.info(" 企业语料处理保存报告异常processItem:{} ", e); + } + } + oldAiAnalysisRequestLogs.setDifyResponse(JSON.toJSONString(execDifyFlow)); + aiAnalysisRequestLogsService.save(oldAiAnalysisRequestLogs); + aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount() + 1); + aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("1"); + aiAnalysisErrorsService.save(aiAnalysisErrors); + } }); - } - } catch (Exception e) { - log.error("processMessageByTask 定时任务补偿处理消息异常",e.getMessage()); + } } - return ResultMsg.ok(); - } + } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/AiAnalysisErrorsMapper.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/AiAnalysisErrorsMapper.java index 5b82eeb..d196aee 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/AiAnalysisErrorsMapper.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/AiAnalysisErrorsMapper.java @@ -3,7 +3,14 @@ package com.volvo.ai.analytic.center.mapper; import com.baomidou.mybatisplus.core.mapper.BaseMapper; import com.volvo.ai.analytic.center.entity.AiAnalysisErrors; import org.apache.ibatis.annotations.Mapper; +import org.apache.ibatis.annotations.Param; + +import java.util.List; @Mapper public interface AiAnalysisErrorsMapper extends BaseMapper { + + List queryAnalysisErrorList(@Param("businessType") String businessType); + + } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/DiFyService.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/DiFyService.java index d561518..c41c2c0 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/DiFyService.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/DiFyService.java @@ -8,5 +8,5 @@ public interface DiFyService { public Object getDiFyObject(DiFyReq diFyReq); - public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType); + public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType, String businessData); } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/TmTelephoneCorpusService.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/TmTelephoneCorpusService.java index c39f839..256518f 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/TmTelephoneCorpusService.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/TmTelephoneCorpusService.java @@ -19,5 +19,5 @@ public interface TmTelephoneCorpusService extends IService { String getCarModelList(); - + void sendMq(String tag, String message); } \ No newline at end of file diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisErrorsServiceImpl.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisErrorsServiceImpl.java index 531f9fe..3e9265f 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisErrorsServiceImpl.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisErrorsServiceImpl.java @@ -32,7 +32,7 @@ public class AiAnalysisErrorsServiceImpl extends ServiceImpl queryAnalysisErrorList(String businessType) { //捞取异常表中属于社区的异常数据 - return null; + return aiAnalysisErrorsMapper.queryAnalysisErrorList(businessType); } } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/DiFyServiceImpl.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/DiFyServiceImpl.java index 3caa372..d930a13 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/DiFyServiceImpl.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/DiFyServiceImpl.java @@ -47,7 +47,7 @@ public class DiFyServiceImpl implements DiFyService{ } @Override - public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType) { + public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType, String businessData) { String aiAnalysisRequestId = AiAnalysisUtils.getAiAnalysisRequestId(businessType); try { Map map = new HashMap<>(); @@ -58,7 +58,7 @@ public class DiFyServiceImpl implements DiFyService{ // 保存请求日志 aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder() .aiAnalysisRequestId(aiAnalysisRequestId) - .businessRequest(JSONObject.toJSONString(diFyReq.getBusinessData())) + .businessRequest(businessData) .difyAgentKey(diFyReq.getFlowId()) .difyRequest(JSON.toJSONString(diFyReq)) .difyResponse(JSON.toJSONString("")) @@ -83,6 +83,7 @@ public class DiFyServiceImpl implements DiFyService{ log.error("dify请求失败",e); aiAnalysisErrorsService.saveAiAnalysisErrors(AiAnalysisErrors.builder() .aiAnalysisRequestId(aiAnalysisRequestId) + .aiAnalysisRequestType(businessType) .aiAnalysisErrorHandlingStatus("0") .aiAnalysisErrorMessage(e.getMessage()) .build()); 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 85c3774..04e7dfe 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 @@ -2,7 +2,6 @@ package com.volvo.ai.analytic.center.service.impl; import cn.hutool.core.date.DatePattern; import cn.hutool.core.date.DateUtil; -import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSONObject; import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl; @@ -25,13 +24,10 @@ import com.volvo.ai.analytic.center.utils.FlowResultSplitUtil; import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections.CollectionUtils; import org.apache.commons.lang3.StringUtils; -import org.apache.rocketmq.client.producer.SendCallback; -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; @@ -87,6 +83,7 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl