当前位置:首页 > 文章列表 > 数据库 > Redis > Redis Pub/Sub断线期间消息丢失时的替代结构

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

来源:17golang原创 2026-09-25 14:45:00 0浏览 收藏

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删除
米坛社区登录入口在哪里?注册、资源与账号安全边界米坛社区登录入口在哪里?注册、资源与账号安全边界
上一篇
米坛社区登录入口在哪里?注册、资源与账号安全边界
模型评测集按任务簇拆分避免数据泄漏的设计
下一篇
模型评测集按任务簇拆分避免数据泄漏的设计
查看更多
最新文章
查看更多
课程推荐
  • 前端进阶之JavaScript设计模式
    前端进阶之JavaScript设计模式
    设计模式是开发人员在软件开发过程中面临一般问题时的解决方案,代表了最佳的实践。本课程的主打内容包括JS常见设计模式以及具体应用场景,打造一站式知识长龙服务,适合有JS基础的同学学习。
    543次学习
  • GO语言核心编程课程
    GO语言核心编程课程
    本课程采用真实案例,全面具体可落地,从理论到实践,一步一步将GO核心编程技术、编程思想、底层实现融会贯通,使学习者贴近时代脉搏,做IT互联网时代的弄潮儿。
    516次学习
  • 简单聊聊mysql8与网络通信
    简单聊聊mysql8与网络通信
    如有问题加微信:Le-studyg;在课程中,我们将首先介绍MySQL8的新特性,包括性能优化、安全增强、新数据类型等,帮助学生快速熟悉MySQL8的最新功能。接着,我们将深入解析MySQL的网络通信机制,包括协议、连接管理、数据传输等,让
    500次学习
  • JavaScript正则表达式基础与实战
    JavaScript正则表达式基础与实战
    在任何一门编程语言中,正则表达式,都是一项重要的知识,它提供了高效的字符串匹配与捕获机制,可以极大的简化程序设计。
    487次学习
  • 从零制作响应式网站—Grid布局
    从零制作响应式网站—Grid布局
    本系列教程将展示从零制作一个假想的网络科技公司官网,分为导航,轮播,关于我们,成功案例,服务流程,团队介绍,数据部分,公司动态,底部信息等内容区块。网站整体采用CSSGrid布局,支持响应式,有流畅过渡和展现动画。
    485次学习
查看更多
AI推荐
  • PubMedQA数据集详解:生物医学问答基准、功能与应用指南
    PubMedQA
    深入了解PubMedQA生物医学问答数据集,涵盖其核心功能、使用方法及在临床决策、药物研发等场景的应用,助力提升NLP模型性能。
    211次使用
  • H2O EvalGPT:开源LLM大模型评估与排行榜工具
    H2O EvalGPT
    H2O EvalGPT是H2O.ai推出的开源LLM评估平台,提供详细的大模型性能排行榜、行业特定基准测试及A/B测试功能,助您快速选择最适合项目的高性能大语言模型。
    265次使用
  • LMArena是什么?伯克利AI模型评估平台使用指南与功能解析
    LMArena
    LMArena是加州大学伯克利分校推出的AI模型匿名评测平台。通过盲测投票机制,用户可对比不同大模型回答并生成实时排行榜,助力开发者优化模型及用户选择最佳AI工具。
    222次使用
  • 斯坦福HELM:大语言模型Holistic Evaluation整体评估框架详解
    HELM
    深入了解斯坦福推出的HELM(Holistic Evaluation of Language Models)大模型评测体系。本文解析其核心功能、安装配置步骤及应用场景,涵盖准确性、公平性、鲁棒性等多维度指标,助力开发者全面优化语言模型性能。
    208次使用
  • CMMLU中文大模型评估基准:功能、使用教程与应用场景解析
    CMMLU
    深入了解CMMLU中文评估基准,涵盖67个学科主题,提供数据集下载、Zero-shot/Five-shot评估方法及排行榜,助力优化中文语言模型性能。
    198次使用
微信登录更方便
  • 密码登录
  • 注册账号
登录即同意 用户协议 和 隐私政策
返回登录
  • 重置密码