Python多进程调用API时如何实现每秒不超过5次的限流控制
解决方案
方案1:最小改动适配现有joblib逻辑
你的场景属于IO密集型任务,用多线程比多进程开销更低,适配逻辑更简单,不需要大改现有代码即可实现限流要求。
- 先安装限流依赖
pip install ratelimit
- 修改后代码
import pandas as pd from joblib import Parallel, delayed from ratelimit import limits, sleep_and_retry import os # 配置参数 BASE_DIR = r"C:\你的实际保存路径" # 替换为你自己的文件保存目录 MAX_REQUESTS_PER_SECOND = 5 # 单次请求平均10秒,最大并发数设为50即可跑满带宽且不超限流阈值 MAX_WORKERS = 50 # 读取id列表 ids = pd.read_csv('data.csv')['Id'].values.tolist() def dump_data(data, idx): filename = os.path.join(BASE_DIR, f"{idx}.csv") data.to_csv(filename, header=True, index=False) # 加限流装饰器:每秒最多调用5次,触发限流时自动休眠等待到可用时间窗口 @sleep_and_retry @limits(calls=MAX_REQUESTS_PER_SECOND, period=1) def get_api(idx): data = call_some_api(idx) # 保留你原有API调用逻辑 dump_data(data, idx) # 改用多线程模式,开销远低于多进程,适配IO密集场景 Parallel(n_jobs=MAX_WORKERS, verbose=50, prefer="threads")( delayed(get_api)(idx) for idx in ids )
注意:已修复原有代码的变量错误:
dump_data入参与实际使用变量不一致的问题,用os.path.join拼接路径适配Windows系统,避免路径符错误。
方案2:更高性能的协程版本(推荐)
API调用属于典型IO密集型任务,用asyncio协程的开销比多线程低很多,同等配置下可以支撑更高并发,处理速度更快。
- 安装依赖
pip install aiohttp pandas aiometer
- 代码示例
import asyncio import aiohttp import pandas as pd from pathlib import Path # 配置参数 BASE_DIR = Path(r"C:\你的实际保存路径") BASE_DIR.mkdir(exist_ok=True) MAX_REQUESTS_PER_SECOND = 5 MAX_CONCURRENT = 50 async def call_some_api_async(session, idx): # 替换为你实际的API调用逻辑,示例为GET请求,POST请求自行调整参数 api_url = f"https://你的API接口地址?id={idx}" async with session.get(api_url) as resp: resp_json = await resp.json() # 根据API实际返回结构调整转DataFrame的逻辑 return pd.DataFrame(resp_json) async def dump_data_async(data, idx): filename = BASE_DIR / f"{idx}.csv" data.to_csv(filename, header=True, index=False) async def process_id(session, idx): data = await call_some_api_async(session, idx) await dump_data_async(data, idx) async def main(): ids = pd.read_csv('data.csv')['Id'].values.tolist() async with aiohttp.ClientSession() as session: # aiometer自带QPS限流+并发数控制,无需额外实现限流逻辑 async with aiometer.amap( lambda idx: process_id(session, idx), ids, max_at_once=MAX_CONCURRENT, max_per_second=MAX_REQUESTS_PER_SECOND, ) as results: # 可在这里添加异常处理逻辑,记录请求失败的id async for _ in results: pass if __name__ == "__main__": # Windows系统asyncio兼容配置,避免运行报错 asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy()) asyncio.run(main())
内容的提问来源于stack exchange,提问作者Mustard Tiger
相关产品推荐
相关产品推荐

