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

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

来源:17golang原创

时间:2026-10-07 13:15:54 220浏览 收藏

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删除
相关阅读
更多>
最新阅读
更多>
课程推荐
更多>