Merge remote-tracking branch 'origin/feature-20250520-release' into feature_20250519_nameplate
# Conflicts: # ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/enums/BusinessTypeEnum.java
This commit is contained in:
@@ -0,0 +1,68 @@
|
||||
package com.volvo.ai.analytic.center.config;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.aop.interceptor.AsyncUncaughtExceptionHandler;
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.scheduling.annotation.AsyncConfigurer;
|
||||
import org.springframework.scheduling.annotation.EnableAsync;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.util.concurrent.*;
|
||||
|
||||
/**
|
||||
* 异步任务线程池装配类
|
||||
* @author gubin
|
||||
* @date 2022-04-14
|
||||
*/
|
||||
@EnableAsync
|
||||
@Slf4j
|
||||
@Component
|
||||
public class AsyncTaskExecutePool implements AsyncConfigurer {
|
||||
|
||||
@Value("${task.pool.corePoolSize}")
|
||||
private int corePoolSize;
|
||||
|
||||
@Value("${task.pool.maxPoolSize}")
|
||||
private int maxPoolSize;
|
||||
|
||||
@Value("${task.pool.queueCapacity}")
|
||||
private int queueCapacity;
|
||||
|
||||
@Value("${task.pool.keepAliveSeconds}")
|
||||
private int keepAliveSeconds;
|
||||
|
||||
|
||||
@Bean
|
||||
@Override
|
||||
public Executor getAsyncExecutor() {
|
||||
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
|
||||
//核心线程池大小
|
||||
executor.setCorePoolSize(corePoolSize);
|
||||
//最大线程数
|
||||
executor.setMaxPoolSize(maxPoolSize);
|
||||
//队列容量
|
||||
executor.setQueueCapacity(queueCapacity);
|
||||
//活跃时间
|
||||
executor.setKeepAliveSeconds(keepAliveSeconds);
|
||||
//线程名字前缀
|
||||
executor.setThreadNamePrefix("async-task-");
|
||||
// setRejectedExecutionHandler:当pool已经达到max size的时候,如何处理新任务
|
||||
// CallerRunsPolicy:不在新线程中执行任务,而是由调用者所在的线程来执行
|
||||
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
|
||||
executor.initialize();
|
||||
return executor;
|
||||
}
|
||||
|
||||
@Override
|
||||
public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() {
|
||||
return (throwable, method, objects) -> {
|
||||
log.error("===="+throwable.getMessage()+"====", throwable);
|
||||
log.error("exception method:"+method.getName());
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
package com.volvo.ai.analytic.center.mapper;
|
||||
|
||||
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
|
||||
import com.volvo.ai.analytic.center.entity.TtAnalysisResultInfo;
|
||||
import org.apache.ibatis.annotations.Mapper;
|
||||
|
||||
/**
|
||||
* @description Ai分析结果明细表
|
||||
* @author BEJSON
|
||||
* @date 2025-03-04
|
||||
*/
|
||||
@Mapper
|
||||
public interface TtAnalysisResultInfoMapper extends BaseMapper<TtAnalysisResultInfo> {
|
||||
|
||||
|
||||
}
|
||||
@@ -2,7 +2,10 @@ package com.volvo.ai.analytic.center.mq;
|
||||
|
||||
import com.alibaba.fastjson.JSONObject;
|
||||
import com.volvo.ai.analytic.center.entity.TmAnalysisResult;
|
||||
import com.volvo.ai.analytic.center.entity.TtAnalysisResultInfo;
|
||||
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
|
||||
import com.volvo.ai.analytic.center.mapper.AiAnalysisRequestLogsMapper;
|
||||
import com.volvo.ai.analytic.center.mapper.TtAnalysisResultInfoMapper;
|
||||
import com.volvo.ai.analytic.center.service.TmAnalysisResultService;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
@@ -12,7 +15,9 @@ import org.apache.rocketmq.spring.core.RocketMQListener;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.util.Arrays;
|
||||
import java.util.Date;
|
||||
import java.util.List;
|
||||
|
||||
@Slf4j
|
||||
@Component
|
||||
@@ -26,6 +31,9 @@ public class CorpushIsLikeConsumer implements RocketMQListener<MessageExt>{
|
||||
|
||||
@Autowired
|
||||
private AiAnalysisRequestLogsMapper aiAnalysisRequestLogsMapper;
|
||||
|
||||
@Autowired
|
||||
private TtAnalysisResultInfoMapper ttAnalysisResultInfoMapper;
|
||||
@Override
|
||||
public void onMessage(MessageExt messageExt) {
|
||||
|
||||
@@ -44,13 +52,20 @@ public class CorpushIsLikeConsumer implements RocketMQListener<MessageExt>{
|
||||
log.info(" 回调的aiAnalysisRequestType为空:{} ", aiAnalysisRequestType);
|
||||
return;
|
||||
}
|
||||
// "isLike": "1" // 1:点赞,2:点踩
|
||||
tmAnalysisResultService.saveTmCorpusReport(TmAnalysisResult.builder()
|
||||
TmAnalysisResult tmAnalysisResult = TmAnalysisResult.builder()
|
||||
.aiAnalysisRequestId(execDifyFlow.getString("analysisRecordId"))
|
||||
.analysisResult(execDifyFlow.toJSONString())
|
||||
.analysisType(aiAnalysisRequestType)
|
||||
.updateTime(new Date())
|
||||
.build());
|
||||
.build();
|
||||
// "isLike": "1" // 1:点赞,2:点踩
|
||||
tmAnalysisResultService.saveTmCorpusReport(tmAnalysisResult);
|
||||
List<String> aiAnalysisRequestIdList = Arrays.asList(BusinessTypeEnum.SMART_ASSISTANT.getCode(),BusinessTypeEnum.SMART_ASSISTANT_QIWEI.getCode(),BusinessTypeEnum.SMART_ASSISTANT_NAMEPLATE.getCode());
|
||||
if(aiAnalysisRequestIdList.contains(aiAnalysisRequestType)){
|
||||
// 保存明细
|
||||
ttAnalysisResultInfoMapper.insert(TtAnalysisResultInfo.builder().analysisResultId(tmAnalysisResult.getId()).analysisType(aiAnalysisRequestType).aiAnalysisRequestId(analysisRecordId).analysisResult(execDifyFlow.toJSONString()).build());
|
||||
|
||||
}
|
||||
} catch (Exception e) {
|
||||
log.info(" corpushIsLikeConsumer AI结果回传处理异常:{} ", e);
|
||||
}
|
||||
|
||||
@@ -24,7 +24,6 @@ public class DataMaskingRuleServiceImpl extends ServiceImpl<DataMaskingRuleMappe
|
||||
//根据适用渠道applicationChannel获取数据脱敏规则List
|
||||
List<DataMaskingRule> dataMaskingRuleList = this.lambdaQuery()
|
||||
.like(DataMaskingRule::getApplicationChannel, applicationChannel)
|
||||
.eq(DataMaskingRule::getApplicationChannel, applicationChannel)
|
||||
.eq(DataMaskingRule::getRuleStatus, YesOrNoConstants.YES)
|
||||
.eq(DataMaskingRule::getIsDeleted, YesOrNoConstants.NO)
|
||||
.list();
|
||||
|
||||
@@ -3,10 +3,12 @@ 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;
|
||||
import com.baomidou.mybatisplus.core.conditions.update.UpdateWrapper;
|
||||
import com.baomidou.mybatisplus.core.toolkit.Wrappers;
|
||||
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
|
||||
import com.volvo.ai.analytic.center.constant.Constant;
|
||||
import com.volvo.ai.analytic.center.dto.req.*;
|
||||
@@ -15,6 +17,7 @@ import com.volvo.ai.analytic.center.entity.*;
|
||||
import com.volvo.ai.analytic.center.enums.*;
|
||||
import com.volvo.ai.analytic.center.mapper.AiAnalysisErrorsMapper;
|
||||
import com.volvo.ai.analytic.center.mapper.AiAnalysisRequestLogsMapper;
|
||||
import com.volvo.ai.analytic.center.mapper.AiAnalyticBusinessConfigMapper;
|
||||
import com.volvo.ai.analytic.center.mapper.MqMessageRecordMapper;
|
||||
import com.volvo.ai.analytic.center.service.DataMaskingRuleService;
|
||||
import com.volvo.ai.analytic.center.service.DiFyService;
|
||||
@@ -35,6 +38,9 @@ import org.springframework.util.CollectionUtils;
|
||||
import java.text.SimpleDateFormat;
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.*;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.Executor;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
@Slf4j
|
||||
@@ -65,8 +71,14 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
||||
@Value("${dify.user}")
|
||||
private String user;
|
||||
|
||||
@Value("${dify.community.targetToken}")
|
||||
private String flowId;
|
||||
@Value("${dify.community.clueAnalysisToken}")
|
||||
private String clueAnalysisToken;
|
||||
|
||||
@Value("${dify.community.caseToken}")
|
||||
private String caseToken;
|
||||
|
||||
@Value("${dify.community.keywordToken}")
|
||||
private String keywordToken;
|
||||
|
||||
@Value("${dify.community.imageToken}")
|
||||
private String imageFlowId;
|
||||
@@ -74,6 +86,15 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
||||
@Value("${rocketmq.producer.topic}")
|
||||
private String topic;
|
||||
|
||||
@Value("${community.maxRetryCount}")
|
||||
private int maxRetryCount;
|
||||
|
||||
@Autowired
|
||||
private AiAnalyticBusinessConfigMapper aiAnalyticBusinessConfigMapper;
|
||||
|
||||
@Autowired
|
||||
private Executor getAsyncExecutor;
|
||||
|
||||
/**
|
||||
* 处理Mq消息
|
||||
*
|
||||
@@ -85,8 +106,11 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
||||
public boolean processMessageByMQ(String message) {
|
||||
log.info("communityProcessMessageByMQ message: {}", message);
|
||||
CommunityTargetDTO communityTargetDTO = JSON.parseObject(message, CommunityTargetDTO.class);
|
||||
if (communityTargetDTO == null) {
|
||||
log.error("communityProcessMessageByMQ message is null");
|
||||
if (communityTargetDTO == null
|
||||
|| StringUtils.isBlank(communityTargetDTO.getTargetId())
|
||||
|| StringUtils.isBlank(communityTargetDTO.getTargetType())
|
||||
|| communityTargetDTO.getContent().isEmpty()) {
|
||||
log.error("社区舆情自动化校验失败入参异常:{}",JSON.toJSONString(communityTargetDTO));
|
||||
return false;
|
||||
}
|
||||
// 生成ai分析请求id
|
||||
@@ -94,7 +118,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()));
|
||||
@@ -112,28 +136,19 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
||||
DifyCommunityTargetDTO difyCommunityTargetDTO = new DifyCommunityTargetDTO();
|
||||
difyCommunityTargetDTO.setTargetId(communityTargetDTO.getTargetId());
|
||||
difyCommunityTargetDTO.setTargetType(communityTargetDTO.getTargetType());
|
||||
difyCommunityTargetDTO.setCommunityRequestId(communityTargetDTO.getCommunityRequestId());
|
||||
//脱敏处理
|
||||
dataMasking(textContent, sb, difyCommunityTargetDTO);
|
||||
|
||||
//根据aiAnalysisRequestId更新请求日志表的difyRequest字段
|
||||
syncUpdateDiFyRequest(difyCommunityTargetDTO, aiAnalysisRequestId);
|
||||
|
||||
DiFyReq diFyReq = new DiFyReq();
|
||||
diFyReq.setUser(user);
|
||||
diFyReq.setFlowId(flowId);
|
||||
diFyReq.setInputs(difyCommunityTargetDTO);
|
||||
//调用舆情文本分析dify工作流
|
||||
difResult = (JSONObject) diFyService.getDiFyObject(diFyReq);
|
||||
//处理结果
|
||||
processingCommunityDifyResponse(difResult);
|
||||
//异步更新请求日志表的difyResponse字段
|
||||
syncUpdateDiFyResponse(difResult, aiAnalysisRequestId);
|
||||
//调用DiFy工作流,并推送到社区
|
||||
processDify(difyCommunityTargetDTO,user,difyResult, aiAnalysisRequestId);
|
||||
} catch (Exception e) {
|
||||
log.error("舆情自动化异常:{}", e.getMessage());
|
||||
//保存错误日志
|
||||
aiAnalysisErrorsMapper.insert(AiAnalysisErrors.builder()
|
||||
.aiAnalysisRequestId(aiAnalysisRequestId)
|
||||
.difyResponse(difResult.toJSONString())
|
||||
.aiAnalysisErrorMessage(e.getMessage())
|
||||
.aiAnalysisRequestType(BusinessTypeEnum.COMMUNITYTARGET.getCode())
|
||||
.build());
|
||||
@@ -141,6 +156,69 @@ 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);
|
||||
|
||||
// 内容主题关键词打标
|
||||
CompletableFuture<JSONObject> keywordWorkFlow = CompletableFuture.supplyAsync(() -> callCommunityWorkFlow(difyCommunityTargetDTO,user, keywordToken),getAsyncExecutor);
|
||||
|
||||
// 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());
|
||||
}
|
||||
|
||||
/**
|
||||
* 调用案件,关键词,线索工作流
|
||||
*/
|
||||
private JSONObject callCommunityWorkFlow(DifyCommunityTargetDTO difyCommunityTargetDTO ,String user, String token) {
|
||||
log.info("开始调用关键词,线索工作流,token: {}",token);
|
||||
DiFyReq diFyReq = new DiFyReq();
|
||||
diFyReq.setUser(user);
|
||||
diFyReq.setInputs(difyCommunityTargetDTO);
|
||||
diFyReq.setFlowId(token);
|
||||
JSONObject difResult = (JSONObject) diFyService.getDiFyObject(diFyReq);
|
||||
return difResult;
|
||||
}
|
||||
|
||||
/**
|
||||
* 调用案件工作流
|
||||
*/
|
||||
private JSONObject callCaseCommunityWorkFlow(DifyCommunityTargetDTO difyCommunityTargetDTO ,String user, String token) {
|
||||
List<AiAnalyticBusinessConfig> aiAnalyticBusinessConfigs = aiAnalyticBusinessConfigMapper.selectList(
|
||||
Wrappers.<AiAnalyticBusinessConfig>lambdaQuery()
|
||||
.eq(AiAnalyticBusinessConfig::getBusinessLine, BusinessTypeEnum.COMMUNITYTARGET.getCode())
|
||||
.eq(AiAnalyticBusinessConfig::getConfigType, BusinessTypeEnum.CASE.getCode())
|
||||
.eq(AiAnalyticBusinessConfig::getIsDeleted, 0)
|
||||
.eq(AiAnalyticBusinessConfig::getConfigVersion, 1)
|
||||
);
|
||||
log.info("开始调用案件工作流,token: {}",token);
|
||||
|
||||
//取出aiAnalyticBusinessConfigs里的所有configData
|
||||
String configDataString = aiAnalyticBusinessConfigs.stream()
|
||||
.map(AiAnalyticBusinessConfig::getConfigData)
|
||||
.collect(Collectors.joining(" "));
|
||||
difyCommunityTargetDTO.setSentiment(configDataString);
|
||||
DiFyReq diFyReq = new DiFyReq();
|
||||
diFyReq.setUser(user);
|
||||
diFyReq.setInputs(difyCommunityTargetDTO);
|
||||
diFyReq.setFlowId(token);
|
||||
JSONObject difResult = (JSONObject) diFyService.getDiFyObject(diFyReq);
|
||||
return difResult;
|
||||
}
|
||||
|
||||
/**
|
||||
* 舆情数据脱敏
|
||||
* @param textContent
|
||||
@@ -187,25 +265,28 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
||||
|
||||
/**
|
||||
* 处理dify返回结果
|
||||
* @param difResult
|
||||
* @param
|
||||
*/
|
||||
private void processingCommunityDifyResponse(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");
|
||||
private void processingCommunityDifyResponse(
|
||||
JSONObject caseResult,
|
||||
JSONObject keywordResult,
|
||||
JSONObject clueAnalysisResult,
|
||||
String aiAnalysisRequestId,
|
||||
String communityRequestId) {
|
||||
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"));
|
||||
caseResult.put("aiAnalysisRequestId", aiAnalysisRequestId);
|
||||
if(StringUtils.isNotBlank(communityRequestId)){
|
||||
caseResult.put("communityRequestId", communityRequestId);
|
||||
}
|
||||
//返回结果推送到社区的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));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -216,19 +297,12 @@ 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()
|
||||
.aiAnalysisRequestId(aiAnalysisRequestId)
|
||||
.businessRequest(message)
|
||||
.difyAgentKey(flowId)
|
||||
.difyAgentKey(clueAnalysisToken+keywordToken+caseToken)
|
||||
.aiAnalysisRequestType(BusinessTypeEnum.COMMUNITYTARGET.getCode())
|
||||
.build());
|
||||
}
|
||||
@@ -434,7 +508,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
||||
|
||||
DiFyReq diFyReq = new DiFyReq();
|
||||
diFyReq.setUser(user);
|
||||
diFyReq.setFlowId(flowId);
|
||||
// diFyReq.setFlowId(flowId);
|
||||
diFyReq.setInputs(record);
|
||||
JSONObject difResult = (JSONObject) diFyService.getDiFyObject(diFyReq);
|
||||
output.setHandleStatus(HandleStatusEnum.ANALYSIS_NORMAL.getCode());
|
||||
@@ -559,7 +633,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
||||
List<AiAnalysisErrors> aiAnalysisErrors = aiAnalysisErrorsMapper.selectList(new LambdaQueryWrapper<AiAnalysisErrors>()
|
||||
.eq(AiAnalysisErrors::getAiAnalysisRequestType, BusinessTypeEnum.COMMUNITYTARGET.getCode())
|
||||
.eq(AiAnalysisErrors::getAiAnalysisErrorHandlingStatus, "0")
|
||||
.lt(AiAnalysisErrors::getRetryCount, 4));
|
||||
.lt(AiAnalysisErrors::getRetryCount, maxRetryCount));
|
||||
if (aiAnalysisErrors != null && aiAnalysisErrors.size() > 0) {
|
||||
//根据ai_analysis_request_id获取AiAnalysisRequestLogs表中的对应的dify_request字段
|
||||
for (AiAnalysisErrors aiAnalysisError : aiAnalysisErrors) {
|
||||
@@ -588,26 +662,25 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
||||
DifyCommunityTargetDTO difyCommunityTargetDTO = new DifyCommunityTargetDTO();
|
||||
difyCommunityTargetDTO.setTargetId(communityTargetDTO.getTargetId());
|
||||
difyCommunityTargetDTO.setTargetType(communityTargetDTO.getTargetType());
|
||||
difyCommunityTargetDTO.setCommunityRequestId(communityTargetDTO.getCommunityRequestId());
|
||||
//脱敏处理
|
||||
dataMasking(textContent, sb, difyCommunityTargetDTO);
|
||||
|
||||
//根据aiAnalysisRequestId更新请求日志表的difyRequest字段
|
||||
syncUpdateDiFyRequest(difyCommunityTargetDTO, aiAnalysisError.getAiAnalysisRequestId());
|
||||
|
||||
DiFyReq diFyReq = new DiFyReq();
|
||||
diFyReq.setUser(user);
|
||||
diFyReq.setFlowId(flowId);
|
||||
diFyReq.setInputs(difyCommunityTargetDTO);
|
||||
//调用舆情文本分析dify工作流
|
||||
JSONObject difResult = (JSONObject) diFyService.getDiFyObject(diFyReq);
|
||||
//处理结果
|
||||
processingCommunityDifyResponse(difResult);
|
||||
JSONArray difyResult = new JSONArray();
|
||||
processDify(difyCommunityTargetDTO,user,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())
|
||||
.set(AiAnalysisErrors::getRetryCount, aiAnalysisError.getRetryCount() + 1)
|
||||
.set(AiAnalysisErrors::getAiAnalysisErrorHandlingStatus, "1"));
|
||||
//根据ai_analysis_request_id更新tt_ai_analysis_request_logs表中的dify_response字段
|
||||
aiAnalysisRequestLogsMapper.update(new AiAnalysisRequestLogs(), new LambdaUpdateWrapper<AiAnalysisRequestLogs>()
|
||||
.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisError.getAiAnalysisRequestId())
|
||||
.set(AiAnalysisRequestLogs::getDifyResponse, difyResult.toJSONString()));
|
||||
}else{
|
||||
aiAnalysisErrorsMapper.update(new AiAnalysisErrors(), new LambdaUpdateWrapper<AiAnalysisErrors>()
|
||||
.eq(AiAnalysisErrors::getAiAnalysisRequestId, aiAnalysisError.getAiAnalysisRequestId())
|
||||
|
||||
@@ -24,6 +24,7 @@ public class TmAnalysisResultServiceImpl extends ServiceImpl<TmAnalysisResultMap
|
||||
if (oldTmCorpusReport == null) {
|
||||
return tmCorpusReportMapper.insert(tmAnalysisResult) > 0;
|
||||
} else {
|
||||
tmAnalysisResult.setId(oldTmCorpusReport.getId());
|
||||
return tmCorpusReportMapper.update(tmAnalysisResult, queryWrapper) > 0;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,18 +1,33 @@
|
||||
package com.volvo.ai.analytic.center.utils;
|
||||
|
||||
import cn.hutool.core.lang.Snowflake;
|
||||
import cn.hutool.core.util.IdUtil;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import javax.annotation.PostConstruct;
|
||||
import java.net.InetAddress;
|
||||
import java.net.UnknownHostException;
|
||||
|
||||
@Slf4j
|
||||
@Component
|
||||
public class AiAnalysisUtils {
|
||||
|
||||
/**
|
||||
* 生成ai分析请求id
|
||||
* @param businessType
|
||||
* @return
|
||||
*/
|
||||
private static Snowflake snowflake;
|
||||
|
||||
@PostConstruct
|
||||
public void init() {
|
||||
try {
|
||||
// 使用 IP 生成唯一的 workerId
|
||||
long workerId = ipToWorkerId(getLocalHostIP());
|
||||
snowflake = IdUtil.getSnowflake(workerId, 0); // datacenterId = 0
|
||||
log.info("Initialized Snowflake with workerId: {}", workerId);
|
||||
} catch (Exception e) {
|
||||
log.error("Failed to initialize Snowflake", e);
|
||||
throw new RuntimeException("Snowflake initialization failed");
|
||||
}
|
||||
}
|
||||
|
||||
public static String getAiAnalysisRequestId(String businessType) {
|
||||
if (businessType == null || businessType.trim().isEmpty()) {
|
||||
log.error("businessType is null or empty, using default value 'unknown'");
|
||||
@@ -20,18 +35,26 @@ public class AiAnalysisUtils {
|
||||
}
|
||||
|
||||
try {
|
||||
long snowflakeId = IdUtil.getSnowflakeNextId();
|
||||
if (snowflakeId == 0) {
|
||||
log.error("Failed to generate Snowflake ID");
|
||||
throw new RuntimeException("Failed to generate Snowflake ID");
|
||||
}
|
||||
// 使用自定义的 Snowflake 实例生成 ID
|
||||
long snowflakeId = snowflake.nextId();
|
||||
String aiAnalysisRequestId = businessType + "-" + snowflakeId;
|
||||
log.info("Generated AI analysis request ID: {}", aiAnalysisRequestId);
|
||||
return aiAnalysisRequestId;
|
||||
} catch (Exception e) {
|
||||
log.error("Error generating AI analysis request ID", e);
|
||||
//生成唯一字符串
|
||||
return businessType + "-"+IdUtil.fastSimpleUUID();
|
||||
return businessType + "-" + IdUtil.fastSimpleUUID();
|
||||
}
|
||||
}
|
||||
|
||||
// 获取本机 IP
|
||||
private static String getLocalHostIP() throws UnknownHostException {
|
||||
return InetAddress.getLocalHost().getHostAddress();
|
||||
}
|
||||
|
||||
// 将 IP 转换为合法的 workerId (0 ~ 31)
|
||||
private static long ipToWorkerId(String ip) {
|
||||
String[] parts = ip.replaceAll("[^\\d.]", "").split("\\.");
|
||||
int lastOctet = Integer.parseInt(parts[parts.length - 1]);
|
||||
return lastOctet % 32; // 限制范围 [0, 31]
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user