Redis 5.0 引入的 Stream 数据类型为 Redis 带来了原生的消息队列能力,支持消息持久化与主从复制,有效解决了服务宕机或网络中断导致的数据丢失问题。本文将拆解 Stream 的核心架构、消费组竞争机制与 ACK 确认流程,并通过完整的命令示例演示消息的发布、读取与状态追踪,帮助开发者快速评估并落地 Redis Stream 方案。
Stream 核心架构与运行机制
Redis Stream 是 Redis 5.0 版本引入的一种新数据类型,也是 Redis 中最为复杂的数据结构之一。它本质上是一个具备消息发布/订阅功能的组件,即常说的消息队列。其设计借鉴了 Kafka 等分布式消息系统,采用生产者/消费者(broker/consumer)模型,但凭借 Redis 自身的特性,在性能和内存利用率上表现优异。
Stream 最大的优势在于提供了消息的持久化存储与主从复制功能。这意味着即使遭遇网络断开或 Redis 宕机重启,存储的消息内容依然不会丢失,保障了数据的高可用性。
消息队列的四大组成部分
Stream 消息队列主要由四部分构成:消息本身、生产者、消费者和消费组。其中,消费组(Consumer Group)是 Stream 实现并发消费与负载均衡的核心机制。
一个 Stream 队列可以拥有多个消费组,每个消费组内又可包含多个消费者。组内消费者之间存在竞争关系:当某个消费者成功消费一条消息后,同组的其他消费者将不会再次消费该消息。被消费的消息 ID 会被记录在等待处理的 pending_ids(官方称为 PEL,Pending Entries List)中。每完成一条消息的消费与确认,消费组的游标 last_delivered_id 就会向前移动,组内消费者继续争抢下一条消息。

图1:Redis Stream 流程处理图
核心概念解析
- Stream direction(数据流):表示消息链,将所有消息按顺序串联。每条消息拥有唯一的标识 ID 和对应的消息内容(Message content)。
- Consumer Group(消费组):拥有唯一的组名,通过
XGROUP CREATE命令创建。一个 Stream 可挂载多个消费组,组内每个消费者也有唯一 ID。 - last_delivered_id(消费组游标):记录消费组当前已递送的消息位置。任意消费者读取消息都会推动该游标向前移动。
- pending_ids(PEL):记录已被客户端读取但尚未 ACK(确认)的消息 ID。若客户端未发送 ACK,该列表会持续增长;一旦消息被 ACK,对应 ID 即从列表中移除。
ACK 确认机制
ACK(Acknowledge character)即确认字符,是数据通信中接收方反馈给发送方的控制信号,表示数据已无误接收。在 Redis Stream 中,消费者处理完消息后必须通过 XACK 命令向服务器发送确认,服务器才会将该消息从 PEL 中清除。这一机制确保了消息处理的可靠性,防止因消费者崩溃导致消息丢失。
常用命令速查表
掌握 Stream 的核心命令是高效使用消息队列的前提。下表汇总了日常开发与运维中最常用的 Stream 指令及其作用:
| 命令 | 说明 |
|---|---|
XADD | 添加消息到 Stream 末尾。 |
XTRIM | 对 Stream 流进行修剪,限制长度或条目数。 |
XDEL | 删除 Stream 中指定的消息。 |
XLEN | 获取 Stream 包含的消息数量。 |
XRANGE | 按 ID 范围正向获取消息列表,自动过滤已删除消息。 |
XREVRANGE | 按 ID 范围反向获取消息列表(ID 从大到小)。 |
XREAD | 以阻塞或非阻塞方式读取消息列表。 |
XGROUP CREATE | 创建消费者组。 |
XREADGROUP GROUP | 读取消费者组中的消息。 |
XACK | 将消息标记为“已处理”,从 PEL 中移除。 |
XGROUP SETID | 为消费者组设置新的最后递送消息 ID。 |
XGROUP DELCONSUMER | 删除指定消费者。 |
XGROUP DESTROY | 删除整个消费者组。 |
XPENDING | 显示待处理(未 ACK)消息的相关信息。 |
XCLAIM | 转移消息的归属权,用于处理超时未确认的消息。 |
XINFO | 查看 Stream、消费者组或消费者的详细信息。 |
XINFO GROUPS | 查看指定 Stream 下所有消费组的信息。 |
XINFO STREAM | 查看 Stream 流的详细元数据。 |
XINFO CONSUMERS key group | 查看指定消费组内所有消费者的状态信息。 |
消息 ID 生成与队列管理
在 Stream 中,每条消息都必须具备唯一且单调递增的 ID。ID 的生成支持系统自动分配与手动指定两种方式。
系统自动生成 ID
使用 * 作为 ID 参数时,Redis 会自动生成基于毫秒时间戳的 ID。格式为 毫秒时间戳-序列号,例如 1610619132674-1 表示在该毫秒内产生的第 1 条消息。
# 添加消息,* 表示由 Redis 自动生成递增 ID 127.0.0.1:6379> XADD mystream * username www.biancheng.net age 10 "1610619132674-1"
自定义 ID
开发者也可手动指定 ID,但必须严格遵守“只增不减”的规则,即新插入的 ID 必须大于当前 Stream 中最大的 ID,否则会报错。
# 自定义 ID 必须为整数格式,且严格递增 127.0.0.1:6379> XADD mystream1 001 name zhangsan addr hebei "1-0" 127.0.0.1:6379> XADD mystream1 002 name lisi addr hunan "2-0" # 插入重复或更小的 ID 将触发错误 127.0.0.1:6379> XADD mystream1 001 name wangwu addr fujian (error) ERR The ID specified in XADD is equal or smaller than the target stream top item 127.0.0.1:6379> XADD mystream1 003 name wangwu addr fujian "3-0"
消息查询与删除
通过 XDEL 可删除指定 ID 的消息,XLEN 用于统计剩余消息数量,XRANGE 支持按范围查询并自动跳过已删除条目。
# 删除指定 ID 的消息
127.0.0.1:6379> XDEL mystream1 001
(integer) 1
# 查看队列长度
127.0.0.1:6379> XLEN mystream1
(integer) 2
# 获取指定范围的消息(- 表示最小,+ 表示最大)
127.0.0.1:6379> XRANGE mystream - +
1) 1) "1610619132674-0"
2) 1) "username"
2) "www.biancheng.net"
3) "age"
4) "10"
2) 1) "1610619178028-0"
2) 1) "username"
2) "c.biancheng.net"
3) "age"
4) "9"
# 使用 count 限制返回数量
127.0.0.1:6379> XRANGE mystream1 - 003 count 1
1) 1) "2-0"
2) 1) "name"
2) "lisi"
3) "addr"
4) "hunan"消费组创建与消息消费
消费组是 Stream 实现多消费者协同工作的核心。通过 XGROUP CREATE 命令可创建消费组,并需指定起始消息 ID 以初始化 last_delivered_id 游标。
创建消费组
创建时可传入具体的起始 ID,或使用 $ 表示从当前 Stream 尾部开始消费(仅接收新消息,忽略历史消息)。
# 创建消费组 ms1,从 ID 0-0 开始消费 127.0.0.1:6379> XGROUP CREATE mystream1 ms1 0-0 OK # 创建消费组 ms3,仅消费新消息(忽略已有数据) 127.0.0.1:6379> XGROUP CREATE mystream1 ms3 $ OK
创建后可通过 XINFO 系列命令查看 Stream 与消费组的详细状态:
# 查看 Stream 元数据
127.0.0.1:6379> XINFO stream mystream1
1) "length"
2) (integer) 2
3) "radix-tree-keys"
4) (integer) 1
5) "radix-tree-nodes"
6) (integer) 2
7) "groups"
8) (integer) 2
9) "last-generated-id"
10) "3-0"
11) "first-entry"
12) 1) "2-0"
2) 1) "name"
2) "lisi"
3) "addr"
4) "hunan"
13) "last-entry"
14) 1) "3-0"
2) 1) "name"
2) "wangwu"
3) "addr"
4) "fujian"
# 查看消费组信息
127.0.0.1:6379> XINFO GROUPS mystream1
1) 1) "name"
2) "ms1"
3) "consumers"
4) (integer) 0
5) "pending"
6) (integer) 0
7) "last-delivered-id"
8) "0-0"
2) 1) "name"
2) "ms3"
3) "consumers"
4) (integer) 0
5) "pending"
6) (integer) 0
7) "last-delivered-id"
8) "3-0"消费组读取与 ACK 确认
消费者通过 XREADGROUP 命令从消费组中拉取消息。该命令支持阻塞等待(BLOCK),读取到的消息会进入消费者的 PEL 列表。业务处理完成后,必须调用 XACK 通知 Redis 移除该消息的待处理状态。

