如何并行化处理DataFrame分片函数?遇属性错误与串行执行问题
ProcessPoolExecutor并行处理DataFrame分片的问题排查与解决
问题重现
- 并行提交任务时报错:
使用executor.submit(my_df_func, my_df_slice)提交任务时,触发AttributeError: Can't get attribute 'my_df_func' on <module '__main__' (built-in)>,无法启动并行进程。 - 修改提交方式后串行执行:
改成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
将业务函数放到独立模块中,避免子进程导入问题:
- 创建
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 )
- 主脚本中导入并调用:
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
相关产品推荐
相关产品推荐

