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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 17:15:03