RAG:基于 Redis 信号量和 Zset 的分布式排队限流设计
前言
几个用户同时在网页系统里点击"发送"按钮,大模型 API 返回了 429 Too Many Requests 。
作为开发,你面临几个选择:
- 限流 → 直接拒绝 :但用户体验很差,前端直接报错
- 不限流 :API 超时、报错、甚至多扣费
- 排队 :暂时处理不了?先排着,有空位了再处理
在一篇讲的那个RAG系统里是选择了第三种。本文分享一个基于 Redis 信号量 + Zset + Lua + Pub/Sub 的排队限流方案,配合 SSE 实时推送排队状态。
一、设计思路
1.1 核心目标
同一时刻最多只能有 N 个请求在处理
超过 N 的请求按时间顺序排队
处理完一个,放行一个
排太久没轮到?告诉用户"系统繁忙"
1.2 为什么选这几个组件?
| Redis 信号量 | 分布式场景统一控制并发数,支持过期自动释放做兜底 |
|---|---|
| Zset | 天然按 score(时间戳)排序,实现公平排队 |
| Lua 脚本 | 判断"是不是队头"和"出队"需要原子操作 |
| Pub/Sub | 减少轮询延迟,处理完一个立刻通知下一个 |
1.3 整体流程
用户 A 发消息
│
▼
AOP 切面拦截 → 调用 ChatQueueLimiter.enqueue()
│
├─① 写入 Zset(score = 当前时间戳)
│
├─② 尝试获取信号量
│ ├─ 成功 → 执行业务逻辑(LLM 回答)
│ └─ 失败 → 进入③
│
├─③ 启动定时轮询(每 200ms)
│ 用 Lua 判断"我是队头吗?"
│ ├─ 是 → 获取信号量 → 执行业务逻辑
│ └─ 否 → 继续等
│
├─④ 排队超时(3 秒)→ SSE 推"系统繁忙"
│
└─⑤ 处理完 → Pub/Sub 通知下一个
二、架构设计
2.1 两层结构
┌─────────────────────────────────────┐
│ @ChatRateLimit 注解(标记入口方法) │
└────────────────┬────────────────────┘
│
┌────────────────▼────────────────────┐
│ ChatRateLimitAspect(AOP 切面) │
│ 负责:拦截请求、提取参数、调用排队器 │
└────────────────┬────────────────────┘
│
┌────────────────▼────────────────────┐
│ ChatQueueLimiter(排队限流器) │
│ 负责:Zset入队/信号量/Lua/Pub/Sub │
│ SSE推送 │
└─────────────────────────────────────┘
2.2 为什么用 AOP?
关键代码只用加一个注解:
@Override
@ChatRateLimit // ← 加这个注解就行
public void streamChat(String question, String conversationId, Boolean
deepThinking, SseEmitter emitter) {
// 业务逻辑完全不用改
}
AOP 切面自动拦截,把请求交给排队器。业务代码零侵入。
三、核心实现(逐段拆解)
3.1 请求入队(Zset + 信号量)
当新的请求进来,先做两件事: 写入 Zset 排队 + 尝试立即获取许可 。
public void enqueue(String question, String conversationId,
SseEmitter emitter, Runnable onAcquire) {
// 1. 生成唯一请求 ID,写入 Zset(按时间戳排序)
String requestId = IdUtil.getSnowflakeNextIdStr();
RScoredSortedSet<String> queue = redissonClient.getScoredSortedSet(QUEUE_KEY);
long seq = nextQueueSeq(); // 生成自增序列号
queue.add(seq, requestId); // score = 序列号,确保公平顺序
// 2. 绑定释放回调(连接断开时自动释放)
Runnable releaseOnce = () -> {
queue.remove(requestId); // 从 Zset 移除
releasePermit(permitRef); // 释放信号量
publishQueueNotify(); // 通知下一个
};
emitter.onCompletion(releaseOnce); // SSE 完成时触发
emitter.onTimeout(releaseOnce); // SSE 超时时触发
emitter.onError(e -> releaseOnce); // SSE 异常时触发
// 3. 立即尝试获取许可
if (tryAcquireIfReady(queue, requestId, permitRef, cancelled, onAcquire)) {
return; // 有空位,直接开始处理
}
// 4. 没空位 → 启动定时轮询
scheduleQueuePoll(queue, requestId, ...);
}
Zset 为什么用自增序列而不是时间戳? 防止同一毫秒多个请求的 score 冲突,序列号保证绝对有序。
3.2 尝试获取许可(信号量 + Lua 脚本)
private boolean tryAcquireIfReady(RScoredSortedSet<String> queue, String requestId,
AtomicReference<String> permitRef,
AtomicBoolean cancelled, Runnable onAcquire) {
// 1. 检查当前还有没有空位
int availablePermits = availablePermits();
if (availablePermits <= 0) return false;
// 2. 用 Lua 判断"我是不是在队头?"
ClaimResult claimResult = claimIfReady(queue, requestId, availablePermits);
if (!claimResult.claimed) return false; // 不是队头,继续等
// 3. 是队头 → 获取信号量许可
String permitId = tryAcquirePermit();
if (permitId == null) {
// 信号量被抢走了?重新入队
queue.add(nextQueueSeq(), requestId);
return false;
}
// 4. 拿到许可了 → 执行业务逻辑
chatEntryExecutor.execute(onAcquire);
return true;
}
这里有一个关键细节: Lua 脚本先判断队头,然后 Java 代码再获取信号量 。为什么不能把获取信号量也放进 Lua?因为 Redisson 的 PermitExpirableSemaphore 是封装好的,写 Lua 里反而复杂。分两步走,中间用 availablePermits 做预检查,把冲突概率降到最低。
3.3 Lua 脚本:原子判队头
-- KEYS[1]: 队列 Zset Key
-- ARGV[1]: 请求 ID
-- ARGV[2]: 最大可进入的 rank(当前空位数)
local rank = redis.call('ZRANK', KEYS[1], ARGV[1])
if not rank then return {0} end
if rank >= tonumber(ARGV[2]) then return {0} end
local score = redis.call('ZSCORE', KEYS[1], ARGV[1])
redis.call('ZREM', KEYS[1], ARGV[1])
return {1, score}
逻辑只有三句话:
- 查排名 : ZRANK 获取请求在队列中的位置
- 判断 :排名 < 空位数 → 轮到你了
- 出队 : ZREM 从队列移除
这些操作在 Lua 脚本里是原子执行的,不会出现两个请求同时认为自己排到了。
3.4 定时轮询 + Pub/Sub 优化
排队中的请求需要等待。两种方式:
方式一:纯轮询(保底方案)
private void scheduleQueuePoll(...) {
scheduler.scheduleAtFixedRate(poller, intervalMs, intervalMs, TimeUnit.
MILLISECONDS);
// 每 200ms 检查一次是否轮到自己
}
方式二:Pub/Sub 通知(优化方案)
一个请求处理完释放信号量时,发一条消息:
private void publishQueueNotify() {
redissonClient.getTopic("rag:global:chat:queue:notify").publish
("permit_released");
}
其他等待的实例订阅了这个频道:
@PostConstruct
public void subscribeQueueNotify() {
RTopic topic = redissonClient.getTopic("rag:global:chat:queue:notify");
topic.addListener(String.class, (channel, msg) -> {
pollNotifier.fire(); // 有人释放了许可,快去检查!
});
}
两种方式共存: Pub/Sub 负责"立刻唤醒",轮询负责"兜底保底"。Pub/Sub 消息可能会丢,但轮询一定会检查。
3.5 超时拒绝 + SSE 推送
排队不是无限的,超过等待时间(默认 3 秒)就拒绝:
if (System.currentTimeMillis() > deadline) {
queue.remove(requestId); // 从队列移除
sendRejectEvents(emitter, context); // 推"系统繁忙"
return;
}
拒绝时通过 SSE 向前端推送:
private void sendRejectEvents(SseEmitter emitter, RejectedContext context) {
SseEmitterSender sender = new SseEmitterSender(emitter);
sender.sendEvent("META", new MetaPayload(context.conversationId, context.
taskId));
sender.sendEvent("REJECT", new MessageDelta("response", "系统繁忙,请稍后再试"));
sender.sendEvent("FINISH", new CompletionPayload(messageId, title));
sender.sendEvent("DONE", "[DONE]");
sender.complete();
}
前端收到 REJECT 事件后,展示"系统繁忙"提示。
四、配置参数
所有参数通过配置文件控制:
rag:
rate-limit:
global:
enabled: true # 是否启用全局限流
max-concurrent: 1 # 最大并发数
max-wait-seconds: 3 # 排队超时时间
lease-seconds: 30 # 信号量自动释放时间(兜底)
poll-interval-ms: 200 # 轮询间隔
五、全过程流程图

