news 2026/7/29 8:11:34

FastAPI + Celery 实战:用任务队列处理耗时操作

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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=info

Worker 启动后会连接 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 response

Celery 常见状态包括:

状态含义
PENDING等待执行,或者结果不存在
STARTEDWorker 已经开始执行
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_iduser_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 请求中持续等待更加稳定。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/7/29 8:08:51

电容与电感工程应用全解析:从基础原理到选型设计实战

1. 项目概述:从“元件”到“伙伴”的认知升级“电容和电感(自总结)”,这个标题一看就是同行在项目复盘或学习沉淀后的产物。它不像教科书目录那样冰冷,更像是一个工程师在调试完一块电源板、或设计完一个滤波电路后&am…

作者头像 李华
网站建设 2026/7/29 8:08:08

第六章系统的配置与性能

软考高级系统架构设计师备考,写着方便自己看,如果有发现不对或少得的地方可以多多指正,谢谢各位兄弟。一、核心性能指标(按应用场景分类)所有指标不需要死记硬背,结合场景理解记忆效率更高:场景…

作者头像 李华
网站建设 2026/7/29 8:07:44

冒泡排序--附图解以及代码

冒泡排序(Bubble Sort)详解 1. 什么是冒泡排序? 冒泡排序是一种基于比较的交换排序算法。 它的核心思想是:重复遍历待排序序列,依次比较相邻的两个元素,若顺序错误(逆序)则交换&…

作者头像 李华
网站建设 2026/7/29 8:05:04

AI会议系统三重防护:提升企业级应用可靠性

1. 项目概述:AI信任危机与企业级解决方案 上周科技圈最热门的新闻莫过于某AI公司公开承诺"出错赔10万"的营销事件,这个看似激进的商业策略背后,折射出当前企业市场对AI技术最核心的焦虑——信任缺失。作为专注智能会议领域的技术从…

作者头像 李华
网站建设 2026/7/29 8:04:56

匿名管道--任务派发程序

整体代码逻辑架构图如下&#xff1a;整体架构图如下&#xff1a;整体代码如下:main.cc:#include "processpool.hpp" #include "Tasks.hpp" int main(int argc, char *argv[]) {if(argc!2){cout<<"Usage: "<<argv[0]<<" &…

作者头像 李华
网站建设 2026/7/29 8:04:01

MANUS Metagloves Pro,让指尖动作精准走进数字世界

在元宇宙、人形机器人、数字孪生、虚拟仿真技术飞速迭代的当下&#xff0c;手部作为人体最灵活、最精细的交互载体&#xff0c;其数字化复刻能力直接决定了沉浸式交互、智能操控、数字内容创作的真实质感。传统动作捕捉设备受限于光学遮挡、数据漂移、精度不足等问题&#xff0…

作者头像 李华