redis pub/sub 必须使用异步连接,同步连接不支持监听;需禁用读超时、持续 poll_msg()、避免 json 解析热路径、手动实现指数退避重连。

pubsub 连接必须用异步客户端,不能复用普通连接
Redis 的 PUB/SUB 机制是单向、无响应的流式协议,redis-rs 的同步连接(get_connection())不支持监听模式,强行调用 subscribe() 会 panic 或卡死。必须使用 get_async_connection() 获取异步连接,并确保整个生命周期在 Tokio runtime 中运行。
常见错误现象:PubSub: Connection closed unexpectedly 或程序静默退出,往往是因为连接被提前 drop,或未在 async fn 内执行 subscribe。
- 订阅前务必调用
conn.set_read_timeout(None),禁用读超时,否则空闲时连接会被断开 - 不要把
AsyncConnection存入结构体长期持有——它不可 Send,且内部状态不可重用 - 每个订阅逻辑建议封装为独立的
tokio::spawn任务,避免阻塞主流程
subscribe() 后必须持续 poll_msg(),不能只调一次
redis::aio::PubSub 不是事件回调模型,而是基于 Future 的拉取式流。调用 subscribe() 只是发送 SUBSCRIBE 命令,真正接收消息靠反复调用 get_message().await —— 它会阻塞直到新消息到达或连接中断。
典型误用:let msg = pubsub.get_message().await? 放在 if 或 match 分支里,只执行一次,导致后续消息全部丢失。
- 必须用
loop { let msg = pubsub.get_message().await?; ... }持续消费 - 若需同时监听多个 channel,用
pubsub.subscribe("ch1", "ch2")一次性注册,不要多次调用 - 注意
msg.get_payload::<string>()</string>可能失败,payload 是Vec<u8></u8>,二进制内容需按业务约定解码
消息反序列化别在 hot loop 里做 JSON 解析
Pub/Sub 流通常是高频、低延迟场景,serde_json::from_slice() 在热路径中会产生明显 GC 压力和 CPU 开销。实测比 rmp-serde(MessagePack)慢 3–4 倍。
Redis 缓存和数据结构管理技能。通过自然语言操作 Redis,支持 String、Hash、List、Set、ZSet、Stream 等数据结构操作。当用户提到 Redis、缓存、消息队列、会话存储时使用此技能。
性能影响:单核下 JSON 解析吞吐易卡在 2–3 万 msg/s,而 MessagePack 可达 10 万+;尤其当 payload 超过 1KB 时差距更明显。
- 发布端统一用
rmp_serde::to_vec(&data).unwrap()编码 - 订阅端用
rmp_serde::from_slice::<myevent>(&msg.get_payload_bytes())</myevent> - 避免在
get_message()后立刻做 heavy work,可先用tokio::task::spawn转移处理
连接断开后自动重连要自己实现,redis-rs 不内置
redis-rs 的 PubSub 类型不提供 reconnect 机制。网络抖动、Redis 重启、timeout 都会导致 get_message().await 返回 Err,此时连接对象已失效,无法继续使用。
容易踩的坑:捕获错误后直接 continue,结果陷入空转,再也收不到消息。
- 必须在
loop外层包一层 retry 循环,例如for _ in 0..=5 { ... } - 每次重连都要重新
Client::open()→get_async_connection()→into_pubsub()→subscribe() - 建议加指数退避:
tokio::time::sleep(Duration::from_millis(100 * 2u64.pow(retry))).await
实际中最容易被忽略的是 Pub/Sub 连接的“一次性”本质:它不像普通命令连接那样可复用,也不像 HTTP 连接有 keep-alive 保活机制。每次断开都意味着整套状态(channel 订阅列表、内部 buffer、waker)全部作废,重连逻辑必须从头构建,且不能漏掉 set_read_timeout(None) 这个关键设置。










