很多人第一次接触 Redis 里的 Stream 时,第一反应往往是:这不就是个消息队列吗?Kafka、RabbitMQ、RocketMQ 哪个不比它强?说实话,我一开始也是这么想的。但真正把它用在日志管道、订单事件流转、甚至轻量级多播通知场景之后,我才意识到 Redis 做这个数据结构的野心不在于替代专业消息中间件,而是提供一套足够轻、足够快、能落地的排队与事件分发方案。这篇就围绕 Stream 这个数据结构本身做一次拆解,把命令用法、消费组机制、底层存储设计和实盘中的坑讲透。
不管你是刚接触 Redis,还是已经用了几年只把它当缓存使,这篇文章读完,你至少能回答出这几个问题:Stream 和 List/PubSub 的本质区别在哪?消费组的 pending 状态到底是个什么机制?为什么 XADD 时要慎重考虑 MAXLEN?以及,Redis 官方为什么要用它做 Kafka 式的数据分区模拟。内容会偏底层一点,但我会尽量用大白话带出来。
1. 为什么 Redis 需要一个全新的"数据结构"
在 Stream 出现之前,Redis 如果要承担消息相关职能,主要靠 List 和 PubSub。这两个结构各有各的问题,而且问题恰好互补:List 能存,但消费模型弱;PubSub 模型灵活,但不落地。Stream 就是冲着这两个短板来的。
1.1 List 做消息队列时到底缺了什么
List 的模型是双向链表,配合 LPUSH、BRPOP 这组命令,可以很轻松地实现后进先出或者先进先出的队列。很多老项目甚至就把 List 当 MQ 用,消息生产方 LPUSH,消费方 BRPOP,阻塞读也很方便。这套方案最大的硬伤在于:一个消费者把消息 POP 出去,消息就从队列里消失了。
如果你只有单消费方,那没问题。但一旦需要多个消费者同时处理同一批消息,比如订单系统要同时把结果通知库存系统和积分系统,List 就抓瞎了——消息被 A 拿走之后 B 就什么都读不到了。你当然可以复制两份数据、建两个 key,但这不是一个数据结构层面解决的问题,纯粹是业务层在打补丁。更麻烦的是,如果消费者处理到一半挂掉,消息已经 POP 出来了但业务没做完,这条消息就彻底丢了,没有任何重试机制。
1.2 PubSub 的广播和丢失,也让人头疼
PubSub 是另一种思路:发布者把消息发到某个频道,所有订阅者都能收到。它解决了多消费方的问题,但又引入了新的麻烦——消息不落盘。发布那一刻没有订阅者在场,这条消息就直接丢了,之后订阅者也补不回来。这种模式适合在线状态广播、简单实时通知,想做可靠交付基本等于痴人说梦。
所以说,List 和 PubSub 各自都只覆盖了消息场景的一半需求。Stream 设计的核心目标,就是我既要能持久化存储,又要支持多个消费者各自维护自己的消费进度,还要能处理消息消费失败后的重试。说白了,它想对标的形态是 Kafka 之于 Redis 的降维版。
1.3 Stream 的定位:一个内存中的追加式日志
Stream 本质上是一个按时间序追加的日志结构。每一条消息有一个全局唯一的 ID,默认由 Redis 自己生成:毫秒级时间戳 + 序号。消费者可以通过游标按顺序读取,也可以在某个范围内随机访问,还可以给消息打标记(消费组 ACK)。它既像 List 一样存储在 Redis 内存中,又比 List 多了完整的消费进度跟踪和消费者组支持。
我之前用它做过一次促销活动的实时行为采集管道:前端行为事件 XADD 到 Stream,多个分析任务各开一个消费者组,互不干扰地读取全量事件。同一个 Stream,不同消费组消费进度完全独立,就像每个组拿了一条消息流的"副本游标"一样。这个能力是 List 给不了的,也是 PubSub 完全不可想象的。
2. 先跑通 Stream 的基础操作:写入与读取
理解了设计动机,下面直接上手。Stream 的命令不算多,核心就 XADD、XREAD、XRANGE、XDEL、XLEN 这几个,先把这几条玩明白,后面的消费组就不会懵。
2.1 XADD 写入:消息 ID 的生成规则
写入一条消息用 XADD,最基本的例子:
XADD order_events * event_type created order_id o_1001 user_id u_8001 amount 99.9这个命令里的*让 Redis 自动生成消息 ID。生成规则是:当前毫秒级 Unix 时间戳作为前半部分,序号作为后半部分;如果同一毫秒内有多条消息,序号自增。两条消息的 ID 形如:
1716220800000-0 1716220800000-1ID 在同一个 Stream 里必须单调递增,这也是它天然有序的原因。你也可以在 XADD 时显式指定 ID,比如从别的系统迁移历史数据时,可以用小一点的 ID,只要它大于 Stream 当前最大 ID 就行。不过日常使用没必要手动造 ID,让 Redis 生成最省心。
XADD 还有一个容易被忽略的选项是 NOMKSTREAM,如果 Stream 不存在,加上这个选项后命令直接返回 nil,不会创建一个空 key。这在纯生产侧脚本里挺有用的,能避免因为笔误写错 key 名而悄悄建出一堆垃圾 key。我从某次事故里学到的教训是:生产环境的写脚本能加 NOMKSTREAM 就加上。
2.2 XREAD 读取:从哪里读、阻塞多久
读消息用 XREAD:
XREAD COUNT 10 STREAMS order_events 0这个0表示从 Stream 里 ID 最小的那条消息开始读。如果只想读新增消息,可以传$,表示只读从现在开始产生的新消息。阻塞模式用 BLOCK:
XREAD BLOCK 5000 COUNT 10 STREAMS order_events $如果 5 秒内没有新消息,命令返回 nil;否则立刻返回新到的消息。BLOCK 0就是无限阻塞,直到有新消息。注意BLOCK之后超过网络空闲超时(比如 TCP 层面的 idle timeout),客户端连接可能被中间设备切断,读到一半抛"stream disconnected"之类的错误。所以生产环境里做阻塞读,至少要在客户端层面上处理重连和续读,不能指望一条长连接永远不断。
2.3 XRANGE 扫范围:排查问题和手工补数据的利器
XRANGE 和 XREVRANGE 用于按 ID 范围扫描消息,是调试 Stream 最舒服的地方。
XRANGE order_events 1716220800000-0 1716220900000-0超大型范围可以写-和+,分别代表最小和最大 ID。这个命令天然支持分页:每次取完最后一条 ID,下次把它当成起始游标继续往后扫就行。之前排查线上数据问题时,我就是靠 XRANGE 把某个时间窗口内的消息全部捞出来,一条一条对字段,很快定位到了脏数据来源。如果没有这个命令,只能写脚本逐条 XPENDING + XCLAIM,麻烦得多。
2.4 XDEL 与 XLEN:删除与统计
XDEL 按 ID 删除指定消息,XLEN 返回 Stream 当前的消息条数。实际使用中 XDEL 不是日常高频命令,因为 Stream 更多是当日志用,靠 XADD 的 MAXLEN 参数做自动裁剪,而不是频繁手工删。真正高频的是 XINFO STREAM,这个命令能告诉你 Stream 当前的条目数、最近消息 ID、消费组数量等,是体检数据结构的首选。
跑通了这五个命令,Stream 的"日志追加写 + 范围读"的基础模型就清楚了。它和 List 最大的区别是:list 的消费会弹出元素,Stream 的读取默认不动数据,消息还稳稳地躺在 Stream 里,直到你主动裁剪或删除。这也是它能支撑多消费组的前提。
3. 消费组:Stream 真正拉开差距的地方
如果说 XADD 和 XREAD 只是给 Stream 加了一层日志的外衣,那消费组就是 Stream 的灵魂。一个 Stream 可以被多个消费者组订阅,每个组维护自己的游标和历史状态,消息发到组里后,组内多个消费者分工处理,处理完还要显式 ACK。
3.1 创建消费组:XGROUP CREATE
消费组从 XGROUP 命令开始:
XGROUP CREATE order_events group_order_center 00表示这个组从头开始消费所有消息。如果想只消费组创建之后的新消息,用$。这个选择一旦定了,后期不容易改,所以建组前想清楚。还有一个 MKSTREAM 选项:当 Stream 不存在时,自动创建一个空 Stream。通常我会建议开发环境用 MKSTREAM,避免因为先建组后建流这种顺序问题报错;生产环境则更谨慎,显式确认目标 Stream 存在更安全。
3.2 XREADGROUP:组内消费的正确姿势
消费消息用 XREADGROUP,注意和 XREAD 的差别:
XREADGROUP GROUP group_order_center worker_1 COUNT 10 BLOCK 2000 STREAMS order_events >关键在于最后那个>。它表示这个消费者要读取的是组里"尚未被投递给任何消费者"的消息。如果不写>,而是写一个具体的消息 ID,那读取的是某个消费者自己的"待处理列表"里的历史消息,通常用于故障恢复后重新拉取未确认消息。这个细节非常容易踩坑,我第一次用的时候因为没搞懂>的作用,配了具体 ID 去消费,结果读出来的全是旧数据判定的逻辑,折磨了半天。
一个组里可以有多个消费者,每个消费者用自己的名称。Redis 的分发规则是:新消息按 round-robin 大致投递给组内不同消费者,不保证绝对的负载均衡,更不可能像 Kafka 那样按分区严格分配。别把 Kafka 的模型硬套到 Redis Stream 上来,否则你会纠结为什么某个消费者一直收不到消息。
3.3 ACK 与 PEL:不确认的消息不会丢
消费者读到消息后,处理业务,然后要主动告诉 Redis 这条消息处理完了:
XACK order_events group_order_center 1716220800000-0如果消费者读了一条消息,但是进程崩溃、超时,或者你故意不 ACK,这条消息就会一直留在消费者的 PEL(Pending Entries List)里。所谓 PEL,就是每个消费者各自的"处理中列表"。Redis 不会自动把超时的消息重新派发,你得通过 XPENDING 查看,再用 XCLAIM 把某个消费者的 pending 消息转移给另一个消费者继续处理。
XPENDING 的典型输出能告诉你:这个组总共有多少待确认、最早最晚待确认 ID、每个消费者的待处理情况。这是排查"消息怎么卡住了"的第一站。
XPENDING order_events group_order_centerXCLAIM 是自主转移:
XCLAIM order_events group_order_center worker_2 60000 1716220800000-0这个命令的含义是:把 group_order_center 组里 worker_1 名下的消息 ID 为 1716220800000-0 的这条,转给 worker_2,只要它"空闲"超过 60000 毫秒。通过这种"超时认领"机制,Stream 实现了一种非常原始但够用的故障转移。注意,XCLAIM 转移后,原消费者的 pending 记录就没了,新消费者要对这条消息负责并最终 ACK。
3.4 消费组模式和生产实践的结合方式
我真刀真枪用过一段时间后,认为 Redis Stream 消费组的适用边界是:消息量日均几十万到几百万、不需要跨机房容灾、不要求秒级延迟、但需要一定可靠性保障的场景。比如订单状态流转提醒、异步发送短信、爬虫任务调度。在这种场景下,一个 Stream 加几个消费者组,比引入一套 Kafka 运维成本低得多,效果也完全够用。
但要提醒的是:不要在单个 Stream 里塞所有业务消息。Redis 单线程模型的底子,决定了写入和读取的吞吐上限,一个实例扛全公司的消息流并不现实。合理姿势是按业务域拆分 key,再配合 Redis Cluster 做横向扩展,每个分片上的 Stream 各自独立。
4. 底层存储设计:listpack 和 Rax 树怎么协作
前面讲的都是命令层面,下面深入一点,看 Stream 在 Redis 内部到底长什么样。这个数据结构从 5.0 引入时就不是简单拿个链表凑合的,它内部由两层结构构成:宏观节点 listpack 和索引树 rax。
4.1 listpack 宏节点:消息的密集存储单元
Redis 6 及之后版本里,Stream 的单条消息不是独立对象,而是打包放在 listpack 里。listpack 是一种内存紧凑型线性结构,和 ziplist 类似,但设计上更干净,访问时不需要连锁更新。每个 listpack 节点大概保存 100 条以内的消息,超过阈值就换一个新的 listpack 节点。
这样做的好处非常直白:内存占用低。消息字段名、字段值都是紧挨着存的,省掉了每个消息对象单独的 dictEntry、robj 等元数据开销。这一点在只有几十个 key 的老 Consumer 上体现不出来,但当你用 Stream 存几百万条日志时,省下来的内存是肉眼可见的。
4.2 rax 索引树:按 ID 快速定位
有了 listpack 节点,怎么快速按 ID 找消息?Redis 用一个 rax(基数树)来做索引,key 是消息 ID(64 位时间戳 + 64 位序号拼接而成的高字节优先的二进制串),value 指向对应的 listpack 节点。
查询逻辑大致是这样的:先按 ID 前缀在 rax 树中定位到离它最近的 listpack 节点,然后在这个宏节点内部做线性扫描。因为单个宏节点内消息量很小,线性扫描成本可以忽略。这种"索引树 + 紧凑块"的组合,相当于在内存中实现了一个迷你版本的 SSTable 索引,兼顾了定位效率和压缩率。
4.3 每个宏节点的内部有序性
宏节点内的消息按 ID 升序排列,节点之间也按最大最小 ID 严格有序。这就让 XRANGE 在扫大范围时非常舒服:rax 树快速跳到起始位置附近,然后一路向后遍历宏节点,中途不需要做任何排序操作。Stream 之所以能支持范围查询和游标遍历,这套结构是根基。
我早期曾经想当然地觉得 Stream 内部就是普通链表,或者直接用有序集合 ZSET 就能模拟。实际上,如果用 ZSET 存储消息,每条消息至少要有 score 和 member 两个对象,内存开销会高出不少;Stream 用 listpack + rax 把消息体紧凑存放,配合专门的遍历逻辑,才是专门为"日志流"优化的形态。这就是"数据结构"的意义所在——它不只是一个 Redis 模块,而是 Redis 内部精心设计的一等公民。
5. 消费组的内部状态:pending 列表与消息流转
如果你只用 Stream 做简单的日志追加和单消费者读取,不需要深究消费组内部。但只要开了消费组,你就必须理解 Redis 在后台维护了哪些状态。因为所有诡异现象,基本都能从这些状态里找到解释。
5.1 消费者组的三个核心元信息
每个消费组维护着三样东西:
- last_delivered_id:组内下一个要投递的消息 ID,只前进,不后退。
- pending_entries:一个按消费者分组的待确认消息集合。
- consumers 列表:组内所有消费者的名字,以及它们各自读到的最后一条消息 ID。
这些元信息也存放在 Stream 内部的 rax 树结构中,Redis 对它们的读写有专门的优化路径。特别是 pending_entries,在 Redis 7.0 之前是一个独立的基数树,每个节点记录消息归属的消费者名、delivery time 和 delivery count 等信息。Redis 7.0 做了一次比较大的重构,把同一个消费者的待处理消息单独抽出,消费组整体的 pending 树不再为每个消费者重复保存完整的消息 ID,节省了大量内存。
5.2 一个消息从投递到 ACK 的完整旅程
严格来说,一条消息在消费组里的生命周期是:
- XADD 写入 Stream,属于未投递状态。
- XREADGROUP 读取后,进入投递状态,记录到 pending_entries。
- XACK 返回后,从 pending_entries 移除,标记完成。
- 如果一直不 ACK,就停留在 pending_entries,等待 XCLAIM 或 XAUTOCLAIM 转移。
实际项目里最常见的问题,就是消费者拉到了消息,处理成功后忘记 XACK,或者在确认之前进程崩溃。结果是 Stream 里消息还在,但消费者永远不再读取它们(因为>只投递新消息),pending 越积越多。我用 XAUTOCLAIM 做过一次大规模清理:把超时 10 分钟的消息全部认领回来重新处理,同时扫描 pending 列表里 delivery count 超过 5 次的消息,直接标记成死信做补偿告警。
5.3 Redis 7.0 对消费者组的大改:不只是在省内存
Redis 7.0 引入的 consumer groups v2,在很多细节上比以前顺滑了。最明显的是内存优化,前面说的 pending 列表不再重复存储消费者名;再者是 XAUTOCLAIM 命令的引入,以前你得先 XPENDING 查出超时 ID 列表,再循环 XCLAIM,现在一条命令就能"扫描并认领"一个批次,操作体验接近专业 MQ 的 dead letter 概念。像这种痛点,官方花了几个大版本才补完,也从侧面说明 Stream 的消费组从设计到成熟经历了相当长的时间。
我个人建议,如果你项目里的 Redis 版本还停留在 5.x 或 6.x,用 Stream 消费组时尽量评估升级到 7.x。不只是为了省内存,更是为了 XAUTOCLAIM 这个运维刚需。
6. 实战踩坑记录与内存控制经验
文章最后一部分,把我实际部署 Stream 过程中踩过的坑和一些内存控制的经验交个底。这些东西教科书上不会写,但遇到一次能耽误你半天。
6.1 大消息的坑:一条消息塞了几百 KB
Stream 的底层 listpack 是为小字段优化的,一条消息塞个几百 KB 的 JSON,listpack 会比较吃力,内存占用会明显偏高。更要命的是大对象可能触发底层节点分裂、重写开销,对单线程 Redis 来说是雪上加霜。我的经验是:Stream 单条消息体尽量控制在 1KB 以内,最多不超过 10KB。真要传大量数据,把内容转存到对象存储,Stream 里只放引用 ID。
6.2 无限增长的问题:必须配 MAXLEN 或 MINID
很多人上线时不会给 Stream 设计裁剪策略,结果几天之后内存被打满。XADD 自带裁剪参数:
XADD order_events MAXLEN ~ 100000 * field value ...MAXLEN ~ 100000表示近似保留约 10 万条消息。那个~符号是关键,它告诉 Redis 可以"大概"裁剪,不需要精确到每条。近似裁剪的实现是:在宏节点粒度上删除过旧消息,一次删一整块,性能和精确裁剪的天壤之别。精确裁剪需要一条一条遍历删除,在大 Stream 上可能导致明显的阻塞;近似裁剪基本只要删除 rax 树左侧的空节点,成本接近 O(1)。
MINID 是另一种裁剪方式,指定保留到哪个消息 ID 之前的所有消息都丢掉:
XADD order_events MINID 1716220800000-0 * field value ...对按时间语义保留日志的场景,MINID 比 MAXLEN 好用得多,因为你可以直接说"保留 24 小时内的消息",不用先换算条数。但无论选哪个,都必须显式配置,Redis 不会帮你自动清理。
6.3 trim 造成的删除开销也不容忽视
就算用了 MAXLEN ~,如果 Stream 的写入速率忽高忽低,裁剪动作可能集中在某个时间点爆发,导致短暂延迟毛刺。我在一次大促流量模拟时见过实例的延迟从 1ms 跳到 300ms,排查半天发现是一批超大写入触发了深裁剪。解决方案是给 XADD 加上 LIMIT 参数,限制一次操作最多删多少节点,把裁剪摊平到多次写入里。
6.4 消费组消息积压与内存配额
开了消费组的 Stream 比普通 Stream 更需要注意:如果某个组长时间不消费,所有新消息都会积压在 Stream 里,同时每个消息都会在 pending 树上留下记录。内存占用往往是"Stream 本体 + 消费组元数据"双份增长。监控时除了看 Redis 总体内存,还要看每个 Stream 的 XINFO STREAM 里的 length,和每个消费组的 XPENDING 数量。我一般会给核心 Stream 建一条定时巡检脚本:发现 pending 超过阈值就触发消费者扩容和死信核查。
6.5 客户端阻塞读与连接断开
还有一个高频问题:客户端用 XREADGROUP BLOCK 长时间阻塞,网络设备把空闲连接断开,客户端抛 stream disconnected 错误。Redis 端不一定有感知,但业务方的重试逻辑必须设计好。一个稳妥做法是:阻塞读的错误处理里区分超时(返回 nil)和断连(抛异常),断连时不要立即重新 XREADGROUP>,而是先用 XAUTOCLAIM 把上一个阻塞周期里可能已投递但未 ACK 的消息回收,避免消息重复消费或永久滞留。这套逻辑我花了不少时间才调顺,真到故障演练那天,才知道它有多值。
个人经验来说,Redis Stream 不是一个能替代一切消息队列的银弹;它更像是一块非常称手的"轻量消息基础设施"。如果你的团队本身就有 Redis 运维能力,消息量在百万级以内,不想再部署和维护一套独立的中间件,那 Stream 绝对值得认真研究。要是哪天项目规模上去了,要严格的事务性消息、跨地域容灾、海量消息堆积,再迁移到 Kafka 或者云上托管 MQ 也不迟——毕竟 Stream 的消费组模型学起来并不亏,很多概念都能平移过去。