Redis 使用方案

作者:old wang 发布时间: 2025-05-28 阅读量:6 评论数:0

Redis 在 RAG 项目中承担了四类职责:业务缓存、分布式协调、并发控制、跨实例通信。选用 Redisson 作为主力客户端,因为它将 Redis 的数据结构(信号量、ZSET、Topic、分布式锁)封装为 Java 原生对象,不需要手动管理连接池和序列化。简单 KV 场景用 StringRedisTemplate 作为轻量补充。

四类职责的可靠性要求和故障影响不同:缓存挂了可以降级,系统变慢但不挂;分布式协调和并发控制挂了,对应功能直接不可用——这是需要补齐的薄弱点。

1. 客户端选型:Redisson + StringRedisTemplate

项目中同时使用两种 Redis 客户端:

Redisson:用于需要高级数据结构的场景——可过期信号量、ZSET 队列、分布式锁、发布订阅。Redisson 将这些 Redis 数据结构和操作封装为 Java 对象RPermitExpirableSemaphoreRScoredSortedSetRTopic),调用方式和操作本地对象一致,不用自己写 Lua 脚本(除非有原子性要求才自己写)。

StringRedisTemplate:用于简单的 KV 读写和自定义 Lua 脚本执行。缓存类场景用它的 opsForValue().get/set,幂等消费用它的 execute(RedisScript) 执行自定义 Lua。

Key 前缀通过 RedisKeySerializer 统一管理—framework.cache.redis.prefix 配置项注入前缀,所有 Key 自动带上前缀,便于多环境区分和运维清理。

2. 业务缓存

2.1 意图树缓存

场景:意图树是每次对话都要加载的基础数据——意图分类器需要把所有叶子节点列给 LLM 打分。如果不加缓存,每次对话都要从 t_intent_node 表查询、构建树结构、填充 fullPath。数据库压力不大,但延时不可忽略。

实现IntentTreeCacheManager 将构建好的意图树 JSON 序列化后存入 Redis,Key 为 ragent:intent:tree,过期时间 7 天。

缓存的读取策略是"每次读 Redis,不命中则查数据库并回写缓存"。之所以不设本地缓存(Caffeine)做二级缓存,是因为意图树数据量不大(几十 KB 到几百 KB),Redis 网络开销可忽略。多实例部署时,本地缓存失效需要额外同步机制,引入的复杂度比省下的延时更不划算。

缓存更新:管理后台修改意图节点时,主动调用 clearIntentTreeCache() 删除 Redis Key。下次对话请求触发重建。7 天过期是安全网——如果删缓存的操作因为网络抖动没执行成功,最坏等 7 天自动纠正。相比永久缓存的风险(永远读到旧数据),7 天是一个安全的兜底。

2.2 术语映射缓存

场景:查询重写时需要将用户口语中的术语映射为标准表述(如"那个打日志的"→"SLS 日志服务")。术语映射规则存储在 t_query_term_mapping 表,读多写少。

实现QueryTermMappingCacheManager,Key 为 ragent:query-term:mappings,过期 7 天。读取和更新策略与意图树缓存一致——查 Redis → 不命中则查数据库 → 回写缓存 → 管理后台变更时主动删除。

3. 分布式并发控制

3.1 聊天全局限流:公平排队

场景:LLM API 有并发配额限制(百炼默认约 10~50 个并发),多个服务实例共享配额,需要跨实例的全局并发控制。

实现FairDistributedRateLimiter,是项目中最复杂的 Redis 使用场景。

四组件协同:

ZSET 公平队列:每个请求入队时分配一个全局递增序列号(通过 Redis 原子计数器 RAtomicLong 递增),作为 ZSET 的 score。ZSET 按 score 排序天然是先到先得。Key 为 rag:global:chat:queue

可过期信号量许可池:Redisson 的 RPermitExpirableSemaphore 管理并发槽位,许可数等于 max-concurrent(默认 10)。每个许可有租约时间 leaseSeconds(默认 30 秒),超时自动归还——即使持有许可的请求因为异常没有显式释放,槽位也不会永久丢失。

Lua 原子抢占:ZSET 排队和信号量获取之间存在竞态——两个请求可能同时判定自己是队头、同时去拿许可。

Lua 脚本 queue_claim_atomic.lua 把三步操作合并为一次 Redis 原子调用:

