登录
推荐 文章 Go 技术 课程 下载 专题 AI
首页 >  数据库 >  Redis

Redis Pub/Sub断线期间消息丢失时的替代结构

来源:17golang原创

时间:2026-09-25 14:45:00 369浏览 收藏

Redis Pub/Sub 断线期间消息丢失,通常不是重连代码写错,而是消息模型本身没有历史记录。Redis 官方文档明确将 Pub/Sub 定义为 at-most-once:消息发送后,如果订阅者因网络断开或处理失败没有接住,就不会再次发送。需要补读、确认和重试时,应把关键事件写入 Redis Stream,再用消费组读取。

官方地址:https://redis.io/

要点速览
  • Pub/Sub 适合实时广播,不提供断线后的历史重放。
  • Stream 配合消费组、XACK 和待确认列表,才能建立可恢复消费。
  • 迁移时要同时设计幂等键、重试上限、消息保留和异常认领。

先判断:丢消息是 Pub/Sub 的交付边界

SUBSCRIBE 建立的是在线推送关系,发布者通过 PUBLISH 把消息推给当时已订阅的客户端。客户端短暂断网、进程重启或消费回调报错时,Redis 不会替它保存一个等待补发的游标。因此,单纯增加重连次数只能缩短空窗期,不能找回空窗期已经发布的消息。

可以把需求拆成三个指标:是否允许丢失、是否需要按顺序补读、是否需要知道某条消息已经处理完成。只要后两个答案为“需要”,就不应把 Pub/Sub 当作唯一承载层。

需求Pub/SubStream + 消费组
在线广播直接推送,延迟低需要额外读取逻辑
断线补读没有历史游标按 ID 继续读取
处理确认没有 ACK 状态XACK 记录确认
Redis Pub/Sub在线推送与断线丢失边界说明图
图1:Redis Pub/Sub 的在线广播与断线丢失边界说明图,不是运行截图。

用 Stream 把事件变成可恢复记录

迁移的关键不是把 PUBLISH 换成另一个命令,而是先改变写入语义:生产者用 XADD 把事件写入 Stream,消费者通过消费组领取消息。Stream 条目有 ID,消费组维护已投递但尚未确认的 Pending Entries List(PEL),进程重启后可以继续从历史 ID 或待确认列表处理。

下面的命令展示最小结构。示例中的订单字段只是说明数据关系,生产环境应补上事件版本、幂等键和产生时间。

# 写入持久化事件;* 让 Redis 生成单调递增的 Stream ID
redis-cli XADD order-events * order_id 1001 status paid event_version 1

# 首次初始化消费组;0-0 表示允许从已有历史开始读取
redis-cli XGROUP CREATE order-events order-workers 0-0 MKSTREAM

# 读取分配给 worker-a 的新消息;BLOCK 只控制等待时长
redis-cli XREADGROUP GROUP order-workers worker-a COUNT 10 BLOCK 5000 STREAMS order-events >

# 业务处理成功后再确认;不要在真正处理前提前 XACK
redis-cli XACK order-events order-workers 1690000000000-0

GROUP 后面是消费组和消费者名称,> 表示读取尚未分配给其他消费者的新消息。处理失败时不要立即确认,否则 Redis 会认为消息已经完成;应把消息留在待确认列表,交给重试或认领逻辑。

用待确认列表处理重试和断线恢复

消费者重启后,先处理自己留下的历史待确认消息,再读取新消息,避免只盯着 > 而把旧任务遗忘。运维侧用 XPENDING order-events order-workers 查看待确认数量、最小和最大 ID 以及消费者分布;对长时间未处理的条目,再结合认领命令转交给健康消费者。

# 查看消费组的待确认概况;用于发现断线或处理卡住的消费者
redis-cli XPENDING order-events order-workers

# 读取某个消费者历史上已经投递但未确认的消息
redis-cli XREADGROUP GROUP order-workers worker-a COUNT 10 STREAMS order-events 0

# 重试成功后确认同一条消息;业务接口必须按 event_id 做幂等
redis-cli XACK order-events order-workers 1690000000000-0

可靠消费不等于“业务只执行一次”。网络超时可能让消费者已经完成外部写入,却来不及发送 XACK,随后消息再次投递。因此要在业务表或下游接口保存稳定的 event_id,用唯一约束、幂等更新或去重记录抵挡重复执行。消息保留时长也要覆盖最长允许的断线恢复窗口,不能刚确认就无条件删除仍可能被审计或补偿的事件。

Redis Stream消费组待确认消息与重试认领关系说明图
图2:Redis Stream 消费组、待确认列表、重试与确认关系说明图,不是运行截图。

迁移时的参数与边界清单

先让生产者双写或在边界层把原事件转换为 Stream,再逐个切换消费者;不要先停止 Pub/Sub 后才尝试补历史,因为旧消息并不存在。压测时记录写入吞吐、消费延迟、Pending 数量和重试次数,分别观察正常运行、消费者重启和 Redis 连接恢复三种场景。

  • 顺序:同一业务分区使用稳定的 Stream 读序;多个消费者组之间不要假设全局处理顺序。
  • 容量:用明确的保留窗口或长度上限治理 Stream,保留策略要大于最长恢复时间。
  • 失败:重试超过上限后写入死信 Stream,并保留原始 ID、错误原因和最后一次时间。
  • 广播:多个独立系统都要收到同一事件时,为每个系统建立独立消费组,而不是让多个消费者共享一个组。

如果消息只是“有人在线就刷新一下”的状态通知,丢失不会产生业务后果,继续用 Pub/Sub 更简单;如果它代表订单状态、库存变化、任务指令或审计事件,就应优先使用 Stream 或专门的消息系统。

相关问题

Redis Pub/Sub 重连后能补回断线消息吗?

不能。重连只能重新订阅,无法获得断线期间没有保存的历史消息;需要补读时应从 Stream、数据库或其他持久化队列恢复。

Stream 消费成功后为什么还会重复处理?

业务完成与发送 XACK 不是一个原子动作,确认前断线就可能再次投递,所以消费处理必须设计幂等。

所有 Pub/Sub 都应该改成 Stream 吗?

不需要。只做在线广播、允许丢失且不需要补偿的通知保留 Pub/Sub;需要可恢复消费的事件才迁移。

声明:本文转载于:17golang原创 如有侵犯,请联系study_golang@163.com删除
相关阅读
更多>
最新阅读
更多>
课程推荐
更多>