You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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 Future attached 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 01:30:21