Files
cs-note/hhs/Redis/15-Stream.md
T
2026-05-24 21:18:14 +08:00

456 lines
16 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
---
tags: [Redis, 缓存, Stream, 消息队列]
create time: 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 能胜任任务队列、事件溯源等对可靠性有要求的场景。
## 一、核心概念
### 基本操作命令
```bash
# === 生产者 ===
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` 正常消费循环 |
### Stream 底层结构
每条消息以 Radix Tree + Listpack 存储,时间复杂度 O(1) 追加。ID 格式为 `<毫秒时间戳>-<序号>`,天然有序。
> [!question] Stream 会像 List 一样消费后删除吗?
> 不会。Stream 是**只追加日志**,消息消费后仍在。需要显式 `XDEL` 或通过 `MAXLEN`/`MINID` 策略裁剪。这带来了「可回溯」的优势,但也意味着必须主动管理容量。
## 二、Consumer Group 机制详解
Consumer Group 是 Stream 的核心抽象,它让多个消费者**协作消费同一条 Stream**,每条消息只会被组内的一个消费者处理。
### 消息生命周期
```mermaid
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` 后才会移除。
```bash
# 查看组内 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 记录了「谁领了哪条消息、领了多久、重试了几次」,这就是「至少一次投递」语义的实现基础。
### 消费者组内部机制
```mermaid
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 手动认领。
## 三、Go 实战
### 生产者
```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
}
```
### 消费者(含重连与重试)
```go
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 实现「至少一次投递」的关键。
### 初始化消费者组
```go
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 —— 检查未确认消息
```bash
# 概览:返回 [最小ID, 最大ID, 总条数, 各消费者pending数]
XPENDING mystream mygroup
# 详细列表:按 ID 范围查询
XPENDING mystream mygroup - + 10
# - = 最小ID,+ = 最大ID,10 = 最多返回10条
```
### XCLAIM —— 转移消息所有权
当消费者宕机后,需要把它 pending 的消息转给存活的消费者处理:
```bash
# 将空闲超过 30000ms 的消息转给 consumer-b
XCLAIM mystream mygroup consumer-b 30000 <msg-id-1> <msg-id-2>
```
```go
// 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,
Idle: 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 合并为一条原子命令:
```bash
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 —— 按数量裁剪
```bash
# 写入时裁剪(推荐,边写边裁)
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+)
```bash
# 删除 ID 小于指定值的消息
XADD mystream MINID 1685000000000-0 * field value
XADD mystream MINID ~ 1685000000000-0 * field value # 近似模式
```
### XTRIM —— 手动裁剪
```bash
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 不可替代
## 七、死信队列(Dead Letter Queue)
当某条消息被重试 N 次仍然失败时,继续重试没有意义。此时应将其移入**死信队列**,等待人工介入或补偿处理。
### 实现思路
Stream 本身没有内置 DLQ,但可以通过 XCLAIM 的 `delivery-count` 自行实现:
```go
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.RetryCount >= 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 {
return
}
// 写入死信 Stream(附带原始 ID 以便溯源)
dlqFields := msgs[0].Values
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 堆积意味着业务逻辑有系统性问题
## 八、消息完整生命周期
```mermaid
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
```
## 关联笔记
- [[hhs/Redis/09-高级特性]] — Pub/Sub 基础、Lua 脚本原子操作
- [[hhs/Redis/03-基本命令速查]] — Redis 命令速查手册
- [[hhs/Redis/README]] — 知识索引总览