定时任务处理舆情异常
This commit is contained in:
@@ -2,13 +2,17 @@ package com.volvo.ai.analytic.center.entity;
|
|||||||
|
|
||||||
import com.baomidou.mybatisplus.annotation.*;
|
import com.baomidou.mybatisplus.annotation.*;
|
||||||
import com.volvo.common.core.base.BaseEntity;
|
import com.volvo.common.core.base.BaseEntity;
|
||||||
|
import lombok.AllArgsConstructor;
|
||||||
import lombok.Builder;
|
import lombok.Builder;
|
||||||
import lombok.Data;
|
import lombok.Data;
|
||||||
|
import lombok.NoArgsConstructor;
|
||||||
|
|
||||||
import java.time.LocalDateTime;
|
import java.time.LocalDateTime;
|
||||||
|
|
||||||
@Data
|
@Data
|
||||||
@Builder
|
@Builder
|
||||||
|
@AllArgsConstructor
|
||||||
|
@NoArgsConstructor
|
||||||
@TableName("tt_ai_analysis_errors")
|
@TableName("tt_ai_analysis_errors")
|
||||||
public class AiAnalysisErrors extends BaseEntity {
|
public class AiAnalysisErrors extends BaseEntity {
|
||||||
|
|
||||||
|
|||||||
@@ -148,26 +148,10 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
|||||||
diFyReq.setInputs(difyCommunityTargetDTO);
|
diFyReq.setInputs(difyCommunityTargetDTO);
|
||||||
//调用舆情文本分析dify工作流
|
//调用舆情文本分析dify工作流
|
||||||
difResult = (JSONObject) diFyService.getDiFyObject(diFyReq);
|
difResult = (JSONObject) diFyService.getDiFyObject(diFyReq);
|
||||||
JSONObject difResultJSONObject = difResult.getJSONObject("text");
|
//处理结果
|
||||||
if (difResultJSONObject != null) {
|
dueCommunityDifyResponse(difResult);
|
||||||
DifyCommunityTargetResult difyCommunityTargetResult = new DifyCommunityTargetResult();
|
//异步更新请求日志表的difyResponse字段
|
||||||
String targetId = difResultJSONObject.getString("targetId");
|
syncUpdateDiFyResponse(difResult, aiAnalysisRequestId);
|
||||||
String targetType = difResultJSONObject.getString("targetType");
|
|
||||||
String commentAnswer = difResultJSONObject.getString("commentAnswer");
|
|
||||||
String keyWords = difResultJSONObject.getString("keyWords");
|
|
||||||
String targetTag = difResultJSONObject.getString("targetTag");
|
|
||||||
|
|
||||||
difyCommunityTargetResult.setTargetType(targetType);
|
|
||||||
difyCommunityTargetResult.setTargetId(targetId);
|
|
||||||
difyCommunityTargetResult.setCommentAnswer(commentAnswer);
|
|
||||||
difyCommunityTargetResult.setTargetTag(targetTag);
|
|
||||||
difyCommunityTargetResult.setKeyWords(keyWords);
|
|
||||||
//返回结果推送到社区的MQ
|
|
||||||
rocketMQTemplate.syncSend(topic, JSON.toJSONString(difyCommunityTargetResult));
|
|
||||||
log.info("舆情分析发送回调MQ完成: {}", JSON.toJSONString(difyCommunityTargetResult));
|
|
||||||
//异步更新请求日志表的difyResponse字段
|
|
||||||
syncUpdateDiFyResponse(difResult, aiAnalysisRequestId);
|
|
||||||
}
|
|
||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
log.error("舆情自动化异常:{}", e.getMessage());
|
log.error("舆情自动化异常:{}", e.getMessage());
|
||||||
//保存错误日志
|
//保存错误日志
|
||||||
@@ -181,6 +165,30 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
|||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 处理dify返回结果
|
||||||
|
* @param difResult
|
||||||
|
*/
|
||||||
|
private void dueCommunityDifyResponse(JSONObject difResult) {
|
||||||
|
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");
|
||||||
|
|
||||||
|
difyCommunityTargetResult.setTargetType(targetType);
|
||||||
|
difyCommunityTargetResult.setTargetId(targetId);
|
||||||
|
difyCommunityTargetResult.setCommentAnswer(commentAnswer);
|
||||||
|
difyCommunityTargetResult.setTargetTag(targetTag);
|
||||||
|
difyCommunityTargetResult.setKeyWords(keyWords);
|
||||||
|
//返回结果推送到社区的MQ
|
||||||
|
rocketMQTemplate.syncSend(topic, JSON.toJSONString(difyCommunityTargetResult));
|
||||||
|
log.info("舆情分析发送回调MQ完成: {}", JSON.toJSONString(difyCommunityTargetResult));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@Async
|
@Async
|
||||||
protected void syncUpdateDiFyRequest(DifyCommunityTargetDTO difyCommunityTargetDTO, String aiAnalysisRequestId) {
|
protected void syncUpdateDiFyRequest(DifyCommunityTargetDTO difyCommunityTargetDTO, String aiAnalysisRequestId) {
|
||||||
aiAnalysisRequestLogsMapper.update(new AiAnalysisRequestLogs(),
|
aiAnalysisRequestLogsMapper.update(new AiAnalysisRequestLogs(),
|
||||||
@@ -526,12 +534,41 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
|||||||
@Override
|
@Override
|
||||||
public void communityMessageByTask() {
|
public void communityMessageByTask() {
|
||||||
|
|
||||||
//捞取异常表中属于社区的异常数据
|
try {
|
||||||
List<AiAnalysisErrors> aiAnalysisErrors = aiAnalysisErrorsMapper.selectList(new LambdaQueryWrapper<AiAnalysisErrors>()
|
//捞取异常表中属于社区的异常数据
|
||||||
.eq(AiAnalysisErrors::getAiAnalysisRequestType, BusinessTypeEnum.COMMUNITYTARGET.getCode())
|
List<AiAnalysisErrors> aiAnalysisErrors = aiAnalysisErrorsMapper.selectList(new LambdaQueryWrapper<AiAnalysisErrors>()
|
||||||
.eq(AiAnalysisErrors::getAiAnalysisErrorHandlingStatus, "0")
|
.eq(AiAnalysisErrors::getAiAnalysisRequestType, BusinessTypeEnum.COMMUNITYTARGET.getCode())
|
||||||
.lt(AiAnalysisErrors::getRetryCount, 4));
|
.eq(AiAnalysisErrors::getAiAnalysisErrorHandlingStatus, "0")
|
||||||
|
.lt(AiAnalysisErrors::getRetryCount, 4));
|
||||||
|
if (aiAnalysisErrors != null && aiAnalysisErrors.size() > 0) {
|
||||||
|
//根据ai_analysis_request_id获取AiAnalysisRequestLogs表中的对应的dify_request字段
|
||||||
|
for (AiAnalysisErrors aiAnalysisError : aiAnalysisErrors) {
|
||||||
|
AiAnalysisRequestLogs aiAnalysisRequestLogs = aiAnalysisRequestLogsMapper.selectOne(new LambdaQueryWrapper<AiAnalysisRequestLogs>()
|
||||||
|
.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisError.getAiAnalysisRequestId()));
|
||||||
|
if (aiAnalysisRequestLogs != null) {
|
||||||
|
//获取dify_request字段
|
||||||
|
String difyRequest = aiAnalysisRequestLogs.getDifyRequest();
|
||||||
|
if (StringUtils.isNotEmpty(difyRequest)) {
|
||||||
|
//调用dify接口
|
||||||
|
DiFyReq diFyReq = new DiFyReq();
|
||||||
|
diFyReq.setUser(user);
|
||||||
|
diFyReq.setFlowId(flowId);
|
||||||
|
diFyReq.setInputs(JSONObject.parseObject(difyRequest));
|
||||||
|
JSONObject difResult = (JSONObject) diFyService.getDiFyObject(diFyReq);
|
||||||
|
this.dueCommunityDifyResponse(difResult);
|
||||||
|
}
|
||||||
|
//根据ai_analysis_request_id更新ai_analysis_errors表中的retry_count字段+1,更新status字段为1
|
||||||
|
aiAnalysisErrorsMapper.update(new AiAnalysisErrors(), new LambdaUpdateWrapper<AiAnalysisErrors>()
|
||||||
|
.eq(AiAnalysisErrors::getAiAnalysisRequestId, aiAnalysisError.getAiAnalysisRequestId())
|
||||||
|
.set(AiAnalysisErrors::getRetryCount, aiAnalysisError.getRetryCount() + 1)
|
||||||
|
.set(AiAnalysisErrors::getAiAnalysisErrorHandlingStatus, "1"));
|
||||||
|
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
} catch (Exception e) {
|
||||||
|
throw new RuntimeException(e);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private String getUserStatus(String oldStr) {
|
private String getUserStatus(String oldStr) {
|
||||||
|
|||||||
Reference in New Issue
Block a user