第 2 章 · 消息队列:Redis Queue 与 Celery Worker
本章目标:
- 理解 Redis 队列原语(
rpush/lpush/brpop)的工作机制- 掌握 Celery + Redis 的完整配置与运行方式
- 学会将 LLM 调用等耗时任务异步化,实现非阻塞响应
- 掌握 Celery 重试策略:
retry_backoff与max_retries- 能搭建一条完整链路:FastAPI 提交任务 → Celery Worker 处理 → 轮询结果
2.1 为什么需要消息队列
在前一章中我们用 FastAPI 的 BackgroundTasks 处理轻量后台任务(参见 FastAPI ch16)。但当任务需要:
- 跨进程执行(主服务崩溃不影响队列)
- 持久化(Worker 重启后任务不丢失)
- 重试与失败隔离(单个任务失败不影响其他任务)
- 水平扩展(多个 Worker 并行消费)
BackgroundTasks 就不够用了。消息队列是生产级 Agent 服务的必备基础设施。
客户端 ──HTTP──▶ FastAPI 服务 ──enqueue──▶ Redis Queue ──dequeue──▶ Celery Worker
(持久化、重试、并发)2.2 Redis 队列原语
Redis 是最常用的消息队列后端。它的核心队列操作只有三个:
| 命令 | 作用 | 阻塞版 |
|---|---|---|
RPUSH queue item | 从右侧推入元素(生产者) | — |
LPUSH queue item | 从左侧推入元素 | — |
RPOP queue | 从右侧弹出元素(消费者) | BRPOP |
LRANGE queue 0 -1 | 查看队列全部内容 | — |
2.2.1 阻塞弹出:BRPOP
RPOP 在队列为空时会立即返回 nil,导致消费者必须不断轮询(忙轮询浪费 CPU)。BRPOP 是阻塞版本:当队列为空时,客户端挂起等待,直到有新元素入队或超时。
import redis
r = redis.Redis(host='localhost', port=6379, db=0)
# 生产者:推入任务
r.rpush('tasks', 'task_001')
r.rpush('tasks', 'task_002')
# 消费者:阻塞等待(最多等 5 秒)
item = r.brpop('tasks', timeout=5)
# item 形如 (b'tasks', b'task_001')
if item:
queue_name, task_id = item
print(f'Got task: {task_id.decode()}')
else:
print('No tasks in queue')2.2.2 多队列广播:BRPOPLPUSH
一个常见需求是"优先级队列":将低优先级任务从普通队列搬到降级队列,等主队列清空后再处理。
# 将任务从 'high' 队列移到 'low' 队列(原子操作)
result = r.brpoplpush('high', 'low', timeout=5)2.3 Celery 入门:配置与 Task 定义
Celery 是一个分布式任务队列框架,支持多种 Broker(Redis、RabbitMQ、Amazon SQS 等)。我们用 Redis 作为 Broker。
2.3.1 安装
pip install -U celery[redis] redis启动 Redis 服务(使用 Docker 最简单):
docker run -d --name redis-queue -p 6379:6379 redis:7-alpine2.3.2 定义 Celery App 和 Task
# app.py
from celery import Celery
app = Celery(
'agent_worker',
broker='redis://localhost:6379/0',
backend='redis://localhost:6379/1',
)
# 设置时区与序列化
app.conf.update(
task_serializer='json',
result_serializer='json',
accept_content=['json'],
timezone='Asia/Shanghai',
enable_utc=True,
)2.3.3 定义一个简单 Task
# tasks.py
from app import app
import time
@app.task(bind=True, max_retries=3)
def call_llm(self, prompt: str, model: str = "gpt-4o-mini") -> dict:
"""模拟 LLM 调用(实际项目中替换为真实 API 调用)"""
try:
# 模拟 API 调用延迟
time.sleep(2)
return {
"model": model,
"response": f"这是针对 '{prompt}' 的模拟回复",
"tokens": 42,
}
except Exception as exc:
# 自动重试,指数退避
raise self.retry(exc=exc, countdown=2 ** self.request.retries)2.3.4 启动 Worker
# 在终端 1 启动 Worker(监听 redis 队列)
celery -A tasks worker --loglevel=info --concurrency=42.3.5 提交任务
# 同步提交,立即返回 AsyncResult
result = call_llm.delay("什么是大语言模型?", model="gpt-4o")
print(f"任务 ID: {result.id}")
# 阻塞等待结果(不推荐在生产中使用,会占用连接)
# response = result.get(timeout=30)
# 非阻塞:定期检查状态
import time
while not result.ready():
time.sleep(1)
print(f"完成!结果: {result.get()}")2.4 Agent 异步任务模式:LLM 调用入队
在 Agent 服务中,LLM 调用通常是最慢的环节(网络延迟 + 生成时间)。把它放入队列可以让 FastAPI 立即返回任务 ID,客户端通过轮询或 WebSocket 获取结果。
# agent_tasks.py
from app import app
from openai import AsyncOpenAI
import asyncio
client = AsyncOpenAI(api_key="sk-...")
@app.task(bind=True, max_retries=5, retry_backoff=True)
def generate_agent_response(
self,
messages: list[dict],
tools: list[dict] | None = None,
) -> dict:
"""
Agent 核心推理任务:
- 将多轮对话送入 LLM
- 解析 tool_calls 并执行(由 Worker 内部完成)
- 返回最终文本响应
"""
try:
# 调用 LLM(此处用 OpenAI SDK,可替换为任意 Provider)
response = asyncio.run(
client.chat.completions.create(
model="gpt-4o",
messages=messages,
tools=tools or [],
)
)
choice = response.choices[0]
return {
"role": choice.message.role,
"content": choice.message.content,
"tool_calls": [
{
"id": tc.id,
"function": {
"name": tc.function.name,
"arguments": tc.function.arguments,
},
}
for tc in choice.message.tool_calls or []
],
}
except Exception as exc:
# retry_backoff=True 时 Celery 自动指数退避重试
raise self.retry(exc=exc)2.4.1 FastAPI 端点:提交任务并立即返回
# main.py
from fastapi import FastAPI, HTTPException
from agent_tasks import generate_agent_response
app = FastAPI(title="Agent 异步服务")
@app.post("/chat")
async def chat_endpoint(prompt: str, model: str = "gpt-4o"):
"""
提交对话任务,立即返回 task_id。
客户端用 task_id 轮询结果。
"""
task = generate_agent_response.delay(
messages=[{"role": "user", "content": prompt}],
model=model,
)
return {"task_id": task.id, "status": "queued"}
@app.get("/chat/{task_id}")
async def get_result(task_id: str):
"""轮询任务结果"""
from celery.result import AsyncResult
result = AsyncResult(task_id)
if result.state == "PENDING":
return {"status": "pending"}
if result.state == "STARTED":
return {"status": "processing"}
if result.state == "FAILURE":
raise HTTPException(status_code=500, detail=str(result.result))
if result.state == "SUCCESS":
return {"status": "done", "result": result.result}
return {"status": result.state}2.5 重试策略详解
生产环境中 LLM API 经常遇到限流(429)或超时。Celery 提供两种重试机制:
2.5.1 任务级别重试(bind=True)
@app.task(bind=True, max_retries=5, retry_backoff=True)
def robust_llm_call(self, prompt: str):
try:
return do_api_call(prompt)
except RateLimitError:
# retry_backoff=True:第1次等1s,第2次等2s,第3次等4s...
raise self.retry(exc=self.__class__.exc, countdown=2 ** self.request.retries)
except ConnectionError:
# ConnectionError 立即重试,不等
raise self.retry(exc=self.__class__.exc, max_retries=10)2.5.2 全局默认配置
app.conf.update(
# 所有 task 默认最多重试 3 次
task_default_max_retries = 3,
# 启用指数退避(需要 bind=True)
task_default_retry_backoff = True,
# 退避基准(秒)
task_default_retry_backoff_base = 2,
# 任务超时(硬超时,Worker 内强制 kill)
task_soft_time_limit = 120, # 软超时:引发 SoftTimeLimitExceeded
task_time_limit = 180, # 硬超时:直接 kill Worker 子进程
)2.6 完整示例:FastAPI + Celery + Redis
创建一个可运行的 Agent 服务项目:
agent-service/
├── app.py # Celery App 配置
├── tasks.py # 所有异步任务
├── main.py # FastAPI 应用
├── requirements.txt
└── docker-compose.yml # Redis + Worker 一键启动requirements.txt
celery[redis]>=5.4.0
fastapi>=0.115.0
uvicorn[standard]>=0.32.0
redis>=5.0.0
openai>=1.50.0
pydantic>=2.0.0docker-compose.yml
version: "3.9"
services:
redis:
image: redis:7-alpine
ports: ["6379:6379"]
api:
build: .
command: uvicorn main:app --host 0.0.0.0 --port 8000
ports: ["8000:8000"]
environment:
- REDIS_URL=redis://redis:6379/0
- OPENAI_API_KEY=${OPENAI_API_KEY}
depends_on: [redis]
worker:
build: .
command: celery -A tasks worker --loglevel=info --concurrency=4
environment:
- REDIS_URL=redis://redis:6379/0
- OPENAI_API_KEY=${OPENAI_API_KEY}
depends_on: [redis]运行
# 启动基础设施
docker compose up -d redis
# 终端 1:启动 Worker
celery -A tasks worker --loglevel=info --broker=redis://localhost:6379/0
# 终端 2:启动 API
uvicorn main:app --reload
# 终端 3:测试
curl -X POST http://localhost:8000/chat \
-H "Content-Type: application/json" \
-d '{"prompt": "解释量子计算"}'
# → {"task_id": "xxx", "status": "queued"}
curl http://localhost:8000/chat/xxx
# → {"status": "done", "result": {"role": "assistant", "content": "..."}}本章小结
- Redis 队列原语:
RPUSH/LPOP实现 FIFO,BRPOP实现阻塞消费,避免忙轮询; - Celery 配置:
broker(任务来源)和backend(结果存储)都指向 Redis 的不同数据库; - Task 定义:用
@app.task装饰函数,bind=True可获得self以调用重试; - 延迟提交:
task.delay(...)返回AsyncResult,不阻塞 FastAPI 请求; - 轮询结果:通过
/chat/{task_id}接口查询AsyncResult.state,处理 PENDING/STARTED/SUCCESS/FAILURE 四种状态; - 重试策略:
retry_backoff=True启用指数退避,max_retries限制最大尝试次数。
📌 更高级的异步交互模式(WebSocket 实时推送结果、Server-Sent Events)在 LangGraph ch08 中讲解,本章先掌握轮询模式即可。
🛠️ 动手实践
本地启动 Redis(
docker run -d --name redis -p 6379:6379 redis:7-alpine),编写一个 Celery Task 完成「将用户输入写入 Redis List,再由另一个 Task 弹出并返回长度」的完整链路,验证delay()+get()流程。在
generate_agent_responseTask 中加入模拟RateLimitError(前 3 次调用抛异常,第 4 次成功),观察 Celery Worker 日志中的自动重试行为,记录每次重试的时间间隔,验证指数退避是否生效。将本章的 FastAPI 端点改为返回 SSE 流式响应(
StreamingResponse),当 Worker 任务完成后一次性推送 JSON 结果,替代轮询接口。提示:使用asyncio.to_thread包装result.get(timeout=...)。