1. 扫描队头窗口内的存活成员(通过 entry 标记判断死活)

2. 清理僵尸(entry 标记过期的成员,说明所属实例已崩溃)

3. 检查当前请求是否在队头窗口内,在则 ZREM 出队并返回原始 score

为什么不直接用信号量的排队能力?

Redisson 的 tryAcquire(waitTime, leaseTime, unit) 确实支持带等待时间的获取,但它是非公平的——哪个实例先抢到算哪个,不保证先到先得。ZSET 排队队列保证了请求按到达顺序分配许可,先请求的用户先拿到许可。

Pub/Sub 跨实例唤醒:任一实例释放许可后,通过 RTopic 广播 permit_changed 消息。其他实例的等待者收到通知后,立即触发一轮轮询检查("许可池有空闲吗?我排到队头了吗?")PollNotifier 内部用 AtomicBoolean firing + 计数器做通知合并——连续到达的多次通知只触发一次批量扫描,避免每个通知都触发一轮全量轮询——本质是防止通知过载。

定时轮询兜底:即使 Pub/Sub 消息丢失,本地调度线程每 200 毫秒pollIntervalMs主动检查一次,保证最终能感知到许可释放。

僵尸清理:如果某个实例 JVM 崩溃,它的排队条目残留在 ZSET 中无人消费。每个排队请求在入队时写了 entry 标记 Keyrag:global:chat:entry:{requestId}),TTL 等于剩余等待时间加 5 秒缓冲。JVM 崩溃后标记过期,后续 Lua 脚本发现队头对应的 entry 已过期,将其作为僵尸清除,队头位置让给下一个活着的请求。

3.2 文档上传限流:信号量 + Filter

场景:文档上传是资源密集型操作(解析 + 分块 + 向量化),并发上传太多会把 Embedding API 打满或消耗过多内存。

实现UploadRateLimitFilter + Redisson RPermitExpirableSemaphore。Filter 注册在最高优先级Ordered.HIGHEST_PRECEDENCE),在 Spring 的 Multipart 解析之前拦截——如果拿不到许可,连临时文件都不会产生,直接返回 HTTP 429。

信号量配置rag:document:upload,最大并发 10,获取等待 30 秒,租约 30 秒。

与聊天限流的区别:上传限流不需要公平排队——用户上传文档对"先到先得"不敏感。直接用 semaphore.tryAcquire(waitTime, leaseTime, unit),简洁且够用SemaphoreInitializer 在启动时通过 @PostConstruct 初始化信号量的许可数。

许可在 Filter 的 finally 块中释放。如果释放失败(如租约已过期),只打 WARN 日志不抛异常——上传已经成功,许可过期说明处理耗时超出预期,自动归还也不影响业务结果。

4. 分布式协调

4.1 幂等控制:防重复提交

场景:用户快速点击两次上传按钮,两个 HTTP 请求几乎同时到达。需要保证只有一个请求执行业务逻辑,另一个直接拒绝。

实现@IdempotentSubmit 注解 + IdempotentSubmitAspect AOP 切面。底层用 Redisson 的分布式锁RLock,锁 Key 为 idempotent-submit:path:{路径}:currentUserId:{用户}:md5:{参数MD5}

锁 Key 包含了用户 ID 和参数 MD5——同一个用户同一批参数短时间内不该提交两次tryLock() 不带超时参数,拿不到立刻返回失败。这是幂等控制场景的正确做法:不需要等,重复请求就该立刻拒绝。

Redisson 的看门狗机制(watchDog)自动续期锁的持有时间,防止业务处理超过锁默认 TTL 导致锁提前释放。但要注意:如果 JVM 进入 FGC,看门狗线程也被暂停,锁可能意外释放——这个场景概率很低,但对支付类场景是严重问题。当前 RAG 场景最多重复创建一条对话或重复上传一个文档,后果可控。

4.2 幂等控制:防重复消费

场景:RocketMQ 在异常情况下可能重复投递消息,消费端必须区分"第一次消费"和"重复投递"。

实现@IdempotentConsume 注解 + IdempotentConsumeAspect AOP 切面。底层用 Redis Lua 脚本 SET key value NX GET PX expire_ms,在一次 Redis 调用中完成"检查是否存在 + 写入标记 + 设过期时间"三步:

- 首次消费:SET NX 成功,返回 nil,执行业务逻辑,完成后 SET CONSUMED

