增加失败补偿

This commit is contained in:
zren25
2025-03-14 17:02:21 +08:00
parent 2127e3e2d5
commit f72537a6c1
8 changed files with 146 additions and 56 deletions

View File

@@ -1,51 +1,147 @@
package com.volvo.ai.analytic.center.job; 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.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.enums.BusinessTypeEnum;
import com.volvo.ai.analytic.center.service.AiAnalysisErrorsService; import com.volvo.ai.analytic.center.enums.CategoryEnum;
import com.volvo.common.core.util.ResultMsg; 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 com.xxl.job.core.handler.annotation.XxlJob;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils; 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.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component; 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.List;
import java.util.Map;
@Slf4j @Slf4j
@Component @Component
@RestController
public class CorpusFailJob { public class CorpusFailJob {
@Autowired @Autowired
private AiAnalysisErrorsService aiAnalysisErrorsService; 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") @XxlJob("corpusFailTask")
public ResultMsg corpusFailTask() { public void corpusFailTask() {
try {
log.info("语料解析失败重试处理"); log.info("语料解析失败重试处理");
List<AiAnalysisErrors> aiAnalysisErrorsListlist = aiAnalysisErrorsService.queryAnalysisErrorList(BusinessTypeEnum.SMART_ASSISTANT.getCode()); List<AiAnalysisErrors> aiAnalysisErrorsListlist = aiAnalysisErrorsService.queryAnalysisErrorList(BusinessTypeEnum.SMART_ASSISTANT.getCode());
if(CollectionUtils.isNotEmpty(aiAnalysisErrorsListlist)) { if(CollectionUtils.isNotEmpty(aiAnalysisErrorsListlist)) {
log.info("语料解析失败重试处理 size:{}", aiAnalysisErrorsListlist.size());
aiAnalysisErrorsListlist.stream().forEach(aiAnalysisErrors -> { aiAnalysisErrorsListlist.stream().forEach(aiAnalysisErrors -> {
if(aiAnalysisErrors.getMaxRetryCount()>= aiAnalysisErrors.getRetryCount()){
log.info("已超过最大重试次数!"); LambdaQueryWrapper<AiAnalysisRequestLogs> 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; return;
} }
}); Map<String, String> 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) { } catch (Exception e) {
log.error("processMessageByTask 定时任务补偿处理消息异常",e.getMessage()); log.info(" 企业语料处理保存报告异常processItem{} ", e);
}
} }
return ResultMsg.ok(); oldAiAnalysisRequestLogs.setDifyResponse(JSON.toJSONString(execDifyFlow));
aiAnalysisRequestLogsService.save(oldAiAnalysisRequestLogs);
aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount() + 1);
aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("1");
aiAnalysisErrorsService.save(aiAnalysisErrors);
} }
});
}
}
} }

View File

@@ -3,7 +3,14 @@ package com.volvo.ai.analytic.center.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper; import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors; import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
import org.apache.ibatis.annotations.Mapper; import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
import java.util.List;
@Mapper @Mapper
public interface AiAnalysisErrorsMapper extends BaseMapper<AiAnalysisErrors> { public interface AiAnalysisErrorsMapper extends BaseMapper<AiAnalysisErrors> {
List<AiAnalysisErrors> queryAnalysisErrorList(@Param("businessType") String businessType);
} }

View File

@@ -8,5 +8,5 @@ public interface DiFyService {
public Object getDiFyObject(DiFyReq diFyReq); public Object getDiFyObject(DiFyReq diFyReq);
public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType); public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType, String businessData);
} }

View File

@@ -19,5 +19,5 @@ public interface TmTelephoneCorpusService extends IService<TmTelephoneCorpus> {
String getCarModelList(); String getCarModelList();
void sendMq(String tag, String message);
} }

View File

@@ -32,7 +32,7 @@ public class AiAnalysisErrorsServiceImpl extends ServiceImpl<AiAnalysisErrorsMa
public List<AiAnalysisErrors> queryAnalysisErrorList(String businessType) { public List<AiAnalysisErrors> queryAnalysisErrorList(String businessType) {
//捞取异常表中属于社区的异常数据 //捞取异常表中属于社区的异常数据
return null; return aiAnalysisErrorsMapper.queryAnalysisErrorList(businessType);
} }
} }

View File

