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

如何并行化处理DataFrame分片函数?遇属性错误与串行执行问题

ProcessPoolExecutor并行处理DataFrame分片的问题排查与解决

问题重现

  1. 并行提交任务时报错:
    使用executor.submit(my_df_func, my_df_slice)提交任务时,触发AttributeError: Can't get attribute 'my_df_func' on <module '__main__' (built-in)>,无法启动并行进程。
  2. 修改提交方式后串行执行:
    改成executor.submit(my_df_func(my_df_slice))后代码能运行,但所有分片都在主线程串行处理,完全没有用到进程池的并行能力。

错误原因解析

1. AttributeError的根源

ProcessPoolExecutor基于多进程实现,子进程通过pickle序列化传递函数和参数。在Windows系统(或部分Python环境)中,子进程会重新导入主模块:

  • 如果函数my_df_func定义在主脚本的__main__模块中,子进程导入时无法正确找到该函数的定义(pickle无法序列化主模块中的函数到子进程)。
  • 若在交互式环境(如Jupyter Notebook)中运行代码,也会触发该错误,因为交互式环境的__main__模块结构特殊。

2. 串行执行的原因

executor.submit(my_df_func(my_df_slice))这种写法是先在主线程执行my_df_func(my_df_slice),然后把函数的返回值(通常是None)传给submit。此时进程池实际接收的是一个已执行完成的空任务,自然不会启动并行进程,所有逻辑都在主线程串行完成。

修复方案

方案1:解决AttributeError

将业务函数放到独立模块中,避免子进程导入问题:

  1. 创建df_utils.py文件,写入函数:
import pandas as pd
import os

def my_df_func(my_df_slice, output_dir):
    # 按col1分组后,整个分片的col1值一致,直接取第一个即可
    col1_val = my_df_slice['col1'].iloc[0]
    output_file = os.path.join(output_dir, f'out_{col1_val}.csv')
    
    # 批量处理数据,替代逐行遍历(效率更高)
    my_df_slice['content'] = 'info_' + my_df_slice['col2'].astype(str) + '&id=' + my_df_slice['col3'] + '&abc=' + my_df_slice['col4']
    
    # 写入文件,自动判断是否加表头
    my_df_slice[['content']].to_csv(
        output_file,
        mode='a',
        header=not os.path.isfile(output_file),
        index=False
    )
  1. 主脚本中导入并调用:
from concurrent.futures import ProcessPoolExecutor
import pandas as pd
from df_utils import my_df_func
import os

if __name__ == '__main__':
    # 假设my_df是已加载的DataFrame,output_dir是输出目录
    my_df = pd.read_csv('your_data.csv')
    output_dir = './output'
    os.makedirs(output_dir, exist_ok=True)
    
    # 按col1分组生成分片
    df_groups = [group for _, group in my_df.groupby('col1')]
    
    # 正确提交并行任务
    with ProcessPoolExecutor(max_workers=4) as executor:
        # 用列表推导式收集Future对象,可后续获取结果或处理异常
        futures = [executor.submit(my_df_func, group, output_dir) for group in df_groups]
        # 可选:等待所有任务完成
        for future in futures:
            future.result()

方案2:修复串行执行问题

必须使用executor.submit(函数对象, 参数1, 参数2...)的格式,绝对不能直接调用函数后传结果。正确的提交方式就是你最初的写法:

executor.submit(my_df_func, my_df_slice)

只要解决了AttributeError,这种写法就能正常启动并行进程。

额外优化:避免逐行遍历

原代码中逐行遍历DataFrame的效率极低,改成批量处理(如上例中的my_df_slice['content'] = ...)能大幅提升性能,同时避免循环中的重复IO操作。

内容的提问来源于stack exchange,提问作者Rajat

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 19:20:32