当前位置:首页 > 文章列表 > 数据库 > Redis > Pub/Sub 与 Streams 不只是是否持久化:订阅模型怎么选

Pub/Sub 与 Streams 不只是是否持久化:订阅模型怎么选

来源:17golang原创 2026-10-07 13:15:54 0浏览 收藏

Redis Pub/Sub 与 Streams 的区别不只是“一个不持久化、一个持久化”。真正决定选型的是谁应该收到消息、订阅者离线后是否补读、多个实例是全员广播还是分摊任务、处理失败后由谁恢复。只给当前在线连接推送可丢通知,选 Pub/Sub;要保存事件并由每个读者维护位置,选 Streams 的 XREAD;要让一组工作实例分摊任务并确认处理,选 Streams Consumer Group。

官方文档:https://redis.io/docs/latest/develop/data-types/streams/

选型速查
  • 在线广播、允许断线丢消息:Pub/Sub。
  • 每个读者都要独立回放同一事件流:Streams + XREAD,应用保存最后 ID。
  • 同一服务的多个实例只需处理一份任务:Streams + XREADGROUP。
  • 多个业务服务都要处理每条事件:为每个业务建立独立消费组,不要都塞进同一个组。

三种模型先按谁能收到消息来区分

维度Pub/SubStreams + XREADStreams + 消费组
接收关系所有当前在线订阅者每个读者按自己的 ID 读取同组消费者分摊,多个组彼此独立
离线补读不支持支持支持,并记录待确认条目
投递语义至多一次由读者位置管理决定通常按至少一次设计
服务端状态频道与在线连接持久 Stream 条目条目、组游标、消费者与 PEL
典型场景实时通知、缓存失效广播回放、审计、独立订阅任务分摊、失败重试、积压治理

最容易选错的是“广播”。Pub/Sub 天然把一条消息推给所有在线订阅者;Streams 消费组则保证一条新消息在同一组内交给一个消费者。想让计费和搜索服务都看到订单事件,应建立 billing 与 search 两个组;如果把它们放进同一个组,事件会被二者分摊。

Redis Pub/Sub、Streams 独立读取与消费组分摊的静态关系图
图1:订阅模型结构图。Pub/Sub 面向当前在线订阅者广播,XREAD 由各读者维护位置,XREADGROUP 则在同一组内分摊条目。

Pub/Sub 解决的是在线广播

Redis 官方说明 Pub/Sub 使用至多一次语义:消息由服务器发出后不会重发,订阅者因网络断开或处理错误错过消息,就永久丢失。它的优势是模型轻、延迟低、发布者无需维护消费进度,适合“此刻在线的人看到即可”的通知。

# 订阅当前在线消息;断线期间的消息不会补发
redis-cli SUBSCRIBE ui:updates

# 向频道广播一次界面刷新通知
redis-cli PUBLISH ui:updates '{"type":"refresh","scope":"orders"}'

缓存失效提示、WebSocket 节点间转发、在线状态变化都可能适合 Pub/Sub,但前提是业务真能容忍丢失。若“缓存失效消息丢了就会长期读到旧值”,则不能只依赖频道通知,还要有过期时间、版本号或可重建状态兜底。

Streams 的 XREAD 是独立位置读取

Stream 是带唯一 ID 的追加日志。消息写入后会保留,直到显式删除或按策略裁剪。XREAD 让读者从某个 ID 之后读取;不同读者可以各自保存最后处理 ID,因此每个读者都能看到同一批条目,也能在重启后继续。

# 写入事件,并用近似 MAXLEN 控制日志不会无限增长
redis-cli XADD orders:events MAXLEN '~' 100000 '*' type paid order_id 9001

# 从指定 ID 之后读取;应用应持久保存最后成功处理的 ID
redis-cli XREAD COUNT 20 BLOCK 5000 STREAMS orders:events 0-0

XREAD 不会替应用保存“谁处理过什么”。如果每个服务都必须获得完整事件历史,可以让它们分别维护游标;但实例很多时,游标存储、故障接管和重复处理都要自己实现。需要服务端协作分摊时,消费组更合适。

消费组把组内实例变成协作消费者

消费组维护组级最后投递位置,并为每个消费者记录待确认条目。使用 > 读取时,新条目只交给组内一个消费者;同一 Stream 可以建立多个消费组,每个组都会独立获得条目。

# 从现有历史起点创建计费组;MKSTREAM 可在键不存在时创建 Stream
redis-cli XGROUP CREATE orders:events billing 0 MKSTREAM

# worker-1 阻塞读取组内尚未投递的新消息
redis-cli XREADGROUP GROUP billing worker-1 COUNT 10 BLOCK 5000 STREAMS orders:events '>'

# 业务处理成功后确认;确认前条目会留在该组的 PEL 中
redis-cli XACK orders:events billing 1740000000000-0

这里的 XACK 只把消息引用从当前消费组的 Pending Entries List 中移除,并不等于删除 Stream 条目。条目是否删除或裁剪,是独立的数据保留决策;多个消费组并存时尤其不能把确认和全局删除混为一谈。

消费组的代价是确认、积压和恢复责任

消费者拿到消息后如果崩溃,消息会继续挂在它名下的 PEL。Redis 不会猜测何时转交;应用需要用 XPENDING 观察未确认数量和空闲时间,再用 XAUTOCLAIM 把长时间未处理的条目交给健康消费者。

