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

如何基于活跃工作线程数动态调整API下载任务的休眠时间?

动态调整多线程下载的休眠时间适配API速率限制

核心思路

用线程安全的计数器跟踪当前活跃工作线程数,每次需要休眠时根据当前活跃数动态计算休眠时间,既保证总请求速率不超过API限制(每分钟200次),又能在部分线程完成任务后,让剩余线程自动缩短休眠时间、提升下载速度。

修改后的完整代码

import os
import pandas as pd
import concurrent.futures
import threading
import time

# 待下载的表列表
list_of_travelperks_tables=['users','trips','invoices','bookings','suppliers']

# 存储结果的容器
dfs = []
dfs_dict = {}

# 配置参数
num_workers = 5
rate_limit = 60/200  # 单请求基准间隔,保证总速率不超200次/分钟

# 线程安全的活跃线程计数器
active_workers = 0
worker_lock = threading.Lock()

def load_and_convert_table(table_name):
    global active_workers
    # 线程启动时增加活跃计数
    with worker_lock:
        active_workers += 1
    try:
        # 传递速率限制、锁和获取活跃数的方法
        data = load_travelperk_table(
            table_name, 
            os.environ['TP_API_KEY'], 
            rate_limit, 
            worker_lock, 
            lambda: active_workers
        )
        df = pd.json_normalize(data)
        return table_name, df
    finally:
        # 线程结束时减少活跃计数
        with worker_lock:
            active_workers -= 1

def load_travelperk_table(table_name: str, api_key: str, rate_limit: float, worker_lock: threading.Lock, get_active_workers):
    # 替换为你的实际API下载逻辑
    data = []
    # 示例:分页请求场景下动态计算休眠时间
    page = 1
    has_more = True
    while has_more:
        # 线程安全地获取当前活跃线程数
        with worker_lock:
            current_active = get_active_workers()
        # 动态计算休眠时间:总速率 = 活跃数 / 休眠时间 = 1/rate_limit(符合200次/分钟限制)
        time_sleep = rate_limit * current_active
        time.sleep(time_sleep)
        
        # 实际API请求代码示例
        # response = requests.get(
        #     f"https://api.example.com/{table_name}",
        #     headers={"Authorization": f"Bearer {api_key}"},
        #     params={"page": page}
        # )
        # page_data = response.json()
        # data.extend(page_data)
        # has_more = page_data.get("has_more", False)
        # page += 1
        
        # 模拟分页结束,避免死循环
        has_more = False
    return data

# 启动线程池执行任务
with concurrent.futures.ThreadPoolExecutor(max_workers=num_workers) as executor:
    futures = executor.map(load_and_convert_table, list_of_travelperks_tables)
    for table_name, df in futures:
        df_name = f'df_{table_name}'
        dfs_dict[df_name] = df
        dfs.append((df_name, df))

关键说明

  1. 线程安全计数:用threading.Lock保护active_workers变量,避免多线程同时读写导致的计数错误。
  2. 动态休眠计算:每次请求前获取当前活跃线程数,休眠时间 = 基准间隔 × 活跃线程数,确保总请求速率稳定在每分钟200次(1/rate_limit = 200/60次/秒)。
  3. 自动适配速率:当某个线程完成任务后,active_workers计数减少,剩余线程的休眠时间自动缩短,单个线程的请求频率提升,从而加快剩余表的下载速度。
  4. 异常安全:用finally块保证线程结束时一定会减少活跃计数,避免计数异常。

内容的提问来源于stack exchange,提问作者user15915737

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 16:27:51