diff --git a/ai-analytic-center-biz/pom.xml b/ai-analytic-center-biz/pom.xml
index a4007da..2c34d50 100644
--- a/ai-analytic-center-biz/pom.xml
+++ b/ai-analytic-center-biz/pom.xml
@@ -226,6 +226,11 @@
2.3.0
+
+
+ org.springframework.boot
+ spring-boot-starter-data-redis
+
diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/AnalysisDifyMqConsumer.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/AnalysisDifyMqConsumer.java
index 7c01412..fbea4cb 100644
--- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/AnalysisDifyMqConsumer.java
+++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/AnalysisDifyMqConsumer.java
@@ -6,6 +6,8 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
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 lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
@@ -43,26 +45,43 @@ public class AnalysisDifyMqConsumer implements RocketMQListener {
@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;
+
private final ObjectMapper objectMapper = new ObjectMapper();
@Override
public void onMessage(MessageExt messageExt) {
+ log.info("analysisDifyMqConsumer 当前线程: {}, 线程ID: {}", Thread.currentThread().getName(), Thread.currentThread().getId());
long startTime = System.currentTimeMillis();
+ log.info("redis计数数量: " + redisCounterRateLimiter.getCurrentCount(ConstantStr.DIFY_COUNTERRATELIMIT));
+ boolean flag = redisCounterRateLimiter.incrementAndCheck(ConstantStr.DIFY_COUNTERRATELIMIT, difyLimit, expire);
+ if(!flag){
+ log.info("analysisDifyMqConsumer 请求dify超过基数 :{},稍后请求: " + redisCounterRateLimiter.getCurrentCount(ConstantStr.DIFY_COUNTERRATELIMIT));
+ throw new RuntimeException("请求dify超过基数 " );
+ }
try {
- log.info("analysisDifyMqConsumer 当前线程: {}, 线程ID: {}", Thread.currentThread().getName(), Thread.currentThread().getId());
String message = new String(messageExt.getBody());
log.info("analysisDifyMqConsumer message: " + message);
DiFyReq difyReq = JSONObject.parseObject(message, DiFyReq.class);
CompletableFuture future = diFyService.asyncExecuteDifyFlow(difyReq);
+ JSONObject difyRequest = JSONObject.parseObject(JSONObject.toJSONString(difyReq.getInputs()), JSONObject.class);
+ String aiAnalysisRequestId = difyRequest.getString("aiAnalysisRequestId");
// JSONObject json = future.get();
future.thenAccept(result -> {
- JSONObject difyRequest = JSONObject.parseObject(JSONObject.toJSONString(difyReq.getInputs()), JSONObject.class);
- String aiAnalysisRequestId = difyRequest.getString("aiAnalysisRequestId");
// 处理异步结果
log.info("异步处理asyncExecuteDifyFlow完成aiAnalysisRequestId: {},处理结果:{}", aiAnalysisRequestId, result);
+ }).exceptionally(ex -> {
+ log.error("异步处理asyncExecuteDifyFlow 失败aiAnalysisRequestId: {} ,{}", aiAnalysisRequestId,ex.getMessage());
+ return null;
});
log.info("analysisDifyMqConsumer处理完成,耗时:{}", System.currentTimeMillis() - startTime);
} catch (Exception e) {
diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisDifyServiceImpl.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisDifyServiceImpl.java
index 8d32758..8848a81 100644
--- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisDifyServiceImpl.java
+++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisDifyServiceImpl.java
@@ -18,7 +18,9 @@ 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.FlowResultSplitUtil;
+import com.volvo.ai.analytic.center.utils.RedisCounterRateLimiter;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
import org.apache.commons.lang3.StringUtils;
@@ -59,9 +61,15 @@ public class AiAnalysisDifyServiceImpl implements AiAnalysisDifyService {
private String callbackTopic;
@Resource
private RocketMQTemplate rocketMqTemplate;
+
+ @Autowired
+ private RedisCounterRateLimiter redisCounterRateLimiter;
@Override
public boolean updateAiDifyResult(String message) {
+ // 计数-1
+ redisCounterRateLimiter.decrement(ConstantStr.DIFY_COUNTERRATELIMIT);
+ log.info("redisdecrement计数数量: " + redisCounterRateLimiter.getCurrentCount(ConstantStr.DIFY_COUNTERRATELIMIT));
if(StringUtils.isNotEmpty(message)){
AnalysisDifyResultDTO analysisResp = JSONObject.parseObject(message, AnalysisDifyResultDTO.class);
diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/utils/ConstantStr.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/utils/ConstantStr.java
index a200176..4e0be67 100644
--- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/utils/ConstantStr.java
+++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/utils/ConstantStr.java
@@ -13,4 +13,6 @@ public class ConstantStr {
public static final String corpus_user = "corpush_user";
public static final String CARMODELLIST_CACHEKEY = "ai_carmodellist";
+
+ public static final String DIFY_COUNTERRATELIMIT = "DIFY_COUNTERRATELIMIT";
}
diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/utils/RedisCounterRateLimiter.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/utils/RedisCounterRateLimiter.java
new file mode 100644
index 0000000..6e9f906
--- /dev/null
+++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/utils/RedisCounterRateLimiter.java
@@ -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 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 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);
+ }
+}
\ No newline at end of file