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

Python concurrent.futures.ProcessPoolExecutor内存耗尽崩溃问题

问题原因
  1. 进程内存叠加:ProcessPoolExecutor 基于多进程实现,每个子进程通过fork机制复制父进程的内存空间。你手动调用large_df.copy()生成副本,再加上子进程默认复制的父进程large_df,等于每个进程持有两份大DataFrame;若CPU为16核,默认max_workers=16,16个进程的内存占用直接叠加,16GB内存很快就会耗尽。
  2. map方法预加载过载:executor.map会一次性生成所有参数对(包括repeat生成的所有DataFrame副本),当parameters数量极大时,这些副本会提前占用大量内存,直接触发内存溢出崩溃。
  3. 写时复制失效:如果func中对传入的DataFrame做了修改操作,会触发写时复制(COW)机制,原本共享的内存页会被实际复制,进一步加剧内存占用。
解决方案

1. 子进程内单次加载DataFrame(最易实现)

放弃在父进程传递DataFrame副本,改为每个子进程启动时仅加载一次DataFrame,避免父进程提前生成大量副本,同时减少重复读取文件的开销。

修改后核心代码:

import concurrent.futures
import pandas as pd

from my_tests import func

parameters = [
    (arg1, arg2, arg3),
    ...
]
csv_path = "your_data.csv"

# 进程初始化函数:每个子进程启动时加载一次DataFrame
def init_worker(path):
    global worker_df
    worker_df = pd.read_csv(path)

# 修改func,直接使用全局的worker_df
def func(params):
    global worker_df
    # 原func的计算逻辑...
    return result

with concurrent.futures.ProcessPoolExecutor(
    initializer=init_worker, 
    initargs=(csv_path,)
) as executor:
    for test_result in executor.map(func, parameters):
        # 处理结果...

2. 用共享内存共享DataFrame(内存效率最高)

利用pandas 1.3.0+支持的shareable_memory特性,将DataFrame存入共享内存,所有子进程直接读取共享内存中的数据,无需复制,内存占用几乎和单进程一致。

核心代码:

import concurrent.futures
import pandas as pd

from my_tests import func

parameters = [
    (arg1, arg2, arg3),
    ...
]
large_df = pd.read_csv(csv_path)

# 将DataFrame存入共享内存
sm = large_df.to_shared_memory()
sm_name = sm.name

def func_shared(params, sm_name):
    # 从共享内存加载DataFrame
    df = pd.read_shared_memory(sm_name)
    # 原func的计算逻辑...
    return result

with concurrent.futures.ProcessPoolExecutor() as executor:
    results = executor.map(func_shared, parameters, repeat(sm_name))
    for test_result in results:
        # 处理结果...

# 任务结束后清理共享内存
sm.unlink()

3. 分批提交任务(降低瞬时内存占用)

如果必须在父进程传递DataFrame,不要一次性提交所有任务,而是分批处理,每批任务完成后释放对应内存。

核心代码:

import concurrent.futures
import pandas as pd

from my_tests import func

parameters = [
    (arg1, arg2, arg3),
    ...
]
large_df = pd.read_csv(csv_path)
batch_size = 8  # 根据内存情况调整每批任务数

with concurrent.futures.ProcessPoolExecutor(max_workers=8) as executor:
    # 拆分参数为多个批次
    for i in range(0, len(parameters), batch_size):
        batch_params = parameters[i:i+batch_size]
        # 为当前批次生成对应数量的DataFrame副本
        dfs = [large_df.copy() for _ in batch_params]
        # 提交批次任务并处理结果
        for test_result in executor.map(func, dfs, batch_params):
            # 处理结果...
        # 批次完成后,dfs会被垃圾回收,释放内存

4. 优化func内部内存使用

检查func内部是否存在不必要的DataFrame副本,比如仅保留计算所需的列,避免生成冗余副本:

def func(df, params):
    # 只保留计算需要的列,大幅减少内存占用
    df = df[['col1', 'col2', 'col3']]
    # 后续计算逻辑...
    return result

内容的提问来源于stack exchange,提问作者janboro

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 13:26:01