Files
cs-note/hhs/Redis/15-Stream.md
T
2026-06-08 23:08:57 +08:00

19 KiB
Raw Blame History

tags, create time
tags create time
Redis
缓存
Stream
消息队列
2026-05-24 10:00

Redis Stream 消息队列

概述

Redis Stream 是 Redis 5.0 引入的日志型数据结构,专门为可靠消息传递设计。它弥补了 Pub/Sub 「发完即忘、不持久化」的先天缺陷,同时又比 Kafka/RabbitMQ 这类外部中间件轻量——不需要额外部署组件,复用已有的 Redis 实例即可。

[!question] Pub/Sub 够用吗?为什么要用 Stream? Pub/Sub 的核心问题:消息不持久化,订阅者下线就丢失。Stream 将每条消息追加写入内存(可配合 AOF 持久化),并通过 Consumer Group 实现消费确认机制——消息只有被 ACK 后才算真正消费完成。这使得 Stream 能胜任任务队列、事件溯源等对可靠性有要求的场景。

一、核心概念

基本操作命令

# === 生产者 ===
XADD mystream * field1 value1 field2 value2
# * 表示自动生成 ID(时间戳-序号),如 1685000000000-0

# === 消费者(无消费者组模式)===
XREAD COUNT 10 BLOCK 5000 STREAMS mystream 0
# BLOCK 5000 = 最多等 5 秒;0 = 从头读取
# 返回一批 XEntry,每条包含 ID + fields

# === 消费者组模式(推荐)===
XGROUP CREATE mystream mygroup $    # $ = 只消费创建后的新消息
XREADGROUP GROUP mygroup consumer1 COUNT 1 BLOCK 5000 STREAMS mystream >
# > 表示「只取未分配的新消息」
# 返回后消息进入该消费者的 pending 列表

XACK mystream mygroup <entry-id>    # 确认消费完成

[!tip] > vs $ vs 0 —— 三个特殊 ID

ID 含义 典型场景
0 Stream 的第一条消息 数据重放、全量回溯
$ 当前最新消息的 ID XGROUP CREATE 时指定起始点
> 只返回组内尚未投递给任何消费者的消息 XREADGROUP 正常消费循环

注意:> 只在 XREADGROUP 中有意义。如果传 0 读取消费者组,则返回该消费者的 pending 消息而非新消息——这在"重启后恢复未完成任务"场景中非常有用。

Stream ID 自动生成机制

ID 格式为 <毫秒时间戳>-<序号>,由 Redis 服务器生成(用 * 时)。规则如下:

  1. 时间戳部分:取当前毫秒级 Unix 时间戳
  2. 序号部分:同一毫秒内递增(从 0 开始);如果时间戳前进到下一毫秒,序号重置为 0
XADD mystream * k v    # 返回 1685000000000-0
XADD mystream * k v    # 同一毫秒内 → 1685000000000-1
XADD mystream * k v    # 跨毫秒后 → 1685000000001-0

[!question] 能手动指定 ID 吗? 可以,但必须严格递增。新 ID 必须大于 Stream 中已有最大 ID,否则报错 (ERR) The ID specified in XADD is equal or smaller than the target stream top item。手动指定 ID 的典型场景:数据迁移、事件溯源回放。

底层存储结构

每条消息以 Radix Tree + Listpack 存储,追加写入时间复杂度 O(1)。同一时间戳内的多条消息会紧凑打包在 Listpack 节点中,极大节省内存。

[!question] Stream 会像 List 一样消费后删除吗? 不会。Stream 是只追加日志,消息消费后仍在。需要显式 XDEL 或通过 MAXLEN/MINID 策略裁剪。这带来了「可回溯」的优势,但也意味着必须主动管理容量。

历史消息读取

消费者组模式适用于"实时消费",但有时你需要按范围回溯历史消息——比如排障、数据对账、事件溯源。这时用 XRANGE/XREVRANGE:

# 按 ID 范围读取:从最小 ID 到最大 ID,最多 10 条
XRANGE mystream - + COUNT 10
# - = 最小 ID,+ = 最大 ID

# 按时间范围读取(ID 只需写时间戳部分,序号默认为 0)
XRANGE mystream 1685000000000 1685000060000 COUNT 10

