如何使用multiprocessing.Pool.map结合生成器处理100GB+大文件?
解决多进程处理大文件时的内存溢出问题
问题根源
multiprocessing.Pool.map会将传入的可迭代对象(比如生成器)一次性转换为列表并加载到内存中,这就是处理100GB+文件时出现Memory Error的核心原因。
解决方案
1. 使用Pool.imap()或Pool.imap_unordered()
这两个方法不会一次性加载所有数据,而是分批迭代处理生成器的输出,同时支持设置chunksize参数优化进程间通信开销:
imap():保持输出结果与输入顺序一致imap_unordered():不保证结果顺序,但处理速度可能更快
示例代码:
import multiprocessing def generate_data(): with open('data.file') as data_file: for line in data_file: yield line.strip() # 预处理减少传输数据量 def process(data): # 替换为你的实际处理逻辑 return len(data) def run(data): return process(data) if __name__ == '__main__': pool = multiprocessing.Pool() # 根据CPU核心数和内存调整chunksize,建议设为1000~10000 chunksize = 1000 # 迭代处理结果,避免缓存所有结果 for result in pool.imap(run, generate_data(), chunksize=chunksize): # 按需处理单个结果,比如写入输出文件 pass pool.close() pool.join()
2. 手动分块处理数据
将大文件按固定行数分块,每次向进程池提交一个数据块而非单行数据,进一步减少进程间通信次数,提升效率。
示例代码:
import multiprocessing def generate_data_chunks(chunk_size=1000): with open('data.file') as data_file: chunk = [] for line in data_file: chunk.append(line.strip()) if len(chunk) == chunk_size: yield chunk chunk = [] # 处理剩余不足一个chunk的数据 if chunk: yield chunk def process_chunk(chunk): # 批量处理整个数据块的逻辑 return [len(line) for line in chunk] if __name__ == '__main__': pool = multiprocessing.Pool() for chunk_results in pool.imap(process_chunk, generate_data_chunks()): # 处理每个数据块的结果 pass pool.close() pool.join()
3. 使用apply_async()手动提交任务
通过循环逐个提交生成器元素到进程池,配合回调函数处理结果,避免缓存大量任务。
示例代码:
import multiprocessing def generate_data(): with open('data.file') as data_file: for line in data_file: yield line.strip() def process(data): # 实际处理逻辑 return len(data) def collect_result(result): # 回调函数:处理单个任务的返回结果 pass if __name__ == '__main__': pool = multiprocessing.Pool() for item in generate_data(): pool.apply_async(process, args=(item,), callback=collect_result) pool.close() pool.join()
注意事项
- 尽量在生成器中完成预处理(如去除空白、过滤无效行),减少进程间传输的数据量
chunksize不宜过大或过小:过小会增加通信开销,过大可能导致单个进程占用过多内存- 如果不需要保持结果顺序,优先使用
imap_unordered()提升处理速度
内容的提问来源于stack exchange,提问作者tsu rugi
相关产品推荐
相关产品推荐

