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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 20:45:53