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 的情况下可以快速启动 Redisdocker 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, brokerredis_url, backendredis_url, include[app.tasks], )然后补充基础配置celery_app.conf.update( task_serializerjson, result_serializerjson, accept_content[json], timezoneAsia/Shanghai, enable_utcTrue, result_expires3600, task_track_startedTrue, )这些配置的作用包括使用 JSON 序列化任务参数只接受 JSON 格式的任务记录任务是否已经开始执行任务结果保存一小时统一处理任务时间。不建议使用能够反序列化任意 Python 对象的格式接收不可信数据否则可能带来安全风险。五、定义第一个异步任务在app/tasks.py中定义任务import time from app.celery_app import celery_app celery_app.task( bindTrue, namedocuments.process, ) def process_document( self, document_id: int, ): self.update_state( statePROGRESS, meta{ progress: 10, message: 开始处理文档, }, ) time.sleep(2) self.update_state( statePROGRESS, meta{ progress: 50, message: 正在分析内容, }, ) time.sleep(2) self.update_state( statePROGRESS, meta{ progress: 90, message: 正在保存结果, }, ) time.sleep(1) return { document_id: document_id, progress: 100, message: 处理完成, }使用bindTrue后任务函数的第一个参数是任务实例本身可以通过self.update_state()更新任务进度。这里使用time.sleep()模拟耗时操作。在真实项目中可以替换为文件解析、模型调用或数据处理逻辑。六、启动 Celery Worker在项目根目录执行celery \ -A app.celery_app.celery_app \ worker \ --loglevelinfoWorker 启动后会连接 Redis并等待新任务。开发环境也可以指定并发数量celery \ -A app.celery_app.celery_app \ worker \ --loglevelinfo \ --concurrency4--concurrency4表示 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_code202, ) 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, appcelery_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等待执行或者结果不存在STARTEDWorker 已经开始执行PROGRESS自定义的处理中状态SUCCESS执行成功FAILURE执行失败RETRY等待重新执行REVOKED任务已被撤销需要注意PENDING不一定表示任务还在排队。如果任务 ID 不存在、结果已经过期Celery 也可能返回PENDING。因此生产环境可以在数据库中单独保存任务记录不要完全依赖 Celery 的结果状态判断任务是否存在。九、增加自动重试机制外部接口超时、网络短暂中断等问题不应该直接导致整个任务永久失败。可以为任务增加自动重试class ExternalServiceError(Exception): pass celery_app.task( bindTrue, namedocuments.process_with_retry, autoretry_for(ExternalServiceError,), retry_backoffTrue, retry_backoff_max60, retry_jitterTrue, max_retries3, ) def process_with_retry( self, document_id: int, ): result call_external_service( document_id ) if not result: raise ExternalServiceError( 外部服务暂时不可用 ) return result这里使用了指数退避策略任务不会立即连续重试而是逐步增加等待时间。retry_jitterTrue会在重试时间中加入随机变化避免大量失败任务在同一时刻重新请求外部服务。并不是所有错误都适合重试网络超时可以重试临时服务错误可以重试请求频率受限可以延迟重试参数格式错误不应该重试用户无权限不应该重试数据本身不存在通常不应该重试。如果不区分错误类型重试机制可能把一次错误放大成多次无效请求。十、设置任务超时时间有些任务可能因为程序错误或外部服务无响应而长时间无法结束。可以设置软超时和硬超时from celery.exceptions import ( SoftTimeLimitExceeded, ) celery_app.task( bindTrue, soft_time_limit50, time_limit60, ) 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_idoperation_id, user_iduser_id, amountamount, )数据库还可以对operation_id建立唯一索引CREATE UNIQUE INDEX idx_operation_id ON payment_operations(operation_id);相比先查询再写入数据库唯一约束能够更可靠地阻止并发情况下的重复处理。十二、不要把大文件直接放进任务消息下面的做法并不推荐process_file.delay( file_binary_data )把完整文件或大段内容放入消息队列会带来以下问题Redis 内存占用增加消息传输变慢序列化和反序列化成本增加任务日志可能意外记录敏感内容Worker 获取任务时需要传输大量数据。更合理的方式是先把文件保存到对象存储或文件系统然后只传递文件 IDprocess_file.delay( file_id )Worker 根据file_id获取文件并执行处理。同样不建议直接传递数据库对象、连接对象或无法使用 JSON 序列化的复杂类型。十三、同言翻译中的异步任务应用对于实时性要求较高的功能系统通常需要快速返回结果但并不是所有操作都必须在当前请求中同步完成。以 同言翻译 为例实时翻译本身可以通过 WebSocket 或流式接口处理而会话结束后的摘要生成、历史记录整理、关键词提取、术语统计和文件导出等任务则可以交给 Celery 异步执行。例如用户结束一段会话后FastAPI 可以立即创建摘要任务app.post( /sessions/{session_id}/summary, status_code202, ) def create_session_summary( session_id: int, user_id: int Depends( get_current_user_id ), ): session get_user_session( user_iduser_id, session_idsession_id, ) if not session: raise HTTPException( status_code404, detail会话不存在, ) task generate_summary.delay( session_idsession_id, user_iduser_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], queuereports, )启动专门处理报表的 Workercelery \ -A app.celery_app.celery_app \ worker \ -Q reports \ --loglevelinfo队列隔离可以避免一个耗时任务占满全部 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, terminateTrue, )但强制终止正在执行的任务存在风险数据可能只写入了一部分临时文件可能没有清理数据库事务可能处于异常状态外部请求可能已经发送资源可能无法正常释放。更稳妥的方式是设计“协作式取消”。接口将任务状态标记为取消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 请求中持续等待更加稳定。

相关新闻