画像增加补偿

This commit is contained in:
zren25
2025-06-12 18:10:19 +08:00
parent 43f9154007
commit ee9fadd622
4 changed files with 94 additions and 2 deletions

View File

@@ -3,6 +3,7 @@ package com.volvo.ai.analytic.center.job;
import com.alibaba.fastjson.JSON; 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.huaweicloud.sdk.eg.v1.model.CloudEvents;
import com.volvo.ai.analytic.center.dto.PageDto; import com.volvo.ai.analytic.center.dto.PageDto;
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO; import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
import com.volvo.ai.analytic.center.dto.corpus.CorpusReportDTO; import com.volvo.ai.analytic.center.dto.corpus.CorpusReportDTO;
@@ -72,6 +73,16 @@ public class CorpusFailJob {
@Autowired @Autowired
private TmNameplateCorpusService tmNameplateCorpusService; private TmNameplateCorpusService tmNameplateCorpusService;
@Autowired
private CorpusPortraitService corpusPortraitService;
@Resource
HuaWeiEGService huaWeiService;
@Value("${huawei.cloud.EG.channel.ltoChannelId}")
private String ltoChannelId;
@Value("${huawei.cloud.EG.channel.ltoSourceId}")
private String sourceId;
/** /**
* 企微语料处理 * 企微语料处理
@@ -247,4 +258,64 @@ public class CorpusFailJob {
} }
} }
} }
/**
* 画像error补偿
*/
@XxlJob("corpusPortraitFailTask")
public void corpusPortraitFailTask() {
log.info(" corpusPortraitFailTask画像解析失败重试处理");
Integer total = aiAnalysisErrorsService.queryCountAnalysisErrorList(Arrays.asList(BusinessTypeEnum.CORPUS_PORTRAIT_DCC.getCode(),BusinessTypeEnum.CORPUS_PORTRAIT_QIWEI.getCode(),BusinessTypeEnum.CORPUS_PORTRAIT_QIWEI.getCode()));
log.info("corpusPortraitFailTask画像语料解析失败重试处理数据量{}", total);
int totalPages = PageDto.getTotalPages(total, pageSize);
log.info("corpusFailTask 语料解析失败重试处理数据量:{},总页数:{}", total, totalPages);
for (int i = 1; i <= totalPages; i++) {
int offset = (i - 1) * pageSize;
List<AiAnalysisErrors> aiAnalysisErrorsListlist = aiAnalysisErrorsService.queryAnalysisErrorList(Arrays.asList(BusinessTypeEnum.CORPUS_PORTRAIT_DCC.getCode(),BusinessTypeEnum.CORPUS_PORTRAIT_QIWEI.getCode(),BusinessTypeEnum.CORPUS_PORTRAIT_QIWEI.getCode()),offset, pageSize);
if(CollectionUtils.isNotEmpty(aiAnalysisErrorsListlist)) {
log.info("corpusFailTask语料解析失败重试处理 size:{}", aiAnalysisErrorsListlist.size());
List<String> requestIdList = aiAnalysisErrorsListlist.stream()
.map(AiAnalysisErrors::getAiAnalysisRequestId)
.collect(Collectors.toList());
List<AiAnalysisRequestLogs> aiAnalysisRequestLogsList = aiAnalysisRequestLogsMapper.queryByAiAnalysisRequestIds(requestIdList);
Map<String, AiAnalysisRequestLogs> requestLogsMap = aiAnalysisRequestLogsList.stream()
.collect(Collectors.toMap( AiAnalysisRequestLogs::getAiAnalysisRequestId,logs -> logs ));
aiAnalysisErrorsListlist.forEach(aiAnalysisErrors -> {
AiAnalysisRequestLogs oldAiAnalysisRequestLogs = requestLogsMap.get(aiAnalysisErrors.getAiAnalysisRequestId());
if (null != oldAiAnalysisRequestLogs && StringUtils.isBlank(oldAiAnalysisRequestLogs.getBusinessResponse())) {
DiFyReq diFyReq = JSONObject.parseObject(oldAiAnalysisRequestLogs.getDifyRequest(), DiFyReq.class);
CorpusReportDTO corpusReportDTO = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class);
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyReq);
log.info(" corpusFailTask runDify execDifyFlow {}", execDifyFlow);
if (null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")) {
JSONObject text = execDifyFlow.getJSONObject("outputs");
// 发送MQ
log.info("sendEvent mq {}", text);
huaWeiService.sendEvent(corpusPortraitService.setCloudEvents(oldAiAnalysisRequestLogs.getAiAnalysisRequestId(), text.toJSONString(), oldAiAnalysisRequestLogs.getAiAnalysisRequestType().substring(oldAiAnalysisRequestLogs.getAiAnalysisRequestType().lastIndexOf("_")+1)), ltoChannelId);
// 保存报告
oldAiAnalysisRequestLogs.setBusinessResponse( text.toJSONString());
oldAiAnalysisRequestLogs.setDifyResponse(JSON.toJSONString(execDifyFlow));
updateDiFyRequest(oldAiAnalysisRequestLogs, oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("1");
}
aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount() + 1);
updateAiAnalysisErrors(aiAnalysisErrors, oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
}
});
}
}
}
} }

