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

多线程下载任务:实现失败重试与父目录校验机制

下载Worker重试与父目录校验方案

1. 扩展Item类

新增重试次数属性,用于跟踪任务重试情况:

class Item:
    url: str
    download_path: str
    is_downloaded: bool
    retry_count: int = 0  # 记录当前任务的重试次数

2. 线程安全的共享状态管理

多线程环境下,用锁保护共享的失败计数和失效目录集合:

import threading

# 记录每个父目录的累计失败次数,key为父目录名,value为失败次数
dir_fail_counts = {}
# 标记已确认不存在的父目录,后续该目录下的任务直接标记失败
invalid_dirs = set()
# 保护共享资源的线程锁
dir_lock = threading.Lock()

3. 辅助工具函数

实现父目录提取和目录存在性检查:

import os

def get_parent_dir(download_path: str) -> str:
    """从下载路径中提取目标父目录名称(如路径为"xxx/bam/xxx.file"则返回"bam")"""
    path_parts = os.path.normpath(download_path).split(os.sep)
    return path_parts[-2] if len(path_parts) >= 2 else ""

def is_parent_dir_valid(download_path: str) -> bool:
    """检查下载路径的父目录是否实际存在"""
    parent_path = os.path.dirname(download_path)
    return os.path.isdir(parent_path)

4. 修改后的Worker核心逻辑

整合重试、父目录校验和失效目录拦截逻辑:

from queue import Queue
import urllib.request as request
import urllib.error as error

SENTINEL = "END"

def download_item(to_download: Queue, downloaded: Queue):
    item: Item = to_download.get()
    while item != SENTINEL:
        parent_dir = get_parent_dir(item.download_path)
        
        # 先拦截已失效目录下的任务
        with dir_lock:
            if parent_dir in invalid_dirs:
                item.is_downloaded = False
                downloaded.put(item)
                item = to_download.get()
                continue
        
        try:
            response = request.urlopen(item.url)
            # 下载成功,写入文件并标记状态
            write_to_file(item.download_path, content=response.read())
            item.is_downloaded = True
            downloaded.put(item)
        except error.HTTPError as e:
            item.retry_count += 1
            # 重试次数未达上限,重新放入队列等待重试
            if item.retry_count < 2:
                to_download.put(item)
            else:
                # 重试超过2次,更新父目录失败计数并检查目录有效性
                with dir_lock:
                    dir_fail_counts[parent_dir] = dir_fail_counts.get(parent_dir, 0) + 1
                    # 累计失败超2次时检查目录是否存在
                    if dir_fail_counts[parent_dir] > 2:
                        if not is_parent_dir_valid(item.download_path):
                            invalid_dirs.add(parent_dir)
                # 标记当前任务失败
                item.is_downloaded = False
                downloaded.put(item)
        finally:
            item = to_download.get()
    
    # 将哨兵放回队列,传递结束信号给其他线程
    to_download.put(SENTINEL)

5. 线程池执行示例

基于ThreadPoolExecutor启动50个Worker线程:

from concurrent.futures import ThreadPoolExecutor

def main():
    to_download_q = Queue()
    downloaded_q = Queue()
    
    # 填充待下载任务(示例,需替换为实际任务)
    # for task_url, task_path in your_task_list:
    #     to_download_q.put(Item(url=task_url, download_path=task_path, is_downloaded=False))
    # 放入哨兵,标记任务队列结束
    to_download_q.put(SENTINEL)
    
    # 启动50个Worker线程
    with ThreadPoolExecutor(max_workers=50) as executor:
        for _ in range(50):
            executor.submit(download_item, to_download_q, downloaded_q)
    
    # 处理下载完成的任务(示例)
    while True:
        result_item = downloaded_q.get()
        if result_item == SENTINEL:
            break
        # 此处可添加结果处理逻辑,如日志记录、状态统计等
        print(f"任务[{result_item.url}] 下载状态: {'成功' if result_item.is_downloaded else '失败'}")

if __name__ == "__main__":
    main()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 20:33:09