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

如何用Python多进程Pool优化DataFrame循环处理?

用多进程Pool优化Pandas DataFrame循环处理

问题背景

有大体积Pandas DataFrame,原代码通过循环遍历result的不同取值(1到MAX_NUMBER),生成对应的mc{i}_aft和mc{i}_bef列,但单进程下groupby、cumcount等操作效率低下,希望通过multiprocessing.Pool利用多核CPU提升处理速度。

原代码逻辑是针对每个i值:

  1. 按con, mac, result分组计算累计计数
  2. 过滤出result == i的计数作为mc{i}_aft列的初始值
  3. 按con, mac分组做前向填充并补0
  4. 生成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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 16:32:51