Python Pandas使用multiprocessing多进程处理DataFrame无新列返回求解
多进程修改DataFrame无效果的解决方案
问题根因
Python multiprocessing 模块默认采用 spawn 模式创建子进程,每个子进程会独立复制父进程的内存空间,子进程内修改的DataFrame只是自身持有的副本,不会同步到父进程的原始DataFrame上,因此运行后父进程的df看不到新增列。
解决方案
不要在子进程内直接修改全局df,将计算好的列结果返回给父进程,由父进程统一合并到原始df中即可,以下是两种可直接运行的实现方案:
方案1:使用Queue接收子进程返回结果
import pandas as pd from io import StringIO import multiprocessing # 构造原始DataFrame df_str = """ ValOption RB test 0 SLA 4 3 1 AC 5 4 2 SLA 5 5 3 AC 2 4 4 SLA 5 5 5 AC 3 4 6 SLA 4 3 """ df = pd.read_csv(StringIO(df_str.strip()), sep='\s+') # 调整函数逻辑:接收df和队列参数,计算完成后将列名和列值传入队列 def func1(df, queue): r1 = df['test'] + 1 queue.put(('r1', r1)) def func2(df, queue): r2 = df['RB'] + 1 queue.put(('r2', r2)) if __name__ == '__main__': res_queue = multiprocessing.Queue() p1 = multiprocessing.Process(target=func1, args=(df, res_queue)) p2 = multiprocessing.Process(target=func2, args=(df, res_queue)) p1.start() p2.start() # 从队列获取子进程返回的结果,合并到原始df for _ in range(2): col_name, col_val = res_queue.get() df[col_name] = col_val p1.join() p2.join() # 输出验证结果 print(df)
方案2:使用进程池Pool简化返回值获取
如果任务逻辑更简单,用multiprocessing.Pool的apply_async方法可以更方便地获取子进程的返回值,不需要手动维护队列:
import pandas as pd from io import StringIO import multiprocessing # 构造原始DataFrame df_str = """ ValOption RB test 0 SLA 4 3 1 AC 5 4 2 SLA 5 5 3 AC 2 4 4 SLA 5 5 5 AC 3 4 6 SLA 4 3 """ df = pd.read_csv(StringIO(df_str.strip()), sep='\s+') def func1(df): return 'r1', df['test'] + 1 def func2(df): return 'r2', df['RB'] + 1 if __name__ == '__main__': # 开2个进程执行任务 with multiprocessing.Pool(processes=2) as pool: res1 = pool.apply_async(func1, args=(df,)) res2 = pool.apply_async(func2, args=(df,)) # 获取计算结果合并到原始df col1, val1 = res1.get() col2, val2 = res2.get() df[col1] = val1 df[col2] = val2 # 输出验证结果 print(df)
内容的提问来源于stack exchange,提问作者William
相关产品推荐
相关产品推荐

