Python concurrent.futures.ProcessPoolExecutor内存耗尽崩溃问题
问题原因
- 进程内存叠加:
ProcessPoolExecutor基于多进程实现,每个子进程通过fork机制复制父进程的内存空间。你手动调用large_df.copy()生成副本,再加上子进程默认复制的父进程large_df,等于每个进程持有两份大DataFrame;若CPU为16核,默认max_workers=16,16个进程的内存占用直接叠加,16GB内存很快就会耗尽。 - map方法预加载过载:
executor.map会一次性生成所有参数对(包括repeat生成的所有DataFrame副本),当parameters数量极大时,这些副本会提前占用大量内存,直接触发内存溢出崩溃。 - 写时复制失效:如果
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
相关产品推荐
相关产品推荐

