如何用多进程解析大型Zip文件?解决重复打开耗时问题
大型Zip文件多进程解析优化方案
问题背景
我有一个包含大量文件的大型Zip文件,解析所有文件耗时很久,因此考虑使用多进程来提速。但不确定实现方式,因为Python的zipfile.ZipFile并非可迭代对象。
我不想先解压所有内容占用额外存储空间,希望直接操作ZipFile对象,也接受其他可行方案。
现有代码的问题
以下代码可运行,但每次执行
get_content()时都会重新打开大型Zip文件,导致单个文件读取耗时长达15秒:
import multiprocessing from zipfile import ZipFile from multiprocessing import Pool import time path = 'zipfile.zip' def get_file_list(zip_path): with ZipFile(zip_path, 'r') as zipObj: listOfiles = zipObj.namelist() return listOfiles def get_content(file_name): start_time = time.time() with ZipFile(path, 'r') as zipObject: with zipObject.open(file_name) as file: content = file.read() end_time = time.time() print(f"It took {end_time - start_time} to open this file") return content def parse_files(): file_list = get_file_list(path) with Pool(multiprocessing.cpu_count()) as p: contents = p.map(get_content, file_list) print(contents) parse_files()
解决方案
方案一:共享文件描述符(避免重复打开Zip)
通过父进程预先打开Zip文件,将文件描述符共享给子进程,子进程复制fd后重新构建ZipFile对象,彻底避免重复打开大文件的开销:
import multiprocessing from zipfile import ZipFile import time import os path = 'zipfile.zip' def get_file_list(zip_path): with ZipFile(zip_path, 'r') as zipObj: return zipObj.namelist() def get_content(args): file_name, fd = args start_time = time.time() # 复制父进程的文件描述符,构建新的文件对象 with os.fdopen(os.dup(fd), 'rb') as f: with ZipFile(f, 'r') as zipObject: with zipObject.open(file_name) as file: content = file.read() end_time = time.time() print(f"读取文件 {file_name} 耗时 {end_time - start_time:.2f} 秒") return content def parse_files(): file_list = get_file_list(path) # 父进程打开Zip文件,获取文件描述符 with open(path, 'rb') as f: fd = f.fileno() with multiprocessing.Pool(multiprocessing.cpu_count()) as p: # 将文件名与共享fd打包为参数传递 contents = p.map(get_content, [(fname, fd) for fname in file_list]) print(f"共解析 {len(contents)} 个文件") if __name__ == "__main__": parse_files()
方案二:多线程优化(IO密集型场景优先)
如果解析工作以IO读取为主、CPU计算量小,多线程是更轻量的选择。虽然ZipFile本身非线程安全,但每个线程独立打开Zip文件的开销远低于多进程重复打开:
import threading from zipfile import ZipFile import time from queue import Queue path = 'zipfile.zip' def worker(q, results): # 每个线程独立打开Zip文件 with ZipFile(path, 'r') as zipObject: while not q.empty(): file_name = q.get() start_time = time.time() with zipObject.open(file_name) as file: content = file.read() end_time = time.time() print(f"读取文件 {file_name} 耗时 {end_time - start_time:.2f} 秒") results.append(content) q.task_done() def parse_files(): with ZipFile(path, 'r') as zipObj: file_list = zipObj.namelist() # 任务队列 q = Queue() for fname in file_list: q.put(fname) results = [] # 线程数取CPU核心数或文件数的较小值 thread_count = min(threading.cpu_count(), len(file_list)) for _ in range(thread_count): t = threading.Thread(target=worker, args=(q, results)) t.daemon = True t.start() q.join() print(f"共解析 {len(results)} 个文件") parse_files()
方案三:内存预读取(小体积Zip优先)
若Zip文件大小在内存可承受范围内,可将整个Zip文件读入内存,子进程直接操作内存中的字节流,彻底消除磁盘IO开销:
import multiprocessing from zipfile import ZipFile import time import io path = 'zipfile.zip' def get_file_list(zip_path): with ZipFile(zip_path, 'r') as zipObj: return zipObj.namelist() def get_content(args): file_name, zip_data = args start_time = time.time() # 从内存字节流初始化ZipFile with ZipFile(io.BytesIO(zip_data), 'r') as zipObject: with zipObject.open(file_name) as file: content = file.read() end_time = time.time() print(f"读取文件 {file_name} 耗时 {end_time - start_time:.2f} 秒") return content def parse_files(): file_list = get_file_list(path) # 父进程读取整个Zip文件到内存 with open(path, 'rb') as f: zip_data = f.read() with multiprocessing.Pool(multiprocessing.cpu_count()) as p: contents = p.map(get_content, [(fname, zip_data) for fname in file_list]) print(f"共解析 {len(contents)} 个文件") if __name__ == "__main__": parse_files()
内容的提问来源于stack exchange,提问作者Aleister Crowley
相关产品推荐
相关产品推荐

