多进程操作Pandas DataFrame后数据丢失,求简便解决方案
解决Pandas DataFrame多进程修改后为空的问题
你遇到的问题本质是多进程的内存隔离机制——每个子进程启动时会复制主进程的内存空间(包括你的df),所以子进程里修改的都是df的副本,主进程里的原始DataFrame完全没被改动,最后导出自然看不到你想要的修改效果。
下面给你两种简便的解决思路:
方案1:避免共享,进程返回结果后主进程统一修改(推荐)
这种方法不用搞复杂的共享机制,每个子进程只负责计算要修改的内容,返回索引和对应的值,主进程最后一次性更新DF。既安全又能保留多进程并行的优势:
import pandas as pd import multiprocessing d = {'col1': [1, 2], 'col2': [3, 4]} df = pd.DataFrame(data=d) task_list = [(0,3,1,5), (1,6,0,7)] # 替换成你的实际任务列表 def process_task(x): # 返回要修改的(行索引, 列名, 新值)元组列表 return [(x[0], "col1", x[1]), (x[2], "col2", x[3])] if __name__ == "__main__": # 多进程必须加这个保护,避免重复初始化 pool = multiprocessing.Pool() results = pool.map(process_task, task_list) pool.close() pool.join() # 主进程统一更新原始DF for result in results: for idx, col, val in result: df.loc[idx, col] = val df.to_csv("test.csv")
为什么推荐这个?因为Pandas的DataFrame本身不是进程安全的,直接共享修改很容易出现数据竞争(比如两个进程同时改同一行),而让进程返回计算结果,主进程统一处理能彻底避免这类问题,代码也更清晰。
方案2:用Manager共享DataFrame(需加锁)
如果你一定要在子进程里直接修改DF,可以用multiprocessing.Manager创建一个可共享的容器存放DF,同时用锁来保证修改的原子性,防止多进程冲突:
import pandas as pd import multiprocessing d = {'col1': [1, 2], 'col2': [3, 4]} task_list = [(0,3,1,5), (1,6,0,7)] def function(x, shared_df, lock): with lock: # 加锁确保同一时间只有一个进程修改DF,避免数据错乱 shared_df.df.loc[x[0],"col1"] = x[1] shared_df.df.loc[x[2],"col2"] = x[3] if __name__ == "__main__": manager = multiprocessing.Manager() shared_df = manager.Namespace() # 创建可共享的命名空间 shared_df.df = pd.DataFrame(data=d) lock = multiprocessing.Lock() # 进程锁 pool = multiprocessing.Pool() # 把共享对象和锁传给每个任务 pool.starmap(function, [(task, shared_df, lock) for task in task_list]) pool.close() pool.join() shared_df.df.to_csv("test.csv")
⚠️ 注意:
- 必须用
if __name__ == "__main__"包裹多进程启动代码,这是Windows系统的硬性要求,也能避免Unix系统的重复初始化问题。 - 锁是必须的,否则多个进程同时修改DF会导致数据错乱甚至抛出异常。
- 这种方法的性能会比方案1差,因为锁会让进程串行修改,失去多进程并行的优势,所以只有在任务必须在子进程里修改DF时才考虑使用。
额外提示:如果你的任务是大规模数据处理,更推荐用dask或者swifter这类专门为Pandas设计的并行处理库,它们会自动处理多进程/多线程的共享问题,代码更简洁高效。
内容的提问来源于stack exchange,提问作者Anthonious Ids Tong
相关产品推荐
相关产品推荐

