Merge branch 'master' into feature-20250603-spokesman‌

This commit is contained in:
lxu75
2025-06-18 15:53:27 +08:00
29 changed files with 1291 additions and 151 deletions

View File

@@ -0,0 +1,81 @@
package com.volvo.ai.analytic.center.controller;
import com.fasterxml.jackson.core.JsonParser;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.DeserializationFeature;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.volvo.ai.analytic.center.dto.req.AnalysisQueryReq;
import com.volvo.ai.analytic.center.dto.req.AnalysisReq;
import com.volvo.ai.analytic.center.dto.resp.AnalysisResp;
import com.volvo.ai.analytic.center.service.AiAnalysisDifyService;
import com.volvo.common.core.util.ResultMsg;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.validation.annotation.Validated;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
@RestController()
@Api(tags = "分析中心接口")
@Slf4j
@RefreshScope
@RequestMapping("analysis")
public class AiAnalysisDifyController {
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Autowired
private AiAnalysisDifyService aiAnalysisDifyService;
@PostMapping("/updateByAiId")
@ApiOperation(value = "更新dify结果")
public ResultMsg<Object> updateByAiId(@RequestBody String message) {
log.info("updateByAiId message: {}", message);
aiAnalysisDifyService.updateAiDifyResult(message);
return ResultMsg.ok("ok");
}
@PostMapping("/aiAnalyze")
@ApiOperation(value = "Ai解析接口")
public AnalysisResp<Object> aiAnalyze(@Validated @RequestBody String message) {
log.info("aiAnalyze data: {}",message);
ObjectMapper objectMapper = new ObjectMapper();
objectMapper.configure(JsonParser.Feature.ALLOW_UNQUOTED_CONTROL_CHARS, true); // 允许未转义的控制字符
objectMapper.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false);
try {
AnalysisReq analysisReq = objectMapper.readValue(message, AnalysisReq.class);
analysisReq.validate();
return aiAnalysisDifyService.aiAnalyze(analysisReq);
} catch (JsonProcessingException e) {
log.info("aiAnalyze data error: {}",e);
return AnalysisResp.failed("解析失败");
}
}
@PostMapping("/query")
@ApiOperation(value = "Ai解析结果查询")
public AnalysisResp<Object> query(@RequestBody @Validated AnalysisQueryReq analysisQueryReq) {
log.info("aiAnalyze query data: {}",analysisQueryReq);
return aiAnalysisDifyService.query(analysisQueryReq);
}
@PostMapping("/callback")
@ApiOperation(value = "Ai解析结果查询")
public AnalysisResp<Object> testCallback(@RequestBody AnalysisResp analysisResp) {
log.info("aiAnalyze testCallback data: {}",analysisResp);
return AnalysisResp.success("ok");
}
}

View File

@@ -1,48 +0,0 @@
package com.volvo.ai.analytic.center.controller;
import com.alibaba.fastjson.JSONObject;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.common.core.util.ResultMsg;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
@RestController
@Api(tags = "AiDifyResult")
@RequestMapping("")
@Slf4j
@RefreshScope
public class AiDifyResultController {
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@PostMapping("/updateByAiId")
@ApiOperation(value = "更新dify结果")
public ResultMsg<Object> updateByAiId(@RequestBody String message) {
JSONObject messageJson = JSONObject.parseObject(message);
AiAnalysisRequestLogs aiAnalysisRequestLogs = new AiAnalysisRequestLogs();
aiAnalysisRequestLogs.setAiAnalysisRequestId(messageJson.getString("aiAnalysisRequestId"));
aiAnalysisRequestLogs.setDifyResponse(messageJson.getString("difyResponse"));
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(aiAnalysisRequestLogs);
return ResultMsg.ok("ok");
}
}

View File

@@ -0,0 +1,81 @@
package com.volvo.ai.analytic.center.controller;
import com.alibaba.fastjson.JSONObject;
import com.volvo.ai.analytic.center.service.AiAnalysisDifyService;
import com.volvo.ai.analytic.center.utils.RedisZSetUtil;
import com.volvo.common.core.util.ResultMsg;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.web.bind.annotation.*;
import java.time.Instant;
import java.util.Set;
@RestController()
@Api(tags = "分析中心接口")
@Slf4j
@RefreshScope
@RequestMapping("redis")
public class RedisTestController {
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Autowired
private AiAnalysisDifyService aiAnalysisDifyService;
@Autowired
private RedisTemplate<String, String> redisTemplate;
@Autowired
private RedisZSetUtil redisZSetUtil;
@PostMapping("/testRedis")
public void test(@RequestBody String message) {
String zsetKey = "myZSet";
JSONObject json = JSONObject.parseObject(message);
// 添加成员并设置过期时间戳
long expireTime1 = Instant.now().plusSeconds(30).toEpochMilli(); // 30秒后过期
long expireTime2 = Instant.now().plusSeconds(60).toEpochMilli(); // 60秒后过期
redisZSetUtil.addWithExpire(zsetKey, "member1", expireTime1);
redisZSetUtil.addWithExpire(zsetKey, "member2", expireTime2);
// 获取未过期的成员
Set<String> validMembers = redisZSetUtil.getValidMembers(zsetKey);
log.info("Valid members: " + validMembers);
}
@PostMapping("/queryRedis")
@ApiOperation(value = "queryRedis")
public ResultMsg<Object> queryRedis(@RequestBody String key) {
log.info("aiAnalyze queryRedis : {}", redisTemplate.opsForZSet().zCard(key));
Set<String> getredisSet = redisTemplate.opsForZSet().range(key, 0, -1);
log.info("queryRedis:{}", getredisSet);
return ResultMsg.ok(getredisSet);
}
@PostMapping("/del")
@ApiOperation(value = "del")
public ResultMsg<Object> del(@RequestParam("key") String key, @RequestParam("nameAuth") String nameAuth) {
log.info("aiAnalyze queryRedis : {}",key);
if(nameAuth.equals("kb2f78sQCUvg")){
redisTemplate.opsForZSet().removeRange(key, 0, -1);
Set<String> getRedisSet = redisTemplate.opsForZSet().range(key, 0, -1);
log.info("removeRange:{}" , getRedisSet);
redisTemplate.delete(key);
Set<String> getredisSet2 = redisTemplate.opsForZSet().range(key, 0, -1);
log.info("delete:{}", getredisSet2);
return ResultMsg.ok(getredisSet2);
}
return ResultMsg.ok(null);
}
}

