如何在FastAPI异步REST端点中执行同步任务且不阻塞事件循环?
问题描述
我用FastAPI和Celery搭建了一个处理CPU密集型同步任务的应用,有一个异步REST端点request_data,用户调用它请求数据,支持可选参数force:
- 当
force=False时,直接返回Cassandra中的数据; - 当
force=True时,需要同步执行CPU密集型任务并在当前请求中返回结果。
当前遇到的问题:这个端点是异步函数,调用CPU密集型任务会阻塞事件循环,影响FastAPI处理其他请求。
我试过以下方案但均不理想:
- 定义同步和异步两个函数,但FastAPI无法通过中间件将请求路由到对应函数;
- 异步提交Celery任务,但无法
await结果,其他等待方式都会阻塞事件循环; - 将REST端点改为同步函数,但会失去Cassandra
execute_async提供的异步IO优势,因此希望避免这种做法。
我需要一种不阻塞事件循环(不影响FastAPI处理其他请求)的执行方式,清楚异步函数中调用同步任务最终仍需处理CPU密集工作,但只要保证事件循环不被阻塞即可。
当前代码示例:
from fastapi import FastAPI import cassandra_wrapper # 封装Cassandra数据库访问的类 app = FastAPI() @app.get("/data") async def request_data(force=False): if force: # 执行计算并返回数据 return cpu_intensive_task() else: return cassandra_wrapper.get_data() # 可以调用这个Celery任务执行计算 @celery.task(bind=True, name="cpu_intensive_task") def cpu_intensive_celery_task(): return cpu_intensive_task # 也可以直接调用这个函数执行计算 def cpu_intensive_task(): import time time.sleep(5)
可行解决方案
方案1:使用asyncio.to_thread将同步任务移到线程池
这是FastAPI官方推荐的异步函数中运行同步CPU密集任务的方式,会把同步函数放到单独线程执行,不阻塞事件循环。
修改后端点代码:
import asyncio from fastapi import FastAPI import cassandra_wrapper app = FastAPI() @app.get("/data") async def request_data(force=False): if force: # 将CPU密集任务放到线程池执行,不阻塞事件循环 result = await asyncio.to_thread(cpu_intensive_task) return result else: # 若get_data是同步方法,也可套to_thread保持异步优势 return await cassandra_wrapper.get_data() def cpu_intensive_task(): import time time.sleep(5) return {"result": "computed data"}
说明:asyncio.to_thread利用Python默认线程池,将同步任务移至线程执行,事件循环可继续处理其他请求,待线程任务完成后再返回结果。
方案2:结合Celery异步获取结果(非阻塞方式)
若必须使用Celery,可将Celery的任务等待逻辑包装为可await的对象,避免阻塞事件循环:
import asyncio from fastapi import FastAPI from celery.result import AsyncResult import cassandra_wrapper import celery app = FastAPI() celery_app = celery.Celery(...) # 初始化Celery实例 @app.get("/data") async def request_data(force=False): if force: task = cpu_intensive_celery_task.delay() # 用线程池包装Celery的wait方法,非阻塞事件循环 result = await asyncio.get_event_loop().run_in_executor(None, task.wait) return result else: return await cassandra_wrapper.get_data() @celery_app.task(bind=True, name="cpu_intensive_task") def cpu_intensive_celery_task(): import time time.sleep(5) return {"result": "celery computed data"}
说明:run_in_executor将Celery的task.wait()移至线程池执行,事件循环可同时处理其他请求,不会被阻塞。
方案3:自定义线程池(进阶)
若需灵活配置线程数(比如根据CPU核心数调整),可自定义线程池:
from concurrent.futures import ThreadPoolExecutor import asyncio from fastapi import FastAPI import cassandra_wrapper app = FastAPI() # 根据业务需求自定义线程池大小 executor = ThreadPoolExecutor(max_workers=4) @app.get("/data") async def request_data(force=False): if force: loop = asyncio.get_event_loop() result = await loop.run_in_executor(executor, cpu_intensive_task) return result else: return await cassandra_wrapper.get_data() def cpu_intensive_task(): import time time.sleep(5) return {"result": "custom pool computed data"}
说明:自定义线程池可根据并发需求调整线程数量,避免默认线程池的限制,适合高并发场景下的CPU密集任务处理。
内容的提问来源于stack exchange,提问作者ATK
相关产品推荐
相关产品推荐

