Python Pandas Groupby大数据聚合任务的多线程优化求助
优化2000万+条数据的Pandas GroupBy聚合效率
问题背景
现有2000万+条记录的数据集,按userid分组后对20个变量执行多统计量聚合(mean/max/min),单线程处理耗时约40分钟,需要通过并行计算提升效率。
数据集结构示例:
userid var1 ... var20 1 323 ... 450 1 443 ... 357 2 467 ... 587 3 235 ... 345 3 578 ... 768 4 354 ... 365
原聚合代码:
varlistdic = { "var1": ["mean", "max", "min"], "var2": ["mean", "max", "min"], # ... 省略var3到var19的配置 "var20": "max" } gr = df.groupby(['userid']) df_agg = gr.agg(varlistdic)
优化方案
1. 单线程基础优化(立竿见影)
Pandas默认groupby会对分组键排序,大数据集下排序会消耗大量时间,先关闭排序:
gr = df.groupby(['userid'], sort=False) # 关闭排序,减少耗时 df_agg = gr.agg(varlistdic)
该改动通常能降低30%-50%的耗时,是成本最低的优化手段。
同时确保数据类型最优:将userid设为整数类型(而非字符串),数值变量用float32/int32(精度允许的情况下),减少内存占用的同时提升计算速度。
2. Dask多进程并行(推荐)
Dask是专门处理大数据集的并行计算框架,自动拆分数据并在多进程中执行groupby操作,适合超大规模数据集:
import dask.dataframe as dd # 将Pandas DataFrame转为Dask DataFrame,分区数设为CPU核心数的1-2倍 ddf = dd.from_pandas(df, npartitions=8) # 执行聚合 varlistdic = { "var1": ["mean", "max", "min"], "var2": ["mean", "max", "min"], # ... 省略其他变量 "var20": "max" } ddf_agg = ddf.groupby('userid').agg(varlistdic) # 计算结果并转回Pandas DataFrame df_agg = ddf_agg.compute()
该方案通常能将耗时压缩至原单线程的1/4到1/8(取决于CPU核心数)。
3. 手动多进程拆分处理(无额外依赖)
若不想引入Dask,可使用Python标准库multiprocessing手动拆分userid分组,并行处理后合并结果:
import pandas as pd from multiprocessing import Pool def process_chunk(user_ids): # 筛选当前chunk的userid数据 chunk_df = df[df['userid'].isin(user_ids)] # 执行聚合 return chunk_df.groupby('userid').agg(varlistdic) # 拆分唯一userid为多个chunk unique_users = df['userid'].unique() chunk_size = len(unique_users) // 8 # 按CPU核心数拆分 user_chunks = [unique_users[i:i+chunk_size] for i in range(0, len(unique_users), chunk_size)] # 多进程执行 with Pool(processes=8) as pool: results = pool.map(process_chunk, user_chunks) # 合并所有结果(每个userid仅在一个chunk中,直接合并即可) df_agg = pd.concat(results)
注意:该方案要求数据集能完全载入内存,内存不足时优先选择Dask。
4. Pandas多线程后端(PyArrow)
Pandas 2.0+支持PyArrow作为计算后端,部分操作自动启用多线程:
# 先安装依赖:pip install pyarrow import pandas as pd # 启用PyArrow后端 pd.set_option('mode.copy_on_write', True) df = df.convert_dtypes(dtype_backend='pyarrow') # 执行聚合,部分操作自动多线程 gr = df.groupby(['userid'], sort=False) df_agg = gr.agg(varlistdic)
该方案无需大幅修改代码,但加速效果可能不如Dask,适合不想引入新框架的场景。
总结
优先尝试关闭groupby排序;若耗时仍不达标,推荐用Dask实现多进程并行以最大化利用CPU资源;受依赖限制时,再考虑手动多进程拆分或PyArrow后端。
内容的提问来源于stack exchange,提问作者Craig Davis
相关产品推荐
相关产品推荐