View File

@@ -0,0 +1,74 @@
package com.volvo.ai.analytic.center.job;
import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
import com.volvo.ai.analytic.center.feign.DiFyFeign;
import com.volvo.ai.analytic.center.mapper.AiAnalysisRequestLogsMapper;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.common.core.util.ResultMsg;
import com.xxl.job.core.context.XxlJobHelper;
import com.xxl.job.core.handler.annotation.XxlJob;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RestController;
import java.util.Arrays;
import java.util.List;
/**
*
*/
@Slf4j
@Component
@RestController
public class AiAnalysisDifyJob {
@Autowired
private AiAnalysisRequestLogsMapper aiAnalysisRequestLogsMapper;
@Autowired
private DiFyFeign diFyFeign;
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
/**
* 失败的查询处理
*/
@XxlJob("workflowRunIdFaile")
@PostMapping("workflowRunIdFaile")
public ResultMsg workflowRunIdFaile(@RequestBody String paramJson) {
try {
// 获取任务参数
String param = XxlJobHelper.getJobParam();
if(StringUtils.isEmpty(param)){
param = paramJson;
}
List<String> workflowRunIdList = Arrays.asList(param.split(","));
LambdaQueryWrapper<AiAnalysisRequestLogs> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.in(AiAnalysisRequestLogs::getWorkflowRunId, workflowRunIdList);
queryWrapper.eq(AiAnalysisRequestLogs::getIsDeleted, "0");
List<AiAnalysisRequestLogs> aiAnalysisRequestLogsList = aiAnalysisRequestLogsMapper.selectList(queryWrapper);
aiAnalysisRequestLogsList.stream().forEach(aiAnalysisRequestLogs -> {
JSONObject jsonResult = diFyFeign.queryWorkFlowById("Bearer "+aiAnalysisRequestLogs.getDifyAgentKey(),aiAnalysisRequestLogs.getWorkflowRunId());
String outputs = jsonResult.getString("outputs");
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder()
.aiAnalysisRequestId(aiAnalysisRequestLogs.getAiAnalysisRequestId())
.difyResponse(outputs)
.build());
});
} catch (Exception e) {
log.error("processMessageByTask 定时任务补偿处理消息异常",e.getMessage());
throw new RuntimeException(e);
}
return ResultMsg.ok();
}
}

View File

@@ -1,56 +0,0 @@
package com.volvo.ai.analytic.center.job;
import com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper;
import com.volvo.ai.analytic.center.service.TmOdsVdqwMessagearchivingService;
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
import com.volvo.common.core.util.ResultMsg;
import com.xxl.job.core.context.XxlJobHelper;
import com.xxl.job.core.handler.annotation.XxlJob;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RestController;
@Slf4j
@Component
@RestController
public class DccCorpusJob {
@Autowired
private TmOdsVdqwMessagearchivingService tmOdsVdqwMessagearchivingService;
@Autowired
private TmTelephoneCorpusService tmTelephoneCorpusService;
@Autowired
private TmTelephoneCorpusMapper tmTelephoneCorpusMapper;
/**
* dcc语料处理
*/
@XxlJob("dccCorpusJob")
public ResultMsg dccCorpusJob(@RequestBody String paramJson) {
try {
// 获取任务参数
String param = XxlJobHelper.getJobParam();
if(StringUtils.isEmpty(param)){
param = paramJson;
}
// 分页查询 过滤已跑批并发送的
tmOdsVdqwMessagearchivingService.runQiWeiCorpusDify(param);
} catch (Exception e) {
log.error("processMessageByTask 定时任务补偿处理消息异常",e.getMessage());
throw new RuntimeException(e);
}
return ResultMsg.ok();
}
}

View File

@@ -0,0 +1,11 @@
package com.volvo.ai.analytic.center.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.volvo.ai.analytic.center.entity.TcBusinessType;
import org.apache.ibatis.annotations.Mapper;
@Mapper
public interface TcBusinessTypeMapper extends BaseMapper<TcBusinessType> {
}

View File

