Skip to content

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)中移除该消息。
bash
# 生产者:追加消息
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-0

MAXLEN ~ 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 → 重试

三个核心机制

  1. 组内分发:一条消息只发给组内一个消费者(类似 Kafka 的 partition),通过 XREADGROUP 自动分配。分配策略是轮询 + 空闲优先,不是按 hash 分片——所以无法保证同一 key 的消息落到同一消费者。

  2. PEL(Pending Entries List):每条投递但未确认的消息都会进入 PEL。如果消费者宕机,其他消费者可以调用 XCLAIM 接管未 ACK 的消息。PEL 是内存中的链表结构,每个 entry 包着消息 ID、消费者名称、投递时间戳、重试次数。

  3. 消息重投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 mystreamXPENDING mystream group 两个指标,设置告警阈值(比如 PEL 超过 10 万就告警)。代码里 XACK 必须和业务逻辑走同一个 try-catch-finally,不能单独放在异步回调里。

与 List、Pub/Sub 的对比

对比维度List(BRPOP/LPUSH)Pub/SubStream(含消费者组)
持久化✅ 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,部分设备上报超时。

排查过程

  1. INFO memory 发现 used_memory_rss 12GB,其中 8GB 是 Stream 数据。
  2. XPENDING mystream group 返回 PEL 长度 860 万。
  3. 进一步查 XINFO STREAM mystream 发现 Stream 实际只保留 2000 条(MAXLEN ~ 2000),但 PEL 有 860 万——问题不在 Stream 本身,在 PEL。
  4. 找到那个不 ACK 的消费者:告警模块的消费者。代码里 XACK 被放在一个异步回调里,回调因为线程池满了被丢弃,XACK 永远没执行。

修复

  • XACK 移到业务处理方法的 finally 块里,和业务事务同生命周期。
  • 加监控:XPENDING 长度超过 5 万就告警。
  • 清理堆积:先 XCLAIM 批量重新分配,再 XACK 已过期数据。写了脚本分批处理,避免单次 XCLAIM 太大阻塞 Redis。

教训

PEL 是悬在 Stream 头上的达摩克利斯之剑XACK 不是可选的,少了它整个内存模型就崩了。

与 Kafka 的能力边界差异

Stream 在中小规模(单机或小集群,消息量 < 10 万/秒)非常好用,但和 Kafka 对比有硬伤:

对比维度Redis StreamKafka
分区扩展单实例,需手动分 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 主从切换 + 异步复制有丢消息风险)

最佳实践总结

  1. XACK 必须和业务逻辑同生命周期:try-catch-finally 的 finally 块里,不是异步回调里。
  2. MAXLEN 必须加:不裁剪的 Stream 会吃掉所有内存。MAXLEN ~ 10000~ 近似裁剪性能更好。
  3. 监控 PEL 长度:PEL 不受 MAXLEN 控制,是独立的内存消耗源。XPENDING 超过 10 万就告警。
  4. 消费者组用 > 消费XREADGROUP> 参数代表"只消费我还没消费过的消息",不要用具体的 ID,否则会重复消费。
  5. 重试机制XCLAIM 前先检查 XPENDING 的 idle 时间,不要短于 30 秒,避免消息刚被获取还没处理就被别人抢走。
  6. 幂等消费:Stream 提供 at-least-once,消费者业务逻辑必须幂等(用业务 ID 去重或唯一约束)。
  7. 不要用 Stream 做广播:广播场景用 Pub/Sub,Stream 做广播需要每个消费者建独立组,运维成本高。

参考

参考:Redis 官方文档 — Redis Streams 参考:Redis 源码 src/t_stream.c — 核心实现 参考:Martin Kleppmann — 《Designing Data-Intensive Applications》第 11 章(流处理)

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。