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

如何用Multiprocessing加速Pandas的groupby/apply操作?

嘿,我明白你想通过多进程加速Pandas里的groupby+apply操作的需求——毕竟大数据量下单进程确实慢得让人头疼!我来给你一步步拆解可行的方案,顺便帮你避开那些容易导致代码失败的坑。

先从测试数据和基础操作说起

先模拟一个和你场景类似的DataFrame,方便后续演示:

import pandas as pd
import numpy as np

# 生成10万行测试数据,包含100个不同的cluster_id
df = pd.DataFrame({
    'cluster_id': np.random.randint(0, 100, size=100000),
    'value1': np.random.randn(100000),
    'value2': np.random.randn(100000)
})

普通单进程的groupby求和操作是这样的:

# 单进程分组求和
result_single = df.groupby('cluster_id')[['value1', 'value2']].sum()
方法一:用multiprocessing手动实现并行分组

核心思路是:把每个cluster的分组任务拆出来,分给不同的进程并行处理,最后合并结果。

完整代码如下:

import multiprocessing as mp

def process_single_cluster(cluster_id):
    # 提取当前cluster对应的子数据集
    group_data = df[df['cluster_id'] == cluster_id]
    # 这里替换成你自己的自定义函数,比如求和
    return pd.Series({
        'cluster_id': cluster_id,
        'value1_sum': group_data['value1'].sum(),
        'value2_sum': group_data['value2'].sum()
    })

if __name__ == '__main__':
    # 获取所有唯一的cluster_id
    unique_clusters = df['cluster_id'].unique()
    
    # 创建进程池,进程数建议设为你的CPU核心数(避免资源浪费)
    with mp.Pool(mp.cpu_count()) as pool:
        # 并行处理每个cluster的任务
        parallel_results = pool.map(process_single_cluster, unique_clusters)
    
    # 将所有进程的结果合并成最终DataFrame
    result_multi = pd.DataFrame(parallel_results).set_index('cluster_id')
方法二:用concurrent.futures更简洁的并行

如果你觉得multiprocessing的语法有点繁琐,可以用concurrent.futures.ProcessPoolExecutor,语法更直观:

from concurrent.futures import ProcessPoolExecutor

def process_single_cluster(cluster_id):
    group_data = df[df['cluster_id'] == cluster_id]
    return pd.Series({
        'cluster_id': cluster_id,
        'value1_sum': group_data['value1'].sum(),
        'value2_sum': group_data['value2'].sum()
    })

if __name__ == '__main__':
    unique_clusters = df['cluster_id'].unique()
    
    with ProcessPoolExecutor(max_workers=mp.cpu_count()) as executor:
        parallel_results = list(executor.map(process_single_cluster, unique_clusters))
    
    result_multi = pd.DataFrame(parallel_results).set_index('cluster_id')
避坑指南:那些容易导致代码失败的问题

你之前代码运行失败,大概率是踩了这些坑:

  • 必须加if __name__ == '__main__'::Windows系统下,多进程启动时会重新导入模块,如果没有这个判断,会无限创建进程导致崩溃;Linux/macOS虽然不会崩溃,但加上更规范。
  • 内存爆炸问题:如果你的DataFrame特别大,每个进程都会复制一份完整的DataFrame,内存会直接飙高。解决方法:可以把DataFrame按cluster拆分后存到文件,让每个进程读取对应文件;或者用dask.dataframe这类专门的并行数据框架。
  • 函数序列化问题:如果你的自定义函数里包含无法被Python序列化的对象(比如某些自定义类的实例),多进程会抛出序列化错误。解决方法:把函数改成纯函数(只依赖输入参数,不引用外部非序列化对象),或者用cloudpickle替代默认的序列化器。
  • 小分组效率问题:如果你的cluster数量极多但每个cluster的数据极少,多进程的启动和数据传递开销会超过并行带来的收益,反而单进程更快。
更省心的工具:用swifter自动并行

如果你不想手动处理多进程的细节,可以用swifter库——它会自动判断你的任务适合单进程还是多进程,语法和普通apply几乎一样:

# 先安装swifter:pip install swifter
import swifter

def custom_sum(group):
    return group[['value1', 'value2']].sum()

# 自动选择最优执行方式(单/多进程)
result_swifter = df.groupby('cluster_id').swifter.apply(custom_sum)

最后可以验证一下并行结果和单进程结果是否一致:

print(result_single.equals(result_multi))  # 输出True说明结果一致

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 09:16:56