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

为何多进程未提升我的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")   

问题根源

你的多进程版本比单进程慢,核心原因有以下几点:

  1. 没有真正实现并行计算
    你用了pool.apply()方法,这个方法是同步阻塞的——它只会把任务交给进程池中的一个进程执行,之后就等待该进程完成,完全没用到多进程并行。本质上相当于把单进程的任务搬到子进程里跑,反而多了进程启动、内存拷贝的额外开销。

  2. 大对象的内存拷贝开销
    多进程模式下(尤其是Windows系统),子进程会拷贝父进程的全部内存空间。你的DataFrame有5000万行数据,这么大的对象拷贝会消耗大量时间和内存,这部分开销远远超过了多进程可能带来的收益。

  3. Pandas本身的优化已经足够高效
    Pandas的groupby和agg操作底层是C扩展实现的,并且做了大量向量化优化,单进程下的计算效率已经很高。盲目套多进程框架,反而会因为额外的进程调度、数据传递成本拖慢整体速度。

  4. 任务没有拆分
    即使你用了正确的并行方法,当前代码也没有把聚合任务拆分成多个子任务。要实现有效的并行,需要将数据按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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 07:55:02