并行处理大型Pandas Dataframe:避免进程复制的新方法咨询
大型Pandas DataFrame并行处理的优化方案
你当前代码的核心问题是每个子进程都会完整复制整个DataFrame,导致内存占用过高。以下是近年实用的解决思路和工具:
1. 拆分数据后传递分片
避免让子进程访问全局的df,提前按列拆分数据,只把单列数据传给对应进程:
from joblib import Parallel, delayed def process(col_data): return col_data.apply(another_function) # 提前拆分列数据,避免子进程复制全量DataFrame column_data_list = [df[col] for col in df.columns] results = Parallel(n_jobs=-2)(delayed(process)(data) for data in column_data_list)
2. 利用共享内存复用数据
让多个进程共享同一份DataFrame内存,避免重复复制:
Python 3.8+ 原生共享内存方案
from multiprocessing import shared_memory import pandas as pd # 创建共享内存块 shm = shared_memory.SharedMemory(create=True, size=df.memory_usage().sum()) # 将DataFrame存入共享内存 df_shared = pd.DataFrame(df.values, index=df.index, columns=df.columns) df_shared_array = df_shared.to_numpy() def process(col_name): # 从共享内存加载目标列 col_idx = df.columns.get_loc(col_name) col_data = pd.Series(df_shared_array[:, col_idx], index=df.index) return col_data.apply(another_function) results = Parallel(n_jobs=-2)(delayed(process)(col) for col in df.columns) # 释放共享内存 shm.close() shm.unlink()
3. 专用并行计算工具包
Dask(专为大数据并行设计)
自动拆分DataFrame并分布式处理,无需手动管理进程:
import dask.dataframe as dd # 将Pandas DataFrame转为Dask DataFrame,按CPU核心数设置分区 ddf = dd.from_pandas(df, npartitions=4) # 并行应用函数,最后计算得到结果 results = ddf.apply(another_function, axis=0).compute()
Swifter(自动选择最优处理方式)
自动判断用Pandas串行或Dask并行处理,简化代码:
import swifter # 直接对列应用函数,Swifter会自动适配并行逻辑 results = df.swifter.apply(another_function, axis=0)
4. 进程池优化传递逻辑
用multiprocessing.Pool的initializer在子进程初始化时加载一次数据,复用内存(Linux/macOS下基于fork机制,内存共享效率更高):
from multiprocessing import Pool def init_worker(shared_df): global df_worker df_worker = shared_df def process(col_name): return df_worker[col_name].apply(another_function) # 子进程初始化时加载DataFrame,后续复用 with Pool(processes=-2, initializer=init_worker, initargs=(df,)) as pool: results = pool.map(process, df.columns)
内容的提问来源于stack exchange,提问作者Jun
相关产品推荐
相关产品推荐

