
很多人第一次接触 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 读取从哪里读、阻塞多久读消息用 XREADXREAD COUNT 10 STREAMS order_events 0这个0表示从 Stream 里 ID 最小的那条消息开始读。如果只想读新增消息可以传$表示只读从现在开始产生的新消息。阻塞模式用 BLOCKXREAD 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这条消息就会一直留在消费者的 PELPending 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 是消息 ID64 位时间戳 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 大消息的坑一条消息塞了几百 KBStream 的底层 listpack 是为小字段优化的一条消息塞个几百 KB 的 JSONlistpack 会比较吃力内存占用会明显偏高。更要命的是大对象可能触发底层节点分裂、重写开销对单线程 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 的消费组模型学起来并不亏很多概念都能平移过去。