首页 / Redis 入门教程 / 流 Stream

Redis 入门教程

流 Stream

本教程共 40 篇 · 第 14 篇 · 更新于 2026-08-02

redisstreamXADDXREADXGROUP消息队列事件溯源消费者组

14. 流 Stream

本节目标

  • 理解 Stream 是「只追加日志(append-only log)」,支持持久化与随机访问
  • 掌握写入与读取:XADD / XLEN / XRANGE / XREVRANGE / XREAD / XTRIM
  • 掌握消费者组:XGROUP CREATE / XREADGROUP / XACK / XPENDING / XCLAIM
  • 理解 Stream 与 Pub/Sub 的本质区别(持久化 vs 瞬时)
  • 用「消息队列 + 至少一次投递」实战巩固

Stream 是 Redis 5.0 引入的数据类型,设计目标是可靠的消息队列与事件流。它是一条「消息链表」:每条消息有全局唯一的 ID 和若干字段(field-value),按追加顺序持久化。相比 Pub/Sub(发布订阅,消息不持久、掉线即丢),Stream 能:

  • 持久化:消息写入即落盘(配合 AOF/RDB),重启不丢;
  • 记住消费位置:每个消费者组有独立的游标,断线重连可从断点继续;
  • 支持消费者组:多个消费者分担同一流,实现负载均衡与「至少一次」投递语义。

14-1 写入与基本读取

XADD 向流追加一条消息。ID 用 * 让 Redis 自动生成(格式 <毫秒时间戳>-<序号>),消息体是一组 field-value:

命令:

127.0.0.1:6379> XADD mystream * sensor "temp" value "26.5"
127.0.0.1:6379> XADD mystream * sensor "humidity" value "60"

输出:

"1754000000000-0"
"1754000000000-1"

获取消息数量用 XLEN,按 ID 区间取消息用 XRANGE- 表示最小、+ 表示最大):

命令:

127.0.0.1:6379> XLEN mystream
127.0.0.1:6379> XRANGE mystream - + COUNT 10

输出:

(integer) 2
1) 1) "1754000000000-0"
   2) 1) "sensor"
      2) "temp"
      3) "value"
      4) "26.5"
2) 1) "1754000000000-1"
   2) 1) "sensor"
      2) "humidity"
      3) "value"
      4) "60"

反向取用 XREVRANGE(参数顺序是 end start)。限制流长度、防止无限增长用 XTRIM

命令:

127.0.0.1:6379> XTRIM mystream MAXLEN 1000

输出:

(integer) 0

MAXLEN 1000 表示最多保留 1000 条(前面的会被裁剪),非常适合「只保留最近 N 条」的日志流。

14-2 独立消费者读取:XREAD

XREAD 以「阻塞或非阻塞」方式从流尾部读取新消息,$ 表示从当前最新位置之后开始读:

命令:

127.0.0.1:6379> XREAD COUNT 2 STREAMS mystream 0

输出:

1) 1) "mystream"
   2) 1) 1) "1754000000000-0"
         2) 1) "sensor"
            2) "temp"
            3) "value"
            4) "26.5"
      2) 1) "1754000000000-1"
         2) 1) "sensor"
            2) "humidity"
            3) "value"
            4) "60"

STREAMS 后跟的 0 是从 ID 0 开始读(即读全部历史);若用 $ 则只等新消息。加上 BLOCK 5000 会阻塞最多 5 秒等待新消息——这就是「等待队列里有新任务」的经典写法。

14-3 消费者组:多消费者协同

当单消费者处理不过来时,用消费者组(Consumer Group) 让多个消费者分担。组维护一个共享游标 last_delivered_id,每条消息只投递给组内的一个消费者,实现负载均衡;已读但未确认(ACK)的消息进入「待处理(pending)」列表,可重新认领,保证「至少一次」投递。

创建组(从最新位置 $ 开始,或 0 从头消费):

命令:

127.0.0.1:6379> XGROUP CREATE mystream cg1 $

输出:

OK

消费者 c1 读取新消息(> 表示「取组内尚未投递过的消息」):

命令:

127.0.0.1:6379> XREADGROUP GROUP cg1 c1 COUNT 1 STREAMS mystream >

输出:

1) 1) "mystream"
   2) 1) 1) "1754000000500-0"
         2) 1) "sensor"
            2) "temp"
            3) "value"
            4) "27.1"

处理完成后用 XACK 确认,消息即从 pending 列表移除:

命令:

127.0.0.1:6379> XACK mystream cg1 1754000000500-0

输出:

(integer) 1

