如何用Python多进程Pool优化DataFrame循环处理?
用多进程Pool优化Pandas DataFrame循环处理
问题背景
有大体积Pandas DataFrame,原代码通过循环遍历result的不同取值(1到MAX_NUMBER),生成对应的mc{i}_aft和mc{i}_bef列,但单进程下groupby、cumcount等操作效率低下,希望通过multiprocessing.Pool利用多核CPU提升处理速度。
原代码逻辑是针对每个i值:
- 按
con, mac, result分组计算累计计数 - 过滤出
result == i的计数作为mc{i}_aft列的初始值 - 按
con, mac分组做前向填充并补0 - 生成
mc{i}_bef列为mc{i}_aft的移位值
多进程改造方案
由于每个i的处理逻辑完全独立,不依赖其他循环的结果,可以将单个i的处理逻辑封装为独立函数,通过Pool并行执行所有任务,最后合并结果。
步骤1:封装单任务处理函数
把单个i的处理逻辑抽成函数,输入原始DataFrame和i值,返回该i对应的两列结果:
import pandas as pd from multiprocessing import Pool import os def process_single_i(args): df, i = args aft_col = f"mc{i}_aft" bef_col = f"mc{i}_bef" group_keys_1 = ['con', 'mac', 'result'] group_keys_2 = ['con', 'mac'] # 复制数据避免多进程间的资源竞争 temp_df = df.copy() # 计算累计计数并生成初始aft列 temp_df['temp_count'] = temp_df.groupby(group_keys_1).cumcount() + 1 temp_df[aft_col] = temp_df['temp_count'].loc[temp_df['result'] == i] # 分组前向填充并补0 temp_df = temp_df.groupby(group_keys_2, group_keys=False).apply( lambda x: x.fillna(method='ffill').fillna(0) ) temp_df.drop(['temp_count'], axis=1, inplace=True) temp_df[aft_col] = temp_df[aft_col].astype(int) # 生成移位后的bef列 temp_df[bef_col] = temp_df.groupby(group_keys_2)[aft_col].shift(1) temp_df.fillna(0, inplace=True) temp_df[bef_col] = temp_df[bef_col].astype(int) # 返回仅包含目标列的结果,保留索引以便合并 return temp_df[[aft_col, bef_col]]
步骤2:启动多进程并行处理
在主程序中准备任务列表,用Pool并行执行,最后合并所有结果:
if __name__ == "__main__": # 初始化原始DataFrame df = pd.DataFrame({ 'mac':['type_a','type_a','type_a','type_a','type_a','type_b','type_b','type_b','type_b','type_b'], 'con':['a','a','a','c','b','a','a','b','a','c'], 'result':[1,1,2,2,3,1,1,3,1,2], }) MAX_NUMBER = 3 # 构造任务:每个任务是(原始df, i)的元组 tasks = [(df, i) for i in range(1, MAX_NUMBER + 1)] # 启动进程池,进程数设为CPU核心数 with Pool(processes=os.cpu_count()) as pool: # 并行执行所有任务 process_results = pool.map(process_single_i, tasks) # 将所有进程返回的列合并到原始DataFrame for result_df in process_results: df = pd.concat([df, result_df], axis=1) # 输出结果 print(df)
关键注意事项
if __name__ == "__main__":Windows系统下必须加这个判断,避免多进程启动时的递归导入问题;Linux/macOS虽然可选,但加了也更规范。- 数据副本:每个进程处理时复制原始df,避免多进程间的共享数据冲突;如果df体积极大,可以只传递必要的列(如
['mac','con','result'])减少内存开销。 - 效率阈值:多进程有启动开销,只有当数据量足够大、MAX_NUMBER足够多时,才能体现出性能优势;小数据量场景下单进程可能更快。
- 替代方案:如果不想手动处理多进程,也可以用
swifter库自动并行化apply,或dask.dataframe处理超大数据集,但multiprocessing.Pool是原生Python方案,无需额外依赖。
内容的提问来源于stack exchange,提问作者ManOnTheMoon
相关产品推荐
相关产品推荐

