Streams 消费者组 · 故障恢复实验台

订单消息从进入流,到领取、滞留、转移和确认的完整旅程
orders:stream
group = billing-workers

消息轨道

新消息只有被消费者组读取后,才会进入 PEL

消费者组

XREADGROUP 把消息交给某个消费者,并留下责任记录

消费者 C1

在线

消费者 C2

在线

待处理列表 PEL

这里不是消息副本,而是“谁拿了、拿了多久、投递几次”的账本
实验时钟:0s
认领阈值:idle ≥ 20s

为什么 Pub/Sub 断线后追不回来,Streams 却能查账?

把 Pub/Sub 想成办公室广播:订阅者在线时能听见,广播完就结束;你断线的那几分钟,没有人替你把广播录下来。Streams 更像带流水号的收件室,每条消息进入后都有 ID,也会保留在流里。消费者组领取消息时,Redis 还会把“这条消息交给了谁”记进待处理列表,所以某个消费者突然退出,其他消费者仍能发现并接手。

观察点Pub/SubRedis Streams
订阅者断线断线期间的消息不会补发消息仍在流中,可按 ID 继续读取
处理责任服务端不记录谁处理过消费者组用 PEL 记录归属、空闲时间和投递次数
适合场景在线通知、实时广播、允许丢失的事件任务分发、异步业务、需要恢复和追踪的事件

一次故障恢复,Redis 实际记了什么

1 · XADD消息获得递增 ID,进入 Stream。此时还没有消费者对它负责。
2 · XREADGROUP组内消费者领取新消息,Redis 同时把它登记进 PEL。
3 · 崩溃消费者退出并不会自动删除责任记录;消息开始积累 idle 时间。
4 · XAUTOCLAIM扫描超过阈值的 PEL 条目,把所有权转给健康消费者。
5 · XACK业务成功后确认,条目从 PEL 移除。消息本体是否裁剪是另一件事。

注意最后一句:XACK 只表示消费者组不再追踪这条待处理责任,并不等于把 Stream 里的原始消息物理删除。流的长度通常通过 MAXLEN 或定期裁剪来控制,这和消费确认是两个维度。

“至少一次”不是免费保险,它把难题交给了业务幂等

故障恰好发生在“业务已经成功、XACK 还没发出”的夹缝里时,Redis 只能看到消息仍在 PEL,于是恢复者会再次投递。为了不丢消息,这种重复是刻意接受的结果。Streams 给你的保证更接近“至少处理一次”,而不是“副作用只发生一次”。

可靠的做法:为每个业务动作准备稳定幂等键,例如 order_id + operation;在数据库唯一约束、幂等记录表或原子脚本里,先判断这个键是否已经完成。不要把 Stream 消息 ID 当成所有业务的天然幂等键——重放、补偿和跨流迁移时,业务身份往往比运输层 ID 更稳定。

上面的实验里,“无幂等”会让每次恢复都重复增加副作用次数;“有业务幂等键”会把已经完成的动作识别出来,只补做确认。重试次数持续升高通常说明它不是偶发故障,而是毒消息、数据不兼容或下游永久拒绝。把它移入死信区后再告警和人工分析,比无限循环占用消费者更可控。