vault backup: 2026-05-24 21:18:14
This commit is contained in:
@@ -0,0 +1,455 @@
|
||||
---
|
||||
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]] — 知识索引总览
|
||||
Reference in New Issue
Block a user