asyncio.sleep配置后仍触发API速率限制429问题排查
问题根因
这不是API服务端故障,是你的异步限流逻辑存在本质设计缺陷,核心问题如下:
asyncio.Semaphore仅能限制同时处于进行中的请求并发数,完全无法控制请求的发送速率。你设置并发数为3,不代表每秒最多发出3个请求——只要请求完成速度够快,1秒内可以连续释放多次信号量,发出远超过4个请求。- 你写的
await asyncio.sleep(1.5)位置完全错误:这段延迟是在收到完整响应、解析完JSON内容之后才执行的,根本不会在请求发送前做间隔,完全无法控制发请求的时间节奏。 - 你日志排查到的“慢请求后抢跑”是这个设计的必然结果:举个例子,3个并发请求中如果有1个请求耗时1.2秒才返回,另外2个请求可能在0.2秒就完成、睡够1.5秒释放了信号量;等慢请求一处理完,信号量会瞬间空出多个位置,排队的任务会立刻批量发出新请求,直接导致1秒时间窗口内到达服务端的请求数超过4个阈值。
- 额外的代码问题:
_run函数启动前的time.sleep(10)是同步阻塞调用,会直接卡停整个事件循环,异步场景下应该用asyncio.sleep;429重试逻辑没有接入速率控制,睡12秒后直接重发很容易再次撞上限流。 - 为什么返回400时不会触发429:400响应通常体积小、处理速度快,请求完成节奏均匀,歪打正着没超过限流阈值;一旦碰到200的大体积响应,JSON解析额外消耗时间打乱了请求节奏,批量待发请求就会直接冲爆限流。
修复方案
要稳定满足4请求/秒的限流要求,不能靠并发信号量+请求后延迟凑数,必须在请求发送前做全局速率控制,推荐用滑动窗口思路实现,核心修改点:
- 新增全局速率限制器,保证每秒最多放行4个请求,不管请求耗时、并发数多少
- 把延迟等待逻辑移到请求发送动作之前
- 替换所有同步阻塞的
time.sleep为异步asyncio.sleep - 429重试时额外增加退避时间,重试请求同样要过速率校验
修正后的核心实现代码:
import asyncio import aiohttp import time from collections import deque class Request: def __init__(self, url: str, method: str="get", payload: str=None): self.url: str = url self.method: str = method self.payload: str or dict = payload or dict() class Response: def __init__(self, url: str, status: int, payload: dict=None, error: bool=False, text: str=None): self.url: str = url self.status: int = status self.payload: dict = payload or dict() self.error: bool = error self.text: str = text or '' class RateLimiter: def __init__(self, max_per_second: int): self.max_per_second = max_per_second self.call_timestamps = deque() async def acquire(self): # 清理1秒窗口外的时间记录 now = time.time() while self.call_timestamps and now - self.call_timestamps[0] >= 1: self.call_timestamps.popleft() # 如果当前窗口请求数达上限,等到最早的请求出窗口 if len(self.call_timestamps) >= self.max_per_second: wait_time = 1 - (now - self.call_timestamps[0]) await asyncio.sleep(wait_time) self.call_timestamps.append(time.time()) def make_requests(headers: dict, requests: list[Request]) -> list[Response]: return asyncio.run(_run(headers, requests)) async def _run(headers: dict, requests: list[Request]) -> list[Response]: semaphore: asyncio.Semaphore = asyncio.Semaphore(3) # 初始化限流器,每秒最多4个请求 rate_limiter = RateLimiter(max_per_second=4) await asyncio.sleep(10) # 替换为异步sleep避免阻塞事件循环 async with aiohttp.ClientSession(headers=headers) as session: tasks: list[asyncio.Task] = [asyncio.create_task(_iterate(semaphore, session, request, rate_limiter)) for request in requests] responses: list[Response] = await asyncio.gather(*tasks) return responses async def _iterate(semaphore: asyncio.Semaphore, session: aiohttp.ClientSession, request: Request, rate_limiter: RateLimiter) -> Response: async with semaphore: return await _fetch(session, request, rate_limiter) async def _fetch(session: aiohttp.ClientSession, request: Request, rate_limiter: RateLimiter) -> Response: try: # 发请求前先过限流校验 await rate_limiter.acquire() async with session.request(request.method, request.url, params=request.payload) as response: print(f"NOW: {time.time()}") print(f"Response Status: {response.status}.") content: dict = await response.json() response.raise_for_status() return Response(request.url, response.status, payload=content, error=False) except aiohttp.ClientResponseError: if response.status == 429: await asyncio.sleep(12) # 重试同样过限流校验 return await _fetch(session, request, rate_limiter) else: return Response(request.url, response.status, error=True)
注:如果不想自己实现限流器,也可以直接用成熟的异步限流库的现成实现,逻辑和上述滑动窗口实现一致,稳定性更高。
内容的提问来源于stack exchange,提问作者Coldchain9
相关产品推荐
相关产品推荐

