适用版本:Redis 7.2.x
一、什么是 Stream
Stream 是 Redis 5.0 引入的一种高级数据结构,用来处理 消息流、事件流、日志流 这类不断追加的数据。
如果只看表面,可以把它理解成:
- 一条可以持续追加消息的日志
- 一个支持消费组的轻量消息队列
- 一个兼顾持久化和多消费者协作的数据流结构
在很多业务里,Stream 经常被用于:
- 异步任务分发
- 订单事件流转
- 用户行为日志收集
- 消费组模式下的消息处理
和传统 List 相比,Stream 最大的特点是:
- 每条消息都有唯一 ID
- 支持按 ID 读取历史消息
- 支持消费组(Consumer Group)
- 支持待确认消息管理
因此它更接近“消息流系统”,而不只是简单队列。
二、Stream 的核心原理
1. 追加式日志模型
Stream 的数据写入方式非常像日志:
- 新消息不断追加到尾部
- 每条消息都有一个递增 ID
- 读取方可以从某个位置继续消费
一条 Stream 消息通常包含:
- 消息 ID
- 多个字段和值
例如:
1717920000000-0 => {
order_id: 10001,
user_id: 2001,
status: created
}
2. 消息 ID 是有序的
Redis Stream 中每条消息都有 ID,通常格式类似:
毫秒时间戳-序号
例如:
1717920000000-0
1717920000000-1
1717920000123-0
这个 ID 决定了消息顺序,也让消费者可以明确地表达:
- 从哪里开始读
- 已经读到哪里
- 哪些消息还没确认
3. 支持消费组协作消费
消费组是 Stream 最重要的特性之一。
在消费组模式下:
- 多个消费者共同消费同一条消息流
- 同一条消息只会分发给组内某一个消费者处理
- 未确认消息会进入待处理列表(PEL)
- 支持重试、转移和确认
这使得 Stream 比 List 更适合做可靠消息消费。
4. 兼顾“可回放”和“可协作消费”
许多消息队列只强调“往前消费”,而 Stream 天然带有日志流的特征,支持:
- 查看历史消息
- 指定位置重放
- 按消费组维护消费进度
这让它在一些轻量事件流场景里很有优势。
三、Stream 适用场景
1. 异步消息队列
例如下单后,不直接同步调用所有后续流程,而是写入 Stream:
- 库存服务消费
- 积分服务消费
- 通知服务消费
这样可以降低耦合,提高系统弹性。
2. 订单或业务事件流
例如:
- 订单创建
- 订单支付
- 订单取消
- 订单退款
这些事件天然具有时间顺序,很适合以 Stream 形式追加记录。
3. 日志与埋点收集
如果需要把应用事件快速写入 Redis,再由多个消费者异步处理,也可以使用 Stream。
4. 多消费者协作处理任务
例如一个任务流需要多个 worker 并发消费,Stream 的消费组可以让多个消费者分摊负载。
四、常用命令与示例
1. XADD:向 Stream 追加消息
XADD order:stream * order_id 10001 user_id 2001 status created
XADD order:stream * order_id 10002 user_id 2002 status created
这里的 * 表示由 Redis 自动生成消息 ID。
执行后会返回类似:
1717920000000-0
也就是说,消息写入成功后,你就拿到了它的唯一标识。
2. XRANGE:按区间读取消息
XRANGE order:stream - +
表示读取整个 Stream 中的所有消息。
如果只读取一段范围:
XRANGE order:stream 1717920000000-0 1717929999999-0 COUNT 10
适合调试、补偿、回放部分历史消息。
3. XREVRANGE:倒序读取
XREVRANGE order:stream + - COUNT 5
适合快速查看最近几条消息。
4. XREAD:非消费组模式读取
XREAD COUNT 2 STREAMS order:stream 0
这个命令表示:
- 从指定 ID 开始读
- 最多读取 2 条消息
如果想阻塞等待新消息:
XREAD BLOCK 5000 STREAMS order:stream $
其中:
$表示从当前末尾开始,只等后续新消息BLOCK 5000表示最多阻塞 5 秒
5. XGROUP CREATE:创建消费组
XGROUP CREATE order:stream order_group 0 MKSTREAM
含义:
- 在
order:stream上创建消费组order_group - 从最早消息开始消费
MKSTREAM表示如果 Stream 不存在则自动创建
6. XREADGROUP:消费组读取消息
XREADGROUP GROUP order_group consumer_1 COUNT 2 STREAMS order:stream >
这里的 > 表示:
- 读取从未投递给本组消费者的新消息
这是消费组模式中最常见的读取方式。
7. XACK:确认消息
消费者处理完消息后,应当显式确认:
XACK order:stream order_group 1717920000000-0
如果不确认,这条消息会保留在待处理列表中。
8. XPENDING:查看待确认消息
XPENDING order:stream order_group
可以查看:
- 当前有多少未确认消息
- 最小和最大待处理 ID
- 哪些消费者持有未确认消息
9. XCLAIM / XAUTOCLAIM:转移长时间未确认消息
如果某个消费者宕机,消息一直不确认,可以让其他消费者接管。
例如:
XAUTOCLAIM order:stream order_group consumer_2 60000 0-0 COUNT 10
表示把空闲超过 60 秒的待处理消息转移给 consumer_2。
这个命令在故障恢复中非常实用。
10. XTRIM:裁剪 Stream 长度
XTRIM order:stream MAXLEN 10000
表示只保留最近 10000 条消息,旧消息自动裁剪。
对于日志流、埋点流这类数据持续增长的场景,长度控制非常重要。
五、一个典型案例:订单创建后的异步处理
假设用户下单后,需要异步完成以下动作:
- 扣减库存
- 发送短信
- 增加积分
1. 生产者写入订单事件
XADD order:event:stream * event order_created order_id 10001 user_id 2001 amount 299
2. 创建消费组
XGROUP CREATE order:event:stream order_worker_group 0 MKSTREAM
3. 消费者读取消息
XREADGROUP GROUP order_worker_group worker_1 COUNT 1 BLOCK 5000 STREAMS order:event:stream >
4. 处理成功后确认
XACK order:event:stream order_worker_group 1717920000000-0
5. 如果消费者异常退出
未确认的消息会保留在待处理列表中,其他消费者可以通过 XPENDING + XAUTOCLAIM 接管处理。
这就是 Stream 相比普通 List 更适合作为消息队列的重要原因。
六、Stream 相比 List 有什么优势
| 维度 | List | Stream |
|---|---|---|
| 消息 ID | 没有天然唯一 ID | 每条消息都有唯一 ID |
| 历史回放 | 不方便 | 支持按 ID 范围读取 |
| 消费组 | 不原生支持 | 原生支持 |
| 待确认机制 | 无 | 有 |
| 消费者协作 | 需要应用层补很多逻辑 | Redis 内置支持 |
| 适合场景 | 简单队列 | 消息流、事件流、可靠消费 |
如果业务只是非常简单的先进先出队列,List 可能够用;但只要涉及:
- 多消费者协作
- 消费确认
- 消息重试
- 历史重放
那么 Stream 更合适。
七、Stream 的优点
1. 同时具备日志流和消息队列特征
它既能追加记录,又能维护消费进度,还能回看历史消息。
2. 原生支持消费组
不用像 List 那样在应用层自行维护复杂分发逻辑。
3. 支持消息确认与重试
有待处理列表,出现消费者故障时可以做恢复和重新分配。
4. 适合构建轻量事件驱动系统
对于很多中小规模业务,Redis Stream 已经足以承担一部分异步事件处理能力。
八、使用 Stream 时的注意事项
1. 消费成功后一定要 ACK
这是最容易遗漏的一点。
如果消费者读取消息后不执行 XACK:
- 消息不会从待处理列表中移除
- 后续会出现积压
- 排障时会发现很多“已经处理但未确认”的消息
所以消费逻辑里要明确区分:
- 读取成功
- 业务处理成功
- 确认完成
2. Stream 不是无限日志仓库
虽然 Stream 可以不断追加,但 Redis 毕竟是内存数据库。若不做裁剪,数据会持续增长。
常见做法:
- 使用
XTRIM控制长度 - 或在
XADD时配合长度限制策略 - 把长期归档交给外部存储系统
3. 要处理消费者宕机后的未确认消息
消费组场景里,真正的难点不只是读消息,而是处理故障后的恢复。
必须考虑:
- 如何查看 pending 消息
- 何时认定某消费者失联
- 如何把消息转给新消费者
- 业务是否支持幂等重试
4. 业务侧要有幂等设计
即使 Stream 提供确认与重分配机制,消息仍然可能因为重试被重复处理。
所以涉及订单、库存、积分等操作时,业务层要具备幂等能力。
5. 不要把 Stream 当成完整 MQ 平台的无差别替代
Redis Stream 很强,但它和专业消息中间件的定位仍有区别。对于:
- 超大规模堆积
- 复杂路由
- 海量消费组
- 强顺序、多副本、高可靠投递
这类复杂场景,仍需结合专门消息系统评估。
6. 关注阻塞读取连接数与消费者模型
如果大量消费者都用阻塞方式读取 Stream,要评估连接资源、线程模型和消费并发策略,避免因为使用方式不当造成服务端压力。
九、实战设计建议
1. 用清晰的命名表达业务语义
例如:
order:event:stream
pay:event:stream
notice:send:stream
消费组也尽量表达角色:
order_worker_group
notice_group
inventory_group
2. 生产者只负责写入,消费者关注处理与确认
职责分离后,系统更容易扩展和排障。
3. 对关键消费流程做幂等保护
例如用业务唯一键、状态机、去重表等方式,保证消息重复处理时不会产生副作用。
4. 监控 pending 与积压长度
在生产环境中,建议重点监控:
- Stream 长度
- 消费组 lag
- pending 数量
- 长时间未 ACK 的消息
这些指标往往能提前暴露消费故障。
十、小结
Stream 是 Redis 高级数据结构中非常重要的一员,它把“追加日志”“消息队列”“消费组协作”几种能力结合到了一起,适合处理实时事件流和轻量消息队列需求。
可以重点记住以下几点:
- Stream 适合消息流、事件流、日志流场景
- 每条消息都有唯一 ID,支持按范围读取历史数据
- 常用命令包括
XADD、XRANGE、XREAD、XGROUP、XREADGROUP、XACK - 消费组是 Stream 的核心特性,支持协作消费与待确认管理
- 使用时一定要关注 ACK、pending、裁剪和幂等
如果你想在 Redis 中实现一个比 List 更可靠、更现代的消息处理方案,那么 Stream 往往是第一选择;但在真正的大规模消息架构里,也要清楚它的边界,结合业务复杂度合理选型。
📝 版权声明:本文为原创技术博客,转载请注明出处。
如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!