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

FastAPI中如何全局共享asyncio.Semaphore限制并发下载

问题原因

asyncio.Semaphore 实例会强绑定到创建它时所在的事件循环。你在模块顶层实例化Classifier时,FastAPI 服务还未正式启动,此时创建的Semaphore绑定的是Python默认生成的预启动事件循环;但Uvicorn/Starlette启动后会创建独立的运行时事件循环处理所有请求,两个循环不互通,跨循环调用异步对象就会触发attached to a different loop报错。
如果把Semaphore初始化逻辑放到run()方法内,每次请求到达都会生成独立的Semaphore实例,计数仅在单个请求内生效,完全无法实现跨请求的全局并发限制。

实现方案

通过FastAPI的生命周期钩子,在服务启动完成、运行时事件循环已经就绪之后,再初始化Semaphore和httpx异步客户端,保证所有请求共享同一个绑定到正确事件循环的Semaphore实例,即可实现全局下载并发数管控。

修正后代码

classifier.py

import asyncio
import httpx
from typing import List, Optional


class Classifier():
    def __init__(self, concurrency_limit: int) -> None:
        self.concurrency_limit = concurrency_limit
        self.client: Optional[httpx.AsyncClient] = None
        self.semaphore: Optional[asyncio.Semaphore] = None

    async def startup(self) -> None:
        # 事件循环运行后再初始化异步对象,绑定到服务实际使用的循环
        self.client = httpx.AsyncClient()
        self.semaphore = asyncio.Semaphore(self.concurrency_limit)

    async def shutdown(self) -> None:
        # 服务关闭时主动释放连接资源
        if self.client:
            await self.client.aclose()

    async def download_async(self, url: str) -> bytes:
        async with self.semaphore:
            response = await self.client.get(url)
        return await response.aread()

    async def run(self, image_urls: List[str]) -> List[bytes]:
        image_list = await asyncio.gather(
            *[self.download_async(url) for url in image_urls]
        )
        # 此处添加图像推理逻辑
        return image_list

api.py

from fastapi import FastAPI
from contextlib import asynccontextmanager
from typing import List
from classifier import Classifier

classifier = Classifier(concurrency_limit=5)

@asynccontextmanager
async def lifespan(app: FastAPI):
    # 服务启动阶段执行初始化
    await classifier.startup()
    yield
    # 服务关闭阶段执行资源清理
    await classifier.shutdown()

server = FastAPI(lifespan=lifespan)

@server.post("/")
async def index(urls: List[str]):
    results = await classifier.run(urls)
    return {"code": 0, "data": results}

兼容旧版本写法

如果使用0.95版本以前的FastAPI(不支持lifespan API),可以用已废弃但仍可用的startup/shutdown事件实现相同效果:

server = FastAPI()
classifier = Classifier(concurrency_limit=5)

@server.on_event("startup")
async def init_async_res():
    await classifier.startup()

@server.on_event("shutdown")
async def clean_async_res():
    await classifier.shutdown()
注意事项
  • 所有绑定事件循环的异步对象(asyncio.Semaphore、asyncio.Queue、httpx.AsyncClient、aiohttp.ClientSession等)都不要在模块顶层、类同步__init__方法内直接创建,必须放到事件循环启动后的异步初始化逻辑中实例化。
  • 该实现下Semaphore为全局单例,无论同时处理多少个请求、每个请求携带多少个图片URL,全局同时执行的下载协程总数永远不会超过设定的concurrency_limit阈值,符合全局限流的需求。

内容的提问来源于stack exchange,提问作者Lezwon Castellino

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 23:27:29