@@ -0,0 +1,91 @@
package com.volvo.ai.analytic.center.mq;
import com.alibaba.fastjson.JSONObject;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.volvo.ai.analytic.center.dto.resp.AnalysisDifyResultDTO;
import com.volvo.ai.analytic.center.dto.resp.AnalysisResp;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.ai.analytic.center.service.DiFyService;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.http.HttpEntity;
import org.springframework.http.HttpHeaders;
import org.springframework.http.ResponseEntity;
import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.client.RestTemplate;
/**
* @ClassName AnalysisDifyMqConsumer
* @Description AI解析MQ-Callback处理
* @Author renzhen
* @Date 2025-03-04 10:18
* @Version 1.0
**/
@Slf4j
@Component
@RefreshScope
@RestController
@RocketMQMessageListener(topic = "${rocketmq.consumer.analysisDify.callbackTopic}",consumerGroup = "${rocketmq.consumer.analysisDify.callbackGroup}",
instanceName = "analysisDifyCallbackMqConsumer",
consumeThreadNumber = 40,
enableMsgTrace = true)
public class AnalysisDifyCallbackMqConsumer implements RocketMQListener<MessageExt> {
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@Value("${dify.corpus.checkDccRepeat}")
private String checkDccRepeat;
@Autowired
private DiFyService diFyService;
@Autowired
private RestTemplate restTemplate;
private final ObjectMapper objectMapper = new ObjectMapper();
@Override
public void onMessage(MessageExt messageExt) {
long startTime = System.currentTimeMillis();
try {
log.info("analysisDifyCallbackMqConsumer 当前线程: {}, 线程ID: {}", Thread.currentThread().getName(), Thread.currentThread().getId());
String message = new String(messageExt.getBody());
log.info("analysisDifyCallbackMqConsumer message: " + message);
AnalysisDifyResultDTO analysisRestDto = JSONObject.parseObject(message, AnalysisDifyResultDTO.class);
AiAnalysisRequestLogs aiAnalysisRequestLogs = aiAnalysisRequestLogsService.queryByAiAnalysisRequestId(analysisRestDto.getAiAnalysisRequestId());
AnalysisResp analysisResp = new AnalysisResp();
analysisResp.setAiAnalysisRequestId(analysisRestDto.getAiAnalysisRequestId());
JSONObject difyJson = JSONObject.parseObject(analysisRestDto.getDifyResponse());
String outputs = difyJson.getString("outputs");
analysisResp.setData(outputs);
HttpHeaders headers = new HttpHeaders();
headers.set("Content-Type", "application/json");
// 封装请求体和请求头
HttpEntity<AnalysisResp> requestEntity = new HttpEntity<>(analysisResp, headers);
ResponseEntity<String> response = restTemplate.postForEntity(aiAnalysisRequestLogs.getCallbackUrl(), requestEntity, String.class); // 响应类型);
if (response.getStatusCode().is2xxSuccessful()) {
log.info("analysisDifyCallbackMqConsumer aiAnalysisRequestId{},回调请求成功url:{}: " ,analysisResp.getAiAnalysisRequestId(), aiAnalysisRequestLogs.getCallbackUrl());
} else {
log.info("analysisDifyCallbackMqConsumer aiAnalysisRequestId{},回调请求失败url:{}: " ,analysisResp.getAiAnalysisRequestId(), aiAnalysisRequestLogs.getCallbackUrl());
}
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder().aiAnalysisRequestId(analysisResp.getAiAnalysisRequestId()).businessResponse(outputs).build());
log.info(" analysisDifyCallbackMqConsumer耗时{}", System.currentTimeMillis() - startTime);
} catch (Exception e) {
log.info(" analysisDifyCallbackMqConsumer mq 处理失败:{}", e.getMessage());
}
}
}

View File

