如何在Python中并行化修改外部DataFrame并存储计算结果?
问题
我需要并行计算并将结果写入共享的DataFrame,单线程下修改外部DataFrame df能得到正确结果,但用multiprocessing多进程时,原DataFrame没有被修改。Julia里用Threads.@threads就能轻松实现,想找Python里类似的简便方案。
附代码示例:
单线程Python代码(正常工作)
import numpy as np import pandas as pd def f(u): df.loc[u] = u**2 df = pd.DataFrame(np.zeros(10), index=range(10)) [f(u) for u in range(10)] print(df.T) # 输出正确结果 # 0 1 2 3 4 5 6 7 8 9 # 0 0.0 1.0 4.0 9.0 16.0 25.0 36.0 49.0 64.0 81.0
多进程Python代码(结果错误)
import multiprocessing as mp import numpy as np import pandas as pd def f(u): df.loc[u] = u**2 df = pd.DataFrame(np.zeros(10), index=range(10)) pool = mp.Pool(2) pool.map(f, range(10)) pool.close() print(df.T) # 输出全0的错误结果 # 0 1 2 3 4 5 6 7 8 9 # 0 0.0 0.0 0.0 0.0 0.0 0.0 0.0 0.0 0.0 0.0
Julia实现(正常并行修改)
using DataFrames function f(u) df[u,:v]=u^2 end df = DataFrame(v=zeros(10)); Threads.@threads for u=1:10 f(u) end
为什么多进程不行
Python的multiprocessing是进程隔离的,每个子进程会复制父进程的内存(包括df),子进程里修改的只是自己的副本,父进程的原始DataFrame不会有任何变化。而Julia的Threads.@threads是共享内存线程,所有线程直接操作同一份数据,所以能直接修改。
解决方案
1. 多线程(最接近Julia的实现,简便)
用concurrent.futures.ThreadPoolExecutor或者threading模块,基于共享内存线程,和Julia的线程模型一致,适合IO密集型任务,或者pandas/numpy内部已经释放GIL的计算场景:
import numpy as np import pandas as pd from concurrent.futures import ThreadPoolExecutor def f(u): df.loc[u] = u**2 df = pd.DataFrame(np.zeros(10), index=range(10)) with ThreadPoolExecutor(max_workers=2) as executor: executor.map(f, range(10)) print(df.T) # 输出正确结果
如果喜欢更贴近Julia的循环写法,也可以用threading手动控制:
import threading import numpy as np import pandas as pd def f(u): df.loc[u] = u**2 df = pd.DataFrame(np.zeros(10), index=range(10)) threads = [] for u in range(10): t = threading.Thread(target=f, args=(u,)) threads.append(t) t.start() for t in threads: t.join() print(df.T)
2. 多进程+结果汇总(CPU密集型任务适用)
如果是纯CPU密集型任务(GIL无法释放),得用多进程的话,不要直接修改共享DataFrame,而是让每个进程返回计算结果,最后在主进程合并:
import multiprocessing as mp import numpy as np import pandas as pd def f(u): # 返回索引和对应计算值 return (u, u**2) df = pd.DataFrame(np.zeros(10), index=range(10)) with mp.Pool(2) as pool: results = pool.map(f, range(10)) # 批量更新到主DataFrame for idx, val in results: df.loc[idx] = val print(df.T)
3. 共享内存(复杂,适合超大数据集)
通过multiprocessing.Array创建共享内存数组,让多进程访问同一份数据,不过实现繁琐,一般不推荐除非处理非常大的DataFrame:
import multiprocessing as mp import numpy as np import pandas as pd def f(args): u, shared_arr, shape = args # 把共享内存数组转为numpy数组 arr = np.frombuffer(shared_arr, dtype=np.float64).reshape(shape) arr[u] = u**2 df = pd.DataFrame(np.zeros(10), index=range(10)) # 创建共享内存数组,类型'd'对应float64 shared_arr = mp.Array('d', df.values.ravel()) shape = df.shape with mp.Pool(2) as pool: pool.map(f, [(u, shared_arr, shape) for u in range(10)]) # 将共享内存数据写回DataFrame df[:] = np.frombuffer(shared_arr, dtype=np.float64).reshape(shape) print(df.T)
总结
- 优先用多线程方案,和Julia的
Threads.@threads体验最接近,代码简洁,适合大部分场景; - CPU密集型任务选多进程+结果汇总,避免共享内存的复杂操作;
- 直接用
multiprocessing修改外部DataFrame不可行,进程隔离导致副本修改不影响主进程数据。
内容的提问来源于stack exchange,提问作者Leo
相关产品推荐
相关产品推荐

