Learn
Redis/09-stream

Stream:可靠消息队列

List 做队列丢消息、不支持多消费组;Pub/Sub 不在线就收不到。Redis 5 引入的 Stream 补齐了这些短板:消息持久保留、消费组分摊、ACK 确认、故障转移认领——一个接近 Kafka 核心语义的日志结构。

1. 数据模型

Stream 是一个只追加的消息日志,每条消息有全局有序的 ID 和一组 field-value:

stream: orders
+---------------------+----------------------------+
| ID                  | fields                     |
+---------------------+----------------------------+
| 1722310000000-0     | action=create sku=2001     |
| 1722310000000-1     | action=pay   order=8801    |   <- 同毫秒第二条, 序号+1
| 1722310005123-0     | action=ship  order=8801    |
+---------------------+----------------------------+
        ID 格式: 毫秒时间戳-序号, 严格递增

2. 生产与独立消费

2.1 XADD 写入

XADD 生产与 XRANGE 回溯
# 星号表示让 Redis 自动生成 ID(毫秒时间戳-序号)
XADD orders * action create sku 2001
XADD orders * action pay order 8801
XLEN orders
# 控制长度:近似裁剪到 1 万条(波浪号表示允许略多,性能好)
XADD orders MAXLEN ~ 10000 * action ship order 8801
# 减号和加号是最小/最大 ID,读取全部消息
XRANGE orders - +

2.2 XRANGE / XREAD 读取

XRANGE key - + 读取全部消息,-/+ 分别代表最小与最大 ID,也可以传具体 ID 做范围查询。想要"只等新消息"则用阻塞读:

# 阻塞读取新消息: $ 表示"只要从现在开始的新消息"
127.0.0.1:6379> XREAD BLOCK 5000 STREAMS orders $
(挂起, 等到新消息或 5 秒超时)

XREAD 模式下每个消费者独立维护自己读到的位置,适合"每个消费者都要全量消息"的广播场景。

3. 消费组:分摊与确认

消费组让多个消费者分摊同一个 Stream 的消息(同组内一条消息只给一个消费者),并提供 ACK 机制。

3.1 创建组与消费

消费组:XGROUP + XREADGROUP
# 这里手写 ID 便于后面 XACK 引用,生产中用星号自动生成
XADD orders 1-1 action create sku 2001
XADD orders 1-2 action pay order 8801
# 从头(0)开始消费;用美元符号则只消费此后的新消息
XGROUP CREATE orders g1 0
# 消费者 c1 领取 2 条新消息,大于号表示"从未投递给本组的消息"
XREADGROUP GROUP g1 c1 COUNT 2 STREAMS orders >
# 再读一次,已经没有新消息了
XREADGROUP GROUP g1 c1 COUNT 2 STREAMS orders >

> 是特殊 ID,表示"给我从未投递给本组的新消息"。消息投递后进入该消费者的 PEL(Pending Entries List,待确认列表),直到被 ACK。

3.2 确认与查看积压

XACK 确认与 XPENDING 查积压
XADD orders 1-1 action create sku 2001
XADD orders 1-2 action pay order 8801
XGROUP CREATE orders g1 0
XREADGROUP GROUP g1 c1 COUNT 2 STREAMS orders >
# 处理完成,确认第一条
XACK orders g1 1-1
# 组内待确认概况:总数、最小/最大 ID、各消费者的挂起数
XPENDING orders g1
# 详细列表:每条的消费者、空闲毫秒数、投递次数
XPENDING orders g1 - + 10

3.3 故障转移:XAUTOCLAIM

消费者 c1 崩溃后,它 PEL 里的消息不会自动重投。其他消费者用 XAUTOCLAIM 把"挂起超过一定时间"的消息认领过来:

XAUTOCLAIM 认领僵尸消息
XADD orders 1-1 action pay order 8801
XGROUP CREATE orders g1 0
# c1 领走消息但不 ACK,模拟处理到一半崩溃
XREADGROUP GROUP g1 c1 COUNT 1 STREAMS orders >
XPENDING orders g1 - + 10
# c2 认领空闲超过 0ms 的消息(生产中填 60000 之类的阈值)
XAUTOCLAIM orders g1 c2 0 0
# 现在这条消息归 c2 了
XPENDING orders g1 - + 10

