FastAPI: 后台任务与 Celery — 异步任务队列
最后更新:2026-08-26
BackgroundTasks 像餐厅叫号——点完菜立刻拿号(API 响应),菜做好了通知你;Celery 像中央厨房——多个厨师同时做不同订单,每个订单可追踪状态、失败重试。
1. 你将学到
- FastAPI 内置
BackgroundTasks:简单场景的快速方案 - Celery 架构:Worker / Broker (Redis) / Backend / Flower 监控
- Celery 与 FastAPI 集成:任务定义、触发、状态查询端点
- 任务重试与错误处理:
@app.task(retry=3, acks_late=True) - Alice 场景:批量爬取价格——API 立即返回任务 ID,Celery Worker 异步处理百万商品价格采集
2. Alice 的真实故事
(1) 痛点:耗时任务阻塞 API 响应
Alice 的 PriceTracker 需要批量爬取百万商品价格,单次爬取需要 30 分钟。如果用同步方式处理,API 请求会超时,Bob 前端等 30 分钟才收到响应,用户体验极差。而 FastAPI 的 BackgroundTasks 只能在当前进程内运行,Worker 重启任务就丢了。
(2) Celery 分布式队列的解法
Celery 将耗时任务放入消息队列(Redis Broker),Worker 进程异步消费,API 立即返回任务 ID。任务失败自动重试,Worker 可水平扩展,进程重启不影响队列中的任务。
PYTHON
from celery import Celery
celery_app = Celery("pricetracker", broker="redis://localhost:6379/0")
@celery_app.task(bind=True, max_retries=3)
def scrape_prices(self, product_ids: list[int]):
# Async price scraping - runs in Celery Worker
...
(3) 收益
百万商品爬取从"阻塞 API 30 分钟"变成"1 秒返回任务 ID + 后台处理"。Worker 可从 1 个扩展到 10 个,爬取时间从 30 分钟降到 3 分钟。任务失败自动重试 3 次,成功率从 95% 提升到 99.9%。
3. BackgroundTasks 轻量方案
(1) 适用场景
▶ 示例:BackgroundTasks 发送通知
PYTHON
from fastapi import FastAPI, BackgroundTasks
from pydantic import BaseModel
app = FastAPI()
class PriceAlertRequest(BaseModel):
product_id: int
target_price: float
email: str
def send_price_alert_email(email: str, product_id: int, price: float):
# Simulate email sending (do NOT use await here)
print(f"Sending alert to {email}: Product {product_id} hit ${price}")
@app.post("/alerts")
async def create_alert(alert: PriceAlertRequest, bg: BackgroundTasks):
# Add task to run after response is sent
bg.add_task(send_price_alert_email, alert.email, alert.product_id, alert.target_price)
return {"message": "Alert created", "product_id": alert.product_id}
输出:
TEXT
📖 仅展示
# 函数定义成功
(2) BackgroundTasks vs Celery 决策树
flowchart TD
Start{Need background task?} --> Time{Takes > 1 min?}
Time -->|No| Simple[Use BackgroundTasks]
Time -->|Yes| Retry{Need retry/resilience?}
Retry -->|No| Simple
Retry -->|Yes| Scale{Need horizontal scaling?}
Scale -->|No| Simple
Scale -->|Yes| Celery[Use Celery]
Simple -->|Pros| P1[Simple, no infra]
Simple -->|Cons| C1[No retry, no scale, lost on restart]
Celery -->|Pros| P2[Retry, scale, persistent, monitor]
Celery -->|Cons| C2[Redis + Worker infrastructure]
| 维度 | BackgroundTasks | Celery |
|---|---|---|
| 复杂度 | 零配置 | 需 Redis + Worker |
| 持久化 | 进程内存 | Redis 持久化 |
| 重试 | 无 | 内置重试机制 |
| 扩展 | 单进程 | Worker 可水平扩展 |
| 监控 | 无 | Flower Dashboard |
| 适用 | < 1 秒的轻量任务 | > 1 分钟的耗时任务 |
4. Celery 架构
(1) 全景架构
flowchart TD
API[FastAPI App] -->|Enqueue Task| Broker[(Redis Broker)]
Broker -->|Consume Task| Worker1[Celery Worker 1]
Broker -->|Consume Task| Worker2[Celery Worker 2]
Broker -->|Consume Task| WorkerN[Celery Worker N]
Worker1 -->|Store Result| Backend[(Redis Backend)]
Worker2 -->|Store Result| Backend
WorkerN -->|Store Result| Backend
API -->|Query Status| Backend
Flower[Flower Monitor] -->|Observe| Broker
Flower -->|Observe| Backend
| 组件 | 作用 | 推荐 |
|---|---|---|
| Broker | 任务消息队列 | Redis |
| Backend | 结果存储 | Redis |
| Worker | 任务执行进程 | celery -A app worker |
| Flower | 监控面板 | celery -A app flower |
▶ 示例:Celery 配置
PYTHON
# app/core/celery_app.py
from celery import Celery
celery_app = Celery(
"pricetracker",
broker="redis://localhost:6379/0",
backend="redis://localhost:6379/1",
)
celery_app.conf.update(
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
task_track_started=True,
task_acks_late=True, # Ack after execution, not before
worker_prefetch_multiplier=4,
result_expires=3600, # Results expire after 1 hour
)
输出:
TEXT
📖 仅展示
# 执行成功
5. Celery 任务定义与触发
(1) 任务定义
▶ 示例:价格爬取任务
PYTHON
# app/tasks/price_scraping.py
from app.core.celery_app import celery_app
import asyncio
from sqlalchemy import select
@celery_app.task(bind=True, max_retries=3, default_retry_delay=60)
def scrape_product_prices(self, product_ids: list[int]):
"""Scrape current prices for given products."""
try:
for pid in product_ids:
# Simulate scraping (in production: HTTP requests to price sources)
price = fetch_price_from_source(pid)
# Store in database
save_price_record(pid, price)
return {"scraped": len(product_ids), "status": "success"}
except Exception as exc:
# Retry with exponential backoff
raise self.retry(exc=exc, countdown=60 * (2 ** self.request.retries))
@celery_app.task(bind=True)
def bulk_scrape_all(self, total_products: int = 1000000, batch_size: int = 1000):
"""Scrape all products in batches - chord pattern."""
batches = [
list(range(i, min(i + batch_size, total_products)))
for i in range(0, total_products, batch_size)
]
# Fan out to individual scrape tasks
for batch in batches:
scrape_product_prices.delay(batch)
return {"total_batches": len(batches), "status": "started"}
输出:
TEXT
📖 仅展示
# 函数定义成功
(2) 任务状态流转
stateDiagram-v2
[*] --> PENDING: Task created
PENDING --> STARTED: Worker picks up
STARTED --> PROGRESS: Running (optional)
PROGRESS --> SUCCESS: Completed
PROGRESS --> FAILURE: Error occurred
STARTED --> FAILURE: Error occurred
FAILURE --> RETRY: max_retries not reached
RETRY --> PENDING: Re-queued
FAILURE --> [*]: max_retries exceeded
SUCCESS --> [*]
▶ 示例:FastAPI 端点触发 Celery 任务
PYTHON
from fastapi import FastAPI, Depends
from app.core.celery_app import celery_app
from app.tasks.price_scraping import scrape_product_prices, bulk_scrape_all
app = FastAPI()
@app.post("/api/v1/scrape/prices")
async def trigger_scrape(
product_ids: list[int],
user=Depends(require_subscription("pro")),
):
# Trigger Celery task - returns task ID immediately
task = scrape_product_prices.delay(product_ids)
return {"task_id": task.id, "status": "pending"}
@app.post("/api/v1/scrape/bulk")
async def trigger_bulk_scrape(
user=Depends(require_subscription("enterprise")),
):
task = bulk_scrape_all.delay(total_products=1000000)
return {"task_id": task.id, "status": "pending"}
输出:
TEXT
📖 仅展示
# 函数定义成功
▶ 示例:任务状态查询端点
PYTHON
from celery.result import AsyncResult
@app.get("/api/v1/tasks/{task_id}")
async def get_task_status(task_id: str):
result = AsyncResult(task_id, app=celery_app)
response = {
"task_id": task_id,
"status": result.status,
}
if result.ready():
if result.successful():
response["result"] = result.result
else:
response["error"] = str(result.result)
elif result.state == "PROGRESS":
response["progress"] = result.info
return response
输出:
TEXT
📖 仅展示
# 函数定义成功
6. Celery Worker 运行与监控
▶ 示例:启动 Worker 和 Flower
BASH
# Start Celery Worker
celery -A app.core.celery_app worker --loglevel=info --concurrency=4
# Start Flower monitoring dashboard
celery -A app.core.celery_app flower --port=5555
# Visit http://localhost:5555 for monitoring dashboard
输出:
TEXT
📖 仅展示
# 命令执行成功
▶ 示例:任务进度报告
PYTHON
from celery import current_task
@celery_app.task(bind=True)
def scrape_with_progress(self, product_ids: list[int]):
total = len(product_ids)
for i, pid in enumerate(product_ids):
# Process each product
price = fetch_price_from_source(pid)
save_price_record(pid, price)
# Report progress
self.update_state(
state="PROGRESS",
meta={"current": i + 1, "total": total, "percent": (i + 1) / total * 100},
)
return {"scraped": total}
输出:
TEXT
📖 仅展示
# 函数定义成功
❓ 常见问题
Q BackgroundTasks 的任务什么时候执行?
A 在响应发送给客户端之后执行。如果任务抛异常,不影响已发送的响应,但会在日志中记录。
Q Celery Worker 和 FastAPI 要在同一个进程吗?
A 不要。Worker 是独立进程,单独部署和扩展。FastAPI 只负责触发任务,Worker 负责执行。
Q Redis 做 Broker 和 Backend 有什么区别?
A Broker 是任务队列(待执行的任务),Backend 是结果存储(已完成的任务结果)。可以用同一个 Redis 不同 DB(如 DB 0 和 DB 1)。
Q task_acks_late=True 有什么用?
A 默认 Worker 收到任务就确认,如果执行中崩溃任务丢失。acks_late=True 在执行完成后才确认,崩溃后任务会被其他 Worker 重新执行。
Q Celery 任务里能用 async/await 吗?
A Celery 任务是同步函数。如需调用异步代码,用
asyncio.run() 包装。或在 Worker 中使用 eventlet/gevent 并发模式。Q 如何处理百万级任务的进度跟踪?
A 用 chord 模式:拆分为千个子任务(每个处理 1000 商品),子任务完成后 chord 回调汇总结果。前端轮询 chord 任务状态。
📖 小节
- BackgroundTasks 适合 < 1 秒的轻量任务(发邮件、写日志),零配置但无重试和持久化
- Celery 是生产级任务队列:Broker(Redis 队列)+ Worker(执行)+ Backend(结果存储)
@app.task(bind=True, max_retries=3)定义可重试任务,self.retry()触发重试- FastAPI 通过
.delay()触发任务、AsyncResult查询状态,API 立即返回任务 ID - Flower 提供可视化监控面板,实时查看 Worker 状态、任务进度、成功率
📝 作业
- 基础题(难度⭐):用 BackgroundTasks 实现一个端点,创建商品后后台发送通知邮件(模拟打印日志),API 立即返回创建结果。提示:
bg: BackgroundTasks+bg.add_task(fn, args) - 进阶题(难度⭐⭐):配置 Celery(Redis Broker + Backend),定义
scrape_prices任务(max_retries=3),FastAPI 端点触发任务并返回 task_id,另一个端点查询任务状态。提示:celery_app.delay()+AsyncResult(task_id) - 挑战题(难度⭐⭐⭐):实现带进度报告的批量爬取任务——
scrape_with_progress用self.update_state(state="PROGRESS")报告进度百分比,前端轮询/tasks/{task_id}展示进度条,Enterprise 用户可触发百万级爬取。提示:self.update_state()+result.info获取进度
---|