# 反向读取(从新到旧)
XREVRANGE mystream + - COUNT 5

# 获取 stream 长度
XLEN mystream

[!tip] XRANGE vs XREAD XREAD 适合持续阻塞消费(配合 BLOCK),是生产者的标准读法。XRANGE 适合一次性按范围查询,不阻塞,常用于排障和数据审计。两者互补,不要混淆。

二、Consumer Group 机制详解

Consumer Group 是 Stream 的核心抽象,它让多个消费者协作消费同一条 Stream,每条消息只会被组内的一个消费者处理。

消息生命周期

flowchart TD
    P["Producer"] -->|"XADD"| S["Stream"]
    S -->|"XREADGROUP"| C1["Consumer A"]
    S -->|"XREADGROUP"| C2["Consumer B"]
    C1 -->|"XACK"| S
    C2 -->|"XACK"| S
    C1 -.->|"crashed"| PDL["Pending List"]
    PDL -->|"XCLAIM"| C2
    C2 -->|"XACK"| S

    style S fill:#dfd,stroke:#090
    style C1 fill:#ddf,stroke:#669
    style C2 fill:#ddf,stroke:#669
    style PDL fill:#fdd,stroke:#f00

Pending List(待确认列表)

当消费者通过 XREADGROUP GROUP ... > 读取消息后,消息进入该消费者的 PEL(Pending Entry List)。只有显式调用 XACK 后才会移除。

# 查看组内 pending 状态
XPENDING mystream mygroup

# 更详细的 pending 信息(包含消费者、空闲时间)
XPENDING mystream mygroup - + 10   # 返回最多 10 条 pending 条目

XPENDING 返回字段:

字段 含义
message-id 消息 ID
consumer 被分配的消费者
idle-ms 自读取以来的空闲时间(毫秒)
delivery-count 被投递的次数(重试计数器)

[!tip] Pending 是 Stream 的「可靠性根基」 Pub/Sub 没有 pending 概念,消息发出去就消失了。Stream 的 PEL 记录了「谁领了哪条消息、领了多久、重试了几次」,这就是「至少一次投递」语义的实现基础。

消费者组内部机制

sequenceDiagram
    participant P as Producer
    participant S as Stream
    participant G as Consumer Group
    participant C1 as Consumer A
    participant C2 as Consumer B

    P->>S: XADD msg-1
    P->>S: XADD msg-2
    P->>S: XADD msg-3
    C1->>G: XREADGROUP COUNT 2
    G->>S: 分配 msg-1, msg-2
    S-->>C1: msg-1, msg-2 (pending)
    C2->>G: XREADGROUP COUNT 2
    G->>S: 分配 msg-3
    S-->>C2: msg-3 (pending)
    C1->>S: XACK msg-1
    C1->>S: XACK msg-2
    C2->>S: XACK msg-3

[!question] 消息是如何分配给消费者的? Redis 采用轮询分配(round-robin):新消息到达时,依次分配给组内活跃的消费者。注意不是负载均衡——如果某个消费者处理慢,它积压的 pending 消息不会自动转移给其他消费者,需要通过 XCLAIM 手动认领。

[!tip] NOACK 模式——高吞吐场景的取舍 XREADGROUP 支持 NOACK 选项,消息读取后跳过 PEL,直接视为已消费。适用于对可靠性要求不高但吞吐要求极高的场景(如实时指标采样)。代价是消费者崩溃时消息会永久丢失。

XREADGROUP GROUP mygroup consumer1 NOACK COUNT 10 BLOCK 5000 STREAMS mystream >

XINFO —— 运维监控利器

排查 Stream 问题时,XINFO 是你的第一选择:

# 查看 stream 元信息(长度、首尾 ID、消费者组数量)
XINFO STREAM mystream

# 查看消费者组列表
XINFO GROUPS mystream

# 查看组内各消费者状态(pending 数、空闲时间)
XINFO CONSUMERS mystream mygroup

[!tip] 运维小贴士 监控脚本可以定期执行 XINFO CONSUMERS,如果某个消费者的 idle 值远超阈值且 pending 不为 0,说明该消费者可能已宕机,需要触发 XCLAIM 或告警。

三、Go 实战

生产者

