先讲一个我实际经历过的场景。有段时间帮朋友排查线上问题用户下单支付后通知服务偶尔会漏发短信后台日志里看不到任何异常但就是有一部分订单没有任何后续动作。他们的实现很朴素——用 Redis List 当消息队列LPUSH 塞消息BRPOP 消费消费完这条消息就没了。问题就出在这消费端进程一旦在 BRPOP 返回之后、业务处理完成之前崩溃这条消息就永久性丢失支付完成了但短信、积分、发票都没动静。这类问题在 Redis 5.0 之前几乎是无解的因为 List 根本没有确认消费的概念。而 Redis 5.0 引入的 Stream 数据类型正是冲着这个长期空白来的。Redis Stream 是一套真正意义上的、内置在 Redis 里的消息队列模型。它支持消息持久化、消费者分组、消费确认ACK、断点续读和按时间回溯把之前用 List、Pub/Sub 拼凑伪队列时踩过的坑一次性填平了。这篇文章我会从设计动机讲起把 Stream 的核心概念、关键命令、实战代码和踩坑记录完整过一遍适合已经在用 Redis、但还没系统接触过 Stream 的开发者也适合准备面试时需要把消息队列这块讲清楚的场景。1. 为什么需要 StreamRedis 以前的伪队列差在哪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 与 ACK2.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 ListPEL。当一条消息被XREADGROUP投递给某个消费者后它的 ID 会进入该组的 PEL。消费者处理完这条消息需要调用XACK把这条消息从 PEL 里移除表示我处理完了。如果迟迟不 ACK消息就一直挂在 PEL 里Redis 就知道这个消费者可能出问题了。PEL 的核心价值在于可追溯、可恢复。它记录了每条 pending 消息的消费者、空闲时间、投递次数。当某个消费者崩溃它的 pending 消息会一直留在 PEL 里其他消费者可以通过XCLAIMRedis 6.2 之后还能用XAUTOCLAIM把这些超过一定空闲时间的消息接管过来继续处理。这样消息从投递到确认全程都有据可查不会因为一个节点挂掉就丢消息。这套机制实现的语义是 at-least-once每条消息至少被处理一次但在极端情况下可能被处理多次——消费者处理完之后还没来得及 ACK 就宕机了消息被其他消费者接管后又处理了一遍。所以你在设计消费逻辑时一定要保证幂等性这个我会在避坑章节专门说。2.4 内存控制MAXLEN 和 MINIDStream 既然持久化存储在 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-0orders是 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) { MapString, 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) { ListMap.EntryStreamEntryID, MapString, String entries jedis.xreadGroup(orderGroup, consumerName, new StreamEntryID(), 100, 5000L, true, orders); if (entries null) { continue; // 阻塞超时返回 null } for (Map.EntryStreamEntryID, MapString, String entry : entries) { StreamEntryID id entry.getKey(); MapString, 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(MapString, 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 之间的本质差异在哪里。