Redis 流 Stream
Stream 是 Redis 5.0 引入的消息队列结构:生产者追加消息、消费者按 ID 顺序读取,消息可持久化(随 RDB/AOF 落盘),还支持消费组与消息确认。相比"发了就忘"的发布订阅,Stream 能可靠地做任务队列、事件流水、日志收集,是在 Redis 上搭消息系统的首选。
追加消息:XADD
XADD 向流追加一条"字段-值"记录并返回消息 ID(毫秒时间戳-自增序号)。ID 传 * 由服务器生成,也可自定更大 ID:
XADD msgq * 发送方 user1 内容 hello
# 输出:"1730000000000-0" # ID 由服务器生成,随运行时刻变化
XADD msgq * 发送方 user2 内容 hi
# 输出:"1730000000001-0"
XLEN msgq
# 输出:(integer) 2
读取消息:XRANGE / XREAD
XRANGE 按 ID 区间拉取(- 最小、+ 最大);XREAD 从指定 ID 之后读取,加 BLOCK 可阻塞等待新消息:
XRANGE msgq - +
# 输出:1) 1) "1730000000000-0"
# 2) 1) "发送方" 2) "user1" 3) "内容" 4) "hello"
# 2) 1) "1730000000001-0"
# 2) 1) "发送方" 2) "user2" 3) "内容" 4) "hi"
XREAD COUNT 1 STREAMS msgq 1730000000000-0
# 输出:1) 1) "msgq"
# 2) 1) 1) "1730000000001-0"
# 2) 1) "发送方" 2) "user2" 3) "内容" 4) "hi"
XDEL 删除单条消息(ID 不复用);XTRIM、XADD 的 MAXLEN 选项把流裁剪到只留最近 N 条,防止无限膨胀。
消费组:多消费者分工
消费组让一条消息只被组内一个消费者处理,实现任务分发。XGROUP CREATE 建组(ID 0 表示从最早消息开始),XREADGROUP 读取:
XGROUP CREATE msgq g1 0
# 输出:OK
XREADGROUP GROUP g1 worker1 COUNT 10 STREAMS msgq >
# 输出:worker1 拿到组内尚未分配给任何人的消息(> 表示只读新消息)
组里再加消费者即可水平扩展;每条消息只会投递给组内的一个消费者。
消息确认:XACK 与 PEL
消费者读到的消息进入其 PEL(待确认列表)。处理成功后要 XACK 确认,否则消息一直记在 PEL,便于异常恢复后重新投递:
XACK msgq g1 1730000000000-0
# 输出:(integer) 1 # 确认成功
XPENDING msgq g1
# 输出:1) (integer) 1 # 还有 1 条未确认
# 2) "1730000000001-0"
# 3) "1730000000001-0"
# 4) 1) 1) "worker1" 2) (integer) 1
消费者崩溃导致消息长期不确认时,可用 XAUTOCLAIM 把超时的 PEL 消息转移给其他消费者重处理。
与 Pub/Sub、Kafka 的对比
| 特性 | Pub/Sub | Stream | Kafka |
|---|---|---|---|
| 消息持久化 | 无,发完即走 | 可(随 RDB/AOF 落盘) | 是,磁盘日志 |
| 消费组与确认 | 无 | 有(XACK/PEL) | 有(offset 提交) |
| 历史回溯 | 不能 | 能(XRANGE 重读) | 能(按 offset) |
| 典型场景 | 实时通知广播 | 应用内可靠任务队列 | 大数据管道 |
使用注意
- Stream 数据默认驻留内存,总量受内存限制;超大吞吐、超长保留选 Kafka,并用 MAXLEN 裁剪控制内存。
- XREADGROUP 的 > 表示"只看未分配的新消息";不带 > 则读自己 PEL 里的待确认消息。
- 消费组要先创建(对不存在的流可用 MKSTREAM 一并创建);XADD 写入不存在的 key 会自动建流,不需要时可加 NOMKSTREAM 禁止。
小结
Stream = 可持久化的消息队列:XADD 生产、XRANGE/XREAD 消费、XGROUP + XREADGROUP + XACK 提供消费组可靠投递。要"不丢消息、能重试、能回溯"就用它而不是 Pub/Sub;要大规模持久队列直接考虑 Kafka。