// StreamProducer 封装了 Stream 的写入逻辑
type StreamProducer struct {
    rdb      *redis.Client
    stream   string
    maxSize  int64 // MAXLEN 裁剪阈值
}

func NewStreamProducer(rdb *redis.Client, stream string) *StreamProducer {
    return &StreamProducer{rdb: rdb, stream: stream, maxSize: 100000}
}

func (p *StreamProducer) Publish(ctx context.Context, fields map[string]interface{}) (string, error) {
    // XADD 的 ApproximateMaxLen 约等于 MAXLEN ~,性能更好
    id, err := p.rdb.XAdd(ctx, &redis.XAddArgs{
        Stream: p.stream,
        ApproximateMaxLen: true, // ~ 模式,不精确计数但更快
        MaxLen: p.maxSize,
        Values: fields,
    }).Result()
    if err != nil {
        return "", fmt.Errorf("XADD failed: %w", err)
    }
    return id, nil
}

消费者(含重连与重试)

type StreamConsumer struct {
    rdb        *redis.Client
    stream     string
    group      string
    consumer   string
    handler    func(ctx context.Context, msg redis.XMessage) error
    batchSize  int64
    blockTime  time.Duration
}

func (c *StreamConsumer) Start(ctx context.Context) {
    for {
        select {
        case <-ctx.Done():
            return
        default:
            // XREADGROUP 阻塞读取
            streams, err := c.rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
                Group:    c.group,
                Consumer: c.consumer,
                Streams:  []string{c.stream, ">"}, // > = 只取新消息
                Count:    c.batchSize,
                Block:    c.blockTime,
            }).Result()

            if err != nil {
                if errors.Is(err, redis.Nil) {
                    continue // 超时无消息,正常
                }
                log.Printf("XReadGroup error: %v, retrying...", err)
                time.Sleep(time.Second) // 退避重试
                continue
            }

            for _, stream := range streams {
                for _, msg := range stream.Messages {
                    if err := c.handler(ctx, msg); err != nil {
                        log.Printf("handler error for %s: %v", msg.ID, err)
                        continue // 不 ACK,消息留在 pending 列表
                    }
                    // 处理成功,确认消费
                    if err := c.rdb.XAck(ctx, c.stream, c.group, msg.ID).Err(); err != nil {
                        log.Printf("XAck error for %s: %v", msg.ID, err)
                    }
                }
            }
        }
    }
}

[!tip] 不 ACK = 进入 pending handler 返回 error 时跳过 XACK,这条消息会留在消费者的 PEL 中。后续通过 XPENDING + XCLAIM 捞回重试——这就是 Stream 实现「至少一次投递」的关键。

初始化消费者组

func EnsureGroup(ctx context.Context, rdb *redis.Client, stream, group string) error {
    // MKSTREAM: 如果 stream 不存在则自动创建
    err := rdb.XGroupCreateMkStream(ctx, stream, group, "$").Err()
    if err != nil && !strings.Contains(err.Error(), "BUSYGROUP") {
        return err
    }
    // BUSYGROUP = 组已存在,忽略即可
    return nil
}

四、消息确认与重试

XPENDING —— 检查未确认消息

# 概览:返回 [最小ID, 最大ID, 总条数, 各消费者pending数]
XPENDING mystream mygroup

# 详细列表:按 ID 范围查询
XPENDING mystream mygroup - + 10
# - = 最小ID,+ = 最大ID,10 = 最多返回10条

XCLAIM —— 转移消息所有权

当消费者宕机后,需要把它 pending 的消息转给存活的消费者处理:

# 将空闲超过 30000ms 的消息转给 consumer-b
XCLAIM mystream mygroup consumer-b 30000 <msg-id-1> <msg-id-2>
// Go 实现:定期扫描超时 pending,认领到自己
func (c *StreamConsumer) ClaimAbandoned(ctx context.Context, minIdle time.Duration) {
    pendings, err := c.rdb.XPendingExt(ctx, &redis.XPendingExtArgs{
        Stream:   c.stream,
        Group:    c.group,
        MinIdle:  minIdle,
        Start:    "-",
        End:      "+",
        Count:    100,
        Consumer: "", // 不限定消费者,扫描所有
    }).Result()
    if err != nil || len(pendings) == 0 {
        return
    }

    ids := make([]string, len(pendings))
    for i, p := range pendings {
        ids[i] = p.ID
    }

    // XCLAIM 认领
    claimed, err := c.rdb.XClaim(ctx, &redis.XClaimArgs{
        Stream:   c.stream,
        Group:    c.group,
        Consumer: c.consumer,
        MinIdle:  minIdle,
        Messages: ids,
    }).Result()
    if err != nil {
        return
    }

    // 对认领到的消息重新执行 handler
    for _, msg := range claimed {
        if err := c.handler(ctx, msg); err == nil {
            c.rdb.XAck(ctx, c.stream, c.group, msg.ID)
        }
    }
}

