欢迎光临
我们一直在努力

Redis Stream 学习笔记

一、Redis Stream 是什么

Redis Stream 是 Redis 5.0 引入的一种消息流数据结构,用于实现可靠消息队列。

特点:

特性

说明

消息持久化

消息不会自动删除

顺序ID

每条消息都有唯一ID

消费者组

支持多个消费者协作

ACK机制

确保消息处理完成

消息可回溯

可以读取历史消息

数据结构示意:

Stream
|
|— 1689345632421-0
| field=value
|
|— 1689345632425-0
| field=value

Stream 本质类似:

Kafka lite

适用于:

  • 任务队列
  • 日志流
  • 异步任务
  • 微服务解耦
  • AI任务调度

二、核心概念

1 Stream

Stream 是一个 消息日志结构。

特点:

  • 按时间顺序存储
  • 自动生成 ID
  • 可以存多个字段

示例:

r.xadd("mystream", {"field": "value"})

生成:

mystream
|
|— 1689345632421-0
field=value

2 Consumer Group(消费者组)

消费者组用于 实现任务分配与负载均衡。

结构:

Stream
|
|— Consumer Group
|
|— consumer1
|— consumer2
|— consumer3

特点:

  • 同组消费者 不会重复消费消息
  • Redis 自动分配任务

示例:

r.xgroup_create("mystream", "group1", id="0", mkstream=True)

参数说明:

参数

含义

mystream

stream名称

group1

消费者组

0

从最早消息开始

mkstream

如果stream不存在就创建

3 Consumer(消费者)

消费者是 具体执行任务的 worker。

例如:

group1
|
|— worker1
|— worker2

任务分配:

msg1 → worker1
msg2 → worker2
msg3 → worker1

Redis 会自动负载均衡。

注意:消费者名称必须唯一,否则会共享 Pending List。

三、核心 API

1 XADD

作用:向 Stream 添加消息

r.xadd("mystream", {"field": "value"})

Redis命令:

XADD mystream * field value

特点:

  • 自动生成 ID
  • 支持多字段
  • 数据持久化

返回:

1689345632421-0

2 XREADGROUP

作用:消费者读取消息

msgs = r.xreadgroup(
"group1",
"consumer1",
{"mystream": ">"},
count=1,
block=2000
)

参数:

参数

含义

group1

消费者组

consumer1

消费者名称

>

只读取新消息

count

一次读取数量

block

阻塞等待时间

返回结构:

[
('mystream',
[
('1689345632421-0', {'field':'value'})
]
)
]

3 XACK

作用:确认消息处理完成

r.xack("mystream", "group1", msg_id)

Redis命令:

XACK mystream group1 msg_id

作用:

  • 从 Pending List 删除消息
  • 表示任务完成

四、Redis Stream 内部机制

核心结构:

Stream
|
|— Entry

消费者组结构:

Stream
|
|— Consumer Group
|
|— Consumer
|
|— Pending Entries List

Pending Entries List (PEL)

PEL 记录:已发送但未 ACK 的消息

结构:

Pending List
|
|— msg_id
|— consumer
|— idle_time
|— delivery_count

流程:

1 生产消息

XADD

消息进入 Stream。

2 消费消息

XREADGROUP

Redis会:1.返回消息  2.将消息放入 PEL

3 ACK

XACK

Redis会:

PEL 删除消息

为什么需要 ACK

解决:worker崩溃导致任务丢失

例子:

msg1 → worker1
worker1 crash

msg1 仍在:

Pending List

其他 worker 可以重新接管任务。

五、多个消费者的行为

同一消费者组

group1
|
|— worker1
|— worker2

消息分配:

msg1 → worker1
msg2 → worker2

特点:不会重复消费

不同消费者组

mystream

group1
group2

消息:msg1

消费情况:

group1 → msg1
group2 → msg1

特点:不同组会各自消费一遍消息

六、count 和 block 参数

count:一次读取多少条消息

count=10

适用:

任务类型

count

CPU任务

1~10

IO任务

10~50

日志流

50~200

block:阻塞等待新消息时间(毫秒)

block=5000

​没有消息时最多等待5秒

常见值:

场景

block

实时任务

1000

普通任务

5000

低频任务

10000

作用:避免 CPU 空轮询。

七、Redis Stream 的真实架构

Redis 不会执行任务。

Redis只是:消息队列

执行任务的是:Worker

典型架构:

API服务器
|
| XADD
|
Redis Stream
|
| XREADGROUP
|
Worker
|
执行任务

八、图片处理任务案例

用户上传图片
|
FastAPI
|
保存图片
|
XADD任务
|
Redis Stream
|
Worker消费
|
执行任务

示例任务:

task = {
"type": "compress",
"path": "/data/img1.png"
}

Worker执行:

compress_image(path)
save_to_db(path)

九、单进程 / Docker 场景如何运行 Worker

如果不能单独启动 worker 进程,可以:

方案1 FastAPI 后台任务

@app.on_event("startup")
async def startup():
asyncio.create_task(stream_worker())

FastAPI
|
|— HTTP API
|
|— Redis Worker

方案2 线程 Worker

threading.Thread(target=worker, daemon=True).start()

方案3 Docker 多进程

container
|
|— api
|— worker

方案4 Kubernetes

API Pod
Worker Pod

十、生产环境注意事项

1 consumer 名称必须唯一

consumer="worker"

consumer=f"worker-{uuid}"

2 Pending消息可能堆积

msg进入Pending

XAUTOCLAIM

回收任务。

3 Stream 会无限增长

需要定期:

XTRIM

或:

xadd(…, maxlen=100000)

4 多进程问题

如果使用:

uvicorn –workers 4

可能启动多个 worker。

解决:Redis 分布式锁

十一、Redis 队列实现方式对比

方式

命令

特点

List

LPUSH / BRPOP

简单

PubSub

PUBLISH / SUBSCRIBE

实时但不可靠

Stream

XADD / XREADGROUP

可靠队列

生产推荐:

Redis Stream

十二、核心理解

Redis Stream 本质:

分布式任务队列

完整系统:

Producer
|
Redis Stream
|
Consumer Group
|
Workers

应用场景:

  • 图片处理
  • 视频处理
  • AI任务
  • RAG构建
  • 日志处理
赞(0)
未经允许不得转载:171主机测评 » Redis Stream 学习笔记
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址