☰
Ragent 流量保护:Redis ZSET 公平排队与分布式并发控制实现原理完整指南
2026/9/25 18:50:07 网站建设 项目流程

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 单线程内一次执行完,避免竞态:

  1. 取队头窗口:用ZRANGE取出前maxRank + slack个条目(slack 是额外余量,方便在僵尸密集时仍能让存活条目推进到窗口内);
  2. 识别僵尸:逐个检查 entry 存活标记,标记已过期(实例崩溃)的条目直接ZREM清理;
  3. 窗口判定:只有自己位于"存活条目的队头 maxRank 窗口内"才允许出队,否则返回失败、继续排队;
  4. 原子出队:成功后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 即时唤醒:许可释放、有人入队/退队时,通过 RedissonRTopic广播一条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 的"生产级特性"章节中。

小结:这套设计好在哪? ✅

  1. 公平:ZSET 自增 score 保证严格 FIFO,抢占失败按原 score 回队,任何人插不了队;
  2. 原子:Lua 脚本把"验位、出队、清僵尸"合成一次 Redis 原子操作,多实例竞争不串位;
  3. 容错:许可带租期自动过期 + 存活标记 TTL,实例崩溃不留"死锁"和"僵尸";
  4. 低延迟:Pub/Sub 即时唤醒 + 通知合并,既不等满轮询周期,也不触发惊群风暴;
  5. 体验友好:排队、拒绝、超时全程经 SSE 推送状态,用户侧始终有确定性反馈。

对于正在建设 RAG / Agent 应用的团队,这套"ZSET 公平排队 + 过期信号量 + 事件唤醒"的组合是一个非常实用的流量保护参考方案。

【免费下载链接】ragent企业级 Agentic RAG 智能体 - 全链路覆盖文档解析、多路检索、意图识别、问题重写、会话记忆、MCP 工具调用与深度思考。面向真实业务场景,从 0 到 1 完整工程实现。项目地址: https://gitcode.com/gh_mirrors/ragent1/ragent

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询