XAUTOCLAIM —— Redis 6.2+ 自动认领

XAUTOCLAIM 把 XPENDING + XCLAIM 合并为一条原子命令:

XAUTOCLAIM mystream mygroup consumer-b 30000 0-0 COUNT 10
# 30000 = 最小空闲时间(ms)
# 0-0   = 起始游标(首次从头开始)
# COUNT = 每次最多返回 10 条

# 返回:next-start-id(下次传入作为游标)+ 消息列表 + 已删除的 ID

[!question] XCLAIM 会不会被两个消费者同时认领同一条消息? 不会。XCLAIM 是单线程原子操作,Redis 保证同一消息在任意时刻只属于一个消费者。如果 consumer-A 已经 ACK 了某条消息,后续 XCLAIM 不会再返回它。

五、容量管理

Stream 是只追加的,消息不会自动删除。生产环境必须配置裁剪策略。

MAXLEN —— 按数量裁剪

# 写入时裁剪(推荐,边写边裁)
XADD mystream MAXLEN 10000 * field value

# 精确裁剪 vs 近似裁剪
XADD mystream MAXLEN 10000 * field value      # 精确,O(N) 开销
XADD mystream MAXLEN ~ 10000 * field value    # 近似,~ 模式,O(1) 开销

MINID —— 按 ID 裁剪(Redis 6.2+)

# 删除 ID 小于指定值的消息
XADD mystream MINID 1685000000000-0 * field value
XADD mystream MINID ~ 1685000000000-0 * field value  # 近似模式

XDEL —— 删除单条消息

# 删除指定 ID 的消息
XDEL mystream 1685000000000-0

# 批量删除
XDEL mystream 1685000000000-0 1685000000000-1

[!warning] XDEL 不会清除 PEL 引用 与 MAXLEN/MINID 一样,XDEL 删除的是 Stream 中的原始消息,但 Pending Entry List 中的引用仍然存在。被删除的消息 ID 在 XPENDING 中仍可见,XCLAIM 尝试读取时会返回 nil。

XTRIM —— 手动裁剪

XTRIM mystream MAXLEN ~ 5000
XTRIM mystream MINID 1685000000000-0

[!warning] 近似裁剪(~)的取舍 ~ 模式不会精确保留 N 条,实际保留数量可能略多(因为底层按 Radix Tree 节点边界裁剪)。对于绝大多数场景,近似裁剪足够,而且性能显著优于精确裁剪。建议生产环境一律使用 ~。

[!question] 裁剪会影响 pending 中的消息吗? 会。MAXLEN/MINID 会删除 stream 中的原始消息,但 pending 列表中的引用不会被清除。这意味着 XPENDING 仍会显示这些 ID,但 XCLAIM 尝试读取时会返回 nil。设计裁剪策略时需确保消费者的处理速度跟得上写入速度。

六、Stream vs Pub/Sub vs Kafka 对比

维度 Stream Pub/Sub Kafka
持久化 内存 + AOF 不持久化 磁盘日志
消费者组 原生支持 不支持 原生支持
投递语义 至少一次(配合 ACK) 最多一次(fire-and-forget) 至少一次 / 精确一次
消息回溯 支持(按 ID 范围) 不支持 支持(按 offset)
吞吐量 中等(10w+ QPS) 高(无持久化开销) 极高(百万级 QPS)
运维复杂度 低(复用 Redis) 低 高(独立集群)
适用规模 中小规模 实时通知 大规模事件流

[!tip] 如何选型?

  • 已有 Redis、消息量不大(< 10w/s):Stream 是最佳选择,零额外运维成本
  • 实时推送、允许丢消息:Pub/Sub 更简单直接
  • 海量事件流、需要精确一次语义:Kafka 不可替代

