一、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构建
- 日志处理