图2:Redis Stream 消费流程示意图
# 消费组 ms1 中的消费者 c1 读取 1 条消息
# > 表示从 last_delivered_id 之后开始读取,读取后游标自动前移
127.0.0.1:6379> XREADGROUP GROUP ms1 c1 COUNT 1 STREAMS mystream1 >
1) 1) "mystream1"
2) 1) 1) "2-0"
2) 1) "name"
2) "lisi"
3) "addr"
4) "hunan"
# 继续读取下一条
127.0.0.1:6379> XREADGROUP GROUP ms1 c1 COUNT 1 STREAMS mystream1 >
1) 1) "mystream1"
2) 1) 1) "3-0"
2) 1) "name"
2) "wangwu"
3) "addr"
4) "fujian"
# BLOCK 1000 表示阻塞等待 1 秒,无新消息则返回 nil
127.0.0.1:6379> XREADGROUP GROUP ms1 c1 COUNT 1 BLOCK 1000 STREAMS mystream1 >
(nil)
# 指定具体 ID 读取(非游标模式)
127.0.0.1:6379> XREADGROUP GROUP ms1 c1 COUNT 1 STREAMS mystream1 1
1) 1) "mystream1"
2) 1) 1) "2-0"
2) 1) "name"
2) "lisi"
3) "addr"
4) "hunan"
# 超出范围返回空列表
127.0.0.1:6379> XREADGROUP GROUP ms1 c1 COUNT 1 STREAMS mystream1 3
1) 1) "mystream1"
2) (empty list or set)
# 添加新消息
127.0.0.1:6379> XADD mystream1 004 name zhangwu age 24
# 消费组 ms3 读取新消息(因创建时使用 $,仅接收新消息)
127.0.0.1:6379> XREADGROUP GROUP ms3 c2 COUNT 2 STREAMS mystream1 >
1) 1) "mystream1"
2) 1) 1) "4-0"
2) 1) "name"
2) "zhangwu"
3) "age"
4) "21"
# 确认消息已处理,从 PEL 中移除
127.0.0.1:6379> XACK mystream1 ms1 002
(integer) 1在实际生产环境中,建议结合 XPENDING 监控未确认消息,并使用 XCLAIM 处理因消费者宕机而长期滞留 PEL 的消息,从而构建高可靠的消息处理闭环。

