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

如何用multiprocessing.Pool高效并行处理大型Pandas DataFrame分组?

并行处理Pandas分组DataFrame的内存优化方案

问题根源

你的内存耗尽问题主要来自两个核心原因:

  1. 全局变量的进程复制:在Windows的spawn模式下,每个子进程会重新导入主模块,若将df_large设为全局变量,每个进程都会加载一份完整的大DataFrame,内存占用直接乘以进程数。即使在Linux/macOS的fork模式下,写时复制机制也会因主/子进程对DataFrame的修改操作触发全量复制,导致内存飙升。
  2. 分组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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 08:10:29