# 查看计费组的未确认摘要,确认是否出现积压
redis-cli XPENDING orders:events billing

# 将空闲超过 60 秒的待确认条目转给恢复消费者
redis-cli XAUTOCLAIM orders:events billing worker-recovery 60000 0-0 COUNT 20

重领意味着同一业务事件可能再次执行,所以消费者必须按业务键幂等。例如用 order_id + event_type 建立处理记录,先判断是否已经生效,再提交副作用并确认。不要把“Redis 只返回一个消费者”误解为端到端恰好一次。

Redis Streams 消费组 PEL、确认、重领与保留策略的静态关系图
图2:Streams 消费责任账本。消费组记录已投递未确认条目,处理成功后确认,停滞条目则需要观测和重领。

从 Pub/Sub 迁移到 Streams 时要补齐什么

把 PUBLISH 换成 XADD 只是迁移的开始。旧模型没有消费进度和积压,Streams 会把这些责任显式化。迁移评审至少包含下面六项。

迁移项要做的决定
事件 ID使用 Stream ID 作为读取位置,业务字段另带幂等键
广播边界每个独立业务建立消费组;同组只放可分摊的实例
起始位置从历史起点、某个 ID 或只读新消息,必须明确
确认时机业务副作用成功后再 XACK
失败恢复设置 pending 告警、最小空闲时间、最大投递次数与隔离策略
容量按最慢消费者和审计需求确定 MAXLEN 或 XTRIM 规则

切换期若同时写 Pub/Sub 与 Stream,要承认它们是两个独立写操作:任一写入失败都可能造成两边短暂不一致。迁移逻辑应指定权威来源、重试策略和切换完成条件,而不是默认“双写就等于可靠”。

回归检查与运行指标

Pub/Sub 主要观察在线订阅数、发布速率和客户端断线;Streams 除了写入和读取吞吐,还要关注 Stream 长度、消费组 lag、PEL 数量、最老 pending 空闲时间和重复处理率。消息系统从无状态广播升级为可恢复日志后,运维成本也随之升级。

# 查看 Stream 长度和组信息,核对积压与消费位置
redis-cli XLEN orders:events
redis-cli XINFO GROUPS orders:events

# 查看 Stream 元数据,确认首尾 ID 与保留结果
redis-cli XINFO STREAM orders:events

压测时不要只比较每秒命令数。Pub/Sub 没有落盘条目、PEL 和确认管理,Streams 为可回放和恢复付出了存储与状态成本;真正应比较的是业务能接受的丢失、重复、恢复时间和积压上限。

常见问题

每个消费者都要收到消息,能使用同一个消费组吗?

不能。同组消费者分摊条目。每个业务都要独立处理时,为每个业务建立不同消费组,或使用各自维护 ID 的 XREAD。

XACK 后消息是否从 Stream 消失?

传统 XACK 只清理当前组的待确认引用,条目仍保留在 Stream 中,直到删除或裁剪。不要用确认代替保留策略。

Streams 是否保证不会重复?

不能把消费组等同于端到端恰好一次。消费者处理后未成功确认、消息被重领等情况都会带来重复机会,业务副作用必须幂等。

只要消息重要就一定选 Streams 吗?

若权威状态已保存在数据库,通知只是提示在线节点刷新,Pub/Sub 仍可能足够;若通知本身就是必须处理的业务事实,才需要 Streams 或其他具备持久与恢复能力的消息系统。

版本声明
本文转载于:17golang原创 如有侵犯,请联系study_golang@163.com删除
go.work 应不应该提交到仓库,团队协作边界如何定go.work 应不应该提交到仓库,团队协作边界如何定
上一篇
go.work 应不应该提交到仓库,团队协作边界如何定
通过 go work sync 对齐工作区构建列表与模块依赖
下一篇
通过 go work sync 对齐工作区构建列表与模块依赖
查看更多
最新文章
查看更多
课程推荐
  • 前端进阶之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模型性能。
    364次使用
  • H2O EvalGPT:开源LLM大模型评估与排行榜工具
    H2O EvalGPT
    H2O EvalGPT是H2O.ai推出的开源LLM评估平台,提供详细的大模型性能排行榜、行业特定基准测试及A/B测试功能,助您快速选择最适合项目的高性能大语言模型。
    420次使用
  • LMArena是什么?伯克利AI模型评估平台使用指南与功能解析
    LMArena
    LMArena是加州大学伯克利分校推出的AI模型匿名评测平台。通过盲测投票机制,用户可对比不同大模型回答并生成实时排行榜,助力开发者优化模型及用户选择最佳AI工具。
    433次使用
  • 斯坦福HELM:大语言模型Holistic Evaluation整体评估框架详解
    HELM
    深入了解斯坦福推出的HELM(Holistic Evaluation of Language Models)大模型评测体系。本文解析其核心功能、安装配置步骤及应用场景,涵盖准确性、公平性、鲁棒性等多维度指标,助力开发者全面优化语言模型性能。
    386次使用
  • MMBench详解:多模态大模型基准测试、功能特点与使用指南
    MMBench
    MMBench是由上海人工智能实验室等机构联合推出的多模态基准测试平台,提供细粒度能力评估、大规模数据集及VLMEvalKit工具。本文详细介绍其核心功能、安装使用方法及应用场景,助力开发者全面评估多模态模型性能。
    213次使用
微信登录更方便
  • 密码登录
  • 注册账号
登录即同意 用户协议 和 隐私政策
返回登录
  • 重置密码