舆情二期

This commit is contained in:
lxu75
2025-04-25 17:40:31 +08:00
parent f11066d7af
commit 778a6ebf3f

View File

@@ -3,6 +3,7 @@ package com.volvo.ai.analytic.center.service.impl;
import cn.hutool.core.date.DateUtil;
import com.alibaba.cloud.commons.lang.StringUtils;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
@@ -95,7 +96,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
// 异步保存请求日志
syncSaveRequestLogs(message, aiAnalysisRequestId);
JSONObject difResult = new JSONObject();
JSONArray difyResult = new JSONArray();
try {
boolean hasImageNodeType = communityTargetDTO.getContent().stream()
.anyMatch(contentNode -> NodeTypeEnum.IMAGE.getCode().equals(contentNode.getNodeType()));
@@ -136,20 +137,17 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
JSONObject caseResult = caseWorkFlow.get();
JSONObject keywordResult = keywordWorkFlow.get();
JSONObject clueAnalysisResult = clueAnalysisWorkFlow.get();
difyResult.add(caseResult);
difyResult.add(keywordResult);
difyResult.add(clueAnalysisResult);
//处理结果
processingCommunityDifyResponse(difResult,caseResult,keywordResult,clueAnalysisResult);
//异步更新请求日志表的difyResponse字段
syncUpdateDiFyResponse(difResult, aiAnalysisRequestId);
//调用舆情文本分析dify工作流
// difResult = (JSONObject) diFyService.getDiFyObject(diFyReq);
processingCommunityDifyResponse(caseResult,keywordResult,clueAnalysisResult,aiAnalysisRequestId);
} catch (Exception e) {
log.error("舆情自动化异常:{}", e.getMessage());
//保存错误日志
aiAnalysisErrorsMapper.insert(AiAnalysisErrors.builder()
.aiAnalysisRequestId(aiAnalysisRequestId)
.difyResponse(difResult.toJSONString())
.difyResponse(difyResult.toJSONString())
.aiAnalysisErrorMessage(e.getMessage())
.aiAnalysisRequestType(BusinessTypeEnum.COMMUNITYTARGET.getCode())
.build());
@@ -215,28 +213,23 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
/**
* 处理dify返回结果
* @param difResult
* @param
*/
private void processingCommunityDifyResponse(JSONObject difResult,
private void processingCommunityDifyResponse(
JSONObject caseResult,
JSONObject keywordResult,
JSONObject clueAnalysisResult) {
if (difResult != null) {
DifyCommunityTargetResult difyCommunityTargetResult = new DifyCommunityTargetResult();
String targetId = difResult.getString("targetId");
String targetType = difResult.getString("targetType");
String commentAnswer = difResult.getString("commentAnswer");
String keyWords = difResult.getString("keyWords");
String targetTag = difResult.getString("targetTag");
JSONObject clueAnalysisResult,
String aiAnalysisRequestId) {
if (caseResult != null && keywordResult != null && clueAnalysisResult != null) {
difyCommunityTargetResult.setTargetType(targetType);
difyCommunityTargetResult.setTargetId(targetId);
difyCommunityTargetResult.setCommentAnswer(commentAnswer);
difyCommunityTargetResult.setTargetTag(targetTag);
difyCommunityTargetResult.setKeyWords(keyWords);
caseResult.put("keyWords", keywordResult.getString("keyWords"));
caseResult.put("targetTag", clueAnalysisResult.getString("targetTag"));
//返回结果推送到社区的MQ
rocketMQTemplate.syncSend(topic, JSON.toJSONString(difyCommunityTargetResult));
log.info("舆情分析发送回调MQ完成: {}", JSON.toJSONString(difyCommunityTargetResult));
rocketMQTemplate.syncSend(topic, caseResult.toString());
log.info("舆情分析发送回调MQ完成: {}", caseResult);
aiAnalysisRequestLogsMapper.update(new AiAnalysisRequestLogs(),
new UpdateWrapper<AiAnalysisRequestLogs>().set("dify_response", caseResult.toJSONString())
.eq("ai_analysis_request_id", aiAnalysisRequestId));
}
}
@@ -247,13 +240,6 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
.eq("ai_analysis_request_id", aiAnalysisRequestId));
}
@Async
protected void syncUpdateDiFyResponse(JSONObject difResult, String aiAnalysisRequestId) {
aiAnalysisRequestLogsMapper.update(new AiAnalysisRequestLogs(),
new UpdateWrapper<AiAnalysisRequestLogs>().set("dify_response", difResult.toJSONString())
.eq("ai_analysis_request_id", aiAnalysisRequestId));
}
@Async
protected void syncSaveRequestLogs(String message, String aiAnalysisRequestId) {
aiAnalysisRequestLogsMapper.insert(AiAnalysisRequestLogs.builder()