加日志,发送mq之前

This commit is contained in:
ZLI263
2025-09-17 12:44:27 +08:00
parent 6e98fa3550
commit 2f4ae4debc

View File

@@ -32,8 +32,8 @@ import java.util.concurrent.Executors;
/** /**
* @description 铭牌语料表-同步表
* @author rz * @author rz
* @description 铭牌语料表-同步表
* @date 2025-03-04 * @date 2025-03-04
*/ */
@RefreshScope @RefreshScope
@@ -68,6 +68,7 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
@Autowired @Autowired
private DataMaskingRuleService dataMaskingRuleService; private DataMaskingRuleService dataMaskingRuleService;
@Override @Override
public void runNameplateCorpusDifyRetry(String paramJson) { public void runNameplateCorpusDifyRetry(String paramJson) {
long startTime = System.currentTimeMillis(); long startTime = System.currentTimeMillis();
@@ -93,12 +94,12 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
} }
} }
Integer total = tmNameplateCorpusMapper.countQueryTmNameplateCorpusRetry(statTime, endTime, customerFlowIds ,retry); Integer total = tmNameplateCorpusMapper.countQueryTmNameplateCorpusRetry(statTime, endTime, customerFlowIds, retry);
int totalPages = PageDto.getTotalPages(total, pageSize); int totalPages = PageDto.getTotalPages(total, pageSize);
// 获取消息列表 // 获取消息列表
int optimalThreadPoolSize = Runtime.getRuntime().availableProcessors() + 1; int optimalThreadPoolSize = Runtime.getRuntime().availableProcessors() + 1;
log.info("获取的线程数:{}",optimalThreadPoolSize); log.info("获取的线程数:{}", optimalThreadPoolSize);
// 创建线程池 // 创建线程池
ExecutorService executor = Executors.newFixedThreadPool(optimalThreadPoolSize); // 根据需求调整线程池大小 ExecutorService executor = Executors.newFixedThreadPool(optimalThreadPoolSize); // 根据需求调整线程池大小
@@ -122,13 +123,13 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
} }
// 关闭线程池 // 关闭线程池
executor.shutdown(); executor.shutdown();
log.info("重跑铭牌语料铭牌数据跑批结束 耗时:{}",System.currentTimeMillis()-startTime); log.info("重跑铭牌语料铭牌数据跑批结束 耗时:{}", System.currentTimeMillis() - startTime);
} }
public void processItem(TmNameplateCorpus item) { public void processItem(TmNameplateCorpus item) {
try { try {
Optional.ofNullable(item).filter(tmNameplateCorpus -> tmNameplateCorpus.getCustomerFlowId()!=null && tmNameplateCorpus.getNameplateContent()!=null).orElseThrow(()->new RuntimeException("铭牌语料为空")); Optional.ofNullable(item).filter(tmNameplateCorpus -> tmNameplateCorpus.getCustomerFlowId() != null && tmNameplateCorpus.getNameplateContent() != null).orElseThrow(() -> new RuntimeException("铭牌语料为空"));
String customerFlowId = item.getCustomerFlowId(); String customerFlowId = item.getCustomerFlowId();
String nameplateContent = item.getNameplateContent(); String nameplateContent = item.getNameplateContent();
@@ -149,8 +150,8 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
inputMap.put("model", carModel); inputMap.put("model", carModel);
inputMap.put("customerFlowId", item.getCustomerFlowId()); inputMap.put("customerFlowId", item.getCustomerFlowId());
inputMap.put("analysisScene", "3"); inputMap.put("analysisScene", "3");
inputMap.put("version",2); inputMap.put("version", 2);
inputMap.put("businessType",BusinessTypeEnum.SMART_ASSISTANT_NAMEPLATE.getCode()); inputMap.put("businessType", BusinessTypeEnum.SMART_ASSISTANT_NAMEPLATE.getCode());
diFyImageReq.setInputs(inputMap); diFyImageReq.setInputs(inputMap);
@@ -168,7 +169,7 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
CompletableFuture<Void> summaryFuture = CompletableFuture.runAsync(() -> { CompletableFuture<Void> summaryFuture = CompletableFuture.runAsync(() -> {
//第一个业务场景 开始: 总结和分类业务场景 //第一个业务场景 开始: 总结和分类业务场景
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT_NAMEPLATE.getCode(), JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT_NAMEPLATE.getCode(),
JSONObject.toJSONString(corpusReportDTO),null); JSONObject.toJSONString(corpusReportDTO), null);
log.info("runDify execDifyFlow ,铭牌语料 ,总结和分类业务,返回: {}", execDifyFlow); log.info("runDify execDifyFlow ,铭牌语料 ,总结和分类业务,返回: {}", execDifyFlow);
}, executorService); }, executorService);
long endTime = System.currentTimeMillis(); long endTime = System.currentTimeMillis();
@@ -181,11 +182,11 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
DiFyReq diFyImageReq2 = new DiFyReq(); DiFyReq diFyImageReq2 = new DiFyReq();
diFyImageReq2.setUser(ConstantStr.corpus_user); diFyImageReq2.setUser(ConstantStr.corpus_user);
Map<String, Object> inputMap2 = new HashMap<>(); Map<String, Object> inputMap2 = new HashMap<>();
inputMap2.put("businessId",item.getCustomerFlowId()); inputMap2.put("businessId", item.getCustomerFlowId());
inputMap2.put("communicateDate",item.getNameplateEndTime().toString()); inputMap2.put("communicateDate", item.getNameplateEndTime().toString());
inputMap2.put("analysisScene", "3"); inputMap2.put("analysisScene", "3");
inputMap2.put("version",2); inputMap2.put("version", 2);
inputMap2.put("chat",JSONObject.toJSONString(corpusReportDTO)); inputMap2.put("chat", JSONObject.toJSONString(corpusReportDTO));
diFyImageReq2.setInputs(inputMap2); diFyImageReq2.setInputs(inputMap2);
// 创建新的DiFyReq对象以避免线程安全问题 // 创建新的DiFyReq对象以避免线程安全问题
@@ -207,39 +208,41 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
} }
String outputsFinal = "outputs";
String outputsFinal="outputs";
String aiAnalysisRequestIdFinal = "aiAnalysisRequestId"; String aiAnalysisRequestIdFinal = "aiAnalysisRequestId";
@Override @Override
public void sendNameplateLto(JSONObject execDifyFlow,String aiAnalysisRequestId, TmNameplateCorpus tmNameplateCorpus,String businessType) { public void sendNameplateLto(JSONObject execDifyFlow, String aiAnalysisRequestId, TmNameplateCorpus tmNameplateCorpus, String businessType) {
try {
log.info("铭牌sendNameplateLto getCustomerFlowId:{}", tmNameplateCorpus.getCustomerFlowId()); log.info("铭牌sendNameplateLto getCustomerFlowId:{}", tmNameplateCorpus.getCustomerFlowId());
log.info("铭牌sendNameplateLto text contnt :{}", execDifyFlow.toString()); log.info("铭牌sendNameplateLto text contnt :{}", execDifyFlow);
log.info("铭牌sendNameplateLto toJSONString :{}", execDifyFlow.toJSONString());
if (null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")) { if (null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")) {
String text = execDifyFlow.getString(outputsFinal); String text = execDifyFlow.getString(outputsFinal);
// 发送MQ // 发送MQ
log.info("铭牌send mq 业务类型是:{}, Json 是:{}", businessType, text); log.info("铭牌send mq 业务类型是:{}, Json 是:{}", businessType, text);
if (BusinessTypeEnum.CORPUS_PORTRAIT_NAMEPLATE.getCode().equals(businessType)) {// 客户画像 if (BusinessTypeEnum.CORPUS_PORTRAIT_NAMEPLATE.getCode().equals(businessType)) {// 客户画像
tmTelephoneCorpusService.sendMq( CategoryEnum.NAMEPLATE_VOICE_PORTRAIT.getCode(), text); tmTelephoneCorpusService.sendMq(CategoryEnum.NAMEPLATE_VOICE_PORTRAIT.getCode(), text);
}else { //一句话总结+分类 } else { //一句话总结+分类
tmTelephoneCorpusService.sendMq( CategoryEnum.NAMEPLATE_VOICE.getCode(), text); tmTelephoneCorpusService.sendMq(CategoryEnum.NAMEPLATE_VOICE.getCode(), text);
} }
try {
ttNameplateRecordMapper.insert(TtNameplateRecord.builder().nameplateCorpusId(tmNameplateCorpus.getId()).customerFlowId(tmNameplateCorpus.getCustomerFlowId()).build()); ttNameplateRecordMapper.insert(TtNameplateRecord.builder().nameplateCorpusId(tmNameplateCorpus.getId()).customerFlowId(tmNameplateCorpus.getCustomerFlowId()).build());
aiAnalysisRequestLogsService.saveOrUpdateAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder().aiAnalysisRequestId(aiAnalysisRequestId).businessResponse(text).build()); aiAnalysisRequestLogsService.saveOrUpdateAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder().aiAnalysisRequestId(aiAnalysisRequestId).businessResponse(text).build());
}
} catch (Exception e) { } catch (Exception e) {
log.info(" 铭牌语料处理保存报告异常processItem{} ", e); log.info(" 铭牌语料处理保存报告异常processItem{} ", e);
} }
} }
}
@Override @Override
public ResultMsg<Object> updateNameplate(String message) { public ResultMsg<Object> updateNameplate(String message) {
try { try {
if(StringUtils.isNotEmpty(message)){ if (StringUtils.isNotEmpty(message)) {
log.info("updateNameplate 回调 message: {}", message); log.info("updateNameplate 回调 message: {}", message);
JSONObject analysisResp = JSONObject.parseObject(message); JSONObject analysisResp = JSONObject.parseObject(message);
String aiAnalysisRequestId = analysisResp.getString(aiAnalysisRequestIdFinal); String aiAnalysisRequestId = analysisResp.getString(aiAnalysisRequestIdFinal);
@@ -248,16 +251,16 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
String businessType; String businessType;
AiAnalysisRequestLogs aiAnalysisRequestLogs = new AiAnalysisRequestLogs(); AiAnalysisRequestLogs aiAnalysisRequestLogs = new AiAnalysisRequestLogs();
JSONObject difyJson =null; JSONObject difyJson = null;
if (StringUtils.isEmpty(customerFlowId)){ // 客户画像场景 if (StringUtils.isEmpty(customerFlowId)) { // 客户画像场景
log.info("customerFlowId 为空,客户画像场景"); log.info("customerFlowId 为空,客户画像场景");
String text = analysisResp.getString(outputsFinal); String text = analysisResp.getString(outputsFinal);
JSONObject outputs = JSONObject.parseObject(text); JSONObject outputs = JSONObject.parseObject(text);
customerFlowId = outputs.getString("businessId"); customerFlowId = outputs.getString("businessId");
aiAnalysisRequestId=analysisResp.getString(aiAnalysisRequestIdFinal); aiAnalysisRequestId = analysisResp.getString(aiAnalysisRequestIdFinal);
businessType=BusinessTypeEnum.CORPUS_PORTRAIT_NAMEPLATE.getCode(); businessType = BusinessTypeEnum.CORPUS_PORTRAIT_NAMEPLATE.getCode();
aiAnalysisRequestLogs.setBusinessResponse(text); aiAnalysisRequestLogs.setBusinessResponse(text);
}else { //总结+分类 } else { //总结+分类
businessType = BusinessTypeEnum.SMART_ASSISTANT_NAMEPLATE.getCode(); businessType = BusinessTypeEnum.SMART_ASSISTANT_NAMEPLATE.getCode();
difyJson = JSONObject.parseObject(difyResponse); difyJson = JSONObject.parseObject(difyResponse);
aiAnalysisRequestLogs.setBusinessResponse(difyJson.getString(outputsFinal)); aiAnalysisRequestLogs.setBusinessResponse(difyJson.getString(outputsFinal));
@@ -276,9 +279,9 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
LambdaQueryWrapper<TmNameplateCorpus> queryWrapper = new LambdaQueryWrapper<>(); LambdaQueryWrapper<TmNameplateCorpus> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(TmNameplateCorpus::getCustomerFlowId, customerFlowId); queryWrapper.eq(TmNameplateCorpus::getCustomerFlowId, customerFlowId);
List<TmNameplateCorpus> tmNameplateCorpusList = tmNameplateCorpusMapper.selectList(queryWrapper); List<TmNameplateCorpus> tmNameplateCorpusList = tmNameplateCorpusMapper.selectList(queryWrapper);
if (CollectionUtils.isNotEmpty(tmNameplateCorpusList)){ //这是补偿机制 if (CollectionUtils.isNotEmpty(tmNameplateCorpusList)) { //这是补偿机制
log.info("66666666666666"); log.info("66666666666666");
sendNameplateLto(difyJson,aiAnalysisRequestId, tmNameplateCorpusList.get(0),businessType); sendNameplateLto(difyJson, aiAnalysisRequestId, tmNameplateCorpusList.get(0), businessType);
LambdaQueryWrapper<AiAnalysisErrors> errorQueryWrapper = new LambdaQueryWrapper<>(); LambdaQueryWrapper<AiAnalysisErrors> errorQueryWrapper = new LambdaQueryWrapper<>();
errorQueryWrapper.eq(AiAnalysisErrors::getAiAnalysisRequestId, aiAnalysisRequestId); errorQueryWrapper.eq(AiAnalysisErrors::getAiAnalysisRequestId, aiAnalysisRequestId);
errorQueryWrapper.eq(AiAnalysisErrors::getAiAnalysisErrorHandlingStatus, "0"); errorQueryWrapper.eq(AiAnalysisErrors::getAiAnalysisErrorHandlingStatus, "0");
@@ -287,7 +290,7 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
if (oldAiAnalysisErrors != null) { if (oldAiAnalysisErrors != null) {
aiAnalysisErrorsMapper.update(AiAnalysisErrors.builder().aiAnalysisRequestId(oldAiAnalysisErrors.getAiAnalysisRequestId()).aiAnalysisErrorHandlingStatus("1").build(), errorQueryWrapper); aiAnalysisErrorsMapper.update(AiAnalysisErrors.builder().aiAnalysisRequestId(oldAiAnalysisErrors.getAiAnalysisRequestId()).aiAnalysisErrorHandlingStatus("1").build(), errorQueryWrapper);
} }
}else { } else {
log.info("没找到语料记录"); log.info("没找到语料记录");
} }
return ResultMsg.ok(); return ResultMsg.ok();
@@ -299,14 +302,15 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
return ResultMsg.failed(); return ResultMsg.failed();
} }
@Override @Override
public ResultMsg<Object> mockInsert(String data) { public ResultMsg<Object> mockInsert(String data) {
if(StringUtils.isNotEmpty(data)){ if (StringUtils.isNotEmpty(data)) {
List<TmNameplateCorpus> tmNameplateCorpusList = new ArrayList<>(); List<TmNameplateCorpus> tmNameplateCorpusList = new ArrayList<>();
for(int i=0;i<100;i++){ for (int i = 0; i < 100; i++) {
TmNameplateCorpus analysisResp = JSONObject.parseObject(data, TmNameplateCorpus.class); TmNameplateCorpus analysisResp = JSONObject.parseObject(data, TmNameplateCorpus.class);
analysisResp.setCustomerFlowId(analysisResp.getCustomerFlowId()+i); analysisResp.setCustomerFlowId(analysisResp.getCustomerFlowId() + i);
tmNameplateCorpusList.add(analysisResp); tmNameplateCorpusList.add(analysisResp);
} }
this.saveBatch(tmNameplateCorpusList); this.saveBatch(tmNameplateCorpusList);
@@ -314,6 +318,7 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
} }
return ResultMsg.failed(); return ResultMsg.failed();
} }
@Override @Override
public List<TmNameplateCorpus> queryTelephoneCorpusByCustomerFlowId(List<String> customerFlowIds) { public List<TmNameplateCorpus> queryTelephoneCorpusByCustomerFlowId(List<String> customerFlowIds) {
return tmNameplateCorpusMapper.queryTmNameplateCorpusByCustomerFlowIds(customerFlowIds); return tmNameplateCorpusMapper.queryTmNameplateCorpusByCustomerFlowIds(customerFlowIds);