为何多进程未提升我的pandas数据分组聚合代码性能?
为什么Pandas分组聚合的多进程版本比单进程更慢?
我实现了数据分组与聚合的代码,数据中存在重复的userid,按userid进行分组,聚合变量为随机数。编写了单进程和多进程两个版本,但多进程版本运行速度更慢,请问这是为什么?
单进程版本
import time start=time.time() import numpy as np import random import pandas as pd def fun_user_id(start, end, step): num = np.linspace(start, end,(end-start) *int(1/step)+1).tolist() return [round(i, 0) for i in num] def fun_rand_num(): return list(map(lambda x: random.randint(300,800), range(1, 25000001))) def aggregate_df(): grvar=df.groupby("userid") dfvar=grvar.agg(varlistdic) dfvar=dfvar.pipe(lambda x: x.set_axis(x.columns.map('_'.join),axis=1)) dfvar.reset_index(inplace=True) return dfvar if __name__=='__main__': userid=fun_user_id(1,25000001,.5) var1=fun_rand_num() var2=fun_rand_num() var3=fun_rand_num() var4=fun_rand_num() var5=fun_rand_num() var6=fun_rand_num() var7=fun_rand_num() var8=fun_rand_num() var9=fun_rand_num() var10=fun_rand_num() df = pd.DataFrame(list(zip(userid,var1, var2,var3,var4,var5,var6,var7,var8,var9,var10)), columns =['userid','var1', 'var2','var3','var4','var5','var6', 'var7','var8','var9','var10']) varlistdic= {"var1" : ["mean","max","min"], "var2" : ["mean","max","min"], "var3" : ["mean","max","min"], "var4" : ["mean","max","min"], "var5" : ["mean","max","min"], "var6" : ["mean","max","min"], "var7" : ["mean","max","min"], "var8" : ["mean","max","min"], "var9" : ["mean","max","min"], "var10" : ["mean","max","min"], } df_sum=aggregate_df() end=time.time() print(f"Elapsed time: {end- start:.2f} seconds")
多进程版本
import time start=time.time() import numpy as np import random import pandas as pd import multiprocessing as mp def fun_user_id(start, end, step): num = np.linspace(start, end,(end-start) *int(1/step)+1).tolist() return [round(i, 0) for i in num] def fun_rand_num(): return list(map(lambda x: random.randint(300,800), range(1, 25000001))) def aggregate_df(): grvar=df2.groupby("userid") dfvar=grvar.agg(varlistdic) dfvar=dfvar.pipe(lambda x: x.set_axis(x.columns.map('_'.join),axis=1)) dfvar.reset_index(inplace=True) return dfvar if __name__=='__main__': userid=fun_user_id(1,25000001,.5) var1=fun_rand_num() var2=fun_rand_num() var3=fun_rand_num() var4=fun_rand_num() var5=fun_rand_num() var6=fun_rand_num() var7=fun_rand_num() var8=fun_rand_num() var9=fun_rand_num() var10=fun_rand_num() df2 = pd.DataFrame(list(zip(userid,var1, var2,var3,var4,var5,var6,var7,var8,var9,var10)), columns =['userid','var1', 'var2','var3','var4','var5','var6', 'var7','var8','var9','var10']) varlistdic= {"var1" : ["mean","max","min"], "var2" : ["mean","max","min"], "var3" : ["mean","max","min"], "var4" : ["mean","max","min"], "var5" : ["mean","max","min"], "var6" : ["mean","max","min"], "var7" : ["mean","max","min"], "var8" : ["mean","max","min"], "var9" : ["mean","max","min"], "var10" : ["mean","max","min"], } with mp.Pool(mp.cpu_count()) as pool: df_sum2=pool.apply(aggregate_df) end=time.time() print(f"Elapsed time: {end- start:.2f} seconds")
问题根源
你的多进程版本比单进程慢,核心原因有以下几点:
没有真正实现并行计算
你用了pool.apply()方法,这个方法是同步阻塞的——它只会把任务交给进程池中的一个进程执行,之后就等待该进程完成,完全没用到多进程并行。本质上相当于把单进程的任务搬到子进程里跑,反而多了进程启动、内存拷贝的额外开销。大对象的内存拷贝开销
多进程模式下(尤其是Windows系统),子进程会拷贝父进程的全部内存空间。你的DataFrame有5000万行数据,这么大的对象拷贝会消耗大量时间和内存,这部分开销远远超过了多进程可能带来的收益。Pandas本身的优化已经足够高效
Pandas的groupby和agg操作底层是C扩展实现的,并且做了大量向量化优化,单进程下的计算效率已经很高。盲目套多进程框架,反而会因为额外的进程调度、数据传递成本拖慢整体速度。任务没有拆分
即使你用了正确的并行方法,当前代码也没有把聚合任务拆分成多个子任务。要实现有效的并行,需要将数据按userid拆分,让每个进程处理一部分分组,最后合并结果,而不是把整个聚合任务丢给一个子进程。
修正方案
如果你想真正用多进程加速分组聚合,可以参考以下思路:
- 先获取所有唯一的userid,将其分成N个批次(N等于CPU核心数)
- 每个进程处理一个批次的userid对应的子DataFrame,完成聚合
- 最后将所有进程的聚合结果合并
示例代码片段:
def aggregate_part(user_ids): # 筛选当前批次的userid数据 part_df = df2[df2['userid'].isin(user_ids)] grvar = part_df.groupby("userid") dfvar = grvar.agg(varlistdic) dfvar = dfvar.pipe(lambda x: x.set_axis(x.columns.map('_'.join), axis=1)) dfvar.reset_index(inplace=True) return dfvar if __name__=='__main__': # ... 数据生成部分省略 ... # 获取唯一userid并拆分批次 unique_users = df2['userid'].unique() num_processes = mp.cpu_count() # 拆分userid为num_processes个批次 user_batches = np.array_split(unique_users, num_processes) with mp.Pool(num_processes) as pool: # 用map并行处理每个批次 results = pool.map(aggregate_part, user_batches) # 合并所有结果 df_sum2 = pd.concat(results, ignore_index=True) end=time.time() print(f"Elapsed time: {end- start:.2f} seconds")
注意:这种方式的加速效果取决于数据规模和分组数量,如果分组数量太少,拆分后的每个任务计算量太小,可能还是不如单进程高效。
内容的提问来源于stack exchange,提问作者Craig Davis
相关产品推荐
相关产品推荐

