使用asyncio调用5请求/秒限制的API时,如何解决速率限制触发问题?
如何正确处理5请求/秒的API速率限制
我来帮你梳理下问题所在,然后给出可行的解决方案——你现在的代码之所以还触发速率限制,主要是同步和异步逻辑混用导致时间控制混乱,再加上没有统一的速率管控机制,咱们一步步来改:
问题分析
你的代码存在几个核心问题:
- 同步与异步sleep混用:同步的
time.sleep()会阻塞整个事件循环,而异步的await asyncio.sleep()只阻塞当前任务,两者混在一起会让请求的时间间隔完全不可控,比如主线程的time.sleep(1)会卡住所有异步任务,但异步任务里的sleep又不会影响主线程的循环速度,很容易导致短时间内发出过多请求。 - 零散的sleep无法精确控速:你在各个地方加的sleep是固定值,但异步任务是并行执行的,相当于同一时间可能有多个请求同时发出去,远远超过5请求/秒的限制。
- 重试逻辑不够完善:遇到异常时直接固定sleep后重试,没有针对API返回的429(速率限制)错误做针对性的退避策略,反而可能加重请求拥堵。
解决方案
1. 统一切换到全异步框架
首先要把所有同步请求(比如requests.post)换成异步HTTP库(比如aiohttp),因为同步请求会阻塞事件循环,让异步逻辑失去意义。
2. 实现统一的速率限制器
用令牌桶算法实现一个速率限制器,精确控制每秒最多发出5个请求,所有请求(包括你的quick-search和crop请求)都必须先获取许可才能发送。
3. 优化重试逻辑
针对API返回的429错误,使用指数退避的方式重试(比如第一次等1秒,第二次等2秒,第三次等4秒),避免反复触发速率限制。
修改后的完整代码示例
第一步:实现速率限制器
import asyncio from collections import deque class RateLimiter: def __init__(self, max_requests: int, time_window: float): self.max_requests = max_requests # 时间窗口内最大请求数 self.time_window = time_window # 时间窗口(秒) self.request_timestamps = deque() # 存储请求时间戳 self.lock = asyncio.Lock() # 确保线程安全 async def acquire(self): async with self.lock: now = asyncio.get_event_loop().time() # 移除时间窗口外的旧请求记录 while self.request_timestamps and now - self.request_timestamps[0] > self.time_window: self.request_timestamps.popleft() # 如果请求数已达上限,计算需要等待的时间 if len(self.request_timestamps) >= self.max_requests: wait_time = self.time_window - (now - self.request_timestamps[0]) await asyncio.sleep(wait_time) # 等待后再清理一次过期记录 now = asyncio.get_event_loop().time() while self.request_timestamps and now - self.request_timestamps[0] > self.time_window: self.request_timestamps.popleft() # 记录当前请求时间 self.request_timestamps.append(now)
第二步:修改异步请求函数和主逻辑
import asyncio from aiohttp import ClientSession, BasicAuth # 初始化速率限制器:5请求/秒 rate_limiter = RateLimiter(max_requests=5, time_window=1.0) async def fire_post(session, payload): # 先获取速率许可,确保不超限 await rate_limiter.acquire() try: async with session.post(url_crop, data=payload) as response: # 专门处理速率限制错误(429) if response.status == 429: retry_delay = 1 # 最多重试3次 for _ in range(3): await asyncio.sleep(retry_delay) await rate_limiter.acquire() async with session.post(url_crop, data=payload) as retry_response: if retry_response.status != 429: return await retry_response.json() retry_delay *= 2 # 指数退避 # 多次重试失败,返回响应文本 return await response.text() # 正常响应返回JSON return await response.json() except Exception as e: print(f"请求出错: {str(e)}") # 异常情况下重试一次,带速率控制 await asyncio.sleep(2) await rate_limiter.acquire() try: async with session.post(url_crop, data=payload) as response: return await response.json() if response.status == 200 else await response.text() except: return await response.text() async def main(object_data, API_KEY, search_request): async with ClientSession() as session: tasks_visual = [] tasks_analytical = [] for object_id, object_coors in object_data.items(): # 处理quick-search请求,同样受速率限制 await rate_limiter.acquire() try: async with session.post( 'https://api.planet.com/data/v1/quick-search', auth=BasicAuth(API_KEY, ''), json=search_request, ssl=False ) as response: search_result = await response.json() except Exception as e: print(e) await asyncio.sleep(5) await rate_limiter.acquire() async with session.post( 'https://api.planet.com/data/v1/quick-search', auth=BasicAuth(API_KEY, ''), json=search_request, ssl=False ) as response: search_result = await response.json() # 添加视觉任务和分析任务 task_visual = asyncio.create_task(fire_post(session, payload_visual)) tasks_visual.append(task_visual) task_analytical = asyncio.create_task(fire_post(session, payload_analytical)) tasks_analytical.append(task_analytical) # 等待所有任务完成 await asyncio.gather(*tasks_visual) await asyncio.gather(*tasks_analytical) # 运行主函数 asyncio.run(main(object_data, API_KEY, search_request))
关键说明
- 全异步化:用
aiohttp替代requests,确保所有网络请求都不会阻塞事件循环,异步任务能真正并行执行又不会失控。 - 统一控速:所有请求(包括quick-search和crop)都要通过
rate_limiter.acquire()获取许可,严格保证每秒最多5个请求。 - 智能重试:针对429错误做指数退避,既避免反复触发限制,又能在API恢复后尽快完成请求。
内容的提问来源于stack exchange,提问作者Manap Shymyr
相关产品推荐
相关产品推荐

