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

Python中多进程函数无法更新DataFrame的问题求助

为什么multiprocessing并行处理时无法写入目标DataFrame?

这个问题的核心原因其实很简单:Python的多进程是内存隔离的。

你在主进程里创建的df_final,当你把它的切片传给子进程时,子进程拿到的只是这个切片的副本——相当于复制了一份数据到自己的内存空间里。子进程里对df_final的所有修改,都只作用在这个副本上,完全不会影响主进程里的原始df_final。这就是为什么单进程执行时能修改成功(因为只有一个进程,修改的是原数据),但多进程时原df_final还是全0的原因。

而且你的代码里还额外绕了临时文件的弯路,其实完全可以不用——我们可以换一种更简洁可靠的方式:让每个子进程处理自己的DataFrame chunk,然后返回处理后的结果,最后在主进程里把所有结果合并起来。

修改后的代码示例

首先调整generate函数,让它只接收待处理的DataFrame chunk,返回处理后的结果:

def generate(df_chunk):
    # 为当前chunk创建结果DF,初始全0
    df_chunk_final = pd.DataFrame(0, index=df_chunk.index, columns=df_chunk.columns)
    print("in func")
    
    for i in Regions:
        for j in range(df_chunk.shape[0]):
            edge_count = df_chunk[i+"_Edges_Count"].iloc[j]
            new_edge_count = edge_count
            
            for k in range(1, edge_count+1):
                if df_chunk[i+"_D"+str(k)].iloc[j] <= 1.4:
                    df_chunk_final[i+"_D"+str(k)].iloc[j] = df_chunk[i+"_D"+str(k)].iloc[j]
                else:
                    new_edge_count = k-1
                    break
            
            for c in range(1, new_edge_count+1):
                for feature in Features:
                    df_chunk_final[i+feature+str(c)].iloc[j] = df_chunk[i+feature+str(c)].iloc[j]
            
            df_chunk_final[i+"_Edges_Count"].iloc[j] = new_edge_count
    
    return df_chunk_final

然后修改主进程的调用逻辑,去掉临时文件相关的操作,直接收集子进程返回的结果并合并:

import multiprocessing as mp
import math
import pandas as pd

# 确定使用的进程数
number_of_CPU = "Max"
if isinstance(number_of_CPU, int) and 0 < number_of_CPU < mp.cpu_count():
    processes = number_of_CPU
else:
    processes = mp.cpu_count()
# 进程数不能超过数据行数
processes = min(processes, df.shape[0])
print(f'Generating the new training data using {processes} processes.')

# 将原DataFrame拆分为多个chunk
chunk_size = math.ceil(df.shape[0] / processes)
list_of_chunks = [
    df.iloc[i*chunk_size : min((i+1)*chunk_size, df.shape[0]), :] 
    for i in range(processes)
]

# 并行处理所有chunk
with mp.Pool(processes) as pool:
    processed_chunks = pool.map(generate, list_of_chunks)

# 合并所有处理后的chunk得到最终结果
df_final = pd.concat(processed_chunks, axis=0)
# 保存结果到CSV
df_final.to_csv('./sample.csv', index=False)

print("done")

为什么这个方案更好?

  1. 避免了内存隔离的问题:每个子进程独立处理自己的chunk,返回结果后由主进程合并,完全不用操心数据共享的问题。
  2. 代码更简洁:去掉了创建临时目录、写临时文件、合并文件这些额外的IO操作,不仅减少了出错概率,还能提升效率。
  3. 更容易调试:每个子进程的逻辑独立,出问题时更容易定位。

如果你非要用共享内存的方式(比如直接修改同一个DataFrame),也可以用multiprocessing.Manager或者共享数组,但这种方式对pandas DataFrame非常不友好——因为DataFrame底层是numpy数组,需要手动处理类型转换和内存映射,很容易出错,所以非常不推荐。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 07:28:39