Python多线程非阻塞下载遇阻塞及API限流问题求助
问题描述
我编写了一个接收订单ID列表的函数,针对每个订单ID检查状态:若状态为success则下载对应文件,否则10秒后重新检查,直至所有订单文件下载完成。API密钥每秒最多允许5次调用,因此使用了5线程的ThreadPoolExecutor。
当前遇到的问题:
- 文件写入时程序阻塞,无法同时处理其他订单的状态检查或下载;
- 处理偏同步,需等前5个订单完成才处理后续订单,希望在遵守API限流的前提下,让所有订单的状态检查和下载并行执行;
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)
解决方案
核心优化思路
- 拆分线程池:将API请求(状态检查、获取下载链接)和文件下载写入分离到两个独立线程池,避免IO操作占用API限流的线程;
- 改用
submit+as_completed:一次性提交所有订单任务,线程空闲时立即处理下一个订单,无需等待前一批完成; - 日志化状态输出:用
logging替代直接print,每条状态更新单独成行,避免覆盖; - 增加异常处理:防止单个订单失败导致整个程序崩溃,同时便于排查问题。
优化后完整代码
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
相关产品推荐
相关产品推荐