六、与普通限流的区别
| 维度 | 普通限流(计数器/令牌桶 | 本方案 |
|---|---|---|
| 超限后的行为 | 直接拒绝 | 排队等待 |
| 是否公平 | 不保证 ,先到不一定先得 | Zset 保证公平 |
| 超时处理 | 无,直接拒绝 | 超时后通知用户 |
| 分布式支持 | 需要额外实现 | 原生支持 (基于 Redis) |
| 用户体验 | 报错 | 排队提示/系统繁忙 |
维度 普通限流(计数器/令牌桶) 本方案 超限后的行为 直接拒绝 排队等待 是否公平 不保证,先到不一定先得 Zset 保证公平 超时处理 无,直接拒绝 超时后通知用户 分布式支持 需要额外实现 原生支持 (基于 Redis) 用户体验 报错 排队提示/系统繁忙
七、总结
这套排队限流方案的核心思想是: 不直接拒绝用户,而是让用户排队等待。 配合 SSE 实时推送,用户体验比直接报错好得多。
如果你也在做类似的大模型 RAG 系统,这个方案值得参考。代码是基于aop切面开发的对业务代码侵入程度低,几乎可以做到低成本改动。
最后想说一句: 限流的本质不是在"拒绝用户"和"挂掉"之间二选一,而是要找到一个让系统和用户都能接受的中间状态。 排队 + 通知,就是那个中间状态。
更多推荐




所有评论(0)