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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 01:27:25