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
相关产品推荐
相关产品推荐

