如何用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
相关产品推荐
相关产品推荐