@@ -0,0 +1,141 @@
package com.volvo.ai.analytic.center.mq;
import com.alibaba.fastjson.JSONObject;
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
import com.volvo.ai.analytic.center.service.AiAnalysisErrorsService;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.ai.analytic.center.service.DiFyService;
import com.volvo.ai.analytic.center.utils.ConstantStr;
import com.volvo.ai.analytic.center.utils.RedisCounterRateLimiter;
import com.volvo.ai.analytic.center.utils.RedisLockService;
import com.volvo.ai.analytic.center.utils.RedisZSetUtil;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.RestController;
import java.time.Instant;
import java.util.concurrent.CompletableFuture;
/**
* @ClassName AnalysisDifyMqConsumer
* @Description AI解析 MQ处理
* @Author renzhen
* @Date 2025-03-04 10:18
* @Version 1.0
**/
@Slf4j
@Component
@RefreshScope
@RestController
@RocketMQMessageListener(topic = "${rocketmq.consumer.analysisDify.topic}",consumerGroup = "${rocketmq.consumer.analysisDify.group}",
instanceName = "analysisDifyMqConsumer",
consumeThreadNumber = 20,
enableMsgTrace = true)
public class AnalysisDifyMqConsumer implements RocketMQListener<MessageExt> {
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@Value("${dify.corpus.checkDccRepeat}")
private String checkDccRepeat;
@Value("${rocketmq.consumer.analysisDify.difyLimit}")
private int difyLimit;
@Value("${rocketmq.consumer.analysisDify.expire}")
private Long expire;
@Autowired
private DiFyService diFyService;
@Autowired
private RedisCounterRateLimiter redisCounterRateLimiter;
@Autowired
private AiAnalysisErrorsService aiAnalysisErrorsService;
@Autowired
private RedisTemplate<String, String> redisTemplate;
@Autowired
private RedisZSetUtil redisZSetUtil;
@Autowired
private RedisLockService redisLockService;
@Override
public void onMessage(MessageExt messageExt) {
log.info("analysisDifyMqConsumer 当前线程: {}, 线程ID: {}", Thread.currentThread().getName(), Thread.currentThread().getId());
long startTime = System.currentTimeMillis();
String lockKey = "LOCK_PREFIX:" + ConstantStr.DIFY_COUNTERRATELIMIT;
// RLock lock = redissonClient.getLock("LOCK_PREFIX:" + ConstantStr.DIFY_COUNTERRATELIMIT);
try {
boolean isLocked = redisLockService.tryLock(lockKey,ConstantStr.DIFY_COUNTERRATELIMIT, expire); // 不等待,立即尝试获取
if (isLocked) {
Long count = redisZSetUtil.zCard(ConstantStr.DIFY_COUNTERRATELIMIT);
log.info("redis计数数量: {}", count);
if(count>=difyLimit){
log.info("analysisDifyMqConsumer 请求dify超过基数 {},稍后请求: " + redisZSetUtil.getValidMembers(ConstantStr.DIFY_COUNTERRATELIMIT));
throw new RuntimeException("请求dify超过基数 " );
}
String message = new String(messageExt.getBody());
log.info("analysisDifyMqConsumer message: " + message);
DiFyReq difyReq = JSONObject.parseObject(message, DiFyReq.class);
JSONObject difyRequest = JSONObject.parseObject(JSONObject.toJSONString(difyReq.getInputs()), JSONObject.class);
String aiAnalysisRequestId = difyRequest.getString("aiAnalysisRequestId");
long expireTime = Instant.now().plusSeconds(expire).toEpochMilli(); // 60秒后过期
redisZSetUtil.addWithExpire(ConstantStr.DIFY_COUNTERRATELIMIT, aiAnalysisRequestId, expireTime);
CompletableFuture<JSONObject> future = diFyService.asyncExecuteDifyFlow(difyReq);
future.thenAccept(result -> {
log.info("异步处理asyncExecuteDifyFlow完成aiAnalysisRequestId: {},处理结果:{}", aiAnalysisRequestId, result);
JSONObject data = result.getJSONObject("data");
if(null == result || data.get("status").equals("failed")){
aiAnalysisErrorsService.saveAiAnalysisErrors(AiAnalysisErrors.builder()
.aiAnalysisRequestId(aiAnalysisRequestId)
.aiAnalysisRequestType(difyRequest.getString("aiAnalysisRequestType"))
.aiAnalysisErrorHandlingStatus("0")
.aiAnalysisErrorMessage(data.getString("error"))
.build());
}
}).exceptionally(ex -> {
aiAnalysisErrorsService.saveAiAnalysisErrors(AiAnalysisErrors.builder()
.aiAnalysisRequestId(aiAnalysisRequestId)
.aiAnalysisRequestType(difyRequest.getString("aiAnalysisRequestType"))
.aiAnalysisErrorHandlingStatus("0")
.aiAnalysisErrorMessage(ex.getMessage())
.build());
log.error("异步处理asyncExecuteDifyFlow 失败aiAnalysisRequestId: {} ,{}", aiAnalysisRequestId,ex.getMessage());
return null;
});
redisLockService.releaseLock(lockKey, ConstantStr.DIFY_COUNTERRATELIMIT);
log.info("analysisDifyMqConsumer处理完成耗时{}", System.currentTimeMillis() - startTime);
}else{
log.info("redis锁未获取到 ");
throw new RuntimeException("锁超时 " );
}
} catch (Exception e){
log.error("analysisDifyMqConsumer 异常:{}", e.getMessage());
throw new RuntimeException("请求dify超过基数 " );
}finally {
redisLockService.releaseLock(lockKey, ConstantStr.DIFY_COUNTERRATELIMIT);
// 建议记录解锁日志
log.info("释放锁成功锁KEY: {}", ConstantStr.DIFY_COUNTERRATELIMIT);
}
}
}

View File

@@ -0,0 +1,19 @@
package com.volvo.ai.analytic.center.service;
import com.volvo.ai.analytic.center.dto.req.AnalysisQueryReq;
import com.volvo.ai.analytic.center.dto.req.AnalysisReq;
import com.volvo.ai.analytic.center.dto.resp.AnalysisDifyResultDTO;
import com.volvo.ai.analytic.center.dto.resp.AnalysisResp;
public interface AiAnalysisDifyService {
public boolean updateAiDifyResult(String message);
AnalysisResp aiAnalyze(AnalysisReq analysisReq);
AnalysisResp query(AnalysisQueryReq analysisQueryReq);
}

View File

@@ -6,6 +6,6 @@ import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
public interface AiAnalysisRequestLogsService extends IService<AiAnalysisRequestLogs> {
boolean saveAiAnalysisRequestLogs(AiAnalysisRequestLogs aiAnalysisRequestLogs);
AiAnalysisRequestLogs queryByAiAnalysisRequestId(String aiAnalysisRequestId);
AiAnalysisRequestLogs queryAiAnalysisRequestLogsByBusinessReponse(String sourceId);
}

View File

