联调企微客户画像场景

This commit is contained in:
ZLI263
2025-09-17 19:30:11 +08:00
parent 4aeb7d22dd
commit 1e98704fdf
2 changed files with 40 additions and 23 deletions

View File

@@ -58,6 +58,9 @@ public class CorpusPortraitServiceImpl implements CorpusPortraitService {
@Value("${dify.corpus.portrait.qiweiToken}")
private String qiweiToken;
@Value("${dify.corpus.portrait.oneToken}")
private String oneTokenPortrait;
@Value("${dify.corpus.portrait.nameplateToken}")
private String nameplateAppKey;
@@ -221,7 +224,7 @@ public class CorpusPortraitServiceImpl implements CorpusPortraitService {
Map<String, Object> inputMap = new HashMap<>();
DiFyReq diFyImageReq = new DiFyReq();
diFyImageReq.setUser(ConstantStr.corpus_user);
diFyImageReq.setFlowId(qiweiToken);
diFyImageReq.setFlowId(oneTokenPortrait);
List<DataMaskingRule> maskingRuleItems = dataMaskingRuleService.getDataMaskingRuleListByApplicationChannel(BusinessTypeEnum.CORPUS_PORTRAIT_QIWEI.getCode());
RunMaskingRuleInput runMaskingRuleInput = new RunMaskingRuleInput();
runMaskingRuleInput.setDataMaskingRules(maskingRuleItems);
@@ -263,6 +266,8 @@ public class CorpusPortraitServiceImpl implements CorpusPortraitService {
inputMap.put("communicateDate", DateUtil.format(maxMsgTimeItem.getMsgTime(), DatePattern.NORM_DATETIME_PATTERN));
inputMap.put("businessId",unionId);
inputMap.put("analysisScene", "3");
inputMap.put("businessType", BusinessTypeEnum.CORPUS_PORTRAIT_QIWEI.getCode());
diFyImageReq.setInputs(inputMap);

View File

@@ -37,6 +37,8 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicInteger;
@@ -127,29 +129,39 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
int offset = (i - 1) * pageSize;
List<OdsVdqwMessageOTD> messageList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByData(statTime, endTime, offset, pageSize, retry);
// 处理查询到的数据
// 使用 CompletableFuture 并行处理
CompletableFuture<?>[] futures = messageList.stream()
.map(item -> CompletableFuture.runAsync(() -> {
long startTime2 = System.currentTimeMillis();
log.info("处理企微语料开始 FromUserId={}, AcceptUserId={}",
item.getFromUserId(), item.getAcceptUserId());
try {
processItem(item, finalStatTime, finalEndTime);
} catch (Exception e) {
log.error("处理企微语料失败: FromUserId={}, AcceptUserId={}, 异常: {}",
item.getFromUserId(), item.getAcceptUserId(), e.getMessage(), e);
}
try {
corpusPortraitService.portraitQiWei(item, finalStatTime, finalEndTime);
} catch (Exception e) {
log.error("corpusPortrait企微画像处理异步任务执行失败", e);
}
log.info("处理企微语料 调用完成: 结束时间:{}, 总用时: {}", System.currentTimeMillis(),System.currentTimeMillis()-startTime2);
}, executor))
.toArray(CompletableFuture[]::new);
// 使用多线程并行执行两个业务场景
ExecutorService executorService = Executors.newFixedThreadPool(2);
// 使用 CompletableFuture 并行处理
messageList.forEach(item -> {
long startTime2 = System.currentTimeMillis();
log.info("处理企微语料开始 FromUserId={}, AcceptUserId={}",
item.getFromUserId(), item.getAcceptUserId());
// 等待所有任务完成
CompletableFuture.allOf(futures).join();
CompletableFuture.runAsync(() -> {
//第一个业务场景 开始: 总结和分类业务场景
try {
processItem(item, finalStatTime, finalEndTime);
} catch (Exception e) {
log.error("处理企微语料失败: FromUserId={}, AcceptUserId={}, 异常: {}",
item.getFromUserId(), item.getAcceptUserId(), e.getMessage(), e);
}
}, executorService);
CompletableFuture.runAsync(() -> {
//第一个业务场景 开始: 总结和分类业务场景
try {
corpusPortraitService.portraitQiWei(item, finalStatTime, finalEndTime);
} catch (Exception e) {
log.error("corpusPortrait企微画像处理异步任务执行失败", e);
}
}, executorService);
log.info("处理企微语料 调用完成: 结束时间:{}, 总用时: {}", System.currentTimeMillis(), System.currentTimeMillis() - startTime2);
});
}
} catch (Exception e) {