多线程下载任务:实现失败重试与父目录校验机制
下载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
相关产品推荐
相关产品推荐

