如何在asyncio中限制依赖型Web请求序列速率且不阻塞响应?
解决带依赖的异步API请求全局速率限制问题
针对你遇到的「分类内请求有依赖、跨分类并行、全局请求速率限制」的需求,asyncio生态里有标准的解决方案:实现全局异步速率限制器,结合串行的分类任务与并行的任务调度。
核心思路
- 全局速率控制:用一个事件循环安全的速率限制器,统一管控所有API请求的发送频率,确保不超过10请求/秒的上限。
- 分类内串行执行:每个分类的
ka→kb→kc请求因存在依赖,需在单个任务内串行执行(用await保证顺序)。 - 跨分类并行调度:多个分类的任务通过
asyncio.TaskGroup并行启动,充分利用异步IO的优势,避免阻塞。
代码实现
1. 全局速率限制器
这个限制器基于时间窗口思想,控制单位时间内的请求数:
import asyncio class RateLimiter: def __init__(self, max_requests: int, period: float): self.max_requests = max_requests # 单位时间内最大请求数 self.period = period # 时间窗口(秒) self.request_timestamps = [] self.lock = asyncio.Lock() async def acquire(self): async with self.lock: now = asyncio.get_event_loop().time() # 清理超出时间窗口的请求记录 self.request_timestamps = [t for t in self.request_timestamps if now - t < self.period] if len(self.request_timestamps) >= self.max_requests: # 计算需要等待的时间,直到时间窗口内有空闲名额 wait_time = self.period - (now - self.request_timestamps[0]) await asyncio.sleep(wait_time) # 等待后重新清理时间窗口 now = asyncio.get_event_loop().time() self.request_timestamps = [t for t in self.request_timestamps if now - t < self.period] self.request_timestamps.append(now) # 支持async with语法 async def __aenter__(self): await self.acquire() return self async def __aexit__(self, *args): pass
2. 带速率限制的API请求函数
所有API请求都需要通过速率限制器获取许可后再发送:
import aiohttp async def collect_a_data(k: int, session: aiohttp.ClientSession, limiter: RateLimiter): async with limiter: async with session.get(f"https://your-api-endpoint/{k}/a") as resp: data = await resp.json() # 返回是否需要继续请求kb的标志 return data.get("need_b", False) async def collect_b_data(k: int, session: aiohttp.ClientSession, limiter: RateLimiter): async with limiter: async with session.get(f"https://your-api-endpoint/{k}/b") as resp: data = await resp.json() # 返回是否需要继续请求kc的标志 return data.get("need_c", False) async def collect_c_data(k: int, session: aiohttp.ClientSession, limiter: RateLimiter): async with limiter: async with session.get(f"https://your-api-endpoint/{k}/c") as resp: data = await resp.json() # 处理kc数据 return data
3. 分类任务与主调度
每个分类的串行逻辑封装在单个任务中,所有分类任务并行启动:
async def collect_abc_data(k: int, session: aiohttp.ClientSession, limiter: RateLimiter): # 按顺序执行ka→kb(按需)→kc(按需) need_b = await collect_a_data(k, session, limiter) need_c = False if need_b: need_c = await collect_b_data(k, session, limiter) if need_c: await collect_c_data(k, session, limiter) return {"category": k, "status": "completed"} async def main(): # 初始化速率限制器:10请求/秒 rate_limiter = RateLimiter(max_requests=10, period=1.0) async with aiohttp.ClientSession() as session: async with asyncio.TaskGroup() as tg: # 并行启动所有分类的任务 tasks = [tg.create_task(collect_abc_data(k, session, rate_limiter)) for k in [1, 2, 3]] # 收集所有任务结果 results = [task.result() for task in tasks] print(results) if __name__ == "__main__": asyncio.run(main())
为什么这个方案能解决你的问题
- 避免速率超标:全局速率限制器统一管控所有请求,不管是哪个分类的
ka/kb/kc,都会被限制在10请求/秒以内,不会出现多个任务并行导致的速率溢出。 - 满足依赖要求:每个分类的请求在单个任务内用
await串行执行,严格保证ka→kb→kc的顺序,且依赖判断逻辑清晰。 - 无阻塞延迟:所有等待逻辑用
asyncio.sleep实现,只会让出当前任务的执行权,不会阻塞整个事件循环。当某个分类在等待API响应时,其他分类的任务可以继续发送请求(只要不超过速率限制),避免了time.sleep导致的全局阻塞问题。 - 高效并行:跨分类的任务通过
TaskGroup并行调度,充分利用异步IO的优势,大幅缩短总执行时间。
内容的提问来源于stack exchange,提问作者SapereAude
相关产品推荐
相关产品推荐

