FastAPI + Celery 实战:用任务队列处理耗时操作
在 Web 项目中,有些任务无法在几百毫秒内完成,例如:
批量处理文件;
生成数据报表;
发送邮件或短信;
调用 AI 模型;
处理音频和视频;
执行大量数据计算。
如果直接在 HTTP 接口中执行这些任务,用户就必须一直等待。任务耗时过长时,还可能触发网关超时,甚至占满 Web 服务的工作进程。
一种常见的解决方案是引入任务队列:接口只负责接收请求并创建任务,真正的耗时操作交给后台 Worker 执行。
本文将使用 FastAPI、Celery 和 Redis,实现一个简单的异步任务系统,并介绍任务状态查询、失败重试、超时控制和幂等性等实际问题。
一、同步接口存在哪些问题?
假设系统有一个文档处理接口:
import time from fastapi import FastAPI app = FastAPI() @app.post("/documents/{document_id}/process") def process_document(document_id: int): time.sleep(20) return { "document_id": document_id, "status": "completed", }客户端调用接口后,需要等待 20 秒才能收到响应。
这种写法存在几个明显问题:
用户等待时间过长;
请求可能被网关提前关闭;
Web 服务进程长时间被占用;
任务失败后不容易自动重试;
服务重启可能导致正在执行的任务丢失;
无法方便地控制任务并发数量。
更合理的方式是让接口立即返回任务 ID:
{ "task_id": "b7e8c1b8-xxxx-xxxx-xxxx-xxxxxxxxxxxx", "status": "queued" }客户端之后通过任务 ID 查询处理进度。
二、任务队列的基本架构
一个基础的异步任务系统通常包含四个部分:
客户端 ↓ FastAPI 接口 ↓ 消息代理 Redis ↓ Celery Worker ↓ 执行耗时任务各部分的职责如下:
FastAPI:接收请求并创建任务;
Redis:保存等待执行的任务消息;
Celery Worker:从队列取出任务并执行;
Result Backend:保存任务状态和执行结果。
FastAPI 和 Celery Worker 是两个独立进程。即使某个任务需要运行几十秒,也不会持续占用原来的 HTTP 请求。
三、安装依赖
安装 FastAPI、Celery 和 Redis 相关依赖:
pip install fastapi uvicorn "celery[redis]"本地已经安装 Docker 的情况下,可以快速启动 Redis:
docker run \ --name celery-redis \ -p 6379:6379 \ -d redis:7示例项目结构如下:
async-task-demo/ ├── app/ │ ├── __init__.py │ ├── celery_app.py │ ├── main.py │ └── tasks.py └── requirements.txt四、创建 Celery 应用
在app/celery_app.py中创建 Celery 实例:
import os from celery import Celery redis_url = os.getenv( "CELERY_REDIS_URL", "redis://localhost:6379/0", ) celery_app = Celery( "async_tasks", broker=redis_url, backend=redis_url, include=["app.tasks"], )然后补充基础配置:
celery_app.conf.update( task_serializer="json", result_serializer="json", accept_content=["json"], timezone="Asia/Shanghai", enable_utc=True, result_expires=3600, task_track_started=True, )这些配置的作用包括:
使用 JSON 序列化任务参数;
只接受 JSON 格式的任务;
记录任务是否已经开始执行;
任务结果保存一小时;
统一处理任务时间。
不建议使用能够反序列化任意 Python 对象的格式接收不可信数据,否则可能带来安全风险。
五、定义第一个异步任务
在app/tasks.py中定义任务:
import time from app.celery_app import celery_app @celery_app.task( bind=True, name="documents.process", ) def process_document( self, document_id: int, ): self.update_state( state="PROGRESS", meta={ "progress": 10, "message": "开始处理文档", }, ) time.sleep(2) self.update_state( state="PROGRESS", meta={ "progress": 50, "message": "正在分析内容", }, ) time.sleep(2) self.update_state( state="PROGRESS", meta={ "progress": 90, "message": "正在保存结果", }, ) time.sleep(1) return { "document_id": document_id, "progress": 100, "message": "处理完成", }使用bind=True后,任务函数的第一个参数是任务实例本身,可以通过self.update_state()更新任务进度。
这里使用time.sleep()模拟耗时操作。在真实项目中,可以替换为文件解析、模型调用或数据处理逻辑。
六、启动 Celery Worker
在项目根目录执行:
celery \ -A app.celery_app.celery_app \ worker \ --loglevel=infoWorker 启动后会连接 Redis,并等待新任务。
开发环境也可以指定并发数量:
celery \ -A app.celery_app.celery_app \ worker \ --loglevel=info \ --concurrency=4--concurrency=4表示 Worker 最多可以同时运行四个任务。
并发数并不是越高越好。如果任务会占用大量内存、GPU 或外部接口配额,过高的并发反而可能导致系统不稳定。
七、通过 FastAPI 创建任务
在app/main.py中创建接口:
from fastapi import FastAPI from pydantic import BaseModel from app.tasks import process_document app = FastAPI() class TaskRequest(BaseModel): document_id: int @app.post( "/tasks", status_code=202, ) def create_task(request: TaskRequest): task = process_document.delay( request.document_id ) return { "task_id": task.id, "status": "queued", }delay()不会直接执行任务,而是将任务发送到 Redis。
接口使用202 Accepted状态码,表示服务器已经接受请求,但任务尚未完成。
请求示例:
curl \ -X POST \ -H "Content-Type: application/json" \ -d '{"document_id": 1001}' \ http://localhost:8000/tasks响应示例:
{ "task_id": "9c402926-xxxx-xxxx-xxxx-xxxxxxxxxxxx", "status": "queued" }八、查询任务状态
客户端拿到任务 ID 后,可以定期查询状态:
from celery.result import AsyncResult from app.celery_app import celery_app @app.get("/tasks/{task_id}") def get_task_status(task_id: str): task = AsyncResult( task_id, app=celery_app, ) response = { "task_id": task_id, "status": task.state, } if task.state == "PROGRESS": response["progress"] = task.info elif task.state == "SUCCESS": response["result"] = task.result elif task.state == "FAILURE": response["error"] = str(task.info) return responseCelery 常见状态包括:
| 状态 | 含义 |
|---|---|
PENDING | 等待执行,或者结果不存在 |
STARTED | Worker 已经开始执行 |
PROGRESS | 自定义的处理中状态 |
SUCCESS | 执行成功 |
FAILURE | 执行失败 |
RETRY | 等待重新执行 |
REVOKED | 任务已被撤销 |
需要注意,PENDING不一定表示任务还在排队。如果任务 ID 不存在、结果已经过期,Celery 也可能返回PENDING。
因此,生产环境可以在数据库中单独保存任务记录,不要完全依赖 Celery 的结果状态判断任务是否存在。
九、增加自动重试机制
外部接口超时、网络短暂中断等问题,不应该直接导致整个任务永久失败。
可以为任务增加自动重试:
class ExternalServiceError(Exception): pass @celery_app.task( bind=True, name="documents.process_with_retry", autoretry_for=(ExternalServiceError,), retry_backoff=True, retry_backoff_max=60, retry_jitter=True, max_retries=3, ) def process_with_retry( self, document_id: int, ): result = call_external_service( document_id ) if not result: raise ExternalServiceError( "外部服务暂时不可用" ) return result这里使用了指数退避策略,任务不会立即连续重试,而是逐步增加等待时间。
retry_jitter=True会在重试时间中加入随机变化,避免大量失败任务在同一时刻重新请求外部服务。
并不是所有错误都适合重试:
网络超时可以重试;
临时服务错误可以重试;
请求频率受限可以延迟重试;
参数格式错误不应该重试;
用户无权限不应该重试;
数据本身不存在通常不应该重试。
如果不区分错误类型,重试机制可能把一次错误放大成多次无效请求。
十、设置任务超时时间
有些任务可能因为程序错误或外部服务无响应而长时间无法结束。
可以设置软超时和硬超时:
from celery.exceptions import ( SoftTimeLimitExceeded, ) @celery_app.task( bind=True, soft_time_limit=50, time_limit=60, ) def process_with_timeout( self, document_id: int, ): try: return run_long_task( document_id ) except SoftTimeLimitExceeded: clean_temporary_files( document_id ) raise两种超时的区别是:
soft_time_limit:触发异常,允许任务清理资源;time_limit:超过时间后强制终止任务。
硬超时应该略大于软超时,为任务释放文件、连接和临时资源留出时间。
任务内部调用外部 API 时,仍然需要给网络请求单独设置超时。Celery 的任务超时不能替代 HTTP 客户端的连接和读取超时。
十一、任务幂等性为什么重要?
任务队列通常采用“至少投递一次”的处理思路。在网络异常、Worker 崩溃或确认消息失败时,同一个任务可能被执行多次。
例如,一个任务负责给用户账户增加 100 元:
def add_balance(user_id): balance = get_balance(user_id) update_balance(user_id, balance + 100)如果任务重复执行,用户余额就会被错误增加多次。
因此,重要任务需要具备幂等性:相同任务执行一次或执行多次,最终结果应该保持一致。
可以为每次业务操作生成唯一编号:
def process_payment( operation_id: str, user_id: int, amount: float, ): if operation_exists(operation_id): return get_operation_result( operation_id ) return create_payment_operation( operation_id=operation_id, user_id=user_id, amount=amount, )数据库还可以对operation_id建立唯一索引:
CREATE UNIQUE INDEX idx_operation_id ON payment_operations(operation_id);相比先查询再写入,数据库唯一约束能够更可靠地阻止并发情况下的重复处理。
十二、不要把大文件直接放进任务消息
下面的做法并不推荐:
process_file.delay( file_binary_data )把完整文件或大段内容放入消息队列,会带来以下问题:
Redis 内存占用增加;
消息传输变慢;
序列化和反序列化成本增加;
任务日志可能意外记录敏感内容;
Worker 获取任务时需要传输大量数据。
更合理的方式是先把文件保存到对象存储或文件系统,然后只传递文件 ID:
process_file.delay( file_id )Worker 根据file_id获取文件并执行处理。
同样,不建议直接传递数据库对象、连接对象或无法使用 JSON 序列化的复杂类型。
十三、同言翻译中的异步任务应用
对于实时性要求较高的功能,系统通常需要快速返回结果;但并不是所有操作都必须在当前请求中同步完成。
以 同言翻译 为例,实时翻译本身可以通过 WebSocket 或流式接口处理,而会话结束后的摘要生成、历史记录整理、关键词提取、术语统计和文件导出等任务,则可以交给 Celery 异步执行。
例如,用户结束一段会话后,FastAPI 可以立即创建摘要任务:
@app.post( "/sessions/{session_id}/summary", status_code=202, ) def create_session_summary( session_id: int, user_id: int = Depends( get_current_user_id ), ): session = get_user_session( user_id=user_id, session_id=session_id, ) if not session: raise HTTPException( status_code=404, detail="会话不存在", ) task = generate_summary.delay( session_id=session_id, user_id=user_id, ) return { "task_id": task.id, "status": "queued", }Worker 完成处理后,可以把结果保存到数据库,再通过轮询、WebSocket 或系统通知告知用户。
对于同言翻译这类可能涉及语音、原文和译文的应用,不建议把完整会话内容直接写入 Redis 消息。更安全的方式是只传递session_id和user_id,由 Worker 在验证数据归属后读取必要内容。
任务执行完成后,还应该及时清理临时音频、缓存文件和不再需要的中间数据,避免敏感信息被长期保留。
十四、如何划分不同任务队列?
当系统任务类型较多时,可以使用不同队列进行隔离。
例如:
default 普通任务 high_priority 高优先级任务 documents 文件处理任务 reports 报表生成任务 notifications 通知任务任务可以指定队列:
task = generate_report.apply_async( args=[report_id], queue="reports", )启动专门处理报表的 Worker:
celery \ -A app.celery_app.celery_app \ worker \ -Q reports \ --loglevel=info队列隔离可以避免一个耗时任务占满全部 Worker。
例如,大量报表任务不应该阻塞登录通知或高优先级业务任务。不同队列还可以设置不同的并发数量和服务器资源。
十五、任务完成后如何通知前端?
客户端获取任务结果通常有三种方式。
1. 定时轮询
客户端每隔几秒查询一次任务状态:
const timer = setInterval(async () => { const response = await fetch( `/tasks/${taskId}` ); const task = await response.json(); if ( task.status === "SUCCESS" || task.status === "FAILURE" ) { clearInterval(timer); } }, 2000);轮询实现简单,适合任务数量较少的系统。
2. WebSocket 推送
客户端与服务器保持 WebSocket 连接。任务完成后,服务器主动推送状态。
这种方式实时性更好,但需要管理连接、重连和消息路由。
3. Webhook 回调
如果任务由另一个系统提交,可以在任务完成后调用对方提供的回调地址。
使用 Webhook 时需要进行签名验证,并防止攻击者伪造回调请求。
十六、任务撤销需要注意什么?
Celery 可以撤销尚未开始的任务:
celery_app.control.revoke( task_id )如果任务已经开始,可以请求终止:
celery_app.control.revoke( task_id, terminate=True, )但强制终止正在执行的任务存在风险:
数据可能只写入了一部分;
临时文件可能没有清理;
数据库事务可能处于异常状态;
外部请求可能已经发送;
资源可能无法正常释放。
更稳妥的方式是设计“协作式取消”。
接口将任务状态标记为取消,Worker 在不同处理阶段主动检查:
def check_task_cancelled(task_id): if is_cancelled(task_id): raise TaskCancelledError() def process_large_document( task_id, document_id, ): check_task_cancelled(task_id) load_document(document_id) check_task_cancelled(task_id) analyze_document(document_id) check_task_cancelled(task_id) save_result(document_id)这样可以在安全位置停止任务,并执行必要的清理操作。
十七、生产环境监控哪些指标?
任务系统上线后,建议重点监控:
队列中等待任务的数量;
任务平均等待时间;
任务平均执行时间;
成功率和失败率;
重试次数;
超时任务数量;
Worker 在线数量;
Redis 内存和连接状态;
不同任务类型的资源消耗。
如果队列长度持续增加,通常说明任务产生速度高于 Worker 处理速度。
此时不应该只考虑增加 Worker,还需要检查:
是否出现大量重复任务;
外部服务是否变慢;
任务代码是否存在性能问题;
是否可以合并批量操作;
是否需要对入口进行限流;
是否应该增加任务优先级和队列隔离。
十八、常见误区
误区一:使用 Celery 后接口一定更快
Celery 只能把耗时工作移出当前请求。任务本身的处理速度并不会自动提高。
误区二:任务发送成功就等于业务成功
接口把消息写入 Redis,只能说明任务已经进入队列,不能说明任务已经完成。
误区三:失败任务应该无限重试
无限重试会持续消耗资源。应该设置最大重试次数,并将最终失败的任务记录下来。
误区四:Redis 可以永久保存任务结果
任务结果应该根据业务需要保存到数据库或对象存储。Redis 更适合保存短期状态和缓存数据。
误区五:增加 Worker 数量可以解决所有积压
如果瓶颈是数据库、第三方 API 或 GPU,增加 Worker 可能让下游服务更快达到极限。
十九、总结
FastAPI、Celery 和 Redis 可以组成一套简单实用的异步任务系统:
FastAPI 接收请求 ↓ Celery 创建任务 ↓ Redis 保存消息 ↓ Worker 执行任务 ↓ 客户端查询或接收结果从演示代码走向生产环境,还需要重点考虑:
自动重试;
超时控制;
任务幂等性;
敏感数据保护;
队列隔离;
任务撤销;
失败补偿;
状态持久化;
系统监控。
任务队列的价值不仅是让接口更快返回,还可以把 Web 请求和耗时处理解耦,让两部分独立扩容、独立失败和独立恢复。
对于执行时间较长、允许稍后完成的业务,异步任务通常比在 HTTP 请求中持续等待更加稳定。