@@ -3,6 +3,8 @@ package com.volvo.ai.analytic.center.service;
import com.alibaba.fastjson.JSONObject;
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
import java.util.concurrent.CompletableFuture;
public interface DiFyService {
@@ -13,4 +15,6 @@ public interface DiFyService {
public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType, String businessData, String aiAnalysisRequestId);
public JSONObject executeDifyFlow(DiFyReq diFyReq);
public CompletableFuture<JSONObject> asyncExecuteDifyFlow(DiFyReq diFyReq);
}

View File

@@ -0,0 +1,204 @@
package com.volvo.ai.analytic.center.service.impl;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.volvo.ai.analytic.center.dto.req.AnalysisQueryReq;
import com.volvo.ai.analytic.center.dto.req.AnalysisReq;
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
import com.volvo.ai.analytic.center.dto.resp.AnalysisDifyResultDTO;
import com.volvo.ai.analytic.center.dto.resp.AnalysisResp;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
import com.volvo.ai.analytic.center.entity.TcBusinessType;
import com.volvo.ai.analytic.center.mapper.TcBusinessTypeMapper;
import com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper;
import com.volvo.ai.analytic.center.service.AiAnalysisDifyService;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
import com.volvo.ai.analytic.center.utils.AiAnalysisUtils;
import com.volvo.ai.analytic.center.utils.ConstantStr;
import com.volvo.ai.analytic.center.utils.RedisCounterRateLimiter;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.stream.Collectors;
@Slf4j
@Service
public class AiAnalysisDifyServiceImpl implements AiAnalysisDifyService {
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@Autowired
private TmTelephoneCorpusMapper tmTelephoneCorpusMapper;
@Autowired
private TmTelephoneCorpusService tmTelephoneCorpusService;
@Autowired
private TcBusinessTypeMapper tcBusinessTypeMapper;
@Value("${rocketmq.producer.analysisDify.topic}")
private String analysisDifyTopic;
@Value("${rocketmq.producer.analysisDify.callbackTopic}")
private String callbackTopic;
@Resource
private RocketMQTemplate rocketMqTemplate;
@Autowired
private RedisCounterRateLimiter redisCounterRateLimiter;
@Autowired
private RedisTemplate<String, String> redisTemplate;
@Override
public boolean updateAiDifyResult(String message) {
if(StringUtils.isNotEmpty(message)){
AnalysisDifyResultDTO analysisResp = JSONObject.parseObject(message, AnalysisDifyResultDTO.class);
// 计数-1
Long removedCount = redisTemplate.opsForZSet().remove(ConstantStr.DIFY_COUNTERRATELIMIT, analysisResp.getAiAnalysisRequestId());
log.info("redisdecrement计数数量: " + removedCount);
AiAnalysisRequestLogs oldAiAnalysisRequestLogs = Optional.ofNullable(aiAnalysisRequestLogsService.queryByAiAnalysisRequestId(analysisResp.getAiAnalysisRequestId()))
.orElseThrow(() -> new IllegalArgumentException("AiAnalysisRequestId查询对象为空"));
AiAnalysisRequestLogs aiAnalysisRequestLogs = new AiAnalysisRequestLogs();
aiAnalysisRequestLogs.setAiAnalysisRequestId(analysisResp.getAiAnalysisRequestId());
aiAnalysisRequestLogs.setDifyResponse(analysisResp.getDifyResponse());
aiAnalysisRequestLogs.setWorkflowRunId(analysisResp.getWorkflowRunId());
aiAnalysisRequestLogs.setWorkflowAppId(analysisResp.getWorkflowRunId());
aiAnalysisRequestLogs.setWorkUserId(analysisResp.getWorkUserId());
JSONObject difyJson = JSONObject.parseObject(analysisResp.getDifyResponse());
aiAnalysisRequestLogs.setBusinessResponse(difyJson.getString("outputs"));
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(aiAnalysisRequestLogs);
if(StringUtils.isNotEmpty(oldAiAnalysisRequestLogs.getCallbackUrl())){
// 发送 mq
sendMq(callbackTopic, analysisResp);
}
return true;
}
return false;
}
@Override
public AnalysisResp aiAnalyze(AnalysisReq analysisReq) {
Optional.ofNullable(analysisReq).orElseThrow(() -> {
log.info("请求Ai解析对象为空");
return new IllegalArgumentException("请求Ai解析对象为空");
});
String aiAnalysisRequestId = StringUtils.isEmpty(analysisReq.getAiAnalysisRequestId())? AiAnalysisUtils.getAiAnalysisRequestId(analysisReq.getAiAnalysisRequestType()):analysisReq.getAiAnalysisRequestId();
Map<String, TcBusinessType> queryTcBusinessType = queryTcBusinessType();
TcBusinessType tcBusinessType = Optional.ofNullable(queryTcBusinessType.get(analysisReq.getAiAnalysisRequestType()))
.filter(businessType -> StringUtils.isNotEmpty(businessType.getWorkflowApiKey()))
.orElseThrow(() -> {
log.info("接入业务类型未配置!");
return new IllegalArgumentException("接入业务类型未配置!");
});
DiFyReq diFyReq = createDiFyReq(analysisReq, tcBusinessType, aiAnalysisRequestId);
saveAiAnalysisRequestLogs(analysisReq, diFyReq, aiAnalysisRequestId);
sendMq(analysisDifyTopic, diFyReq);
return AnalysisResp.success(analysisReq.getData(),aiAnalysisRequestId);
}
private DiFyReq createDiFyReq(AnalysisReq analysisReq, TcBusinessType tcBusinessType, String aiAnalysisRequestId) {
DiFyReq diFyReq = new DiFyReq();
diFyReq.setUser(StringUtils.isEmpty(tcBusinessType.getWorkflowUser()) ? analysisReq.getAiAnalysisRequestType().concat("_USER") : tcBusinessType.getWorkflowUser());
diFyReq.setFlowId(tcBusinessType.getWorkflowApiKey());
JSONObject difyRequest = JSONObject.parseObject(JSONObject.toJSONString(analysisReq.getData()), JSONObject.class);
difyRequest.put("aiAnalysisRequestId", aiAnalysisRequestId);
difyRequest.put("aiAnalysisRequestType", tcBusinessType.getBusinessRequestType() );
diFyReq.setInputs(difyRequest);
return diFyReq;
}
private void saveAiAnalysisRequestLogs(AnalysisReq analysisReq, DiFyReq diFyReq, String aiAnalysisRequestId) {
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder()
.aiAnalysisRequestId(aiAnalysisRequestId)
.businessRequest(JSONObject.toJSONString(analysisReq.getData()))
.difyAgentKey(diFyReq.getFlowId())
.difyRequest(JSON.toJSONString(diFyReq))
.aiAnalysisRequestType(analysisReq.getAiAnalysisRequestType())
.callbackUrl(analysisReq.getCallbackUrl())
.build());
}
@Override
public AnalysisResp query(AnalysisQueryReq analysisQueryReq) {
if(null == analysisQueryReq){
log.info("请求Ai查询对象为空");
return AnalysisResp.failed("请求Ai查询对象为空");
}
Optional<AiAnalysisRequestLogs> aiAnalysisRequestLogsOpt = Optional.ofNullable(
aiAnalysisRequestLogsService.queryByAiAnalysisRequestId(analysisQueryReq.getAiAnalysisRequestId())
);
return aiAnalysisRequestLogsOpt.map(logs -> {
if ("true".equals(analysisQueryReq.getRetryAnalyze())) {
log.info("请求Ai查询 需要重新生成AI解析 aiAnalysisRequestId:{}, retryAnalyze{}", analysisQueryReq.getAiAnalysisRequestId(), analysisQueryReq.getRetryAnalyze());
return aiAnalyze(AnalysisReq.builder()
.data(logs.getBusinessRequest())
.aiAnalysisRequestType(logs.getAiAnalysisRequestType())
.callbackUrl(logs.getCallbackUrl())
.build());
}
return AnalysisResp.success(logs.getBusinessResponse(), logs.getAiAnalysisRequestId());
}).orElseGet(() -> {
log.info("查询的AI解析不存在aiAnalysisRequestId:{}", analysisQueryReq.getAiAnalysisRequestId());
return AnalysisResp.failed("查询的AI解析不存在");
});
}
public Map<String, TcBusinessType> queryTcBusinessType() {
LambdaQueryWrapper<TcBusinessType> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(TcBusinessType::getIsDeleted, "0");
List<TcBusinessType> tcBusinessTypeList =tcBusinessTypeMapper.selectList(queryWrapper);
return tcBusinessTypeList.stream()
.collect(Collectors.toMap(
TcBusinessType::getBusinessRequestType,
tcBusinessType -> tcBusinessType
));
}
private void sendMq(String topic, Object message){
rocketMqTemplate.asyncSend(topic, MessageBuilder.withPayload(message).build(),
new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
log.info("请求AI解析发送MQ成功消息体:{}", message);
}
@Override
public void onException(Throwable e) {
log.error("请求AI解析发送MQ异常消息体:{}, 异常:", message, e);
}
}, 10000);
}
}

