在gunicorn/uvicorn运行的FastAPI中实现非阻塞后台长任务
解决方案
针对你的需求,这里提供两种安全可靠的实现方式,分别适配轻量场景和生产环境:
一、轻量方案:基于asyncio后台任务(不依赖第三方服务)
利用asyncio.create_task启动后台协程监控子进程,既不阻塞主事件循环,又能在任务完成后自动处理结果。
1. 定义进程结果处理逻辑
async def handle_process_result(process: asyncio.subprocess.Process, task_meta): # 等待子进程完成,收集输出 stdout, stderr = await process.communicate() result = ProcessResultModel( returncode=process.returncode, stdout=stdout.decode('utf-8') if stdout else '', stderr=stderr.decode('utf-8') if stderr else '' ) # 执行后续业务逻辑:文件操作、数据库写入等 if result.returncode == 0: # 任务成功:更新数据库状态、通知前端等 print(f"任务[{task_meta['id']}]完成:{result}") # 示例:await db.update_video_status(task_meta['id'], "completed") else: # 任务失败:记录错误日志、触发告警等 print(f"任务[{task_meta['id']}]失败:{result.stderr}") # 示例:await db.update_video_status(task_meta['id'], "failed", error=result.stderr)
2. FastAPI路由中调用
from fastapi import FastAPI import asyncio import ffmpeg app = FastAPI() @app.post("/start-encode") async def start_encode(input_path: str, output_path: str, task_id: int): # 构造ffmpeg处理流 stream = ffmpeg.input(input_path).output(output_path) # 启动子进程,不等待完成 process = await stream.run_async_async(run=False) # 启动后台协程监控进程,脱离请求上下文执行 asyncio.create_task(handle_process_result(process, {"id": task_id})) # 立即返回响应,不阻塞主循环 return {"status": "started", "task_id": task_id}
注意事项
- 确保
handle_process_result内的数据库操作、文件操作都是异步实现(如用asyncpg、aiofiles),避免阻塞事件循环。 - 多worker部署时,后台任务仅在当前worker进程内运行,进程重启会丢失未完成任务,适合短周期、非核心任务。
二、生产环境方案:Celery + 消息队列(推荐)
对于需要持久化、重试、分布式处理的核心任务,用Celery异步任务框架是更健壮的选择,完全隔离FastAPI主进程与任务执行进程。
1. 配置Celery
# celery_config.py from celery import Celery # 用Redis作为消息队列和结果存储(也可替换为RabbitMQ) app = Celery( "video_tasks", broker="redis://localhost:6379/0", backend="redis://localhost:6379/0" )
2. 定义编码任务
# tasks.py from celery_config import app import ffmpeg from pydantic import BaseModel class ProcessResultModel(BaseModel): returncode: int = None stdout: str = '' stderr: str = '' @app.task(bind=True) def encode_video_task(self, input_path, output_path, task_id): try: # 同步执行ffmpeg命令(Celery任务在独立进程运行,不影响FastAPI) stdout, stderr = ffmpeg.run( ffmpeg.input(input_path).output(output_path), capture_stdout=True, capture_stderr=True ) # 后续同步业务操作:文件移动、数据库写入等 # 示例:db.session.execute("UPDATE tasks SET status='completed' WHERE id=?", (task_id,)) # db.session.commit() return ProcessResultModel( returncode=0, stdout=stdout.decode('utf-8'), stderr=stderr.decode('utf-8') ).dict() except Exception as e: # 任务失败自动重试(最多3次) self.retry(exc=e, max_retries=3) return ProcessResultModel( returncode=-1, stderr=str(e) ).dict()
3. FastAPI中调用任务
from fastapi import FastAPI from tasks import encode_video_task app = FastAPI() @app.post("/start-encode") async def start_encode(input_path: str, output_path: str, task_id: int): # 发送任务到Celery队列,立即返回 task = encode_video_task.delay(input_path, output_path, task_id) return {"status": "started", "task_id": task_id, "celery_task_id": task.id} # 可选:查询任务状态接口 @app.get("/task-status/{celery_task_id}") async def get_task_status(celery_task_id: str): task = encode_video_task.AsyncResult(celery_task_id) match task.state: case 'PENDING': return {"status": "pending", "message": "任务等待执行"} case 'SUCCESS': return {"status": "completed", "result": task.result} case _: return {"status": task.state, "error": str(task.info)}
启动Celery Worker
celery -A celery_config worker --loglevel=info
方案优势
- 任务持久化:FastAPI进程重启不影响任务执行。
- 支持重试、超时、并发控制,适合核心业务场景。
- 可分布式部署多个Worker,横向扩展处理能力。
内容的提问来源于stack exchange,提问作者s0mbre
相关产品推荐
相关产品推荐