查看待处理消息用 XPENDING,把卡住(如消费者崩溃)的消息转移给别的消费者用 XCLAIM

命令:

127.0.0.1:6379> XPENDING mystream cg1

输出:

1) (integer) 0
2) (nil)
3) (nil)
4) (nil)

当前没有未确认消息(因为上面已 XACK)。真实场景中 XPENDING 会列出「哪些 ID 被谁读取、闲置多久」,据此决定 XCLAIM 转移。

14-4 实战:可靠消息队列

结合上述命令可以搭建一个轻量但可靠的队列:

  1. 生产者 XADD 投递任务;
  2. 多个 worker 用 XREADGROUP GROUP ... > 竞争消费;
  3. 处理成功 XACK,失败则靠 XPENDING + XCLAIM 重试;
  4. XTRIM MAXLEN 控制流长度,避免无限膨胀。

提示:Stream 自 Redis 5.0 起内置。它与 Pub/Sub 的关键区别是——Pub/Sub 消息「发完即忘」、不持久、无消费位点,适合广播通知;Stream 持久化、有消费组与位点,适合「不能丢」的任务队列与事件溯源(event sourcing)。Redis 8.x 中 Stream 还对复制流压缩(BUILD_COMPRESSION,需编译支持)等做了增强。本教程正文统一以 Redis 8.x(最新稳定版)为准。

14-5 常见误区

  • 误区:XADD 必须自己生成 ID。* 让 Redis 自动生成即可;若自定义 ID 必须单调递增,否则会报错。
  • 误区:消费者组里一条消息会发给所有消费者。 不会,组内每条消息只投递给一个消费者(竞争模式);想广播请用 Pub/Sub 或多个组。
  • 误区:XREADGROUP 用 0 和用 > 一样。 > 取「组内新消息」;用具体 ID(如 0)是「重读该消费者已读但未 ACK 的 pending 消息」,二者语义不同。
  • 误区:XACK 不影响数据。 XACK 只把消息从 pending 列表移除,原始消息仍在流里(可用 XDEL 真正删除)。

14.x 消息 ID、投递语义与监控

先看清 Stream 消息 ID 的结构:自动生成的 ID 形如 <毫秒时间戳>-<序号>,前半段是 Redis 服务器本地时间的毫秒数,后半段是同一毫秒内的递增序号。这个设计保证 ID 全局单调递增,天然支持「按时间排序」「从某时间点之后读」。如果你用 * 让 Redis 生成 ID,无需操心;若自定义 ID(如用外部逻辑时钟),必须保证严格递增,否则会报错。

关于「投递语义」要心里有数:Stream 的消费者组保证的是**至少一次(at-least-once)**投递——一条消息会被投递给组内某个消费者,在它 XACK 之前一直留在 pending 列表里;若消费者崩溃未 ACK,别的消费者可通过 XCLAIM 重新认领重投。但这不保证「恰好一次」:如果消费者处理成功、在 XACK 前崩溃,消息会被重投,可能出现重复处理。要做到「恰好一次」,需要业务侧做幂等(如用消息 ID 去重)。这与 Pub/Sub 的「至多一次、发完即忘」形成对比,也是 Stream 适合「不能丢的任务」的原因。

监控与运维上,XINFO STREAM mystream 可看流的长度、首尾 ID、消费者组数;XINFO GROUPS / XINFO CONSUMERS 可看各组的待处理量与消费者空闲时间,是排查「消息堆积、某消费者卡死」的必备手段。最后提醒:XTRIM MAXLEN 的裁剪是「近似」的(为性能允许略超阈值),不要依赖它做精确长度控制;真正要彻底删消息用 XDEL,但被删消息在范围查询里只是被跳过、空间由后续 reclaim 逐步回收。

14-6 小结

  • Stream 是持久化的只追加消息流,每条消息有唯一 ID 和 field-value 体。
  • 写/读:XADD / XLEN / XRANGE / XREVRANGE / XREAD(可 BLOCK)/ XTRIM
  • 消费者组:XGROUP CREATE / XREADGROUP(用 >) / XACK / XPENDING / XCLAIM,实现负载均衡与至少一次投递。
  • 与 Pub/Sub:Stream 持久化、有位点、可靠;Pub/Sub 瞬时、广播。
  • 用武之地:消息队列、事件溯源、日志流、任务分发。

提示:关于版本——本教程正文统一以 Redis 8.x(最新稳定版)为准;官方 redis.io 下载页另标 8.8 为 “Latest stable”,而 GitHub 上 redis/redis 的最新发布 tag 为 8.10.0,二者同属 8.x,命令与类型差异对教材影响极小。