View File

@@ -32,6 +32,13 @@ public class AiAnalysisRequestLogsServiceImpl extends ServiceImpl<AiAnalysisRequ
}
}
@Override
public AiAnalysisRequestLogs queryByAiAnalysisRequestId(String aiAnalysisRequestId) {
LambdaQueryWrapper<AiAnalysisRequestLogs> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisRequestId);
return aiAnalysisRequestLogsMapper.selectOne(queryWrapper);
}
@Override
public AiAnalysisRequestLogs queryAiAnalysisRequestLogsByBusinessReponse(String sourceId) {
return aiAnalysisRequestLogsMapper.queryAiAnalysisRequestLogsByBusinessReponse(sourceId);

View File

@@ -19,6 +19,7 @@ import org.springframework.stereotype.Service;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
@Slf4j
@Service
@@ -123,4 +124,13 @@ public class DiFyServiceImpl implements DiFyService{
return data;
}
@Override
public CompletableFuture<JSONObject> asyncExecuteDifyFlow(DiFyReq diFyReq) {
Map<String, Object> map = new HashMap<>();
map.put("inputs",diFyReq.getInputs());
map.put("user",diFyReq.getUser());
return CompletableFuture.supplyAsync(() -> diFyFeign.runWorkflows("Bearer "+diFyReq.getFlowId(),map));
}
}

View File

