Redis Stream 消息队列
提出问题
面试中经常被问到:"Redis 能不能当消息队列用?" 这个问题背后是对技术选型边界和 trade-off 的考察。用 List 做队列太简陋,Pub/Sub 丢消息,很多人在选型时要么一刀切上 Kafka,要么用 Redis 当万能缓存又硬塞 MQ 功能。Redis Stream 是 Redis 5.0 引入的原生消息队列模型,真正解决了"消息持久化、消费者组、ACK 确认"等一系列问题。理解 Stream 的定位和能力边界,能帮你在中小规模场景下省掉一套 Kafka 集群。
与其他方案的对比
为什么不直接用 List 当队列?
BRPOP/LPUSH 组合是"简陋的队列"——消费者 A 取走消息后内存里就没了,一旦消费者 A 处理到一半宕机,这条消息永远丢失。没有 ACK,没有重投,没有消费者组。线上最常见的事故:List 队列的消费者 OOM 重启,重启后 List 空了,但重启前消费的那批消息全丢了,业务方追问"那条支付回调去哪了"。
为什么不直接用 Pub/Sub?
Pub/Sub 的 fire-and-forget 模式极端危险:如果消费者不在线(订阅尚未建立或网络闪断),消息直接丢弃。Redis 官方文档明确写了 Pub/Sub 不保证消息可达。有团队用它做订单状态变更通知,结果消费者重启窗口期正好赶上促销高峰,30% 的订单通知丢了,排查半天才找到原因。
核心数据结构与命令
Redis Stream 是一个只追加的日志结构,每个消息有唯一的 ID(通常是 timestamp-sequence 格式,保证全局有序)。核心命令:
- XADD:往 Stream 追加消息。支持
MAXLEN裁剪历史,避免内存暴涨。 - XREAD:按 ID 范围读取消息,支持阻塞等待(
BLOCK),类似 Kafka 的简单消费。 - XREADGROUP:通过消费者组消费,支持组内负载均衡。
- XACK:确认消息已被处理,Stream 才会从 PEL(Pending Entries List)中移除该消息。
# 生产者:追加消息
XADD mystream MAXLEN ~ 1000 * sensor-id 1234 temperature 19.8
# 消费者组:创建组并消费
XGROUP CREATE mystream mygroup 0
XREADGROUP GROUP mygroup consumer1 COUNT 1 BLOCK 5000 STREAMS mystream >
# 确认处理完成
XACK mystream mygroup 1600000000000-0MAXLEN ~ 1000 的 ~ 表示近似裁剪——不是精确保留 1000 条,而是在内存里按节点裁剪,性能更高。别写成精确的 MAXLEN 1000,那会 O(N) 遍历删除。实测:MAXLEN ~ 1000 延迟约 0.5μs,MAXLEN 1000 在 Stream 有 10 万条时延迟飙到 50μs。
Consumer Group 工作原理
这是 Stream 和 List/Pub-Sub 最本质的区别。消费者组引入了一个组内负载均衡和消息确认的机制:
投递流程时序
Producer ──XADD──→ Stream (mystream)
│
┌──────────┼──────────┐
▼ ▼ ▼
Consumer-1 Consumer-2 Consumer-3 ← 同组
(msg-1) (msg-2) (msg-3)
│ │ │
▼ ▼ ▼
处理成功 处理成功 处理失败
│ │ │
XACK→PEL移除 XACK→PEL移除 │
│
┌─────────────┘
▼
PEL 中保留 msg-3
│ (pending)
▼
Consumer-2 调用 XCLAIM
接管 msg-3 → 重试三个核心机制
组内分发:一条消息只发给组内一个消费者(类似 Kafka 的 partition),通过
XREADGROUP自动分配。分配策略是轮询 + 空闲优先,不是按 hash 分片——所以无法保证同一 key 的消息落到同一消费者。PEL(Pending Entries List):每条投递但未确认的消息都会进入 PEL。如果消费者宕机,其他消费者可以调用
XCLAIM接管未 ACK 的消息。PEL 是内存中的链表结构,每个 entry 包着消息 ID、消费者名称、投递时间戳、重试次数。消息重投:
XPENDING查看 PEL 状态,XCLAIM转移消息归属权,实现 at-least-once 语义。重投的关键参数是MIN-IDLE-TIME:设置一个闲置超时,只有当消息在 PEL 中停留超过该时间才允许被其他消费者认领,避免短时卡顿就频繁重投。
PEL 踩坑实录
问题:某个消费者写完业务逻辑后忘记调 XACK,PEL 越积越多。线上 Redis 内存从 2GB 涨到 12GB,触发 OOM。
根因:PEL 是 Redis 内存中的数据结构,没有上限,不受 MAXLEN 控制。一个 1000 条/秒的 Stream,如果消费者一直不 ACK,24 小时 PEL 堆积 8640 万条 entry,每条至少 80 字节(消息 ID + 消费者名 + 元信息),直奔 7GB。
解法:监控 XLEN mystream 和 XPENDING mystream group 两个指标,设置告警阈值(比如 PEL 超过 10 万就告警)。代码里 XACK 必须和业务逻辑走同一个 try-catch-finally,不能单独放在异步回调里。
与 List、Pub/Sub 的对比
| 对比维度 | List(BRPOP/LPUSH) | Pub/Sub | Stream(含消费者组) |
|---|---|---|---|
| 持久化 | ✅ RDB/AOF | ❌ 不持久 | ✅ RDB/AOF |
| 消息确认 | ❌ 无 ACK | ❌ 无 ACK | ✅ XACK + PEL |
| 多消费者 | ❌ 竞争消费 | ✅ 广播(fan-out) | ✅ 组内负载均衡 + 通过独立组实现广播 |
| 消息重投 | ❌ 无 | ❌ 无 | ✅ XCLAIM + XPENDING |
| 阻塞读取 | ✅ BRPOP 可阻塞 | ✅ SUBSCRIBE 阻塞 | ✅ XREAD/XREADGROUP BLOCK |
| 消息回溯 | ❌ 消费即删除 | ❌ 不存储 | ✅ 按 ID 范围读取 |
| 内存占用可控 | ✅ 队列长度可控 | ❌ 无存储 | ✅ MAXLEN 近似裁剪 |
| 延迟(P99) | ~0.1ms | ~0.1ms | ~0.5ms(多了 PEL 维护) |
| 适用场景 | 简单任务队列 | 实时通知、实时聊天 | 可靠消息队列、事件流 |
Stream 主要补齐了消息确认和消费者组两个短板,同时保留了持久化。但 Stream 的广播能力需要每个消费者独立创建组(XGROUP CREATE stream $ 后用 XREADGROUP 消费),不像 Pub/Sub 那样一条 SUBSCRIBE 就能接收。
一个真实的生产事故复盘
背景
某 IoT 平台用 Redis Stream 做设备上报数据缓冲。设备 5000 台,每台每 10 秒上报一次,峰值 500 msg/s。四个消费者处理数据清洗 → 入库 → 告警 → 归档。
事故现象
上线两周后,Redis 内存使用率从 30% 暴涨到 85%,Stream 写入延迟从 0.5ms 升到 15ms,部分设备上报超时。
排查过程
INFO memory发现 used_memory_rss 12GB,其中 8GB 是 Stream 数据。XPENDING mystream group返回 PEL 长度 860 万。- 进一步查
XINFO STREAM mystream发现 Stream 实际只保留 2000 条(MAXLEN ~ 2000),但 PEL 有 860 万——问题不在 Stream 本身,在 PEL。 - 找到那个不 ACK 的消费者:告警模块的消费者。代码里
XACK被放在一个异步回调里,回调因为线程池满了被丢弃,XACK永远没执行。
修复
- 把
XACK移到业务处理方法的 finally 块里,和业务事务同生命周期。 - 加监控:
XPENDING长度超过 5 万就告警。 - 清理堆积:先
XCLAIM批量重新分配,再XACK已过期数据。写了脚本分批处理,避免单次XCLAIM太大阻塞 Redis。
教训
PEL 是悬在 Stream 头上的达摩克利斯之剑。XACK 不是可选的,少了它整个内存模型就崩了。
与 Kafka 的能力边界差异
Stream 在中小规模(单机或小集群,消息量 < 10 万/秒)非常好用,但和 Kafka 对比有硬伤:
| 对比维度 | Redis Stream | Kafka |
|---|---|---|
| 分区扩展 | 单实例,需手动分 key 到不同 Stream | 原生 partition 机制,线性扩展 |
| 吞吐上限 | 单核约 10 万 msg/s(受限于 PEL 维护) | 单 partition 约 100 万 msg/s |
| 存储介质 | 全内存,受物理内存限制 | 磁盘,可无限保留 |
| 消息时效 | MAXLEN 裁剪,旧消息自动删除 | 可配置保留时间/大小 |
| offset 管理 | 无持久化 offset,重启后从 > 消费 | 持久化 offset,可任意重置 |
| rebalance | ❌ 手动 XCLAIM | ✅ 自动 rebalance |
| 消息回放 | 手动 XREAD 按 ID 范围 | 随意重置 offset |
| 运维复杂度 | 零(Redis 已有) | 高(需要 ZK/KRaft + 集群管理) |
关键差异一:无分区副本与横向扩展
Stream 的每个 shard 就是单个 Redis 实例的主从复制,没有像 Kafka partition 那样跨多副本的物理分区机制。Kafka 的 partition 可以并行提高吞吐,Stream 的吞吐受限于单实例。如果需要横向扩展,只能手动把不同 key 分散到不同 Stream 实例,运维上比 Kafka 原生 partition 麻烦得多。
关键差异二:无 offset 持久化
Kafka 可以随意重置消费者 offset 到任意时间点回放,Stream 的消费者组通过 XREADGROUP 的 > 参数只能从最新未消费开始,要回放得手动 XREAD 按 ID 范围读。而且如果消费者组被删除重建,XGROUP CREATE stream 0 是从头开始,$ 是从最新开始,没有中间态。
关键差异三:无自动 rebalance
Kafka 消费者加入/退出会触发 partition 重新分配,Redis Stream 的消费者组需要手动 XCLAIM 处理消费者宕机,没有自动 rebalance。如果 4 个消费者挂了 1 个,剩下的 3 个不会自动接过挂掉消费者的消息,必须由外部守护进程定时扫描 PEL 并调用 XCLAIM。
选型决策树
什么时候用 Stream
- 团队已经用了 Redis,不想引入 Kafka/ZooKeeper 运维成本
- 消息量 < 5 万/秒,消息保留时间 < 24 小时
- 可以接受 at-least-once 语义(重复消费需要业务幂等)
- 延迟要求 < 5ms
- 典型场景:任务队列、事件流、日志收集、IoT 数据缓冲
什么时候不该用 Stream
- 需要严格 partition 顺序(同一 key 必须落到同一消费者)
- 消息量百万级/秒
- 需要无限回放历史
- 消费者数量动态变化频繁,需要自动 rebalance
- 不能容忍消息丢失(虽然 Stream 有持久化,但 Redis 主从切换 + 异步复制有丢消息风险)
最佳实践总结
- XACK 必须和业务逻辑同生命周期:try-catch-finally 的 finally 块里,不是异步回调里。
- MAXLEN 必须加:不裁剪的 Stream 会吃掉所有内存。
MAXLEN ~ 10000的~近似裁剪性能更好。 - 监控 PEL 长度:PEL 不受 MAXLEN 控制,是独立的内存消耗源。
XPENDING超过 10 万就告警。 - 消费者组用
>消费:XREADGROUP的>参数代表"只消费我还没消费过的消息",不要用具体的 ID,否则会重复消费。 - 重试机制:
XCLAIM前先检查XPENDING的 idle 时间,不要短于 30 秒,避免消息刚被获取还没处理就被别人抢走。 - 幂等消费:Stream 提供 at-least-once,消费者业务逻辑必须幂等(用业务 ID 去重或唯一约束)。
- 不要用 Stream 做广播:广播场景用 Pub/Sub,Stream 做广播需要每个消费者建独立组,运维成本高。
参考
参考:Redis 官方文档 — Redis Streams 参考:Redis 源码
src/t_stream.c— 核心实现 参考:Martin Kleppmann — 《Designing Data-Intensive Applications》第 11 章(流处理)