如何用aiohttp+asyncio实现每秒5次请求限流以解决429错误?
正确的异步限流轮询实现方案
针对你的需求——基于订单ID列表每10秒轮询状态,同时遵守每秒最多5次请求的API速率限制,以下是经过验证的实现方案:
核心问题分析
之前的方案(Semaphore、自定义RateLimiter等)失效的原因大多是只控制了并发请求数,没控制时间维度的请求速率。比如Semaphore设为5时,5个请求可能在0.1秒内全部完成,紧接着又会发出新的5个请求,瞬间突破每秒5次的限制。
推荐实现:使用aiolimiter精确控速
aiolimiter是专门为异步场景设计的速率限制库,能精准控制单位时间内的请求数,完美匹配你的需求。
完整代码示例
import asyncio import aiohttp from aiolimiter import AsyncLimiter async def fetch_order_status(session, limiter, order_id): # 每次请求前获取限流许可,确保不超过每秒5次 await limiter.acquire() api_url = f"https://api.planet.com/orders/{order_id}" async with session.get(api_url) as response: # 处理429错误:简单重试(需结合限流器,避免加重限流) if response.status == 429: await asyncio.sleep(2) return await fetch_order_status(session, limiter, order_id) # 处理其他错误(可选) response.raise_for_status() return await response.json() async def poll_order_statuses(activated_order_ids, poll_interval=10): # 初始化限流器:1秒内最多允许5次请求 rate_limiter = AsyncLimiter(5, 1) async with aiohttp.ClientSession() as session: while True: # 为每个订单创建带限流的请求任务 tasks = [ fetch_order_status(session, rate_limiter, order_id) for order_id in activated_order_ids ] # 并发执行所有任务(受限于流器控速) status_results = await asyncio.gather(*tasks, return_exceptions=True) # 处理返回结果(示例:打印状态) for order_id, result in zip(activated_order_ids, status_results): if isinstance(result, Exception): print(f"Order {order_id} query failed: {str(result)}") else: print(f"Order {order_id} current state: {result.get('state', 'unknown')}") # 等待下一轮轮询 await asyncio.sleep(poll_interval) if __name__ == "__main__": # 替换为你的订单ID列表 ACTIVATED_ORDERS = ["order_001", "order_002", "order_003"] asyncio.run(poll_order_statuses(ACTIVATED_ORDERS))
关键细节说明
- 精确速率控制:
AsyncLimiter(5, 1)确保每1秒窗口内最多通过5个请求,不管任务并发量多大,都会按速率放行。 - 异步友好:全程使用
asyncio.sleep而非time.sleep,不会阻塞事件循环。 - 错误处理:加入429重试逻辑,同时重试也受限于流器,避免加剧限流问题。
- 轮询逻辑:每一轮所有请求完成后,再等待10秒进入下一轮,符合你的轮询周期要求。
其他注意事项
- 不要手动拆分订单列表分批请求(比如每5个一批加sleep),
aiolimiter会自动处理速率控制,更高效。 - 如果API有分钟级限流(比如每分钟最多100次),可以叠加一个分钟级的限流器,或者在代码中额外控制。
内容的提问来源于stack exchange,提问作者saving_space
相关产品推荐
相关产品推荐