@@ -77,6 +77,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;
@@ -89,7 +92,6 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
@Autowired
private AiAnalyticBusinessConfigMapper aiAnalyticBusinessConfigMapper;
/**
* 处理Mq消息
*
@@ -105,7 +107,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
@@ -137,8 +139,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());
//保存错误日志
@@ -151,28 +155,101 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
return true;
}
private void processDify(DifyCommunityTargetDTO difyCommunityTargetDTO ,String user, JSONArray difyResult, String aiAnalysisRequestId) {
//舆情案件分析
JSONObject caseResult = callCaseCommunityWorkFlow(difyCommunityTargetDTO,user, caseToken);
@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);
//aiAnalysisRequestId 根据查询日志表获取dify_response设置到aiAnalysisRequestLogs表中
updaterDifyResponse(aiAnalysisRequestId, difFeedQualityResult, BusinessTypeEnum.FEEDQUALITY.getCode());
}
} catch (Exception e) {
log.error("Feed流好内容异常:{}", e.getMessage());
saveOrUpdateError(aiAnalysisRequestId, e, BusinessTypeEnum.FEEDQUALITY.getCode());
}
}
private void updaterDifyResponse(String aiAnalysisRequestId, JSONObject difFeedQualityResult, String businessType) {
AiAnalysisRequestLogs aiAnalysisRequestLogs = aiAnalysisRequestLogsMapper.selectOne(
new LambdaQueryWrapper<AiAnalysisRequestLogs>()
.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisRequestId)
.last("for update"));
// 内容主题关键词打标
JSONObject keywordResult = callCommunityWorkFlow(difyCommunityTargetDTO,user, keywordToken);
if (aiAnalysisRequestLogs == null) {
log.error("未找到对应的 AiAnalysisRequestLogs请求ID: {}", aiAnalysisRequestId);
throw new RuntimeException("未找到对应的请求日志");
}
// litecrm线索分析
JSONObject clueAnalysisResult = callCommunityWorkFlow(difyCommunityTargetDTO,user, clueAnalysisToken);
String difyResponseStr = aiAnalysisRequestLogs.getDifyResponse();
JSONObject difyResponse = StringUtils.isNotBlank(difyResponseStr)
? JSONObject.parseObject(difyResponseStr)
: new JSONObject();
difyResult.add(caseResult);
difyResult.add(keywordResult);
difyResult.add(clueAnalysisResult);
//处理结果
processingCommunityDifyResponse(caseResult,keywordResult,clueAnalysisResult, aiAnalysisRequestId,difyCommunityTargetDTO.getCommunityRequestId());
difyResponse.put(businessType, difFeedQualityResult);
aiAnalysisRequestLogsMapper.update(null,
new UpdateWrapper<AiAnalysisRequestLogs>()
.set("dify_response", difyResponse.toJSONString())
.eq("ai_analysis_request_id", aiAnalysisRequestId));
}
@Async
protected void processPublicOpinionAutomationDify(DifyCommunityTargetDTO difyCommunityTargetDTO, String user, JSONArray difyResult, String aiAnalysisRequestId) {
try {
//舆情案件分析
JSONObject caseResult = callCaseCommunityWorkFlow(difyCommunityTargetDTO, user, caseToken);
// 内容主题关键词打标
JSONObject keywordResult = callCommunityWorkFlow(difyCommunityTargetDTO, user, keywordToken);
// litecrm线索分析
JSONObject clueAnalysisResult = callCommunityWorkFlow(difyCommunityTargetDTO, user, clueAnalysisToken);
difyResult.add(caseResult);
difyResult.add(keywordResult);
difyResult.add(clueAnalysisResult);
//处理结果
processingCommunityDifyResponse(caseResult, keywordResult, clueAnalysisResult, aiAnalysisRequestId, difyCommunityTargetDTO.getCommunityRequestId());
} catch (Exception e) {
log.error("舆情自动化异常:{}", e.getMessage());
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);
@@ -184,7 +261,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())
@@ -192,7 +269,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()
@@ -209,6 +286,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
/**
* 舆情数据脱敏
*
* @param textContent
* @param sb
* @param difyCommunityTargetDTO
@@ -224,6 +302,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
/**
* 处理图片信息
*
* @param hasImageNodeType
* @param communityTargetDTO
* @param sb
@@ -253,28 +332,28 @@ 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);
aiAnalysisRequestLogsMapper.update(new AiAnalysisRequestLogs(),
new UpdateWrapper<AiAnalysisRequestLogs>().set("dify_response", caseResult.toJSONString())
.eq("ai_analysis_request_id", aiAnalysisRequestId));
updaterDifyResponse(aiAnalysisRequestId, caseResult, BusinessTypeEnum.PUBLICOPINIONAUTOMATION.getCode());
}
}
@@ -290,7 +369,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());
}
@@ -619,7 +698,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) {
@@ -658,25 +738,34 @@ 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())) {
log.info("补偿社区全部workflow", aiAnalysisError.getAiAnalysisRequestId());
processPublicOpinionAutomationDify(difyCommunityTargetDTO, user, difyResult, aiAnalysisError.getAiAnalysisRequestId());
processFeedDify(difyCommunityTargetDTO, difyResult, aiAnalysisError.getAiAnalysisRequestId());
}
if (BusinessTypeEnum.PUBLICOPINIONAUTOMATION.getCode().equals(aiAnalysisError.getAiAnalysisRequestType())) {
log.info("补偿舆情自动化全部workflow", aiAnalysisError.getAiAnalysisRequestId());
processPublicOpinionAutomationDify(difyCommunityTargetDTO, user, difyResult, aiAnalysisError.getAiAnalysisRequestId());
}
if (BusinessTypeEnum.FEEDQUALITY.getCode().equals(aiAnalysisError.getAiAnalysisRequestType())) {
log.info("补偿Feed流好内容全部workflow", aiAnalysisError.getAiAnalysisRequestId());
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())
.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{
} 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));
@@ -684,7 +773,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
}
}
} catch (Exception e) {
log.error("处理社区异常消息异常:{}", e);
log.error("处理社区异常消息异常:{}", e);
}
}

View File

@@ -17,4 +17,6 @@ public class ConstantStr {
public static final String INTELLIGENT_CUSTOMER_4IN1 = "INTELLIGENT_CUSTOMER_4IN1";
public static final String DIFY_COUNTERRATELIMIT = "DIFY_COUNTERRATELIMIT";
}

View File

@@ -0,0 +1,80 @@
package com.volvo.ai.analytic.center.utils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.data.redis.core.script.DefaultRedisScript;
import org.springframework.stereotype.Service;
import java.util.Collections;
@Service
public class RedisCounterRateLimiter {
private final StringRedisTemplate redisTemplate;
@Autowired
public RedisCounterRateLimiter(StringRedisTemplate redisTemplate) {
this.redisTemplate = redisTemplate;
}
/**
* 增加计数并检查是否超过限制
* @param key 限流key
* @param limit 最大限制数
* @param expire 过期时间(秒)
* @return true-允许请求; false-超过限制
*/
public boolean incrementAndCheck(String key, int limit, long expire) {
// 使用Lua脚本保证原子性
String luaScript =
"local current = redis.call('GET', KEYS[1]) or '0'\n" +
"local num = tonumber(current)\n" +
"if num >= tonumber(ARGV[1]) then\n" +
" return 0\n" +
"else\n" +
" redis.call('INCR', KEYS[1])\n" +
" if num == 0 then\n" +
" redis.call('EXPIRE', KEYS[1], ARGV[2])\n" +
" end\n" +
" return 1\n" +
"end";
DefaultRedisScript<Long> script = new DefaultRedisScript<>();
script.setScriptText(luaScript);
script.setResultType(Long.class);
Long result = redisTemplate.execute(script, Collections.singletonList(key),
String.valueOf(limit), String.valueOf(expire));
return result != null && result == 1L;
}
/**
* 减少计数
* @param key 限流key
*/
public void decrement(String key) {
// 使用Lua脚本防止减到负数
String luaScript =
"local current = redis.call('GET', KEYS[1]) or '0'\n" +
"if tonumber(current) > 0 then\n" +
" redis.call('DECR', KEYS[1])\n" +
"end\n" +
"return 1";
DefaultRedisScript<Long> script = new DefaultRedisScript<>();
script.setScriptText(luaScript);
script.setResultType(Long.class);
redisTemplate.execute(script, Collections.singletonList(key));
}
/**
* 获取当前计数
* @param key 限流key
* @return 当前计数值
*/
public int getCurrentCount(String key) {
String count = redisTemplate.opsForValue().get(key);
return count == null ? 0 : Integer.parseInt(count);
}
}

