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

如何在asyncio中限制依赖型Web请求序列速率且不阻塞响应?

解决带依赖的异步API请求全局速率限制问题

针对你遇到的「分类内请求有依赖、跨分类并行、全局请求速率限制」的需求,asyncio生态里有标准的解决方案:实现全局异步速率限制器,结合串行的分类任务与并行的任务调度。

核心思路

  1. 全局速率控制:用一个事件循环安全的速率限制器,统一管控所有API请求的发送频率,确保不超过10请求/秒的上限。
  2. 分类内串行执行:每个分类的ka→kb→kc请求因存在依赖,需在单个任务内串行执行(用await保证顺序)。
  3. 跨分类并行调度:多个分类的任务通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 18:44:50