增加redis计数
This commit is contained in:
@@ -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<MessageExt> {
|
||||
@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<JSONObject> 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) {
|
||||
|
||||
@@ -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);
|
||||
|
||||
|
||||
@@ -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";
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user