feed流好内容

This commit is contained in:
lxu75
2025-05-23 20:00:34 +08:00
parent 726a58df13
commit 334862149c
3 changed files with 111 additions and 38 deletions

View File

@@ -80,6 +80,9 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
@Value("${dify.community.keywordToken}")
private String keywordToken;
@Value("${dify.community.feedToken}")
private String feedToken;
@Value("${dify.community.imageToken}")
private String imageFlowId;
@@ -110,7 +113,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|| StringUtils.isBlank(communityTargetDTO.getTargetId())
|| StringUtils.isBlank(communityTargetDTO.getTargetType())
|| communityTargetDTO.getContent().isEmpty()) {
log.error("社区舆情自动化校验失败入参异常:{}",JSON.toJSONString(communityTargetDTO));
log.error("社区舆情自动化校验失败入参异常:{}", JSON.toJSONString(communityTargetDTO));
return false;
}
// 生成ai分析请求id
@@ -142,8 +145,10 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
//根据aiAnalysisRequestId更新请求日志表的difyRequest字段
syncUpdateDiFyRequest(difyCommunityTargetDTO, aiAnalysisRequestId);
//调用DiFy工作流并推送到社区
processDify(difyCommunityTargetDTO,user,difyResult, aiAnalysisRequestId);
//异步调用feed流好内容workflow
processFeedDify(difyCommunityTargetDTO, difyResult, aiAnalysisRequestId);
//异步调用舆情自动化工作流,并推送到社区
processPublicOpinionAutomationDify(difyCommunityTargetDTO, user, difyResult, aiAnalysisRequestId);
} catch (Exception e) {
log.error("舆情自动化异常:{}", e.getMessage());
//保存错误日志
@@ -156,35 +161,84 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
return true;
}
private void processDify(DifyCommunityTargetDTO difyCommunityTargetDTO ,String user, JSONArray difyResult, String aiAnalysisRequestId) throws InterruptedException, ExecutionException {
//舆情案件分析
CompletableFuture<JSONObject> caseWorkFlow = CompletableFuture.supplyAsync(() -> callCaseCommunityWorkFlow(difyCommunityTargetDTO,user, caseToken),getAsyncExecutor);
@Async
protected void processFeedDify(DifyCommunityTargetDTO difyCommunityTargetDTO, JSONArray difyResult, String aiAnalysisRequestId) {
try {
DiFyReq diFyFeedQualityReq = new DiFyReq();
diFyFeedQualityReq.setUser(user);
diFyFeedQualityReq.setFlowId(feedToken);
diFyFeedQualityReq.setInputs(difyCommunityTargetDTO);
JSONObject difFeedQualityResult = (JSONObject) diFyService.getDiFyObject(diFyFeedQualityReq);
if (difFeedQualityResult != null) {
difFeedQualityResult.put("communityType",BusinessTypeEnum.FEEDQUALITY.getCode());
//返回结果推送到社区的MQ
rocketMQTemplate.syncSend(topic, difFeedQualityResult.toString());
log.info("Feed流好内容发送回调MQ完成: {}", difFeedQualityResult);
aiAnalysisRequestLogsMapper.update(new AiAnalysisRequestLogs(),
new UpdateWrapper<AiAnalysisRequestLogs>().set("dify_response", difFeedQualityResult.toJSONString())
.eq("ai_analysis_request_id", aiAnalysisRequestId));
}
} catch (Exception e) {
saveOrUpdateError(aiAnalysisRequestId, e, BusinessTypeEnum.FEEDQUALITY.getCode());
}
}
// 内容主题关键词打标
CompletableFuture<JSONObject> keywordWorkFlow = CompletableFuture.supplyAsync(() -> callCommunityWorkFlow(difyCommunityTargetDTO,user, keywordToken),getAsyncExecutor);
@Async
protected void processPublicOpinionAutomationDify(DifyCommunityTargetDTO difyCommunityTargetDTO, String user, JSONArray difyResult, String aiAnalysisRequestId) {
try {
//舆情案件分析
CompletableFuture<JSONObject> caseWorkFlow = CompletableFuture.supplyAsync(() -> callCaseCommunityWorkFlow(difyCommunityTargetDTO, user, caseToken), getAsyncExecutor);
// litecrm线索分析
CompletableFuture<JSONObject> clueAnalysisWorkFlow = CompletableFuture.supplyAsync(() -> callCommunityWorkFlow(difyCommunityTargetDTO,user, clueAnalysisToken),getAsyncExecutor);
// 内容主题关键词打标
CompletableFuture<JSONObject> keywordWorkFlow = CompletableFuture.supplyAsync(() -> callCommunityWorkFlow(difyCommunityTargetDTO, user, keywordToken), getAsyncExecutor);
CompletableFuture<Void> allFutures = CompletableFuture.allOf(caseWorkFlow, keywordWorkFlow, clueAnalysisWorkFlow);
// 等待所有API调用完成
allFutures.get();
// 获取各个API的结果
JSONObject caseResult = caseWorkFlow.get();
JSONObject keywordResult = keywordWorkFlow.get();
JSONObject clueAnalysisResult = clueAnalysisWorkFlow.get();
difyResult.add(caseResult);
difyResult.add(keywordResult);
difyResult.add(clueAnalysisResult);
//处理结果
processingCommunityDifyResponse(caseResult,keywordResult,clueAnalysisResult, aiAnalysisRequestId,difyCommunityTargetDTO.getCommunityRequestId());
// litecrm线索分析
CompletableFuture<JSONObject> clueAnalysisWorkFlow = CompletableFuture.supplyAsync(() -> callCommunityWorkFlow(difyCommunityTargetDTO, user, clueAnalysisToken), getAsyncExecutor);
CompletableFuture<Void> allFutures = CompletableFuture.allOf(caseWorkFlow, keywordWorkFlow, clueAnalysisWorkFlow);
// 等待所有API调用完成
allFutures.get();
// 获取各个API的结果
JSONObject caseResult = caseWorkFlow.get();
JSONObject keywordResult = keywordWorkFlow.get();
JSONObject clueAnalysisResult = clueAnalysisWorkFlow.get();
difyResult.add(caseResult);
difyResult.add(keywordResult);
difyResult.add(clueAnalysisResult);
//处理结果
processingCommunityDifyResponse(caseResult, keywordResult, clueAnalysisResult, aiAnalysisRequestId, difyCommunityTargetDTO.getCommunityRequestId());
} catch (Exception e) {
saveOrUpdateError(aiAnalysisRequestId, e, BusinessTypeEnum.PUBLICOPINIONAUTOMATION.getCode());
}
}
private void saveOrUpdateError(String aiAnalysisRequestId, Exception e, String businnessType) {
//根据aiAnalysisRequestId查询是否存在错误日志无则新增有则更新retryCount+1
AiAnalysisErrors aiAnalysisErrors = aiAnalysisErrorsMapper.selectOne(
new LambdaQueryWrapper<AiAnalysisErrors>()
.eq(AiAnalysisErrors::getAiAnalysisRequestId, aiAnalysisRequestId)
.eq(AiAnalysisErrors::getAiAnalysisRequestType, businnessType)
.last("for update") // 添加悲观锁
);
if (aiAnalysisErrors != null) {
aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount() + 1);
aiAnalysisErrorsMapper.updateById(aiAnalysisErrors);
log.error("处理错误日志失败:{}", e.getMessage());
} else {
aiAnalysisErrorsMapper.insert(AiAnalysisErrors.builder()
.aiAnalysisRequestId(aiAnalysisRequestId)
.aiAnalysisErrorMessage(e.getMessage())
//舆情自动化标识
.aiAnalysisRequestType(businnessType)
.build());
}
}
/**
* 调用案件,关键词,线索工作流
*/
private JSONObject callCommunityWorkFlow(DifyCommunityTargetDTO difyCommunityTargetDTO ,String user, String token) {
log.info("开始调用关键词,线索工作流,token: {}",token);
private JSONObject callCommunityWorkFlow(DifyCommunityTargetDTO difyCommunityTargetDTO, String user, String token) {
log.info("开始调用关键词,线索工作流,token: {}", token);
DiFyReq diFyReq = new DiFyReq();
diFyReq.setUser(user);
diFyReq.setInputs(difyCommunityTargetDTO);
@@ -196,7 +250,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
/**
* 调用案件工作流
*/
private JSONObject callCaseCommunityWorkFlow(DifyCommunityTargetDTO difyCommunityTargetDTO ,String user, String token) {
private JSONObject callCaseCommunityWorkFlow(DifyCommunityTargetDTO difyCommunityTargetDTO, String user, String token) {
List<AiAnalyticBusinessConfig> aiAnalyticBusinessConfigs = aiAnalyticBusinessConfigMapper.selectList(
Wrappers.<AiAnalyticBusinessConfig>lambdaQuery()
.eq(AiAnalyticBusinessConfig::getBusinessLine, BusinessTypeEnum.COMMUNITYTARGET.getCode())
@@ -204,7 +258,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
.eq(AiAnalyticBusinessConfig::getIsDeleted, 0)
.eq(AiAnalyticBusinessConfig::getConfigVersion, 1)
);
log.info("开始调用案件工作流,token: {}",token);
log.info("开始调用案件工作流,token: {}", token);
//取出aiAnalyticBusinessConfigs里的所有configData
String configDataString = aiAnalyticBusinessConfigs.stream()
@@ -221,6 +275,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
/**
* 舆情数据脱敏
*
* @param textContent
* @param sb
* @param difyCommunityTargetDTO
@@ -236,6 +291,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
/**
* 处理图片信息
*
* @param hasImageNodeType
* @param communityTargetDTO
* @param sb
@@ -265,22 +321,24 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
/**
* 处理dify返回结果
*
* @param
*/
private void processingCommunityDifyResponse(
JSONObject caseResult,
JSONObject keywordResult,
JSONObject clueAnalysisResult,
String aiAnalysisRequestId,
String communityRequestId) {
JSONObject caseResult,
JSONObject keywordResult,
JSONObject clueAnalysisResult,
String aiAnalysisRequestId,
String communityRequestId) {
if (caseResult != null && keywordResult != null && clueAnalysisResult != null) {
caseResult.put("keyWords", keywordResult.getString("keyWords"));
caseResult.put("targetTag", clueAnalysisResult.getString("targetTag"));
caseResult.put("aiAnalysisRequestId", aiAnalysisRequestId);
if(StringUtils.isNotBlank(communityRequestId)){
if (StringUtils.isNotBlank(communityRequestId)) {
caseResult.put("communityRequestId", communityRequestId);
}
caseResult.put("communityType",BusinessTypeEnum.PUBLICOPINIONAUTOMATION.getCode());
//返回结果推送到社区的MQ
rocketMQTemplate.syncSend(topic, caseResult.toString());
log.info("舆情分析发送回调MQ完成: {}", caseResult);
@@ -302,7 +360,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
aiAnalysisRequestLogsMapper.insert(AiAnalysisRequestLogs.builder()
.aiAnalysisRequestId(aiAnalysisRequestId)
.businessRequest(message)
.difyAgentKey(clueAnalysisToken+keywordToken+caseToken)
.difyAgentKey(clueAnalysisToken + keywordToken + caseToken)
.aiAnalysisRequestType(BusinessTypeEnum.COMMUNITYTARGET.getCode())
.build());
}
@@ -631,7 +689,8 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
try {
//捞取异常表中属于社区的异常数据
List<AiAnalysisErrors> aiAnalysisErrors = aiAnalysisErrorsMapper.selectList(new LambdaQueryWrapper<AiAnalysisErrors>()
.eq(AiAnalysisErrors::getAiAnalysisRequestType, BusinessTypeEnum.COMMUNITYTARGET.getCode())
.in(AiAnalysisErrors::getAiAnalysisRequestType, BusinessTypeEnum.COMMUNITYTARGET.getCode(),
BusinessTypeEnum.PUBLICOPINIONAUTOMATION.getCode(), BusinessTypeEnum.FEEDQUALITY.getCode())
.eq(AiAnalysisErrors::getAiAnalysisErrorHandlingStatus, "0")
.lt(AiAnalysisErrors::getRetryCount, maxRetryCount));
if (aiAnalysisErrors != null && aiAnalysisErrors.size() > 0) {
@@ -670,8 +729,18 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
syncUpdateDiFyRequest(difyCommunityTargetDTO, aiAnalysisError.getAiAnalysisRequestId());
JSONArray difyResult = new JSONArray();
processDify(difyCommunityTargetDTO,user,difyResult,aiAnalysisError.getAiAnalysisRequestId());
//如果AiAnalysisRequestType是舆情自动化则调用舆情自动化方法processPublicOpinionAutomationDify如果是feed流好内容则调用feed流好内容方法processFeedDify如果是community则两个都调用
if (BusinessTypeEnum.COMMUNITYTARGET.getCode().equals(aiAnalysisError.getAiAnalysisRequestType())) {
processPublicOpinionAutomationDify(difyCommunityTargetDTO, user, difyResult, aiAnalysisError.getAiAnalysisRequestId());
processFeedDify(difyCommunityTargetDTO, difyResult, aiAnalysisError.getAiAnalysisRequestId());
}
if (BusinessTypeEnum.PUBLICOPINIONAUTOMATION.getCode().equals(aiAnalysisError.getAiAnalysisRequestType())) {
processPublicOpinionAutomationDify(difyCommunityTargetDTO, user, difyResult, aiAnalysisError.getAiAnalysisRequestId());
}
if (BusinessTypeEnum.FEEDQUALITY.getCode().equals(aiAnalysisError.getAiAnalysisRequestType())) {
processFeedDify(difyCommunityTargetDTO, difyResult, aiAnalysisError.getAiAnalysisRequestId());
}
//根据ai_analysis_request_id更新ai_analysis_errors表中的retry_count字段+1,更新status字段为1
aiAnalysisErrorsMapper.update(new AiAnalysisErrors(), new LambdaUpdateWrapper<AiAnalysisErrors>()
.eq(AiAnalysisErrors::getAiAnalysisRequestId, aiAnalysisError.getAiAnalysisRequestId())
@@ -681,14 +750,14 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
aiAnalysisRequestLogsMapper.update(new AiAnalysisRequestLogs(), new LambdaUpdateWrapper<AiAnalysisRequestLogs>()
.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisError.getAiAnalysisRequestId())
.set(AiAnalysisRequestLogs::getDifyResponse, difyResult.toJSONString()));
}else{
} else {
aiAnalysisErrorsMapper.update(new AiAnalysisErrors(), new LambdaUpdateWrapper<AiAnalysisErrors>()
.eq(AiAnalysisErrors::getAiAnalysisRequestId, aiAnalysisError.getAiAnalysisRequestId())
.set(AiAnalysisErrors::getAiAnalysisErrorHandlingStatus, "2"));
}
}
} catch (Exception e) {
log.error("补偿社区消息,AIID:{},异常:{}",aiAnalysisError.getAiAnalysisRequestId(), e);
log.error("补偿社区消息,AIID:{},异常:{}", aiAnalysisError.getAiAnalysisRequestId(), e);
aiAnalysisErrorsMapper.update(new AiAnalysisErrors(), new LambdaUpdateWrapper<AiAnalysisErrors>()
.eq(AiAnalysisErrors::getAiAnalysisRequestId, aiAnalysisError.getAiAnalysisRequestId())
.set(AiAnalysisErrors::getRetryCount, aiAnalysisError.getRetryCount() + 1));
@@ -696,7 +765,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
}
}
} catch (Exception e) {
log.error("处理社区异常消息异常:{}", e);
log.error("处理社区异常消息异常:{}", e);
}
}