Python multiprocessing并行API请求限流最优实现方案
问题根因
- 进程数和API限频不匹配:你开了8个工作进程,API单秒限频仅6次,进程启动瞬间会同时发起8个请求,直接超出阈值触发429
- 退避逻辑存在设计缺陷:所有进程独立执行固定10秒的休眠,休眠结束后会同时发起重试,再次瞬间打满请求量,陷入「触发429-休眠-集体重试-再次触发429」的死循环
- 场景选型错误:API请求属于IO密集型任务,多进程的优势是处理CPU密集型计算,用来发请求不仅会带来额外的进程调度、通信开销,还会提升全局速率管控的复杂度
最优实现方案
核心逻辑是从请求发起的源头做全局速率管控,而不是触发429之后再被动退避,具体要点:
- 并发数不用匹配CPU核心数,IO场景下3-4个工作并发完全足够,过高并发没有收益
- 配置全局统一的速率限制器,把总请求速率控制在5-5.5次/秒,预留10%左右的冗余抵消服务端计数窗口的误差,从根源上避免触发429
- 退避逻辑采用指数退避+随机抖动的策略,不要用固定时长休眠,避免多个工作端同时重试再次撞限
- 优先用线程池而非进程池处理这类IO任务:Python线程在IO等待时会释放GIL,性能和多进程基本持平,还可以直接共享内存中的限流器实例,不需要额外做跨进程通信,实现成本极低
修正后的可运行代码
import time import random import requests from concurrent.futures import ThreadPoolExecutor from threading import Lock # 全局配置 API_RATE_LIMIT = 5.5 # 预留0.5次冗余,抵消服务端计数窗口误差 API_URL = "替换为实际API请求地址" MAX_WORKERS = 4 # IO场景4个并发足够,不需要开到和CPU核心数一致 MAX_RETRY = 5 # 令牌桶全局限流器 class RateLimiter: def __init__(self, rate_per_sec): self.rate = rate_per_sec self.tokens = rate_per_sec self.last_update = time.time() self.lock = Lock() def acquire(self): with self.lock: now = time.time() time_passed = now - self.last_update self.tokens = min(self.rate, self.tokens + time_passed * self.rate) self.last_update = now if self.tokens < 1: wait_time = (1 - self.tokens) / self.rate time.sleep(wait_time) self.tokens = 0 self.last_update = time.time() else: self.tokens -= 1 rate_limiter = RateLimiter(API_RATE_LIMIT) def api_call(params_tuple): api_key, user_id = params_tuple query = {'api_key': api_key, 'user_id': user_id} retry_count = 0 while True: # 先获取请求令牌再发请求,从根源避免超量 rate_limiter.acquire() resp = requests.get(API_URL, params=query) if resp.status_code == 200: data = resp.json() print(data) return data elif resp.status_code == 429: retry_count += 1 if retry_count > MAX_RETRY: raise Exception(f"请求{query}超过最大重试次数,仍返回429") # 指数退避+随机抖动,避免集体同时重试 backoff = min(2 ** retry_count, 10) + random.uniform(0, 1) time.sleep(backoff) else: resp.raise_for_status() if __name__ == "__main__": # 替换为实际的请求参数列表 iterable = [("你的api_key", f"user_{i}") for i in range(100)] with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: try: res = list(executor.map(api_call, iterable)) except KeyboardInterrupt: print("收到中断信号,停止执行") executor.shutdown(wait=False, cancel_futures=True)
如果必须使用
multiprocessing实现,需要通过multiprocessing.Manager实现跨进程共享的限流器状态,或者单独启动一个管控进程统一分配请求令牌,但这类实现复杂度高,对于纯API请求场景没有额外收益,不推荐使用。
内容的提问来源于stack exchange,提问作者Coldchain9
相关产品推荐
相关产品推荐

