FastAPI中如何在新线程执行协程作为后台任务并返回结果
FastAPI后台异步任务执行问题解决
问题代码
import asyncio import threading from sqlalchemy.ext.declarative import declarative_base SQLAlchemyModel = declarative_base() class Service: async def publish(self, file_obj): """Sending file_obj to s3""" def render_pdf(self, data: dict): """CPU task.""" async def get_stuff(self) -> bool: """Get data from aio.Redis""" async def get_data_db(self) -> SQLAlchemyModel: """Get data from DB (like postgres)""" async def request_to_3rd_api(self) -> dict: """Working with 3rd party API""" async def task(self, data): """Task which should be execute as background task. I don't want to get a result of this.""" # Here is some blocking cpu bound task and also io bound task. file_obj = self.render_pdf(data) await self.publish(file_obj) def sync_task(self, data): asyncio.run(self.task(data)) async def do_smth_useful(self) -> dict | None: """Interface for working with service. Make request to 3rd-party API. If some stuff exists in Redis then render file with data from 3rd API and save it in S3: I don't need a result here. Method should return result to caller before this task will be completed. """ result = await self.request_to_3rd_api() db_data = await self.get_data_db() data = (result, db_data) if await self.get_stuff(): threading.Thread(target=self.sync_task, args=(data,)).start() return result
触发的错误
RuntimeError: Task <Task pending
name='anyio.from_thread.BlockingPortal._call_func'
coro=<BlockingPortal._call_func() running at
/lib/python3.11/site-packages/anyio/from_thread.py:217>
cb=[TaskGroup._spawn..task_done()
/lib/python3.11/site-packages/anyio/_backends/_asyncio.py:661]>
got Future
attached to a different loop got Futureattached to a different loop
核心问题
- 如何正确在新线程中执行协程作为后台任务,让主线程无需等待后台任务完成即可向客户端返回结果?
- 是否应该为每个线程共享事件循环,还是为新线程创建独立的事件循环?
解决方案
方案1:使用FastAPI内置的BackgroundTasks(推荐)
FastAPI原生支持后台任务,无需手动管理线程和事件循环,是最可靠的实现方式:
from fastapi import FastAPI, BackgroundTasks, Depends app = FastAPI() # 假设Service实例通过依赖注入获取 def get_service(): return Service() @app.post("/execute-useful") async def execute_useful(background_tasks: BackgroundTasks, service: Service = Depends(get_service)): result = await service.do_smth_useful(background_tasks) return result class Service: # ... 保留原有其他方法 ... async def do_smth_useful(self, background_tasks: BackgroundTasks) -> dict | None: result = await self.request_to_3rd_api() db_data = await self.get_data_db() data = (result, db_data) if await self.get_stuff(): # 将后台任务交给FastAPI管理 background_tasks.add_task(self._handle_background_job, data) return result async def _handle_background_job(self, data): # CPU密集任务用asyncio.to_thread放到线程池,避免阻塞事件循环 file_obj = await asyncio.to_thread(self.render_pdf, data) await self.publish(file_obj)
关键说明:
BackgroundTasks会在响应返回给客户端后,自动在后台线程中执行任务。- 把CPU密集的
render_pdf用asyncio.to_thread剥离到线程池,防止阻塞异步事件循环。
方案2:手动管理线程与事件循环(仅作原理说明,不推荐)
如果必须手动创建线程,需严格保证事件循环的线程隔离:
class Service: # ... 保留原有其他方法 ... def _background_thread_entry(self, data): # 为新线程创建独立的事件循环 loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) try: loop.run_until_complete(self._handle_background_job(data)) finally: loop.close() async def _handle_background_job(self, data): file_obj = await asyncio.to_thread(self.render_pdf, data) await self.publish(file_obj) async def do_smth_useful(self) -> dict | None: result = await self.request_to_3rd_api() db_data = await self.get_data_db() data = (result, db_data) if await self.get_stuff(): threading.Thread(target=self._background_thread_entry, args=(data,)).start() return result
关键说明:
- 原错误根源:你尝试在新线程中使用绑定了主线程事件循环的异步对象(如aio.Redis、异步DB连接),这些对象无法跨循环使用。
- 每个线程必须创建独立的事件循环,绝对不能共享主线程的循环(事件循环不是线程安全的)。
事件循环选择原则
- 禁止跨线程共享事件循环:事件循环与线程绑定,跨线程操作必然引发异常。
- 线程独立创建循环:手动管理线程时,必须为每个线程单独初始化、运行、关闭事件循环。
- 优先用FastAPI内置方案:
BackgroundTasks已经封装了线程池和循环隔离逻辑,无需手动处理,稳定性更高。
内容的提问来源于stack exchange,提问作者Leo
相关产品推荐
相关产品推荐

