【发布时间】:2020-02-23 15:19:33
【问题描述】:
我正在使用 mac book,因此,多处理将使用 fork 系统调用而不是产生一个新进程。另外,我正在使用 Python(带有多处理或 Dask)。
我有一个非常大的熊猫数据框。我需要让许多并行子流程与这个大数据帧的一部分一起工作。假设我有 100 个需要并行处理的表分区。我想避免需要制作 100 个这个大数据帧的副本,因为这会压倒内存。所以我目前采取的方法是对其进行分区,将每个分区保存到磁盘,并让每个进程读取它们以处理它们各自负责的部分。但是这种读/写对我来说非常昂贵,我想避免它。
但是,如果我为这个数据帧创建一个全局变量,那么由于 COW 行为,每个进程都可以从这个数据帧中读取,而无需制作它的实际物理副本(只要它不修改它)。现在我的问题是,如果我创建一个全局数据框并命名它:
global my_global_df
my_global_df = one_big_df
然后在我做的一个子流程中:
a_portion_of_global_df_readonly = my_global_df.iloc[0:10]
a_portion_of_global_df_copied = a_portion_of_global_df_readonly.reset_index(drop=True)
# reset index will make a copy of the a_portion_of_global_df_readonly
do something with a_portion_of_global_df_copied
如果我执行上述操作,我会创建整个my_global_df 的副本还是只创建a_portion_of_global_df_readonly 的副本,从而避免复制100 个one_big_df?
另一个更普遍的问题是,当(假设人们使用 UNIX)将数据设置为全局变量时,为什么人们必须处理 Pickle 序列化和/或读/写磁盘以跨多个进程传输数据?如此轻松地有效地使其在所有子进程中可用?使用 COW 作为使任何数据可用于一般子流程的手段是否存在危险?
[来自下面线程的可重现代码]
from multiprocessing import Process, Pool
import contextlib
import pandas as pd
def my_function(elem):
return id(elem)
num_proc = 4
num_iter = 10
df = pd.DataFrame(np.asarray([1]))
print(id(df))
with contextlib.closing(Pool(processes=num_proc)) as p:
procs = [p.apply_async(my_function, args=(df, )) for elem in range(num_iter)]
results = [proc.get() for proc in procs]
p.close()
p.join()
print(results)
【问题讨论】:
-
由于 COW,您可以在每个子进程中修改数据帧,而不会影响父副本或兄弟副本。
-
我知道从子进程修改
my_global_df不会修改父进程中的my_global_df。但我想知道的是,复制my_global_df的一部分是否会有效地复制整个my_global_df或仅复制我正在复制/修改的那部分(例如.iloc[0:10])?如果它会复制整个内容,我将耗尽内存,因为每个子进程都会复制整个my_global_df。 -
它只复制您要求的内容。您将获得一个新的数据框,其中仅包含您要求的行/列。您可以通过在原始 DF 和您创建的新 DF 上调用
.info()来验证。 -
您甚至不需要使用全局变量。只需将数据框作为参数传递给子进程工作函数即可。
-
是的,但如果你这样做,Python 将使用 Pickle 序列化数据帧,这对于大数据帧来说可能非常昂贵
标签: python pandas multiprocessing fork dask