如何在API速率限制下并行处理DataFrame行值并存入结果列
最优实现方案
process_file()是同步阻塞函数的情况下,不要硬套asyncio,强行把同步函数塞到事件循环里只会阻塞整个循环,没有任何性能收益。针对外部API调用这种IO密集型场景,用线程池+令牌桶限流是性价比最高的方案:既能并行发起请求充分利用IO等待时间,又能严格卡准速率上限不触发API封禁,代码改动量也最小。
核心实现代码
import pandas as pd import time from concurrent.futures import ThreadPoolExecutor, as_completed from threading import Lock # 无需修改原有外部函数的源码 def process_file(file_path): # 原有业务逻辑保持不变 pass # 线程安全的令牌桶限流器,请求分发更平滑,无瞬时突刺 class RateLimiter: def __init__(self, max_qps): self.max_qps = max_qps self.remain_tokens = max_qps self.last_refill_ts = time.time() self.lock = Lock() def acquire(self): with self.lock: now = time.time() time_gap = now - self.last_refill_ts # 补充时间窗口内的可用令牌 self.remain_tokens = min(self.max_qps, self.remain_tokens + time_gap * self.max_qps) self.last_refill_ts = now if self.remain_tokens < 1: # 令牌不足时等待对应时长 wait_time = (1 - self.remain_tokens) / self.max_qps time.sleep(wait_time) self.remain_tokens = 0 self.last_refill_ts = time.time() else: self.remain_tokens -= 1 if __name__ == "__main__": # 示例DataFrame初始化 d1 = {1:['Test','Test1','Test2'], 2:['file1','file2','file3'],3:[pd.NA,pd.NA,pd.NA]} df = pd.DataFrame(data=d1) # 初始化限流器,建议设为比上限低1-2的数值留冗余,避免网络波动导致瞬时超限 limiter = RateLimiter(max_qps=19) # 线程数和QPS上限匹配即可,无需开太大 worker_num = 20 res_map = {} def task_wrapper(row_idx, file_val): limiter.acquire() return row_idx, process_file(file_val) with ThreadPoolExecutor(max_workers=worker_num) as executor: all_tasks = [ executor.submit(task_wrapper, idx, row[2]) for idx, row in df.iterrows() ] # 按完成顺序收集结果 for future in as_completed(all_tasks): idx, res = future.result() res_map[idx] = res # 按原始行索引回填结果,避免并行返回顺序错乱导致值错位 df[3] = df.index.map(res_map)
关键说明
- 不要用串行循环加固定间隔sleep的方案:这种方式必须等上一个请求返回才会发下一个,实际QPS会远低于20的上限,根本跑不满配额,效率极差。
- 选线程池而非进程池的原因:外部API调用属于IO等待型操作,线程在等待API响应时会释放GIL,不会触发Python多线程的GIL性能瓶颈;相比多进程,线程不需要做数据跨进程序列化,启动和运行开销低很多。
- 选令牌桶限流而非固定窗口限流的原因:固定窗口容易出现窗口边界瞬时发超2倍请求的问题,也会出现窗口前几十毫秒发完所有配额、剩下时间空等的情况,令牌桶可以把请求均匀分散在整秒内,速率控制更平滑。
- 结果回填必须用原始索引做映射:并行任务的返回顺序和提交顺序不一定一致,直接按返回顺序赋值会出现列值和行不匹配的问题。
可选调整方向
- 如果
process_file()内部包含大量CPU密集计算逻辑,可以把ThreadPoolExecutor换成ProcessPoolExecutor,进程数设置为和CPU核心数一致即可,限流逻辑不需要改动。 - 如果需要加失败重试逻辑,重试时必须重新调用
limiter.acquire()获取令牌,避免重试请求瞬间打满触发API限流。 - 如果目标API支持批量请求,可以调整逻辑把单条入参改成批量入参,同步降低限流阈值即可,能进一步提升处理效率。
内容的提问来源于stack exchange,提问作者splotsh
相关产品推荐
相关产品推荐

