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

Python多线程非阻塞下载遇阻塞及API限流问题求助

问题描述

我编写了一个接收订单ID列表的函数,针对每个订单ID检查状态:若状态为success则下载对应文件,否则10秒后重新检查,直至所有订单文件下载完成。API密钥每秒最多允许5次调用,因此使用了5线程的ThreadPoolExecutor。

当前遇到的问题:

  1. 文件写入时程序阻塞,无法同时处理其他订单的状态检查或下载;
  2. 处理偏同步,需等前5个订单完成才处理后续订单,希望在遵守API限流的前提下,让所有订单的状态检查和下载并行执行;
  3. check_state函数中的print(f"Order ID: {order_id} Status: {state}", end="\r")语句会覆盖之前的订单ID,希望能动态更新每个订单的状态且不互相覆盖。

原始代码:

def download_orders(api_key, headers, activated_orders_ids):
        def check_state(order_id):
            url = f"https://api.com/{order_id}"

            headers = {"Authorization": f"api-key {api_key}"}
            state = requests.get(url, headers=headers).json()["state"]
            print(f"Order ID: {order_id} Status: {state}", end="\r")
            return state

        def get_download_link(order_id):
            download_url = f"https://api.com/{order_id}"

            headers = {"Authorization": f"api-key {api_key}"}

            response = requests.get(download_url, headers=headers)
            return response.json()["_links"]["results"][0]["location"], response.json()["name"]

        def save_file(download_link, order_name):
            with requests.Session() as s:
                r = s.get(download_link, stream=True)
                total_size = int(r.headers.get("content-length", 0))
                block_size = 100 * 1024  # You can adjust this to your preferred block size.

                with open(f"{order_name}.zip", mode="wb") as file, tqdm(
                    desc=f"Downloading {order_name}.zip",
                    total=total_size,
                    unit="B",
                    unit_scale=True,
                    unit_divisor=1024,
                ) as pbar:
                    for data in r.iter_content(chunk_size=block_size):
                        file.write(data)
                        pbar.update(len(data))

        def download_logic(order_id):
            state = check_state(order_id)
            while state != "success":
                state = check_state(order_id)
                time.sleep(10)

            download_link, order_name = get_download_link(order_id)
            save_file(download_link, order_name)
            time.sleep(1)

        # Use multi-threading for parallel execution
        with concurrent.futures.ThreadPoolExecutor(5) as executor:
            executor.map(download_logic, activated_orders_ids)


download_orders(api_key, headers, order_ids)
解决方案

核心优化思路

  1. 拆分线程池:将API请求(状态检查、获取下载链接)和文件下载写入分离到两个独立线程池,避免IO操作占用API限流的线程;
  2. 改用submit+as_completed:一次性提交所有订单任务,线程空闲时立即处理下一个订单,无需等待前一批完成;
  3. 日志化状态输出:用logging替代直接print,每条状态更新单独成行,避免覆盖;
  4. 增加异常处理:防止单个订单失败导致整个程序崩溃,同时便于排查问题。

优化后完整代码

import logging
from concurrent.futures import ThreadPoolExecutor, as_completed
import requests
import time
from tqdm import tqdm

# 配置日志,每条状态更新单独输出一行
logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s - Order %(message)s",
    handlers=[logging.StreamHandler()]
)

def download_orders(api_key, headers, activated_orders_ids):
    # 维护订单状态字典,用于跟踪每个订单的最新状态
    order_states = {}

    def check_state(order_id):
        url = f"https://api.com/{order_id}"
        headers = {"Authorization": f"api-key {api_key}"}
        try:
            response = requests.get(url, headers=headers)
            response.raise_for_status()  # 捕获HTTP请求错误
            state = response.json()["state"]
            order_states[order_id] = state
            logging.info(f"{order_id}: Status updated to {state}")
            return state
        except Exception as e:
            logging.error(f"{order_id}: Failed to check state - {str(e)}")
            return "error"

    def get_download_link(order_id):
        download_url = f"https://api.com/{order_id}"
        headers = {"Authorization": f"api-key {api_key}"}
        try:
            response = requests.get(download_url, headers=headers)
            response.raise_for_status()
            data = response.json()
            download_link = data["_links"]["results"][0]["location"]
            order_name = data["name"]
            return download_link, order_name
        except Exception as e:
            logging.error(f"{order_id}: Failed to get download link - {str(e)}")
            raise

    def save_file(download_link, order_name):
        try:
            with requests.Session() as s:
                r = s.get(download_link, stream=True)
                r.raise_for_status()
                total_size = int(r.headers.get("content-length", 0))
                block_size = 100 * 1024

                with open(f"{order_name}.zip", mode="wb") as file, tqdm(
                    desc=f"Downloading {order_name}.zip",
                    total=total_size,
                    unit="B",
                    unit_scale=True,
                    unit_divisor=1024,
                    leave=True  # 下载完成后保留进度条,避免被覆盖
                ) as pbar:
                    for data in r.iter_content(chunk_size=block_size):
                        file.write(data)
                        pbar.update(len(data))
            logging.info(f"{order_name}.zip downloaded successfully")
        except Exception as e:
            logging.error(f"Failed to save {order_name}.zip - {str(e)}")
            raise

    def download_logic(order_id):
        # 循环检查状态直到成功
        state = check_state(order_id)
        while state != "success":
            time.sleep(10)
            state = check_state(order_id)
        
        # 获取下载链接后,提交到文件线程池处理
        download_link, order_name = get_download_link(order_id)
        file_executor.submit(save_file, download_link, order_name)

    # 初始化两个线程池:API池遵守限流,文件池处理IO操作
    api_executor = ThreadPoolExecutor(max_workers=5)
    file_executor = ThreadPoolExecutor(max_workers=10)

    # 一次性提交所有订单的API处理任务
    futures = [api_executor.submit(download_logic, order_id) for order_id in activated_orders_ids]

    # 等待所有API相关任务完成
    for future in as_completed(futures):
        try:
            future.result()
        except Exception as e:
            logging.error(f"Order processing failed: {str(e)}")

    # 等待所有文件下载任务完成后关闭线程池
    file_executor.shutdown(wait=True)
    api_executor.shutdown(wait=True)

# 调用示例
# download_orders(api_key, headers, order_ids)

优化点说明

  • 线程池拆分:api_executor(5线程)专门处理API请求,严格遵守限流;file_executor(10线程)处理文件下载写入,不占用API线程,避免阻塞其他订单的状态检查;
  • 任务提交方式:用submit替代map,所有订单任务一次性提交,线程空闲时立即处理下一个,无需等待前一批完成;
  • 状态输出优化:用logging输出状态,每条日志带时间戳且单独成行,彻底解决覆盖问题;
  • 异常处理增强:增加response.raise_for_status()捕获HTTP错误,同时捕获其他异常并打印日志,避免单个订单失败影响全局;
  • 进度条保留:设置tqdm的leave=True,下载完成后进度条保留,便于查看历史下载记录。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 09:40:57