Redis 5.0 引入的 Stream 数据结构,是 Redis 作者 antirez 在多年观察各类消息中间件使用模式后,为 Redis 量身打造的原生消息队列实现。与早期依赖 List、Pub/Sub 或第三方组件拼凑出来的”伪队列”不同,Stream 从底层就内置了消费组、消息确认、死信处理、持久化等完整能力,足以在中小流量场景下替代 RabbitMQ 或 Kafka 的一部分职责。本文将系统拆解 Stream 的核心概念、底层结构、命令用法、消费者组语义以及生产环境下的可靠性设计,帮助你在真实业务中正确选型和使用。
一、为什么 Redis 需要 Stream:从 List 队列到原生消息中间件
在 Stream 出现之前,社区用 Redis 实现消息队列主要有三种方式,每种都有明显短板:
- List + BRPOPLPUSH:可以做到阻塞拉取和”处理完再删”的可靠语义,但消息一旦被消费就消失,没有历史可回溯,也没有消费组概念,多消费者只能轮询同一份消息。
- Pub/Sub:纯内存广播,订阅者不在线就丢消息,没有持久化,无法做消费确认,更谈不上积压监控。
- Sorted Set + 时间戳:能实现延迟队列,但消费确认、去重、分组消费都得自己造轮子,运维成本高。
Stream 的设计目标是把这些能力一次补齐:每条消息都有 Redis 自动分配的递增 ID,天然有序且可回溯;消费组(Consumer Group)让多个消费者分摊同一份流;PEL(Pending Entry List)记录未确认消息,支持 ACK 与重投;消息持久化随 RDB/AOF 一起落地,宕机不丢。这种”内建消息中间件语义”的定位,是 Stream 区别于其他 Redis 数据结构的关键。
二、Stream 的底层结构与消息 ID
Stream 在 Redis 内部并非复用 listpack 或 skiplist,而是有一套独立的实现:底层由一个名为
1 | stream |
的 C 结构体维护,核心是基数树(rax tree)加 listpack 节点的组合。每条消息以 ID 为键、字段-值对为内容存储,多个消息被打包进 listpack 节点,节点之间用基数树按 ID 有序串联。这种设计既支持 O(log N) 的 ID 查找,又能批量压缩存储,在百万级消息规模下内存占用依然可控。
消息 ID 是 Stream 的灵魂,它由两部分组成:
1
2 <毫秒时间戳>-<同一毫秒内序号>
例如:1718937600000-0、1718937600000-1
当客户端使用
1 | * |
让 Redis 自动生成 ID 时,Redis 会以服务器当前时间为基准分配,保证严格递增。也可以由客户端显式指定 ID,但必须大于流中已有的最大 ID,否则会报错。这种”时间有序”的 ID 让 Stream 天然支持按时间范围查询,是它区别于 Kafka 偏移量语义的重要特征。
三、基础命令实战:XADD、XLEN、XRANGE、XREAD
下面通过一组实战命令演示 Stream 的基本用法。先创建一条订单事件流:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29 # 向订单流追加一条消息,字段为 order_id 和 amount
127.0.0.1:6379> XADD orders * order_id 1001 amount 199.00
"1718937600000-0"
# 再追加一条
127.0.0.1:6379> XADD orders * order_id 1002 amount 88.50
"1718937600001-0"
# 查看流中消息总数
127.0.0.1:6379> XLEN orders
(integer) 2
# 按范围读取全部消息
127.0.0.1:6379> XRANGE orders - +
1) 1) "1718937600000-0"
2) 1) "order_id"
2) "1001"
3) "amount"
4) "199.00"
2) 1) "1718937600001-0"
2) 1) "order_id"
2) "1002"
3) "amount"
4) "88.50"
# 阻塞式读取新消息,最多等 5 秒
127.0.0.1:6379> XREAD COUNT 10 BLOCK 5000 STREAMS orders $
(nil)
(5.00s)
几个关键点值得注意:
-
1XADD ... *
中的
1*表示让 Redis 自动分配 ID,这是最常用的方式。
-
1XRANGE - +
中的
1-和
1+分别代表最小和最大 ID,用于全量扫描。
-
1XREAD
的
1$表示只读取当前最大 ID 之后的新消息,配合
1BLOCK实现长轮询;如果传具体 ID 则从该 ID 之后开始读。
- 所有读命令都不会删除消息,消息一直在流中保留,直到被
1XTRIM
或
1XDEL主动清理。
四、消费者组:多消费者分摊与 ACK 机制
消费者组是 Stream 最核心的能力。它把一条流逻辑上划分给若干个组,每个组内的消费者分摊消息,互不重复;不同组之间相互独立,各自维护自己的消费进度。这是 Kafka 消费组概念在 Redis 中的精简实现。
4.1 创建消费者组
1
2
3
4
5
6
7 # 创建消费组 order-processors,从流的起始位置开始消费
127.0.0.1:6379> XGROUP CREATE orders order-processors 0
OK
# 也可以从最新位置开始,只消费创建组之后的新消息
127.0.0.1:6379> XGROUP CREATE orders order-processors $ MKSTREAM
OK
第三个参数是消费组起始 ID:
1 | 0 |
表示从头消费历史消息,
1 | $ |
表示只消费组创建后的新消息。
1 | MKSTREAM |
选项允许在流不存在时自动创建空流,避免报错。
4.2 消费消息与 ACK
1
2
3
4
5
6
7
8
9
10
11
12 # 消费者 consumer-1 从 order-processors 组读取消息
127.0.0.1:6379> XREADGROUP GROUP order-processors consumer-1 COUNT 1 STREAMS orders >
1) 1) "orders"
2) 1) 1) "1718937600000-0"
2) 1) "order_id"
2) "1001"
3) "amount"
4) "199.00"
# 处理完成后确认(ACK)
127.0.0.1:6379> XACK orders order-processors 1718937600000-0
(integer) 1
这里的
1 | > |
是一个特殊 ID,含义是”只取该消费者尚未交付过的消息”。一旦消息被交付,它就进入 PEL(Pending Entry List),必须显式
1 | XACK |
才会从 PEL 移除。如果消费者处理过程中崩溃,消息会一直停留在 PEL 中,这正是可靠投递的基础。
4.3 转移未确认消息:XPENDING 与 XCLAIM
当某个消费者宕机后,它未 ACK 的消息会一直留在 PEL 中。可以通过
1 | XPENDING |
查看,再用
1 | XCLAIM |
把消息转移到其他消费者重投:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15 # 查看消费组中所有未确认消息
127.0.0.1:6379> XPENDING orders order-processors
1) (integer) 1 # PEL 总数
2) "1718937600001-0" # 最小 ID
3) "1718937600001-0" # 最大 ID
4) 1) 1) "consumer-1"
2) "1"
# 把 idle 超过 60 秒的消息转移给 consumer-2
127.0.0.1:6379> XCLAIM orders order-processors consumer-2 60000 1718937600001-0
1) 1) "1718937600001-0"
2) 1) "order_id"
2) "1002"
3) "amount"
4) "88.50"
生产环境通常会写一个定时任务,周期性地扫描各消费者 PEL,把 idle 时间超过阈值的消息
1 | XCLAIM |
给健康消费者,这就是死信重投的雏形。Redis 6.2 引入的
1 | XAUTOCLAIM |
命令进一步简化了这个流程:
1
2 # 自动领取 idle 超过 60 秒的消息,每次最多 10 条
127.0.0.1:6379> XAUTOCLAIM orders order-processors consumer-2 60000 0 COUNT 10
五、消息积压治理:XTRIM 与 XDEL
Stream 是持久化结构,消息不会因为被消费而消失,长此以往内存会膨胀。Redis 提供两种清理手段:
| 命令 | 作用 | 是否删除 PEL 中的消息 | 典型场景 | ||
|---|---|---|---|---|---|
|
保留最近 10000 条,其余丢弃 | 是(PEL 中的也可能被删) | 按数量控制内存 | ||
|
删除 ID 小于指定值的所有消息 | 是 | 按时间保留窗口 | ||
|
删除指定 ID 的单条消息 | 否(仅从流中移除) | 精确清理脏数据 |
需要注意:
1 | XTRIM MAXLEN |
在 6.0 之后支持近似裁剪(加
1 | ~ |
符号),即
1 | XTRIM orders MAXLEN ~ 10000 |
,Redis 会以 listpack 节点为粒度批量删除,性能远高于精确裁剪。生产环境几乎都应该用近似模式,避免每次追加消息都触发精确裁剪。
一个常见误区是用
1 | XDEL |
来”确认消费”。这是错的:
1 | XDEL |
只是从流中移除消息本体,并不会从 PEL 中清除,正确做法始终是
1 | XACK |
。两者职责完全不同。
六、生产级可靠投递:客户端实现要点
真正在生产环境跑起来,光会命令还不够。下面给出一段基于 Python redis-py 的消费者骨架,演示幂等消费、心跳上报和死信重投的完整模式。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64 import redis
import time
import json
import logging
logging.basicConfig(level=logging.INFO)
log = logging.getLogger("order-consumer")
r = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
STREAM = "orders"
GROUP = "order-processors"
CONSUMER = "worker-1"
IDLE_THRESHOLD_MS = 60_000 # 60 秒未确认视为挂了
CLAIM_BATCH = 100
def ensure_group():
try:
r.xgroup_create(STREAM, GROUP, id="0", mkstream=True)
except redis.exceptions.ResponseError as e:
if "BUSYGROUP" not in str(e):
raise
def process(msg_id, fields):
# 业务侧必须自己做幂等:用唯一键(如 order_id)查重
order_id = fields.get("order_id")
log.info("processing %s order=%s", msg_id, order_id)
# ... 真实业务逻辑 ...
# 处理成功后 ACK
r.xack(STREAM, GROUP, msg_id)
def consume_loop():
while True:
# 1. 先尝试领取空闲消费者遗留的消息
claimed = r.xautoclaim(
STREAM, GROUP, CONSUMER,
min_idle_time=IDLE_THRESHOLD_MS,
start_id="0", count=CLAIM_BATCH,
)
# redis-py 4.x 返回 (next_id, [(id, fields), ...], deleted_ids)
for msg_id, fields in claimed[1]:
try:
process(msg_id, fields)
except Exception as e:
log.error("claim process failed %s: %s", msg_id, e)
# 2. 再正常拉取新消息
resp = r.xreadgroup(
GROUP, CONSUMER,
streams={STREAM: ">"},
count=10, block=2000,
)
if not resp:
continue
for _, messages in resp:
for msg_id, fields in messages:
try:
process(msg_id, fields)
except Exception as e:
log.error("process failed %s: %s", msg_id, e)
# 不 ACK,留在 PEL 等下一轮 XAUTOCLAIM 重投
if __name__ == "__main__":
ensure_group()
consume_loop()
这段代码体现了几个生产要点:
- 幂等消费:消息可能被重投,业务侧必须以
1order_id
等业务唯一键查重,避免重复发货、重复扣款。
- XAUTOCLAIM 优先:每轮循环先扫一遍 PEL,把宕机消费者遗留的消息接管过来,再拉取新消息,保证不会”新消息源源不断,老消息无人处理”。
- 异常不 ACK:处理失败时不要 ACK,让消息留在 PEL 中,下一轮自然进入重投流程。配合业务侧的重试次数计数可以实现”超过 N 次进死信队列”。
- block 而非忙轮询:
1XREADGROUP
一定要带
1BLOCK,否则空流时会打满 CPU。
七、Stream 与 Kafka、RabbitMQ 的边界
Stream 很强,但不是万能。选型时需要清楚它的能力边界:
| 维度 | Redis Stream | Kafka | RabbitMQ |
|---|---|---|---|
| 单机吞吐 | 10 万级 msg/s | 百万级 msg/s | 万级~十万级 |
| 持久化 | RDB/AOF,宕机有窗口 | 副本 + 日志,强持久 | 队列镜像/仲裁队列 |
| 消费组 | 支持,但单分区顺序 | 原生,分区有序 | 通过 prefetch + ACK |
| 消息回溯 | 按 ID 范围,可回放 | 按 offset,强回放 | 弱,依赖 dead-letter |
| 延迟队列 | 不原生,需 Sorted Set 配合 | 不原生 | 插件支持 |
| 运维成本 | 极低,复用 Redis | 高,需 ZK/KRaft | 中等 |
经验法则:单分区 QPS 在 10 万以内、消息需要持久化但可接受毫秒级宕机窗口、不想引入额外中间件时,Stream 是最佳选择。一旦涉及到跨节点分区顺序、千万级 TPS、严格不丢消息(金融级),就应该上 Kafka 或 RabbitMQ,Redis 作为旁路缓存或轻量队列使用。
八、监控指标与容量规划
把 Stream 上生产环境,需要持续关注三个指标:
- PEL 长度:
1XPENDING orders order-processors
返回的第一项就是 PEL 总数。持续增长说明消费者跟不上或反复失败,是积压的最早信号。
- 流长度:
1XLEN orders
。如果配了
1XTRIM MAXLEN ~但长度仍持续增长,说明裁剪跟不上写入速度,需要调小保留量或增加消费者。
- 消费者 idle 时间:
1XPENDING ... IDLE
子命令可以列出每条 PEL 消息的 idle 时长,用于判断是哪个消费者卡住了。
容量规划上,单条消息假设 200 字节,100 万条消息约 200MB。但 Stream 的 listpack 节点有压缩,实际内存占用会更低。建议把 Stream 的
1 | MAXLEN |
控制在业务可回溯窗口的 2 倍以内,既保证回放能力,又控制内存上限。配合 Redis 的
1 | maxmemory-policy |
设为
1 | noeviction |
,避免 Stream 被动淘汰导致消息丢失。
最后提醒:Stream 的消息确认是”至少一次”(at-least-once)语义,不是”恰好一次”(exactly-once)。业务侧必须接受消息可能重复,并通过幂等设计来兜底。这是所有基于 ACK 机制的消息中间件的共同特征,并非 Redis 独有。理解这一点,才能在设计消费者时把幂等放在第一位,而不是把”不丢不重”寄托在队列本身。
汤不热吧