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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 20:43:12