- 重复投递,消息还在处理中:返回 CONSUMING,抛异常让 MQ 稍后重试

- 重复投递,消息已处理完:返回 CONSUMED,直接跳过

Key 的过期时间必须大于消息的最大处理时间。如果处理还在进行中 Key 就过期了,重复消息会绕过 SET NX 再次处理。这是 Redis 做幂等的固有局限——TTL 和业务处理时间的错配可能导致双写。数据库唯一索引是最后一道防线。

4.3 分布式 Snowflake ID 初始化

场景:Snowflake 算法要求每个实例有唯一的 workerId 和 datacenterId,多实例部署时需要协调分配。

实现SnowflakeIdInitializer@PostConstruct 时通过 Redis Lua 脚本原子分配 workerId 和 datacenterId,然后注册到 Hutool 的 IdUtil。这个方案的局限性是服务启动依赖 Redis 可用——Redis 挂了服务起不来。Snowflake 本身是去中心化的 ID 方案,用 Redis 分配 workerId 相当于在去中心化方案前面加了一个中心化依赖。更稳健的做法是用 K8s StatefulSet 的 Pod 序号或宿主机 IP 后几位来派生 workerId,完全不需要 Redis。

---

5. 跨实例通信

5.1 流取消信号广播

场景:用户点击"停止生成",取消请求可能落在实例 A,但 LLM 推理跑在实例 B。取消信号必须跨实例传递。

实现StreamTaskManager 用 Redis 做两层协调:

取消标记RBucket):Key 为 ragent:stream:cancel:{taskId},值为 true,TTL 30 分钟。写标记后立即广播,所有实例收到广播后本地取消。注册时回查 Redis——处理"取消请求先于 EventHandler 注册到达"的竞态。

取消广播RTopic):发布到 ragent:stream:cancel 通道,消息体为 taskId。所有实例在 @PostConstruct 时订阅。本地取消通过 AtomicBoolean.compareAndSet 保证幂等。

5.2 限流唤醒

聊天限流器在释放许可后通过 RTopic 广播 permit_changed 消息,唤醒其他实例上正在排队的等待者。详见 3.1。

6. Redis 客户端的 Key 序列化

RedisKeySerializer 实现 Spring 的 RedisSerializer<String> 接口,在 Key 读写时自动拼接 framework.cache.redis.prefix 前缀。这个设计让多环境(开发/测试/生产)共用同一套 Redis 时不会出现 Key 冲突——每个环境配置不同的前缀即可。

7. 边界情况与容错

Redis 不可用——缓存场景:意图树和术语映射的缓存读取失败时IntentTreeCacheManagerQueryTermMappingCacheManagergetXxxFromCache() 方法 catch 异常后返回 null,触发数据库加载。系统可用但响应变慢。这是缓存的标准降级模式——挂了不影响核心功能。

Redis 不可用——幂等场景IdempotentConsumeAspect 的 Lua 脚本执行失败会抛异常,MQ 消费中断IdempotentSubmitAspecttryLock() 在 Redis 不可用时的行为取决于 Redisson 的配置——可能抛异常或返回 false。这两种情况下,幂等保护失效。对于上传和消费场景,最坏结果是产生重复数据,不致命但需要事后清理。

Redis 不可用——限流场景:聊天限流器的新请求无法入队ChatQueueLimiter 中开关关闭时走直通分支绕过限流;开关开启时请求无法获取许可。上传限流 Filter 的 getPermitExpirableSemaphore 访问 Redis 失败会抛异常,Filter 返回 500。这两种情况下,限流失效——所有请求都能通过,LLM API 可能被打爆。

Redis 不可用——跨实例通信StreamTaskManager 的取消广播和标记写入失败。最坏情况:用户点了取消,但 LLM 推理没被停掉,继续消耗 API 配额生成没人看的回复。

Snowflake 启动依赖SnowflakeIdInitializer@PostConstruct 时依赖 Redis,Redis 挂会导致服务启动失败。

当前 Redis 已经是架构中的事实单点——核心对话链路的限流、幂等、会话管理、取消协调全部绑在上面。关键保护措施:

- 限流器有开关,紧急时可以关闭限流绕过 Redis 依赖

- 缓存有数据库兜底,Redis 挂了只是变慢

- 幂等保护的数据库唯一索引是最后一道防线

评论