如何用os.scandir遍历大目录并将文件加入多进程队列并行处理
实现递归遍历目录+多进程并行处理文件的正确方式
我来帮你搞定这个并行处理超大目录文件的需求~下面是完整的实现方案,结合你的初始代码和多进程队列的最佳实践:
完整代码实现
import os import multiprocessing import hashlib def doWork(queue): """处理文件的工作函数:计算MD5及其他校验和""" while True: file_path = queue.get() # 收到结束信号时退出循环 if file_path is None: break try: # 计算MD5示例(大文件分块读取,避免内存占用过高) md5_hash = hashlib.md5() with open(file_path, "rb") as f: # 以4KB为块读取文件 for chunk in iter(lambda: f.read(4096), b""): md5_hash.update(chunk) print(f"Processed {file_path}, MD5: {md5_hash.hexdigest()}") # 在这里添加其他校验和的计算逻辑(比如SHA-1、SHA-256等) # sha1_hash = hashlib.sha1() # ... 类似MD5的处理逻辑 ... except Exception as e: # 捕获处理文件时的异常,避免单个文件失败导致进程崩溃 print(f"Error processing {file_path}: {str(e)}") def scan_directory(path, queue): """递归遍历目录,将符合条件的文件路径加入队列""" with os.scandir(path) as it: for entry in it: # 跳过隐藏文件/目录 if entry.name.startswith('.'): continue if entry.is_file(): # 将文件路径放入队列 queue.put(entry.path) elif entry.is_dir(): # 递归遍历子目录 scan_directory(entry.path, queue) def main(): PATH = "/mnt/large_directory" print(f"Starting the scanner in root {PATH}") # 初始化进程安全的队列 file_queue = multiprocessing.Queue() # 确定工作进程数量:建议使用CPU核心数,平衡性能和资源占用 num_workers = multiprocessing.cpu_count() workers = [] # 启动所有工作进程 for _ in range(num_workers): worker_process = multiprocessing.Process(target=doWork, args=(file_queue,)) worker_process.start() workers.append(worker_process) # 开始扫描目录(包裹在try-finally里,确保无论扫描是否出错都能正确关闭进程) try: scan_directory(PATH, file_queue) finally: # 给每个工作进程发送结束信号(None) for _ in range(num_workers): file_queue.put(None) # 等待所有工作进程完成任务 for worker in workers: worker.join() print("All files have been processed successfully!") if __name__ == "__main__": # Windows系统下必须加这个判断,避免进程重复创建 main()
关键细节解释
- 进程安全队列:用
multiprocessing.Queue而非普通的queue.Queue,因为前者是专门为多进程通信设计的,保证了数据传递的安全性。 - 结束信号机制:当目录扫描完成后,给每个工作进程发送一个
None,让工作进程知道没有更多任务,避免无限阻塞在queue.get()上。 - 高效目录遍历:使用
os.scandir()比os.listdir()更高效,它能直接获取文件的属性(是否是文件/目录),不需要额外调用os.path.isfile()或os.path.isdir()。 - 大文件处理:计算校验和时采用分块读取的方式,避免一次性加载大文件到内存,防止内存溢出。
- 异常处理:在
doWork函数中添加异常捕获,确保单个文件处理失败不会导致整个工作进程崩溃,其余文件仍能正常处理。 - 进程数量设置:使用
multiprocessing.cpu_count()获取CPU核心数作为工作进程数量,能最大化利用系统资源,同时避免进程过多导致的上下文切换开销。
内容的提问来源于stack exchange,提问作者ensnare
相关产品推荐
相关产品推荐

