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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 07:06:14