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")
为什么这个方案更好?
- 避免了内存隔离的问题:每个子进程独立处理自己的chunk,返回结果后由主进程合并,完全不用操心数据共享的问题。
- 代码更简洁:去掉了创建临时目录、写临时文件、合并文件这些额外的IO操作,不仅减少了出错概率,还能提升效率。
- 更容易调试:每个子进程的逻辑独立,出问题时更容易定位。
如果你非要用共享内存的方式(比如直接修改同一个DataFrame),也可以用multiprocessing.Manager或者共享数组,但这种方式对pandas DataFrame非常不友好——因为DataFrame底层是numpy数组,需要手动处理类型转换和内存映射,很容易出错,所以非常不推荐。
内容的提问来源于stack exchange,提问作者Mee
相关产品推荐
相关产品推荐

