Skip to content

第 2 章 · 消息队列:Redis Queue 与 Celery Worker

本章目标:

  • 理解 Redis 队列原语(rpush/lpush/brpop)的工作机制
  • 掌握 Celery + Redis 的完整配置与运行方式
  • 学会将 LLM 调用等耗时任务异步化,实现非阻塞响应
  • 掌握 Celery 重试策略:retry_backoffmax_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 是阻塞版本:当队列为空时,客户端挂起等待,直到有新元素入队或超时。

python
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

一个常见需求是"优先级队列":将低优先级任务从普通队列搬到降级队列,等主队列清空后再处理。

python
# 将任务从 'high' 队列移到 'low' 队列(原子操作)
result = r.brpoplpush('high', 'low', timeout=5)

2.3 Celery 入门:配置与 Task 定义

Celery 是一个分布式任务队列框架,支持多种 Broker(Redis、RabbitMQ、Amazon SQS 等)。我们用 Redis 作为 Broker。

2.3.1 安装

bash
pip install -U celery[redis] redis

启动 Redis 服务(使用 Docker 最简单):

bash
docker run -d --name redis-queue -p 6379:6379 redis:7-alpine

2.3.2 定义 Celery App 和 Task

python
# 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

python
# 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

bash
# 在终端 1 启动 Worker(监听 redis 队列)
celery -A tasks worker --loglevel=info --concurrency=4

2.3.5 提交任务

python
# 同步提交,立即返回 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 获取结果。

python
# 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 端点:提交任务并立即返回

python
# 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)

python
@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 全局默认配置

python
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

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.0

docker-compose.yml

yaml
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]

运行

bash
# 启动基础设施
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": "..."}}

本章小结

  1. Redis 队列原语RPUSH/LPOP 实现 FIFO,BRPOP 实现阻塞消费,避免忙轮询;
  2. Celery 配置broker(任务来源)和 backend(结果存储)都指向 Redis 的不同数据库;
  3. Task 定义:用 @app.task 装饰函数,bind=True 可获得 self 以调用重试;
  4. 延迟提交task.delay(...) 返回 AsyncResult,不阻塞 FastAPI 请求;
  5. 轮询结果:通过 /chat/{task_id} 接口查询 AsyncResult.state,处理 PENDING/STARTED/SUCCESS/FAILURE 四种状态;
  6. 重试策略retry_backoff=True 启用指数退避,max_retries 限制最大尝试次数。

📌 更高级的异步交互模式(WebSocket 实时推送结果、Server-Sent Events)在 LangGraph ch08 中讲解,本章先掌握轮询模式即可。


🛠️ 动手实践

  1. 本地启动 Redis(docker run -d --name redis -p 6379:6379 redis:7-alpine),编写一个 Celery Task 完成「将用户输入写入 Redis List,再由另一个 Task 弹出并返回长度」的完整链路,验证 delay() + get() 流程。

  2. generate_agent_response Task 中加入模拟 RateLimitError(前 3 次调用抛异常,第 4 次成功),观察 Celery Worker 日志中的自动重试行为,记录每次重试的时间间隔,验证指数退避是否生效。

  3. 将本章的 FastAPI 端点改为返回 SSE 流式响应StreamingResponse),当 Worker 任务完成后一次性推送 JSON 结果,替代轮询接口。提示:使用 asyncio.to_thread 包装 result.get(timeout=...)


🧪 随堂测验

点击你认为正确的选项。答错时会展示正确答案与原因解析。

1. BRPOP 与 RPOP 的主要区别是什么?

2. Celery 的 retry_backoff=True 配合 bind=True 时,重试间隔如何变化?

3. AsyncResult.get() 阻塞等待结果时有什么风险?

4. 以下哪种场景最适合用 Celery 而非 FastAPI BackgroundTasks?