舆情事件配置化

This commit is contained in:
lxu75
2025-05-13 15:01:52 +08:00
parent fbe657d770
commit 6aa26915a4
3 changed files with 44 additions and 3 deletions

View File

@@ -8,6 +8,7 @@ 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.*;
@@ -16,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;
@@ -87,6 +89,9 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
@Value("${community.maxRetryCount}")
private int maxRetryCount;
@Autowired
private AiAnalyticBusinessConfigMapper aiAnalyticBusinessConfigMapper;
@Autowired
private Executor getAsyncExecutor;
@@ -150,7 +155,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
private void processDify(DifyCommunityTargetDTO difyCommunityTargetDTO ,String user, JSONArray difyResult, String aiAnalysisRequestId) throws InterruptedException, ExecutionException {
//舆情案件分析
CompletableFuture<JSONObject> caseWorkFlow = CompletableFuture.supplyAsync(() -> callCommunityWorkFlow(difyCommunityTargetDTO,user, caseToken),getAsyncExecutor);
CompletableFuture<JSONObject> caseWorkFlow = CompletableFuture.supplyAsync(() -> callCaseCommunityWorkFlow(difyCommunityTargetDTO,user, caseToken),getAsyncExecutor);
// 内容主题关键词打标
CompletableFuture<JSONObject> keywordWorkFlow = CompletableFuture.supplyAsync(() -> callCommunityWorkFlow(difyCommunityTargetDTO,user, keywordToken),getAsyncExecutor);
@@ -176,7 +181,33 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
* 调用案件,关键词,线索工作流
*/
private JSONObject callCommunityWorkFlow(DifyCommunityTargetDTO difyCommunityTargetDTO ,String user, String token) {
log.info("开始调用案件,关键词,线索工作流,token: {}",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);