AI 应用异步任务队列实战:用 Celery / Cloud Tasks 处理长任务与高并发
2026-09-07 · 约 11 分钟阅读
---
title: "AI 应用异步任务队列实战:用 Celery / Cloud Tasks 处理长任务与高并发"
description: "AI应用异步任务队列实战:详解Celery/Cloud Tasks/SQS三种主流方案,覆盖长任务(视频生成/批量推理)、高并发削峰、失败重试与优先级队列。附可直接复用的Python/Node代码示例、生产级Worker配置与监控告警设置,让你的AI应用不被慢请求拖垮,适合后端工程师与SRE,含中转站实测数据,附2026年最新API价格,含可直接复用的代码片段,覆盖中转站稳定性对比,是开发者必读的实战指南。"
date: "2026-09-07"
tags: ["异步任务", "Celery", "高并发", "架构设计"]
---
# AI 应用异步任务队列实战:用 Celery / Cloud Tasks 处理长任务与高并发
用户提交了一个"批量分析 1 万条评论"的需求,前端按钮转了 90 秒没响应,最后超时 504——这是典型的"把同步接口当成万能锤子"的反模式。AI 应用里有大量任务是 慢、贵、不确定时长 的(视频生成、批量 embedding、长文档总结、多步 Agent 推理),这些任务必须用异步队列和前端解耦。
一、为什么 AI 应用特别需要异步队列
传统 Web 请求一般几百毫秒返回,AI 应用有几类典型"慢任务":
| 任务类型 | 典型耗时 | 为什么必须异步 |
|---|---|---|
| 视频生成(Sora 2 / Seedance 2.5) | 1-10 分钟 | 单次生成就要等几十秒到几分钟 |
| 批量文档总结(100 篇 PDF) | 5-30 分钟 | 串行跑要半小时起 |
| 长文档多轮问答(百万字上下文) | 30 秒-3 分钟 | 用户等不了 |
| 大批量 embedding(10 万条文本) | 10-60 分钟 | 限速下必须排队 |
| 多步 Agent 任务(订机票+写报告) | 2-15 分钟 | 多轮工具调用自然就慢 |
这些任务的共同特征是:用户提交后不需要立即拿到结果,可以在稍后查询进度或接收通知。这正是异步队列的甜蜜点。
二、什么时候用同步、什么时候用异步
简单的决策树:
```
用户请求
├─ 延迟 < 5 秒 ────→ 同步 HTTP 响应
└─ 延迟 ≥ 5 秒
├─ 用户在线等待 ──→ SSE 流式响应(不阻塞连接)
└─ 用户可以离开 ──→ 异步任务队列 + 回调/WebSocket 通知
```
流式响应(SSE)和异步队列不是替代关系,是互补关系:
- 流式:用户在页面上看着 AI 一点点打字,适合 5-60 秒的对话类任务
- 异步队列:用户提交后关掉页面都行,适合 1 分钟以上的离线任务
三、三种主流方案对比
| 维度 | Celery + Redis/RabbitMQ | Cloud Tasks (GCP) / SQS (AWS) | 自建轻量队列(RQ/Dramatiq) |
|---|---|---|---|
| 部署成本 | 需要维护 Broker + Worker | 托管服务,几乎零运维 | 比 Celery 轻量,仍需自托管 |
| 适用规模 | 中大规模生产首选 | 大规模生产首选 | 小团队/原型阶段 |
| 优先级队列 | ✅ 完善 | ✅ SQS FIFO / Cloud Tasks 支持 | 部分支持 |
| 延迟队列 | ✅ Celery Countdown | ✅ SQS DelaySeconds | 有限支持 |
| 失败重试 | ✅ 完善 | ✅ 内置 DLQ | 基础支持 |
| 监控 | Flower / PromQL | Cloud Console / CloudWatch | 弱 |
推荐:小团队起步用 Celery + Redis,上规模后切到 Cloud Tasks/SQS 这类托管服务省运维。
四、Celery 实战:批量文档总结
一个典型的"提交一批 PDF → 异步处理 → 完成后通知用户"流程:
```python
# tasks.py
from celery import Celery
import time
app = Celery('ai_tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/1')
# 队列路由:不同任务分到不同 worker 池
app.conf.task_routes = {
'tasks.summarize_doc': {'queue': 'summary'},
'tasks.generate_video': {'queue': 'video', 'priority': 5},
}
@app.task(bind=True, max_retries=3, default_retry_delay=30)
def summarize_doc(self, doc_url: str, user_id: str):
"""异步总结一个文档"""
try:
# 1. 下载文档
content = download_pdf(doc_url)
# 2. 调 LLM(用中转站接入多模型)
from openai import OpenAI
client = OpenAI(
api_key="sk-xxx",
base_url="https://你的中转站域名/v1"
)
response = client.chat.completions.create(
model="gpt-5.6-sol",
messages=[
{"role": "system", "content": "你是文档总结助手,输出 300 字摘要"},
{"role": "user", "content": content[:50000]} # 截断避免超限
],
timeout=120,
)
summary = response.choices[0].message.content
# 3. 持久化结果 + 通知用户
save_result(user_id, doc_url, summary)
notify_user(user_id, f"文档总结已完成:{doc_url}")
return {"status": "ok", "summary": summary}
except Exception as exc:
# 自动重试,3 次后进 DLQ
raise self.retry(exc=exc)
```
提交任务的 API:
```python
# api.py(FastAPI)
from fastapi import FastAPI
from tasks import summarize_doc
app = FastAPI()
@app.post("/api/summarize")
async def submit_summarize(doc_url: str, user_id: str):
# 立即返回 task_id,不等待实际处理
task = summarize_doc.delay(doc_url, user_id)
return {"task_id": task.id, "status": "queued"}
```
Worker 启动:
```bash
# 启动 worker(生产建议用 supervisor 或 systemd)
celery -A tasks worker -Q summary -l info --concurrency=4
# 视频生成用单独的 worker(更慢、需要更长超时)
celery -A tasks worker -Q video -l info --concurrency=2 --soft-time-limit=900
```
五、生产级 Worker 配置
直接 `celery worker` 跑生产会翻车,务必配置这几项:
```python
# celery_config.py
app.conf.update(
# 任务超时(视频生成要长一些)
task_time_limit=300, # 硬超时 5 分钟
task_soft_time_limit=240, # 软超时 4 分钟(抛 SoftTimeLimitExceeded)
# 限速(避免打爆上游 LLM API)
task_annotations={
'tasks.summarize_doc': {'rate_limit': '10/m'}, # 每分钟 10 个
'tasks.generate_video': {'rate_limit': '2/m'}, # 视频更慢
},
# 失败处理
task_acks_late=True, # 任务执行完才 ack,崩溃不丢任务
worker_prefetch_multiplier=1, # 公平调度
worker_max_tasks_per_child=100, # 每处理 100 个任务重启,防止内存泄漏
)
```
六、进度查询与回调
异步任务对用户来说不能是"黑盒",必须能查进度:
```python
@app.get("/api/task/{task_id}")
async def get_task_status(task_id: str):
task = summarize_doc.AsyncResult(task_id)
return {
"task_id": task_id,
"status": task.status, # PENDING / STARTED / SUCCESS / FAILURE
"result": task.result if task.ready() else None,
"progress": task.info.get("progress", 0) if task.info else 0,
}
# 任务内部更新进度(通过 self.update_state)
@app.task(bind=True)
def summarize_doc(self, doc_url, user_id):
self.update_state(state="PROGRESS", meta={"progress": 0, "step": "下载文档"})
content = download_pdf(doc_url)
self.update_state(state="PROGRESS", meta={"progress": 30, "step": "调 LLM"})
summary = call_llm(content)
self.update_state(state="PROGRESS", meta={"progress": 80, "step": "保存结果"})
save_result(user_id, doc_url, summary)
return {"status": "ok"}
```
任务完成时的通知方式选择:
| 方式 | 适用场景 |
|---|---|
| WebSocket | 用户在线,需要实时推送进度 |
| Webhook 回调 | 客户端是另一个服务(非浏览器) |
| 邮件 / 短信 | 长任务,用户可能离场 |
| 站内消息 / 通知中心 | 移动端 App 场景 |
七、托管方案:Cloud Tasks / SQS
当 Celery 自托管成为运维负担时,切到云厂商托管:
```python
# Cloud Tasks(GCP)
from google.cloud import tasks_v2
client = tasks_v2.CloudTasksClient()
parent = client.queue_path("your-project", "us-central1", "ai-tasks")
task = {
"app_engine_http_request": {
"http_method": "POST",
"relative_uri": "/worker/summarize",
"body": json.dumps({"doc_url": doc_url, "user_id": user_id}).encode(),
},
"schedule_time": {"seconds": int(time.time()) + 60}, # 延迟 60 秒
}
response = client.create_task(parent=parent, task=task)
```
托管方案优势:自动扩缩容、DLQ、可视化监控、零运维,代价是要付一点服务费但比自己搭 Redis + Worker 服务器便宜得多。
八、常见反模式
- ❌ 长任务走同步 HTTP:前端必然超时,浪费连接资源
- ❌ 任务队列没有优先级:付费用户和免费用户混在一起排队
- ❌ 没有 DLQ:失败任务默默消失,永远找不到原因
- ❌ Worker 没有限速:任务一多就把上游 LLM API 打到 429
- ❌ 结果只存内存:Worker 重启结果就丢,要存 Redis 或 DB
九、总结
异步任务队列是 AI 应用从"能跑"到"能扛流量"的必经一步。短任务用 SSE 流式、长任务用 Celery/Cloud Tasks,配合 Worker 限速、自动重试、DLQ、进度查询,才能撑住生产级负载。
想了解流式响应的实现细节,可以看 [AI API 流式响应 SSE 实战指南](/blog/ai-api-streaming-sse-implementation-guide);要选稳定的 LLM API 来支撑异步 Worker,可以去 [openairouter.net](https://openairouter.net) 比较各家平台的中转延迟和并发上限。