如何用multiprocessing.Pool高效并行处理大型Pandas DataFrame分组?
并行处理Pandas分组DataFrame的内存优化方案
问题根源
你的内存耗尽问题主要来自两个核心原因:
- 全局变量的进程复制:在Windows的
spawn模式下,每个子进程会重新导入主模块,若将df_large设为全局变量,每个进程都会加载一份完整的大DataFrame,内存占用直接乘以进程数。即使在Linux/macOS的fork模式下,写时复制机制也会因主/子进程对DataFrame的修改操作触发全量复制,导致内存飙升。 - 分组DataFrame的隐式引用:
groupby迭代出的分组DataFrame默认是原大DataFrame的视图,传递给子进程时,序列化过程可能附带原大DataFrame的引用,导致子进程间接持有全量数据;同时主进程若同时保留所有分组的视图,也会占用大量内存。
可行的优化方案
方案1:子进程从磁盘按需加载分组(推荐,适合超大型DataFrame)
让每个子进程仅加载自己需要的分组数据,完全避免主进程传递大对象或全局变量复制的问题。
步骤:
- 主进程仅提取所有唯一的
series_id,不保留全量DataFrame - 子进程根据
series_id从原始数据文件(CSV/Parquet)中过滤加载对应分组
示例代码:
import pandas as pd from multiprocessing import Pool def process_series(series_id): # 以Parquet为例,支持直接过滤加载,效率远高于CSV df_group = pd.read_parquet("large_data.parquet", filters=[("series_id", "=", series_id)]) # 处理逻辑:清洗、去除异常值等 df_group = df_group.dropna(subset=["value_column"]) df_group = df_group[df_group["value_column"].between(0, 100)] # 保存结果 df_group.to_parquet(f"processed_{series_id}.parquet", index=False) if __name__ == "__main__": # 主进程仅读取series_id列,获取唯一值列表 id_df = pd.read_parquet("large_data.parquet", columns=["series_id"]) unique_ids = id_df["series_id"].unique().tolist() id_df = None # 释放主进程内存 num_workers = 4 with Pool(num_workers) as pool: pool.map(process_series, unique_ids)
如果使用CSV,可通过加载后过滤实现:
def process_series(series_id): df_group = pd.read_csv("large_data.csv", dtype={"series_id": str}) df_group = df_group[df_group["series_id"] == series_id].copy() # 后续处理逻辑...
方案2:主进程预生成独立分组副本,再传递给进程池
若主进程内存足以容纳所有分组的独立副本,可先将每个分组转换为独立的DataFrame(避免引用原大DataFrame),再传递给子进程。
示例代码:
import pandas as pd from multiprocessing import Pool def process_series(series_id, df_group): # 处理逻辑 df_clean = df_group.dropna() df_clean.to_csv(f"processed_{series_id}.csv", index=False) if __name__ == "__main__": df_large = pd.read_csv("large_data.csv") # 预生成所有分组的独立副本,切断与原大DataFrame的引用 groups = [(sid, df.copy()) for sid, df in df_large.groupby("series_id")] df_large = None # 立即释放原大DataFrame的内存 num_workers = 4 with Pool(num_workers) as pool: pool.starmap(process_series, groups)
方案3:使用concurrent.futures简化代码
ProcessPoolExecutor的语法更简洁,功能与multiprocessing.Pool一致:
import pandas as pd from concurrent.futures import ProcessPoolExecutor def process_series(series_id, df_group): df_clean = df_group.dropna() df_clean.to_csv(f"processed_{series_id}.csv", index=False) if __name__ == "__main__": df_large = pd.read_csv("large_data.csv") groups = [(sid, df.copy()) for sid, df in df_large.groupby("series_id")] df_large = None num_workers = 4 with ProcessPoolExecutor(max_workers=num_workers) as executor: executor.map(lambda args: process_series(*args), groups)
关键注意事项
- 所有主程序逻辑必须放在
if __name__ == "__main__":块中,避免Windows下spawn模式重复加载代码导致的异常。 - 优先选择从磁盘按需加载的方案,这是控制内存占用最可靠的方式,尤其适合超大型数据集。
- 若使用
fork模式(Linux/macOS),即使使用全局变量,也要确保主进程在创建进程池后不再修改df_large,否则会触发写时复制,导致内存暴涨。
内容的提问来源于stack exchange,提问作者Surányi Olivér
相关产品推荐
相关产品推荐