View File

@@ -0,0 +1,43 @@
package com.volvo.ai.analytic.center.utils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Service;
import java.util.concurrent.TimeUnit;
@Service
public class RedisLockService {
@Autowired
private StringRedisTemplate stringRedisTemplate;
/**
* 尝试获取分布式锁
*
* @param lockKey 锁的键名
* @param requestId 请求标识(用于解锁时验证)
* @param expireTime 锁的过期时间(秒)
* @return 是否获取锁成功
*/
public boolean tryLock(String lockKey, String requestId, long expireTime) {
return stringRedisTemplate.opsForValue()
.setIfAbsent(lockKey, requestId, expireTime, TimeUnit.SECONDS);
}
/**
* 释放分布式锁
*
* @param lockKey 锁的键名
* @param requestId 请求标识(用于验证)
* @return 是否释放锁成功
*/
public boolean releaseLock(String lockKey, String requestId) {
String currentValue = stringRedisTemplate.opsForValue().get(lockKey);
if (currentValue != null && currentValue.equals(requestId)) {
stringRedisTemplate.delete(lockKey);
return true;
}
return false;
}
}

View File

@@ -0,0 +1,49 @@
package com.volvo.ai.analytic.center.utils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import java.time.Instant;
import java.util.Set;
@Service
public class RedisZSetUtil {
// @Autowired
private RedisTemplate<String, String> redisTemplate;
@Autowired
public RedisZSetUtil(RedisTemplate<String, String> redisTemplate) {
this.redisTemplate = redisTemplate;
}
/**
* 添加成员到 ZSet并设置过期时间戳
*/
public void addWithExpire(String zsetKey, String member, long expireTimeMillis) {
redisTemplate.opsForZSet().add(zsetKey, member, expireTimeMillis);
}
/**
* 清理过期的成员
*/
@Scheduled(fixedRate = 1000)
public void cleanupExpiredMembers() {
long now = Instant.now().toEpochMilli();
redisTemplate.opsForZSet().removeRangeByScore(ConstantStr.DIFY_COUNTERRATELIMIT, 0, now);
}
/**
* 获取未过期的成员
*/
public Set<String> getValidMembers(String zsetKey) {
long now = Instant.now().toEpochMilli();
return redisTemplate.opsForZSet().rangeByScore(zsetKey, now, Double.MAX_VALUE);
}
public Long zCard(String zsetKey) {
return redisTemplate.opsForZSet().zCard(zsetKey);
}
}