异步任务Celery高并发架构设计

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/RabbitMQCloud Tasks (GCP) / SQS (AWS)自建轻量队列(RQ/Dramatiq)
部署成本需要维护 Broker + Worker托管服务,几乎零运维比 Celery 轻量,仍需自托管
适用规模中大规模生产首选大规模生产首选小团队/原型阶段
优先级队列✅ 完善✅ SQS FIFO / Cloud Tasks 支持部分支持
延迟队列✅ Celery Countdown✅ SQS DelaySeconds有限支持
失败重试✅ 完善✅ 内置 DLQ基础支持
监控Flower / PromQLCloud 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) 比较各家平台的中转延迟和并发上限。

找到最适合你的 AI API 中转站

收录 125+ 服务商,按价格、模型、标签一键筛选

查看所有中转站 →