联调企微客户画像场景
This commit is contained in:
@@ -78,6 +78,10 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
|
||||
@Value("${dify.corpus.qiweiToken}")
|
||||
private String qiweiToken;
|
||||
|
||||
@Value("${dify.corpus.portrait.oneToken}")
|
||||
private String oneTokenPortrait;
|
||||
|
||||
|
||||
@Value("${batch.size}")
|
||||
public int pageSize = 100;
|
||||
@Autowired
|
||||
@@ -122,6 +126,7 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
|
||||
String finalEndTime = endTime;
|
||||
|
||||
Integer total = tmOdsVdqwMessagearchivingMapper.countOdsVdqwMessageByData(statTime, endTime,retry);
|
||||
log.info("runQiWeiCorpusDify total {}", total);
|
||||
int totalPages = PageDto.getTotalPages(total, pageSize);
|
||||
|
||||
try {
|
||||
@@ -130,7 +135,8 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
|
||||
List<OdsVdqwMessageOTD> messageList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByData(statTime, endTime, offset, pageSize, retry);
|
||||
// 处理查询到的数据
|
||||
// 使用多线程并行执行两个业务场景
|
||||
ExecutorService executorService = Executors.newFixedThreadPool(2);
|
||||
ExecutorService executorService = Executors.newFixedThreadPool(1);
|
||||
log.info("待处理企微聊天列表messageList",messageList.toString());
|
||||
// 使用 CompletableFuture 并行处理
|
||||
messageList.forEach(item -> {
|
||||
long startTime2 = System.currentTimeMillis();
|
||||
@@ -147,18 +153,6 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
|
||||
}
|
||||
}, executorService);
|
||||
|
||||
|
||||
CompletableFuture.runAsync(() -> {
|
||||
//第一个业务场景 开始: 总结和分类业务场景
|
||||
try {
|
||||
corpusPortraitService.portraitQiWei(item, finalStatTime, finalEndTime);
|
||||
} catch (Exception e) {
|
||||
log.error("corpusPortrait企微画像处理异步任务执行失败", e);
|
||||
}
|
||||
}, executorService);
|
||||
|
||||
|
||||
|
||||
log.info("处理企微语料 调用完成: 结束时间:{}, 总用时: {}", System.currentTimeMillis(), System.currentTimeMillis() - startTime2);
|
||||
|
||||
});
|
||||
@@ -187,6 +181,8 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
|
||||
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);
|
||||
@@ -219,6 +215,7 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
|
||||
log.info("没有客户回复的语料,无需解析");
|
||||
return;
|
||||
}
|
||||
log.info("一条完整的企微信息的聊天内容,掩码后:{}", chatList.toString());
|
||||
|
||||
inputMap.put("chat", chatList.toString());
|
||||
inputMap.put("analysisScene", "1");
|
||||
@@ -233,22 +230,69 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
|
||||
corpusReportDTO.setUnionId(unionId);
|
||||
corpusReportDTO.setUserId(userId);
|
||||
corpusReportDTO.setAnalysisScene(1l);
|
||||
// 获取配置
|
||||
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT_QIWEI.getCode(), JSONObject.toJSONString(corpusReportDTO),null);
|
||||
log.info("runDify execDifyFlow , 企微场景,总结需求的dify返回 :{}", execDifyFlow);
|
||||
if (null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")) {
|
||||
JSONObject text = execDifyFlow.getJSONObject("outputs");
|
||||
|
||||
|
||||
|
||||
|
||||
// 使用多线程并行执行两个业务场景
|
||||
ExecutorService executorService = Executors.newFixedThreadPool(2);
|
||||
log.info("第1个业务场景, 优先执行。 ");
|
||||
long startTime = System.currentTimeMillis();
|
||||
CompletableFuture.runAsync(() -> {
|
||||
// 获取配置
|
||||
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT_QIWEI.getCode(), JSONObject.toJSONString(corpusReportDTO),null);
|
||||
log.info("runDify execDifyFlow , 企微场景,总结需求的dify返回 :{}", execDifyFlow);
|
||||
if (null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")) {
|
||||
JSONObject text = execDifyFlow.getJSONObject("outputs");
|
||||
// 发送MQ
|
||||
log.info("send mq企微总结 {}", text);
|
||||
tmTelephoneCorpusService.sendMq( CategoryEnum.ENTERPRISE_WECHAT.getCode(), text.toJSONString());
|
||||
|
||||
try {
|
||||
ttVdqwRecordMapper.insert(TtVdqwRecord.builder().acceptUserId(item.getFromUserId()).fromUserId(item.getAcceptUserId()).msgTime(maxMsgTimeItem.getMsgTime()).from_accept_user_id(item.getFromAcceptUserId()).build());
|
||||
aiAnalysisRequestLogsService.saveOrUpdateAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder().aiAnalysisRequestId(execDifyFlow.getString("aiAnalysisRequestId")).businessResponse(text.toJSONString()).build());
|
||||
} catch (Exception e) {
|
||||
log.info(" 企业语料处理保存报告异常processItem:{} ", e);
|
||||
}
|
||||
}
|
||||
log.info("runDify execDifyFlow ,企微语料 ,总结和分类业务,返回: {}", execDifyFlow);
|
||||
}, executorService);
|
||||
long endTime1 = System.currentTimeMillis();
|
||||
log.info("第一个业务场景(总结和分类)执行时间: {} ms", (endTime1 - startTime));
|
||||
//第一个业务场景, 结束
|
||||
|
||||
long startTime2 = System.currentTimeMillis();
|
||||
CompletableFuture.runAsync(() -> {
|
||||
//"analysisScene": "分析类型", // 1、企微会话 2、AI通话录音 3、AI铭牌(客流) 4、AI铭牌(试驾)
|
||||
DiFyReq diFyImageReq2 = new DiFyReq();
|
||||
diFyImageReq2.setUser(ConstantStr.corpus_user);
|
||||
Map<String, Object> inputMap2 = new HashMap<>();
|
||||
inputMap2.put("businessId", unionId);
|
||||
inputMap2.put("communicateDate", DateUtil.format(maxMsgTimeItem.getMsgTime(), DatePattern.NORM_DATETIME_PATTERN));
|
||||
inputMap2.put("analysisScene", "1");
|
||||
inputMap2.put("version", 2);
|
||||
inputMap2.put("chat", chatList.toString());
|
||||
diFyImageReq2.setInputs(inputMap2);
|
||||
// 创建新的DiFyReq对象以避免线程安全问题
|
||||
diFyImageReq2.setFlowId(oneTokenPortrait);
|
||||
log.info("runDify execDifyFlow ,客户画像场景 ,token-{} , 对象: {}", oneTokenPortrait, diFyImageReq2);
|
||||
JSONObject execDifyFlowForPortrait = diFyService.executeDifyFlow(diFyImageReq2, BusinessTypeEnum.CORPUS_PORTRAIT_NAMEPLATE.getCode(),
|
||||
JSONObject.toJSONString(corpusReportDTO), null);
|
||||
log.info("runDify execDifyFlow ,客户画像场景 ,返回-{} ", execDifyFlowForPortrait);
|
||||
|
||||
JSONObject text = execDifyFlowForPortrait.getJSONObject("outputs");
|
||||
// 发送MQ
|
||||
log.info("send mq企微总结 {}", text);
|
||||
log.info("send mq企微客户画像 {}", text);
|
||||
tmTelephoneCorpusService.sendMq( CategoryEnum.ENTERPRISE_WECHAT.getCode(), text.toJSONString());
|
||||
|
||||
try {
|
||||
ttVdqwRecordMapper.insert(TtVdqwRecord.builder().acceptUserId(item.getFromUserId()).fromUserId(item.getAcceptUserId()).msgTime(maxMsgTimeItem.getMsgTime()).from_accept_user_id(item.getFromAcceptUserId()).build());
|
||||
aiAnalysisRequestLogsService.saveOrUpdateAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder().aiAnalysisRequestId(execDifyFlow.getString("aiAnalysisRequestId")).businessResponse(text.toJSONString()).build());
|
||||
} catch (Exception e) {
|
||||
log.info(" 企业语料处理保存报告异常processItem:{} ", e);
|
||||
}
|
||||
}
|
||||
aiAnalysisRequestLogsService.saveOrUpdateAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder().aiAnalysisRequestId(execDifyFlowForPortrait.getString("aiAnalysisRequestId")).businessResponse(text.toJSONString()).build());
|
||||
|
||||
log.info("runDify execDifyFlow ,企微语料 ,客户画像场景,返回: {}", execDifyFlowForPortrait);
|
||||
|
||||
}, executorService);
|
||||
long endTime2 = System.currentTimeMillis();
|
||||
log.info("第二个业务场景(用户画像)执行时间: {} ms", (endTime2 - startTime2));
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user