Ragent 流量保护:Redis ZSET 公平排队与分布式并发控制实现原理完整指南
【免费下载链接】ragent企业级 Agentic RAG 智能体 - 全链路覆盖文档解析、多路检索、意图识别、问题重写、会话记忆、MCP 工具调用与深度思考。面向真实业务场景,从 0 到 1 完整工程实现。项目地址: https://gitcode.com/gh_mirrors/ragent1/ragent
Ragent是一个企业级 Agentic RAG 智能体项目,全链路覆盖文档解析、多路检索、意图识别、会话记忆与 MCP 工具调用。在真实业务中,LLM 问答接口又慢又贵,突发流量很容易把模型服务压垮。Ragent 用Redis ZSET 公平排队加上分布式并发控制实现了一套完整的流量保护机制,让排队请求按先来后到依次执行,同时严格限制全局并发数。本文用通俗的语言拆解这套机制的实现原理。
为什么 RAG 系统需要流量保护 🛡️
RAG 问答的每一次请求背后,往往串联着改写、检索、Rerank、大模型生成等多个耗时步骤,单个请求可能持续数十秒,且 Token 成本不低。如果放任所有请求直接打到模型:
- 模型服务过载:并发一高,延迟雪崩,所有用户都变慢;
- 突发流量击穿:某一时段用户集中提问,瞬间压垮上游 LLM;
- 先到的请求没有保障:后来的请求可能插队,先到的用户一直等待。
Ragent 的解决方案是:先排队,再放行——用一个全局并发上限(信号量)控制"同时能跑多少条",用一个公平的 FIFO 队列(ZSET 有序集合)保证"谁先到谁先跑"。
全局排队入口:ChatQueueLimiter 怎么接入 👇
在核心问答流水线的接入层,就有一个"全局流控与排队"节点。每个 SSE 聊天请求进来后,并不直接执行,而是先交给 ChatQueueLimiter.java 的enqueue方法:
- 如果全局限流关闭,请求直通线程池立即执行;
- 如果开启,请求携带"最长等待时间"进入排队器,等待期间通过 SSE 向前端推送排队状态;
- 拿到许可后执行业务,超时未拿到则优雅拒绝,并向前端返回"系统繁忙,请稍后再试",同时把这次拒绝记入会话记忆。
限流器的 Bean 装配在 ChatRateLimiterConfig.java 中,限流器名字为rag:global:chat,并发数、许可租期、轮询间隔都来自 RAGRateLimitProperties.java,并且可以在管理后台热更新。
公平排队:为什么选择 Redis ZSET 而不是普通队列?
核心实现是 FairDistributedRateLimiter.java,它是一个分布式限流器——并发数控制在 Redis 里,多实例部署时全局共享同一个"闸口"。
它用 RedisZSET(有序集合)当排队队列,巧妙之处在于:
| 设计 | 做法 | 解决什么问题 |
|---|---|---|
| 位次即公平 | 每个请求入队时用全局自增序号(queueSeqKey)作为 score | 先到的请求 score 更小、排在前面,天然 FIFO |
| 存活标记 | 每个请求额外写一个带 TTL 的 entry 标记(L384-L392) | JVM 崩溃后标记自然过期,条目变成"僵尸"可被识别清理 |
| 原子出队 | 用 Lua 脚本判断"我是否位于队头窗口内"并原子地把自己摘出队列 | 多个实例同时抢同一个槽位时,只有一个能成功 |
入队顺序在 acquire 方法中:先写存活标记、再入队,然后先尝试立即抢占;抢不到才注册定时轮询等待放行。
Lua 脚本:一次原子操作完成"验位 + 出队 + 扫僵尸" 🔧
抢占逻辑全部封装在 queue_claim_atomic.lua 中,在 Redis 单线程内一次执行完,避免竞态:
- 取队头窗口:用
ZRANGE取出前maxRank + slack个条目(slack 是额外余量,方便在僵尸密集时仍能让存活条目推进到窗口内); - 识别僵尸:逐个检查 entry 存活标记,标记已过期(实例崩溃)的条目直接
ZREM清理; - 窗口判定:只有自己位于"存活条目的队头 maxRank 窗口内"才允许出队,否则返回失败、继续排队;
- 原子出队:成功后
ZREM自己并删除存活标记,同时返回原始 score,供后续失败时按原位次重新入队——绝不插队。
-- 简化示意:判断存活位次并原子出队 local headEntries = redis.call('ZRANGE', queueKey, 0, maxRank + slack - 1) -- 遍历:标记缺失的僵尸条目 ZREM 清理;统计自己的存活位次 liveRank if liveRank < 0 or liveRank >= maxRank then return {0} end local score = redis.call('ZSCORE', queueKey, requestId) redis.call('ZREM', queueKey, requestId) -- 出队 redis.call('DEL', entryPrefix .. requestId) -- 删存活标记 return {1, score}Java 侧调用入口是 claimIfReady:返回 1 表示成功出队,返回的 score 用于抢占失败时"原样回队"。
分布式并发控制:过期信号量防止死锁 🔑
排队解决"顺序"问题,信号量解决"并发上限"问题。Ragent 使用的是 Redisson 的RPermitExpirableSemaphore(tryAcquirePermit 中的tryAcquire(0, leaseSeconds)):
- 非阻塞抢许可:出队成功后立即
tryAcquire,抢不到说明槽位已被其他实例拿走,此时按原 score 重新入队,等待下一轮——公平性由此闭环; - 许可自动过期:每个 permit 有租期(lease 秒数),如果拿到许可的进程中途崩溃,许可到期自动释放,不会像普通信号量那样"死锁"卡死整个队列;
- 释放即广播:业务执行完毕在
finally中释放许可(grant 方法里用 try/finally 包装回调),并立即发布通知。
除聊天队列外,同一套过期信号量思路也用于文档上传限流:SemaphoreInitializer.java 在启动时初始化rag:document:upload信号量,参数可在 RagSemaphoreProperties.java 配置(默认并发 10、等待 30 秒、租期 30 秒)。
Pub/Sub 唤醒与轮询:低延迟又不惊群 ⚡
"轮询"太费 Redis,"纯等待"又没人叫你。Ragent 采用定时轮询 + 事件唤醒的混合驱动:
- 定时轮询兜底:每个排队中的"票券(Ticket)"注册一个固定间隔(默认最小 50ms)的轮询任务,检查是否到期、是否轮到自己(scheduleQueuePoll);
- Pub/Sub 即时唤醒:许可释放、有人入队/退队时,通过 Redisson
RTopic广播一条permit_changed消息,其他实例收到后立刻触发本地轮询,不用等下一个轮询周期; - 通知合并防风暴:本机的 PollNotifier 用"firing 标志 + 待处理计数"把连续到达的多条通知合并成一次扫描,避免释放高峰时所有轮询任务被重复唤醒。
票券状态机:一个请求的完整生命周期 🎫
每个排队请求在本机对应一个 Ticket 对象,内部是四态状态机(Ticket),所有状态迁移只经过一个 CAS 协调点,终态互斥,保证业务回调最多触发一次:
| 状态 | 触发条件 | 处理 | |:---|:---|:| |PENDING| 初始状态,正在排队 | 定时轮询 + 事件唤醒 | |GRANTED| 拿到 permit | 交给线程池执行;执行完在 finally 释放许可 | |TIMED_OUT| 等待超过 maxWait | 出队清理,走"系统繁忙"拒绝流程 | |CANCELLED| SSE 连接断开/出错 | 出队清理,静默结束 |
几个容易踩坑的细节都被显式处理了:先写存活标记再入队(防止刚入队的条目被并发 claim 当僵尸清掉);先设 permitRef 再 CAS 状态(防止 grant 与 cancel 竞争导致许可泄漏);GRANTED 后取消不释放许可(避免把正在使用的槽位让给别人)。
如何配置这套流量保护 📝
聊天队列的关键参数集中在 RAGRateLimitProperties.java(rag.ratelimit.global-*):
| 参数 | 含义 | 调优建议 |
|---|---|---|
globalEnabled | 是否启用全局排队 | 生产环境建议开启 |
globalMaxConcurrent | 全局最大并发数 | 按上游 LLM 的承压能力设置 |
globalMaxWaitSeconds | 排队最长等待时间 | 太短用户容易被拒,太长响应慢 |
globalLeaseSeconds | 许可租期 | 应大于单请求最大耗时,崩溃后自动回收 |
globalPollIntervalMs | 轮询间隔 | 延迟与 Redis 开销的平衡点 |
这套机制的完整说明也记录在项目发布文档 docs/releases/v1.0.0.md 与 README.md 的"生产级特性"章节中。
小结:这套设计好在哪? ✅
- 公平:ZSET 自增 score 保证严格 FIFO,抢占失败按原 score 回队,任何人插不了队;
- 原子:Lua 脚本把"验位、出队、清僵尸"合成一次 Redis 原子操作,多实例竞争不串位;
- 容错:许可带租期自动过期 + 存活标记 TTL,实例崩溃不留"死锁"和"僵尸";
- 低延迟:Pub/Sub 即时唤醒 + 通知合并,既不等满轮询周期,也不触发惊群风暴;
- 体验友好:排队、拒绝、超时全程经 SSE 推送状态,用户侧始终有确定性反馈。
对于正在建设 RAG / Agent 应用的团队,这套"ZSET 公平排队 + 过期信号量 + 事件唤醒"的组合是一个非常实用的流量保护参考方案。
【免费下载链接】ragent企业级 Agentic RAG 智能体 - 全链路覆盖文档解析、多路检索、意图识别、问题重写、会话记忆、MCP 工具调用与深度思考。面向真实业务场景,从 0 到 1 完整工程实现。项目地址: https://gitcode.com/gh_mirrors/ragent1/ragent
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考