消费者应定期跑 XAUTOCLAIM 作为"补偿任务"。投递次数(delivery count)随每次认领递增,超过阈值(如 5 次)仍失败的消息应转入死信队列(另一个 Stream)人工处理,避免毒消息无限循环。

3.4 完整消费循环(伪代码)

while True:
    # 1. 优先处理新消息
    msgs = XREADGROUP GROUP g1 c1 COUNT 10 BLOCK 2000 STREAMS orders >
    for id, fields in msgs:
        try:
            process(fields)
            XACK orders g1 id
        except Exception:
            pass          # 不 ACK, 留在 PEL 等待重试/认领
    # 2. 周期性认领僵尸消息(也可由独立任务做)
    XAUTOCLAIM orders g1 c1 60000 0

4. 三种消息方案对比

能力ListPub/SubStream
消息持久化是(弹出即删)否,即发即失是,可回溯
离线消息保留丢失保留
ACK 确认无无有
消费组分摊无无有
多组独立消费(广播)无有有(多个 group)
阻塞消费BLPOPSUBSCRIBEXREAD BLOCK
适用简单任务队列实时通知,可容忍丢失可靠队列、事件流

选型直觉:能容忍丢消息的实时广播用 Pub/Sub;一切"不能丢"的异步任务用 Stream;List 只在极简单场景且不想引入消费组概念时使用。

5. 管理命令

127.0.0.1:6379> XINFO STREAM orders
 1) "length"
 2) (integer) 3
 3) "radix-tree-keys"
 4) (integer) 1
 ...
127.0.0.1:6379> XINFO GROUPS orders
1) 1) "name"
   2) "g1"
   3) "consumers"
   4) (integer) 2
   5) "pending"
   6) (integer) 1
   7) "last-delivered-id"
   8) "1722310000000-1"
127.0.0.1:6379> XDEL orders 1722310000000-0     # 删除指定消息
(integer) 1
127.0.0.1:6379> XTRIM orders MAXLEN ~ 1000      # 手动裁剪
(integer) 0
⚠️Stream 不裁剪会吃光内存

Stream 消息 ACK 之后也不会被删除(这是它能多组消费、可回溯的原因)。必须通过 XADD 的 MAXLEN/MINID 参数或定期 XTRIM 控制长度,否则一个高吞吐 Stream 会无限增长直至 OOM。带波浪号的 MAXLEN ~ N 是近似裁剪,按宏节点整块删除,性能远好于精确裁剪。

ℹ️和 Kafka 的差距

Stream 有消费组与 ACK,但没有分区(单 key 单线程处理,吞吐有单机上限)、没有多副本 ISR 语义(可靠性依赖 Redis 自身持久化与主从复制,主从切换仍可能丢尾部消息)。十万级 QPS 以下的业务队列足够;金融级或超大吞吐场景仍需 Kafka/Pulsar。

小结

  • Stream 是只追加日志:XADD 生产,XRANGE 回溯,XREAD 独立消费
  • 消费组三件套:XREADGROUP 领取(进入 PEL)、XACK 确认、XAUTOCLAIM 认领僵尸消息
  • 投递次数超限的消息应转死信 Stream,防止毒消息循环
  • 必须用 MAXLEN/XTRIM 控制长度;ACK 不等于删除
  • 下一章回到 key 的生命周期:过期与内存淘汰 →
🎯练习
  1. 创建一个 Stream 和消费组,用两个 redis-cli 窗口模拟两个消费者,验证消息被分摊而不是都收到。
  2. 一个消费者领取消息后不 ACK,用 XPENDING 观察其挂起状态,再用 XAUTOCLAIM 从另一个消费者身份认领它。
  3. 对比实验:Pub/Sub 的订阅者断开期间发布消息,重连后收不到;Stream 消费组的消费者离线期间的消息重连后仍能领取。