@@ -47,7 +47,7 @@ public class DiFyServiceImpl implements DiFyService{
} }
@Override @Override
public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType) { public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType, String businessData) {
String aiAnalysisRequestId = AiAnalysisUtils.getAiAnalysisRequestId(businessType); String aiAnalysisRequestId = AiAnalysisUtils.getAiAnalysisRequestId(businessType);
try { try {
Map<String, Object> map = new HashMap<>(); Map<String, Object> map = new HashMap<>();
@@ -58,7 +58,7 @@ public class DiFyServiceImpl implements DiFyService{
// 保存请求日志 // 保存请求日志
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder() aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder()
.aiAnalysisRequestId(aiAnalysisRequestId) .aiAnalysisRequestId(aiAnalysisRequestId)
.businessRequest(JSONObject.toJSONString(diFyReq.getBusinessData())) .businessRequest(businessData)
.difyAgentKey(diFyReq.getFlowId()) .difyAgentKey(diFyReq.getFlowId())
.difyRequest(JSON.toJSONString(diFyReq)) .difyRequest(JSON.toJSONString(diFyReq))
.difyResponse(JSON.toJSONString("")) .difyResponse(JSON.toJSONString(""))
@@ -83,6 +83,7 @@ public class DiFyServiceImpl implements DiFyService{
log.error("dify请求失败",e); log.error("dify请求失败",e);
aiAnalysisErrorsService.saveAiAnalysisErrors(AiAnalysisErrors.builder() aiAnalysisErrorsService.saveAiAnalysisErrors(AiAnalysisErrors.builder()
.aiAnalysisRequestId(aiAnalysisRequestId) .aiAnalysisRequestId(aiAnalysisRequestId)
.aiAnalysisRequestType(businessType)
.aiAnalysisErrorHandlingStatus("0") .aiAnalysisErrorHandlingStatus("0")
.aiAnalysisErrorMessage(e.getMessage()) .aiAnalysisErrorMessage(e.getMessage())
.build()); .build());

View File

@@ -2,7 +2,6 @@ package com.volvo.ai.analytic.center.service.impl;
import cn.hutool.core.date.DatePattern; import cn.hutool.core.date.DatePattern;
import cn.hutool.core.date.DateUtil; import cn.hutool.core.date.DateUtil;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject; import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl; 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 lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils; import org.apache.commons.collections.CollectionUtils;
import org.apache.commons.lang3.StringUtils; 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.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope; import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import javax.annotation.Resource; import javax.annotation.Resource;
@@ -87,6 +83,7 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
private TmCorpusReportService tmCorpusReportService; private TmCorpusReportService tmCorpusReportService;
@Override @Override
public void runQiWeiCorpusDify(String paramJson) { public void runQiWeiCorpusDify(String paramJson) {
log.info("runQiWeiCorpusDify paramJson {}", paramJson); log.info("runQiWeiCorpusDify paramJson {}", paramJson);
@@ -198,9 +195,8 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
corpusReportDTO.setUnionId(unionId); corpusReportDTO.setUnionId(unionId);
corpusReportDTO.setUserId(userId); corpusReportDTO.setUserId(userId);
corpusReportDTO.setAnalysisScene(1l); corpusReportDTO.setAnalysisScene(1l);
diFyImageReq.setBusinessData(corpusReportDTO);
// 获取配置 // 获取配置
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode()); JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode(), JSONObject.toJSONString(corpusReportDTO));
log.info("runDify execDifyFlow {}", execDifyFlow); log.info("runDify execDifyFlow {}", execDifyFlow);
if (null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")) { if (null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")) {
String text = execDifyFlow.getJSONObject("outputs").getString("text"); String text = execDifyFlow.getJSONObject("outputs").getString("text");
@@ -220,19 +216,7 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
ltoMap.put("analysisDetail", resultStrTwo); ltoMap.put("analysisDetail", resultStrTwo);
// 发送MQ // 发送MQ
log.info("send mq {}", ltoMap); log.info("send mq {}", ltoMap);
rocketMqTemplate.asyncSend(topic+":"+ CategoryEnum.ENTERPRISE_WECHAT.getCode(), MessageBuilder.withPayload(JSON.toJSONString(ltoMap)).build(), tmTelephoneCorpusService.sendMq( CategoryEnum.ENTERPRISE_WECHAT.getCode(), JSONObject.toJSONString(ltoMap));
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);
try { try {
// 保存报告 // 保存报告
tmCorpusReportService.saveTmCorpusReport(TmCorpusReport.builder() tmCorpusReportService.saveTmCorpusReport(TmCorpusReport.builder()

View File

@@ -1,7 +1,6 @@
package com.volvo.ai.analytic.center.service.impl; package com.volvo.ai.analytic.center.service.impl;
import cn.hutool.core.date.DateUtil; import cn.hutool.core.date.DateUtil;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONArray; import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject; import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl; import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
@@ -141,9 +140,8 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
CorpusReportDTO corpusReportDTO = new CorpusReportDTO(); CorpusReportDTO corpusReportDTO = new CorpusReportDTO();
corpusReportDTO.setCorpusTime(formattedDateStartTime); corpusReportDTO.setCorpusTime(formattedDateStartTime);
corpusReportDTO.setRecordId(aicorpusTelephone.getSourceId()); corpusReportDTO.setRecordId(aicorpusTelephone.getSourceId());
diFyImageReq.setBusinessData(corpusReportDTO);
// 获取配置 // 获取配置
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode()); JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode(), JSONObject.toJSONString(corpusReportDTO));
log.info("runDify execDifyFlow {}",execDifyFlow); log.info("runDify execDifyFlow {}",execDifyFlow);
if(null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")){ if(null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")){
@@ -166,17 +164,7 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
// 发送MQ // 发送MQ
log.info("send mq {}",ltoMap); log.info("send mq {}",ltoMap);
rocketMqTemplate.asyncSend(topic+":"+CategoryEnum.PHONE_VOICE.getCode(), MessageBuilder.withPayload(JSON.toJSONString(ltoMap)).build(), sendMq( CategoryEnum.PHONE_VOICE.getCode(), JSONObject.toJSONString(ltoMap));
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);
try { try {
// 保存报告 // 保存报告
Date corpusTime = DateUtil.parseUTC(jsonObject.getString("start_time")); Date corpusTime = DateUtil.parseUTC(jsonObject.getString("start_time"));
@@ -227,6 +215,20 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
return "C40 RECHARGE、EM90、EX30、S60、S90、V60、V90、XC40、XC40 RECHARGE、XC60、XC90"; return "C40 RECHARGE、EM90、EX30、S60、S90、V60、V90、XC40、XC40 RECHARGE、XC60、XC90";
} }
@Override
public void sendMq(String tag, String message) {
rocketMqTemplate.asyncSend(topic+":"+ tag, MessageBuilder.withPayload(message).build(),
new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
log.info("发送MQ成功 消息体:{}", message);
}
@Override
public void onException(Throwable e) {
log.error("送MQ异常 消息体:{}, 异常:", message, e);
}
}, 10000);
}
} }