View File

@@ -2,6 +2,7 @@ package com.volvo.ai.analytic.center.job;
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO; import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
import com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper; import com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper;
import com.volvo.ai.analytic.center.service.CorpusPortraitService;
import com.volvo.ai.analytic.center.service.TmOdsVdqwMessagearchivingService; import com.volvo.ai.analytic.center.service.TmOdsVdqwMessagearchivingService;
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService; import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
import com.volvo.common.core.util.ResultMsg; import com.volvo.common.core.util.ResultMsg;
@@ -10,6 +11,7 @@ import com.xxl.job.core.handler.annotation.XxlJob;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
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.PostMapping; import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestBody;
@@ -33,6 +35,13 @@ public class QiWeiCorpusJob {
@Autowired @Autowired
private TmTelephoneCorpusMapper tmTelephoneCorpusMapper; private TmTelephoneCorpusMapper tmTelephoneCorpusMapper;
@Autowired
private CorpusPortraitService corpusPortraitService;
@Value("${huawei.cloud.EG.channel.ltoChannelId}")
private String ltoChannelId;
/** /**
* 企微语料处理 * 企微语料处理
*/ */
@@ -68,9 +77,14 @@ public class QiWeiCorpusJob {
param = paramJson; param = paramJson;
} }
// 执行业务逻辑 // 执行业务逻辑
XxlJobHelper.log("任务参数: {}", param); XxlJobHelper.log("dccCorpusFailRetry任务参数: {}", param);
List<AicorpusTelephoneDTO> corpusDtoList = tmTelephoneCorpusMapper.queryTelephoneCorpusBySourceIds( Arrays.asList(param.split(","))); List<AicorpusTelephoneDTO> corpusDtoList = tmTelephoneCorpusMapper.queryTelephoneCorpusBySourceIds( Arrays.asList(param.split(",")));
corpusDtoList.stream().forEach(item-> tmTelephoneCorpusService.runTelephoneCorpusDify(item)); corpusDtoList.stream().forEach(item->
{
tmTelephoneCorpusService.runTelephoneCorpusDify(item);
corpusPortraitService.portraitDcc(item);
}
);
} catch (Exception e) { } catch (Exception e) {
log.error("processMessageByTask 定时任务补偿处理消息异常",e.getMessage()); log.error("processMessageByTask 定时任务补偿处理消息异常",e.getMessage());
throw new RuntimeException(e); throw new RuntimeException(e);

View File

@@ -2,6 +2,7 @@ package com.volvo.ai.analytic.center.service;
import com.alibaba.fastjson.JSONObject; import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.extension.service.IService; import com.baomidou.mybatisplus.extension.service.IService;
import com.huaweicloud.sdk.eg.v1.model.CloudEvents;
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO; import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
import com.volvo.ai.analytic.center.dto.corpus.OdsVdqwMessageOTD; import com.volvo.ai.analytic.center.dto.corpus.OdsVdqwMessageOTD;
import com.volvo.ai.analytic.center.entity.TmNameplateCorpus; import com.volvo.ai.analytic.center.entity.TmNameplateCorpus;
@@ -22,5 +23,6 @@ public interface CorpusPortraitService {
void portraitNameplate(TmNameplateCorpus item); void portraitNameplate(TmNameplateCorpus item);
CloudEvents setCloudEvents(String aiAnalysisRequestId, String text, String eventType);
} }

View File

@@ -80,6 +80,9 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
@Resource(name = "threadPoolTaskExecutor") @Resource(name = "threadPoolTaskExecutor")
private ThreadPoolTaskExecutor executor; private ThreadPoolTaskExecutor executor;
@Autowired
private CorpusPortraitService corpusPortraitService;
@Override @Override
public void runNameplateCorpusDifyRetry(String paramJson) { public void runNameplateCorpusDifyRetry(String paramJson) {
long startTime = System.currentTimeMillis(); long startTime = System.currentTimeMillis();
@@ -114,6 +117,8 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
CompletableFuture.runAsync(() -> { CompletableFuture.runAsync(() -> {
try { try {
processItem(tmNameplateCorpus); processItem(tmNameplateCorpus);
// 增加 画像手动补偿
corpusPortraitService.portraitNameplate(tmNameplateCorpus);
} catch (Exception e) { } catch (Exception e) {
log.error("重跑铭牌语料失败: customerFlowId={}, AcceptUserId={}, 异常: {}", log.error("重跑铭牌语料失败: customerFlowId={}, AcceptUserId={}, 异常: {}",
tmNameplateCorpus.getCustomerFlowId(), e.getMessage(), e); tmNameplateCorpus.getCustomerFlowId(), e.getMessage(), e);