先讲一个我实际经历过的场景。有段时间帮朋友排查线上问题:用户下单支付后,通知服务偶尔会漏发短信,后台日志里看不到任何异常,但就是有一部分订单没有任何后续动作。他们的实现很朴素——用 Redis List 当消息队列,LPUSH 塞消息,BRPOP 消费,消费完这条消息就没了。问题就出在这:消费端进程一旦在 BRPOP 返回之后、业务处理完成之前崩溃,这条消息就永久性丢失,支付完成了但短信、积分、发票都没动静。这类问题在 Redis 5.0 之前几乎是无解的,因为 List 根本没有"确认消费"的概念。而 Redis 5.0 引入的 Stream 数据类型,正是冲着这个长期空白来的。
Redis Stream 是一套真正意义上的、内置在 Redis 里的消息队列模型。它支持消息持久化、消费者分组、消费确认(ACK)、断点续读和按时间回溯,把之前用 List、Pub/Sub 拼凑伪队列时踩过的坑一次性填平了。这篇文章我会从设计动机讲起,把 Stream 的核心概念、关键命令、实战代码和踩坑记录完整过一遍,适合已经在用 Redis、但还没系统接触过 Stream 的开发者,也适合准备面试时需要把消息队列这块讲清楚的场景。
1. 为什么需要 Stream:Redis 以前的"伪队列"差在哪
1.1 用 List 实现队列的核心痛点
用 List 做队列是 Redis 社区最古老的玩法:生产者 LPUSH 塞消息,消费者 BRPOP 弹出消息。简单场景下它能跑,但深入一想全是漏洞。最致命的一点是,BRPOP 把消息弹出来之后,这条消息就从 Redis 里消失了,消费端有没有处理成功,Redis 完全不关心。消费端拿到消息后,可能在写数据库之前宕机,可能在调用外部 API 时超时重试导致重复,也可能处理逻辑抛了异常但消息已经没了。丢了就是丢了,没有任何"重试"或"补偿"的入口。
第二个痛点是多消费者组的问题。假设一条消息既要同步到搜索服务,又要同步到报表系统,还要推送给用户。用 List 只能做到多个消费者"竞争"同一条消息——谁抢到谁消费,别的消费者就永远看不到。要想让多个下游各自拿到完整数据流,你得维护 N 份 List,生产端往每个 List 里各写一份,消费端各自维护消费进度。数据一致性、重复写入、消费位点存储全是问题。
第三个痛点是消费位点的管理。List 一旦弹出就没有了,无法回答"昨天下午消费到哪一条了""这三天积压了多少消息"这类问题。虽然可以用 LLEN 看到堆积数量,但做不到按 ID 回溯历史数据。运维排障的时候,这条消息是什么时候产生的、内容是什么、有没有被消费过,List 模式下完全是一团黑。
1.2 Pub/Sub 为什么只能当"广播放大器"
有人会问:要广播效果,不是有 Pub/Sub 吗?确实,PUBLISH/SUBSCRIBE 能做到一对多广播,但它的问题更明显:消息发出去之后 Redis 不保存,订阅者不在线,消息就错过了。订阅者消费速度跟不上发布速度,消息直接丢弃,没有背压一说。官方文档都写得很明白,Pub/Sub 的定位是"即发即弃"的广播通道,适合实时在线推送这种场景,不适合做业务消息队列。订单消息丢了可没法跟用户说"不好意思刚才你下线了,所以订单没创建成功"。
1.3 那个长期存在的空白位
所以在 Stream 出现之前,Redis 生态里一直缺一样东西:一个持久化、可确认、支持多消费组、能回溯的消息队列。业务团队要么凑合着用 List 加数据库表手动维护消费状态,要么直接引一个 Kafka 或 RabbitMQ,为一个偶尔才用到的消息队列引入一套重资产。Kafka 再轻量也得部署 Zookeeper(新版本虽然去掉了,但集群运维成本依然不低),对于很多中小团队来说,这是一个不小的负担。
Redis Stream 的出现恰好补上了这个中间档。它不要求你新部署任何东西,Redis 5.0 以上自带,API 风格和 Redis 其他命令一样简单直接。它在内存里做存储,天然速度快;通过消费者组机制,实现"组与组之间广播、组内竞争消费"的完整消息模型。更关键的是,它内置了 ACK 确认机制,消息处理完需要显式确认,没确认的消息可以重新投递,从机制上保证了消息不丢。
2. Stream 核心设计拆解:消息 ID、消费者组、PEL 与 ACK
2.1 消息 ID 不是简单的自增数
在 Stream 里,每条消息都有一个全局唯一的 ID,格式统一为<毫秒时间戳>-<序列号>,比如1689233456789-0。这个 ID 不是随便设计的:前半段是消息产生的毫秒时间戳,后半段是同一毫秒内的自增序号,Redis 保证在同一个 Stream 内,新生成的 ID 一定比旧 ID 大。也就是说,消息天然按时间排序,按 ID 读取就是按时间顺序读取。这让"从某条消息开始继续消费"变得极其简单——记住自己上次消费到的 ID 就行了。
这个设计还带来一个额外好处:你可以只看 ID 就知道消息大约是什么时候产生的,做按时间范围的裁剪或者回溯都很方便。我自己排查线上问题时,XINFO STREAM看 last-generated-id 的毫秒时间戳,就能直接判断出这个 Stream 是不是已经很久没有生产消息了。
需要注意的是,虽然 ID 包含服务器时间,但多实例同时写同一个 Stream 时,如果不同服务器的时钟有偏差,ID 的先后顺序和实际业务发生的先后顺序可能对不上。我后面在避坑章节会细说这个问题。
2.2 Consumer Group:组与组广播,组内竞争
消费者组是 Stream 相比 List 的最大突破。一个 Stream 下可以创建多个消费组,比如orderNoticeGroup、searchSyncGroup、reportGroup。每一条新消息,都会完整地投递给每一个消费组——这是组与组之间的广播关系。但在同一个组内部,消息只会被投递给其中一个消费者,组内成员竞争消费,避免同一条消息被重复处理。
这个模型和 Kafka 的消费者组几乎一致,理解起来不费劲。那么它底层是怎么实现的?关键在于每个消费组独立记录自己的消费位点(last-delivered-id),互不干扰。消费者 A 在组里读到 ID 为1689233456789-0的消息,组位点推进了;而另一个组的位点还停在更早的位置,不影响它从旧消息开始读。List 做不到这一点,因为 List 只有一个"弹出的栈顶";Stream 则通过组位点把每组的消费进度彻底解耦了。
另外,组里的消费者是"虚拟"的概念,不需要像创建用户那样显式注册。你执行XREADGROUP GROUP orderGroup consumer-1时,consumer-1这个消费者就自动出现在组里了。这个设计极大简化了消费者节点的上线和下线——新起的消费进程只要指定组名和消费者名就能立即开始工作,不用手工维护成员列表。
2.3 PEL 和 ACK:消息可靠的灵魂
每个消费组内部,还有一张"待确认消息列表",官方叫法 Pending Entries List(PEL)。当一条消息被XREADGROUP投递给某个消费者后,它的 ID 会进入该组的 PEL。消费者处理完这条消息,需要调用XACK把这条消息从 PEL 里移除,表示"我处理完了"。如果迟迟不 ACK,消息就一直挂在 PEL 里,Redis 就知道这个消费者可能出问题了。
PEL 的核心价值在于"可追溯、可恢复"。它记录了每条 pending 消息的消费者、空闲时间、投递次数。当某个消费者崩溃,它的 pending 消息会一直留在 PEL 里,其他消费者可以通过XCLAIM(Redis 6.2 之后还能用XAUTOCLAIM)把这些超过一定空闲时间的消息"接管"过来继续处理。这样,消息从投递到确认,全程都有据可查,不会因为一个节点挂掉就丢消息。
这套机制实现的语义是 at-least-once:每条消息至少被处理一次,但在极端情况下可能被处理多次——消费者处理完之后还没来得及 ACK 就宕机了,消息被其他消费者接管后又处理了一遍。所以你在设计消费逻辑时,一定要保证幂等性,这个我会在避坑章节专门说。
2.4 内存控制:MAXLEN 和 MINID
Stream 既然持久化存储在 Redis 里,就必须考虑内存。给 Stream 加"只保留最近 N 条消息"的约束非常简单:XADD的时候带上MAXLEN参数,或者用XTRIM命令来裁剪。XTRIM mystream MAXLEN 1000表示只保留最新的 1000 条;MINID模式则按消息 ID 裁剪,比如XTRIM mystream MINID 1689000000000,删除所有 ID 早于指定时间戳的消息,这个模式适合按时间保留数据的场景。
这里有一个实用技巧:MAXLEN可以加~符号,写成XTRIM mystream MAXLEN ~ 1000。意思是"近似裁剪到 1000 条",Redis 在最方便裁剪的节点批量删除,性能更好。对大多数场景来说,差几条完全无感。从 Redis 6.2 开始,还支持给XADD的MAXLEN加LIMIT参数,控制单次裁剪的条数,避免大 Stream 裁剪时阻塞过久。真实业务中,我建议直接把裁剪策略设计在写入路径上,而不是等内存报警了再手动补救。
3. 它到底解决了什么问题:三个真实场景复盘
3.1 场景一:秒杀场景的削峰填谷
高并发场景下,瞬时流量直接把订单服务打挂是常事。用 Stream 做削峰填谷的思路是:用户下单请求进来,写一条 Stream 消息到seckillOrders,立刻给用户返回"排队中";真正的下单逻辑放到消费端,按自己的处理能力从 Stream 里顺序拉取消息处理。因为 Stream 有持久化能力,即使消费端在峰值期间不够快,消息也会积压在 Stream 里,不会丢失;等流量波峰过去,消费端继续把积压的消息消费完。
有一个细节值得注意:在这个场景里,消息的"消费进度"由消费组的位点管理,而不是由调用方管理。即使所有消费者都宕机了,重启后从位点继续消费即可,不需要额外维护"这个用户提交过没有"的状态。相比于 List 模式,省掉了很多手工补偿逻辑。而且如果某个消费者处理一条消息时挂了,这条消息会一直留在 PEL 里,被其他消费者接管后重试,避免"订单已经在支付中被漏处理"的问题。
3.2 场景二:异步任务和事件驱动的标准姿势
很多系统里有大量异步任务:注册后发欢迎邮件、订单完成后发票推送、上传文件后的图片压缩。这类任务的共同特点是:对结果没有强实时性要求,但对可靠性有要求。用 Stream 做事件总线非常顺手——业务代码只管XADD写事件,异步 worker 用XREADGROUP消费,处理完XACK。Producer 和 Consumer 完全解耦,Producer 不需要关心现在有几个 worker、worker 挂没挂。
这里我强烈推荐用阻塞读取。XREADGROUP支持BLOCK参数,和 BRPOP 一样可以让消费者阻塞等待新消息,而不是空转轮询。设置一个合理的超时时间(比如 5000 毫秒),消费者在没有新消息时进入阻塞,有消息了立即返回,既保证了实时性,又不会让 Redis 被无效的轮询请求打满。从实际效果看,空轮询的空耗远高于阻塞读,改用 BLOCK 之后,我们那台 Redis 的 CPU 使用率直接降了 30% 以上。
3.3 场景三:数据同步与日志类的顺序处理
再往上走一个台阶,是典型的"多条数据需要广播给多个下游"的场景。比如订单数据要同步给数据仓库、搜索索引、推荐系统。每个下游的消费能力和故障恢复节奏都不一样,不可能共用同一个消费进度。Stream 的多消费组模型,让每个下游建一个消费组,各读各的,互不影响。数据仓库那边挂了两天,重启后从自己的消费组位点继续补读即可,不影响搜索索引的实时性。
更妙的是,Stream 的消息天然有序,按 ID 递增排列,这让顺序处理变得异常简单。对"必须先处理订单创建、再处理订单状态变更"这类有顺序依赖的事件流,Stream 不需要像多线程队列那样靠锁去保证顺序,只要消费者按 ID 顺序处理就行。而如果消息量特别大,一个组里可以挂多个消费者,Redis 会按 ID 顺序轮转投递,保证同一时刻组内消费的消息依然是有序的。
3.4 Stream 的边界:它不适合干什么
Stream 不是银弹,有些场景我不建议硬上。首先是超大批量消息存储:Stream 的数据在内存里,Redis 的内存容量决定了它能存的消息总量。如果你一天要消化几十亿条消息,或者要求消息保留几个月甚至几年,Stream 不合适,Kafka 这类磁盘型消息队列才是正确选择。其次是真正意义上的分布式事务,Stream 没有事务消息、没有死信队列、没有消息路由规则,这些能力在 RabbitMQ 里是开箱即用的,Stream 需要你自己在消费端实现。
我还经常被问到"Stream 能不能替代 Redis 分布式锁",这是两个完全不同的东西。分布式锁解决的是互斥问题,用 SETNX 加过期时间实现;Stream 解决的是消息流转问题。它们可以配合使用,但不会互相替代。
4. 实操:从生产到消费,把 Stream 链路完整跑通
4.1 生产端:XADD 写消息
先看最基本的写入命令。假设我们有一个订单事件流:
> XADD orders * orderId 1024 userId 88 status PAID "1689233456789-0"orders是 Stream 的 key,*让 Redis 自动生成消息 ID,后面跟着的是键值对形式的消息内容。Stream 的一个 entry 本质是一组字段和值的映射,类似一个小的 Hash。如果你不想让 Redis 自动生成 ID,也可以自己指定一个 ID,只要保证后写入的 ID 比已有 ID 大就行。
写入时推荐顺手带上裁剪策略:
> XADD orders MAXLEN ~ 10000 * orderId 1025 userId 99 status PAID这条命令把 Stream 的长度近似控制在 10000 条以内,不用再单独执行 XTRIM。生产环境我一般都会加这个参数,因为 Stream 如果只写不裁,内存增长会非常隐蔽,等发现时可能已经占用几个 GB 了。
4.2 消费端:XREAD 和 XREADGROUP 怎么选
直接读用XREAD,适合简单的"我从某个位置开始读"场景:
# 从 ID 0 开始读前 10 条(即从头开始) > XREAD COUNT 10 STREAMS orders 0 # 只读最新的消息,阻塞等待 5000 毫秒 > XREAD COUNT 10 BLOCK 5000 STREAMS orders $$表示"从最新的消息开始",也就是只读未来的新消息。注意XREAD不带消费者组,读过的消息不会记录到任何 PEL 中,下次还可以从同样的位置再读。这种模式适合日志收集、手动排查这类不需要分配协作的场景。
多消费者协作必须用XREADGROUP。假设我们已经在 orders 上建了一个orderGroup:
> XGROUP CREATE orders orderGroup 00表示这个消费组从第一条消息开始消费。如果想忽略历史消息,只处理创建组之后的新消息,给$即可。
然后消费者开始工作:
> XREADGROUP GROUP orderGroup consumer-1 COUNT 10 BLOCK 5000 STREAMS orders >注意这里的>是一个特殊标识,意思是"只给我从未被投递给任何消费者的新消息";如果换成具体的 ID,比如0,则是从 PEL 里读取这个消费者自己还没确认的消息,用于故障恢复。这俩的区别很重要,面试时也经常被问到:>走的是正常投递路径,具体 ID 走的是 pending 恢复路径。批量消费、多消费者部署时,>和具体 ID 要分清,别在恢复逻辑里误用了>。
4.3 消费确认与故障恢复:XACK、XPENDING、XCLAIM 组合拳
消费者处理完一条消息,必须显式确认:
> XACK orders orderGroup 1689233456789-0不 ACK 会怎样?看下 PEL 就知道了:
> XPENDING orders orderGroup 1) (integer) 1 2) "1689233456789-0" 3) "1689233456789-0" 4) 1) 1) "consumer-1" 2) "1"输出的四个字段分别是:pending 消息总量、最早一条 pending 的 ID、最晚一条 pending 的 ID、以及每个消费者各自持有的 pending 数量。如果你发现某个消费者名下的 pending 数量一直增长且不减少,基本可以断定这个消费者处理完消息后忘了 ACK,或者已经卡死了。
消费者卡死后,消息不能永远躺在 PEL 里,需要接管。XCLAIM就是干这个的:
> XCLAIM orders orderGroup consumer-2 60000 1689233456789-0这个命令的意思是:把 ID 为1689233456789-0的消息,从原来的消费者手里转交给consumer-2,前提是这条消息已经 pending 超过 60000 毫秒(1 分钟)。consumer-2拿到消息后重新处理,处理完同样执行XACK。
从 Redis 6.2 开始,我推荐用XAUTOCLAIM替代手动 XCLAIM,它比 XCLAIM 更聪明,会自动扫描符合超时条件的多条 pending 消息并批量转移:
> XAUTOCLAIM orders orderGroup consumer-2 60000 0最后一个参数0表示从第一条 pending 开始扫描。XAUTOCLAIM 的返回结果里会带上一个光标,可以用于分页处理,避免一次处理过多消息阻塞太久。我们线上目前所有故障接管逻辑已经全部切到 XAUTOCLAIM 了,代码明显更简洁。
4.4 一套完整的 Java 示例代码
用代码把上面的命令串起来。下面用 Jedis 做客户端,逻辑清晰,方便你直接照着改写。
import redis.clients.jedis.*; import java.util.*; public class OrderStreamDemo { // 生产:写入订单事件 public static String produce(Jedis jedis, String orderId, String userId, String status) { Map<String, String> message = new HashMap<>(); message.put("orderId", orderId); message.put("userId", userId); message.put("status", status); // MAXLEN 近似保留最近 5000 条 return jedis.xadd("orders", StreamEntryID.NEW_ENTRY, message, 5000L, true); } // 消费:从消费组读,处理完 ACK public static void consume(Jedis jedis, String consumerName) { while (true) { List<Map.Entry<StreamEntryID, Map<String, String>>> entries = jedis.xreadGroup("orderGroup", consumerName, new StreamEntryID(), 100, 5000L, true, "orders"); if (entries == null) { continue; // 阻塞超时返回 null } for (Map.Entry<StreamEntryID, Map<String, String>> entry : entries) { StreamEntryID id = entry.getKey(); Map<String, String> msg = entry.getValue(); try { // 业务处理:下单、通知、同步…… handleOrder(msg); // 处理成功,确认消息 jedis.xack("orders", "orderGroup", id); } catch (Exception e) { // 处理失败:先不 ACK,让消息留在 PEL,稍后由其他消费者接管 log.error("handle order failed, id=" + id, e); } } } } private static void handleOrder(Map<String, String> msg) { // 幂等处理 String orderId = msg.get("orderId"); // 查数据库:orderId 是否已处理过,处理过则直接返回 } public static void main(String[] args) { try (Jedis jedis = new Jedis("localhost", 6379)) { // 初始化消费组;如果已存在会报错,捕获忽略即可 try { jedis.xgroupCreate("orders", "orderGroup", new StreamEntryID(), true); } catch (JedisDataException ignored) { } consume(jedis, "consumer-" + System.currentTimeMillis()); } } }这段代码有几个关键点。xreadGroup的最后一个布尔参数true对应命令里的>,表示只读新消息。消费失败时不执行xack,消息留在 PEL 中,之后由超时接管机制重新投递给其他消费者。handleOrder里必须做幂等检查,这是 at-least-once 语义下的硬性要求。
5. 常见问题与避坑指南
5.1 PEL 无限增长,内存被吃掉怎么办
PEL 挂在消费组上,每条 pending 消息都会占用内存。如果消费者处理完消息不 ACK,或者消费者长期宕机,PEL 会越积越大。我见过有团队把 Stream 当"拿到即处理、从不确认"来用,结果 PEL 里堆了几百万条 pending,内存暴涨,还找不到原因。
应对方案分两层。第一层是预防:消费代码里一定把XACK放在业务成功之后,用 try-catch 保证异常路径下不 ACK 而留给接管机制。第二层是治理:写一个定时任务,周期性扫描各组的XPENDING指标,总量超过阈值时告警;对超时的 pending 消息执行XAUTOCLAIM转移,并设置"重试次数上限"——同一条消息接管并处理了 N 次仍然失败,直接XACK掉并记录到日志/死信表里,人工介入。没有死信处理的 Stream 消费是不完整的。
5.2 消息 ID 依赖服务器时间,时钟会有坑吗
ID 自动生成依赖 Redis 服务器的毫秒时间戳。如果多台 Redis 实例时钟有偏差,或者同一台机器时钟回拨,生成的消息 ID 可能打乱。Redis 内部有一个处理逻辑:如果计算出的新 ID 小于等于当前记录的最大 ID,它会把新 ID 强制设为"最大 ID + 1",保证同一个 Stream 里 ID 仍然单调递增。所以时钟回拨确实不会导致重复 ID,但会让 ID 的时间部分失真——消息 ID 里的时间戳就不再准确反映真实写入时间,MINID裁剪的语义也会受影响。
实操层面,我建议:不要让多台物理机的高可用 Redis 同时作为同一个 Stream 的写入端;优先使用单实例 Redis(加好 AOF 持久化)承载 Stream,跨实例的写并发交给业务层路由到不同 key。对绝大多数 Stream 业务场景,单实例完全够用,没必要为了它去搭一个多写集群。
5.3 消费者重复收到消息怎么办
at-least-once 语义决定了重复消息必然存在。消费者处理完一条消息,刚准备 ACK 就宕机了,这条消息在超时后被其他消费者接管,于是又被处理一遍。这种重复无法从机制上消除,只能靠消费端幂等兜底。我的经验是,消息内容里一定要带业务幂等键,比如orderId、taskId,消费逻辑开头先查状态或查去重表,已经处理过就直接 ACK 跳过。数据库层面也建议给幂等键加唯一索引,双保险。
有一个容易被忽略的点:XACK确认的是"这条 I D 被处理完了",不是"这个消息有效"。如果业务上判断消息非法(比如字段缺失),也要显式XACK把它移出 PEL,否则它会永远卡在 pending 队列里。把"处理失败需要重试"和"处理成功但消息无效"区分开,前者不 ACK,后者必须 ACK。
5.4 消息体太大影响性能
Stream 的每个 entry 本质是存在内存里的一个 map,一条消息塞一个几 MB 的 JSON,会直接拖慢XADD和XREADGROUP,还会造成大 key 问题——Redis 在操作大 key 时可能阻塞其他请求。我的建议:Stream 里只放轻量消息体,比如 ID、时间戳、事件类型 + 一个任务 ID;真正的业务负载(如文件路径、完整数据)存在外部存储,消费者根据 ID 自行加载。单条消息控制在 1KB 以内,Redis 的处理性能基本可以拉满。
另外要注意,MAXLEN ~ 1000这种近似裁剪在消息非常大时,裁剪的"块"也可能包含大对象,导致裁剪本身耗时增加。这就是为什么我更推荐在写入时用MAXLEN配合LIMIT,把裁剪分批做掉,而不是攒到几千条再一次性裁。
5.5 消费组和消费者的几个理解误区
误区一:以为消费者必须先创建才能用。前面说过,XREADGROUP里出现的消费者名会自动注册,不需要任何前置操作。但有一点要提醒:一个消费者的 pending 消息只会被该消费者自己读到(通过指定 ID 的方式);如果这个消费者永远不回来了,其他消费者必须靠XCLAIM/XAUTOCLAIM才能接管它的 pending 消息,不能天然地抢过去。
误区二:以为消费者数量越多消费越快。在单 Stream 下,同一组内多个消费者虽然可以并行,但每条消息只会发给一个消费者,所以消费者数量并不是线性提升吞吐。如果积压严重,更合理的做法是给 Stream 拆分成多个分区(比如按 userId 哈希成多个 Stream key),每个分区独立消费。
误区三:XINFO STREAM只能看消息总量——其实它还能看到所有消费组的位点和 pending 概览,一条命令就能判断整个链路是否健康:
> XINFO STREAM orders > XINFO GROUPS orders我很喜欢用这两个命令做日常巡检,比写一堆脚本查状态直观多了。
写在最后:我的实际使用体会
接触 Stream 三年多,我最大的感受是它的可靠性边界比想象中更清晰。它不是万能的,内存容量限制让它撑不起海量消息流,但它把"轻量可靠消息队列"这件事做得足够好:协议简单、消费模型符合直觉、故障恢复路径明确。目前我们线上两个核心业务(订单事件流和导出任务队列)都跑在 Stream 上,日常维护基本只关心两件事——PEL 是否积压、Stream 是否超过裁剪阈值。如果你正卡在"List 丢消息、Kafka 太笨重"的中间地带,Stream 大概率值得你花一个下午把它彻底搞明白。动手搭一个测试 Stream,跑一遍 XADD、XREADGROUP、XACK、XAUTOCLAIM 的完整流程,很快就能体会到它和 List 之间的本质差异在哪里。