语料信息不完整,dcc加日志

This commit is contained in:
ZLI263
2025-09-18 15:08:12 +08:00
parent 40c7d471b6
commit 6ba44cb3d0

View File

@@ -164,96 +164,105 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
}
private synchronized void processItem(OdsVdqwMessageOTD item, String statTime, String endTime) {
private void processItem(OdsVdqwMessageOTD item, String statTime, String endTime) {
log.info("企微语料内容FromUserId{}, AcceptUserId{}", item.getFromUserId(), item.getAcceptUserId());
// 1vdqw_workuserinfo 这个表对应是 B端认证中心userId
// 2vdqw_externalcontact 这个表对应是 企微客户unionId
TmOdsVdqwExternalcontact tmOdsVdqwExternalcontact = getUnionId(Arrays.asList(item.getFromUserId(), item.getAcceptUserId()));
TmOdsVdqwWorkuserinfo tmOdsVdqwWorkuserinfo = getUserId(Arrays.asList(item.getFromUserId(), item.getAcceptUserId()));
if (null == tmOdsVdqwExternalcontact || null == tmOdsVdqwWorkuserinfo) {
log.info("企微查询信息为空 ");
return;
}
String unionId = tmOdsVdqwExternalcontact.getUnionId();
String userId = tmOdsVdqwWorkuserinfo.getMiddleUserId().toString();
log.info("企微查询信息unionId{}, userId{}", unionId, userId);
if (StringUtils.isNotBlank(item.getFromUserId()) && StringUtils.isNotBlank(item.getAcceptUserId())) {
List<OdsVdqwMessageOTD> contetnList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByFromUserIdAndAcceptUserId(statTime, endTime, Arrays.asList(item.getFromUserId(), item.getAcceptUserId()) );
log.info("一条完整的企微信息的聊天内容,未加密:{}", contetnList.toString());
OdsVdqwMessageOTD maxMsgTimeItem = contetnList.stream()
.max((o1, o2) -> o1.getMsgTime().compareTo(o2.getMsgTime()))
.orElse(null);
List<DataMaskingRule> maskingRuleItems = dataMaskingRuleService.getDataMaskingRuleListByApplicationChannel(BusinessTypeEnum.SMART_ASSISTANT_QIWEI.getCode());
RunMaskingRuleInput runMaskingRuleInput = new RunMaskingRuleInput();
runMaskingRuleInput.setDataMaskingRules(maskingRuleItems);
StringBuffer chatList = new StringBuffer();
AtomicInteger externalcontactCount = new AtomicInteger();
contetnList.forEach(contentItem -> {
String title = "";
if(tmOdsVdqwExternalcontact.getExternalUserId().equals(contentItem.getFromUserId())){
title="客户:";
externalcontactCount.addAndGet(1);
}else{
title="顾问:";
}
JSONObject contentJson = JSONObject.parseObject(contentItem.getContent());
String content = contentJson.getString("content");
// 拼接 role 和 text
String chat = title.concat(content);
runMaskingRuleInput.setOldStr(chat);
String corpusChat = dataMaskingRuleService.runMaskingRule(runMaskingRuleInput);
chatList.append(corpusChat).append("\n");
});
if(externalcontactCount.get()<1){
log.info("没有客户回复的语料,无需解析");
DiFyReq diFyImageReq = new DiFyReq();
CorpusReportDTO corpusReportDTO = new CorpusReportDTO();
OdsVdqwMessageOTD finalMaxMsgTimeItem = new OdsVdqwMessageOTD();
StringBuffer chatList = new StringBuffer();
String unionId = "";
synchronized (this) {// 1vdqw_workuserinfo 这个表对应是 B端认证中心userId
// 2vdqw_externalcontact 这个表对应是 企微客户unionId
TmOdsVdqwExternalcontact tmOdsVdqwExternalcontact = getUnionId(Arrays.asList(item.getFromUserId(), item.getAcceptUserId()));
TmOdsVdqwWorkuserinfo tmOdsVdqwWorkuserinfo = getUserId(Arrays.asList(item.getFromUserId(), item.getAcceptUserId()));
if (null == tmOdsVdqwExternalcontact || null == tmOdsVdqwWorkuserinfo) {
log.info("企微查询信息为空 ");
return;
}
log.info("一条完整的企微信息的聊天内容,循环次数:{} ,应该循环次数:{}, 掩码后:{}", externalcontactCount.get(),contetnList.size(),chatList.toString());
unionId = tmOdsVdqwExternalcontact.getUnionId();
String userId = tmOdsVdqwWorkuserinfo.getMiddleUserId().toString();
log.info("企微查询信息unionId{}, userId{}", unionId, userId);
if (StringUtils.isNotBlank(item.getFromUserId()) && StringUtils.isNotBlank(item.getAcceptUserId())) {
List<OdsVdqwMessageOTD> contetnList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByFromUserIdAndAcceptUserId(statTime, endTime, Arrays.asList(item.getFromUserId(), item.getAcceptUserId()));
Map<String, Object> inputMap = new HashMap<>();
DiFyReq diFyImageReq = new DiFyReq();
diFyImageReq.setUser(ConstantStr.corpus_user);
diFyImageReq.setFlowId(qiweiToken);
inputMap.put("chat", chatList.toString());
inputMap.put("analysisScene", "1");
inputMap.put("unionId", unionId);
inputMap.put("consultantId", userId);
inputMap.put("communicateDate", DateUtil.format(maxMsgTimeItem.getMsgTime(), DatePattern.NORM_DATETIME_PATTERN));
inputMap.put("version",2);
diFyImageReq.setInputs(inputMap);
log.info("一条完整的企微信息的聊天内容,未加密:{}", contetnList.toString());
finalMaxMsgTimeItem = contetnList.stream()
.max((o1, o2) -> o1.getMsgTime().compareTo(o2.getMsgTime()))
.orElse(null);
CorpusReportDTO corpusReportDTO = new CorpusReportDTO();
corpusReportDTO.setCorpusTime(DateUtil.format(maxMsgTimeItem.getMsgTime(), DatePattern.NORM_DATETIME_PATTERN));
corpusReportDTO.setUnionId(unionId);
corpusReportDTO.setUserId(userId);
corpusReportDTO.setAnalysisScene(1l);
List<DataMaskingRule> maskingRuleItems = dataMaskingRuleService.getDataMaskingRuleListByApplicationChannel(BusinessTypeEnum.SMART_ASSISTANT_QIWEI.getCode());
RunMaskingRuleInput runMaskingRuleInput = new RunMaskingRuleInput();
runMaskingRuleInput.setDataMaskingRules(maskingRuleItems);
AtomicInteger externalcontactCount = new AtomicInteger();
contetnList.forEach(contentItem -> {
String title = "";
if (tmOdsVdqwExternalcontact.getExternalUserId().equals(contentItem.getFromUserId())) {
title = "客户:";
externalcontactCount.addAndGet(1);
} else {
title = "顾问:";
}
JSONObject contentJson = JSONObject.parseObject(contentItem.getContent());
String content = contentJson.getString("content");
// 拼接 role 和 text
String chat = title.concat(content);
runMaskingRuleInput.setOldStr(chat);
String corpusChat = dataMaskingRuleService.runMaskingRule(runMaskingRuleInput);
chatList.append(corpusChat).append("\n");
});
if (externalcontactCount.get() < 1) {
log.info("没有客户回复的语料,无需解析");
return;
}
log.info("一条完整的企微信息的聊天内容,循环次数:{} ,应该循环次数:{}, 掩码后:{}", externalcontactCount.get(), contetnList.size(), chatList.toString());
Map<String, Object> inputMap = new HashMap<>();
diFyImageReq.setUser(ConstantStr.corpus_user);
diFyImageReq.setFlowId(qiweiToken);
inputMap.put("chat", chatList.toString());
inputMap.put("analysisScene", "1");
inputMap.put("unionId", unionId);
inputMap.put("consultantId", userId);
inputMap.put("communicateDate", DateUtil.format(finalMaxMsgTimeItem.getMsgTime(), DatePattern.NORM_DATETIME_PATTERN));
inputMap.put("version", 2);
diFyImageReq.setInputs(inputMap);
corpusReportDTO.setCorpusTime(DateUtil.format(finalMaxMsgTimeItem.getMsgTime(), DatePattern.NORM_DATETIME_PATTERN));
corpusReportDTO.setUnionId(unionId);
corpusReportDTO.setUserId(userId);
corpusReportDTO.setAnalysisScene(1l);
// 使用多线程并行执行两个业务场景
log.info("第1个业务场景 优先执行。 ");
long startTime = System.currentTimeMillis();
CompletableFuture<Void> summaryTask = CompletableFuture.runAsync(() -> {
executeSummaryTask(diFyImageReq, corpusReportDTO, item, maxMsgTimeItem);
}, executor);
long startTime2 = System.currentTimeMillis();
CompletableFuture<Void> portraitTask = CompletableFuture.runAsync(() -> {
executePortraitTask(unionId, maxMsgTimeItem, chatList, corpusReportDTO);
}, executor);
// 等待两个任务完成
CompletableFuture.allOf(summaryTask, portraitTask).join();
long endTime1 = System.currentTimeMillis();
log.info("第一个业务场景(总结和分类)执行时间: {} ms", (endTime1 - startTime));
long endTime2 = System.currentTimeMillis();
log.info("第二个业务场景(用户画像)执行时间: {} ms", (endTime2 - startTime2));
}
}
// 使用多线程并行执行两个业务场景
log.info("第1个业务场景 优先执行。 ");
long startTime = System.currentTimeMillis();
OdsVdqwMessageOTD finalMaxMsgTimeItem1 = finalMaxMsgTimeItem;
CompletableFuture<Void> summaryTask = CompletableFuture.runAsync(() -> {
executeSummaryTask(diFyImageReq, corpusReportDTO, item, finalMaxMsgTimeItem1);
}, executor);
long startTime2 = System.currentTimeMillis();
OdsVdqwMessageOTD finalMaxMsgTimeItem2 = finalMaxMsgTimeItem;
String finalUnionId = unionId;
CompletableFuture<Void> portraitTask = CompletableFuture.runAsync(() -> {
executePortraitTask(finalUnionId, finalMaxMsgTimeItem2, chatList, corpusReportDTO);
}, executor);
// 等待两个任务完成
CompletableFuture.allOf(summaryTask, portraitTask).join();
long endTime1 = System.currentTimeMillis();
log.info("第一个业务场景(总结和分类)执行时间: {} ms", (endTime1 - startTime));
long endTime2 = System.currentTimeMillis();
log.info("第二个业务场景(用户画像)执行时间: {} ms", (endTime2 - startTime2));
}
public TmOdsVdqwExternalcontact getUnionId(List<String> userIds){