[!warning] Redis Cluster 模式下的 Stream 在 Redis Cluster 中,一个 Stream 被分配到单个 slot,所有写入和读取都路由到同一节点,无法水平分片。Consumer Group 的所有操作(XREADGROUP、XACK、XPENDING)也必须在同一节点完成。如果单个 Stream 成为热点,考虑将其拆分为多个 Stream(如 orders:0、orders:1……)并使用客户端负载均衡。

七、死信队列(Dead Letter Queue)

当某条消息被重试 N 次仍然失败时,继续重试没有意义。此时应将其移入死信队列,等待人工介入或补偿处理。

实现思路

Stream 本身没有内置 DLQ,但可以通过 XCLAIM 的 delivery-count 自行实现:

const maxRetries = 5

func (c *StreamConsumer) ProcessPendingWithDLQ(ctx context.Context) {
    pendings, _ := c.rdb.XPendingExt(ctx, &redis.XPendingExtArgs{
        Stream:  c.stream,
        Group:   c.group,
        Start:   "-",
        End:     "+",
        Count:   100,
        MinIdle: 30 * time.Second,
    }).Result()

    for _, p := range pendings {
        if p.TimesDelivered >= maxRetries {
            // 超过重试次数 → 写入死信 Stream
            c.moveToDLQ(ctx, p.ID)
            continue
        }
        // 认领并重试
        claimed, _ := c.rdb.XClaim(ctx, &redis.XClaimArgs{
            Stream:   c.stream,
            Group:    c.group,
            Consumer: c.consumer,
            MinIdle:  30 * time.Second,
            Messages: []string{p.ID},
        }).Result()
        for _, msg := range claimed {
            if err := c.handler(ctx, msg); err == nil {
                c.rdb.XAck(ctx, c.stream, c.group, msg.ID)
            }
        }
    }
}

func (c *StreamConsumer) moveToDLQ(ctx context.Context, msgID string) {
    // 读取原始消息内容(XCLAIM 的返回可能为空,用 XRANGE 兜底)
    msgs, _ := c.rdb.XRange(ctx, c.stream, msgID, msgID).Result()
    if len(msgs) == 0 {
        // 消息可能已被裁剪删除,仅 ACK 清理 pending
        c.rdb.XAck(ctx, c.stream, c.group, msgID)
        return
    }
    // 拷贝一份,避免污染原始消息的 Values map
    dlqFields := make(map[string]interface{}, len(msgs[0].Values)+2)
    for k, v := range msgs[0].Values {
        dlqFields[k] = v
    }
    dlqFields["original_id"] = msgID
    dlqFields["original_stream"] = c.stream
    c.rdb.XAdd(ctx, &redis.XAddArgs{
        Stream: c.stream + ":dlq",
        Values: dlqFields,
    })
    // 从原 Stream 确认移除
    c.rdb.XAck(ctx, c.stream, c.group, msgID)
}

[!tip] DLQ 最佳实践

  • 死信 Stream 命名建议 <stream>:dlq,保持关联性
  • 在 DLQ 消息中记录 original_id、original_stream、error_reason,方便排查
  • 监控 DLQ 长度,超过阈值立即告警——DLQ 堆积意味着业务逻辑有系统性问题

八、消息完整生命周期

flowchart TD
    P["Producer"] -->|"XADD"| S["Stream"]
    S -->|"XREADGROUP"| C1["Consumer A"]
    S -->|"XREADGROUP"| C2["Consumer B"]
    C1 -->|"XACK success"| ACK["Confirmed"]
    C1 -->|"handler fail"| PEL["Pending Entry List"]
    C2 -->|"XACK success"| ACK
    C2 -->|"handler fail"| PEL
    PEL -->|"retry count < max"| C1
    PEL -->|"retry count >= max"| DLQ["Dead Letter Queue"]
    ACK -->|"MAXLEN trim"| TRIM["Trimmed"]
    DLQ -->|"人工处理"| FIX["Fixed"]

    style S fill:#dfd,stroke:#090
    style PEL fill:#fdd,stroke:#f00
    style DLQ fill:#ffd,stroke:#f90
    style ACK fill:#dff,stroke:#099
    style TRIM fill:#eee,stroke:#999

关联笔记