【问题标题】:Taking advantage of fork system call to avoid read/writing or serializing altogether?利用 fork 系统调用来完全避免读/写或序列化?
【发布时间】: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


【解决方案1】:

总结 cmets,在 Mac 或 Linux 等分叉系统上,子进程具有父地址空间的写时复制 (COW) 视图,包括它可能拥有的任何DataFrames。在子进程中使用和修改数据框是安全的,而无需更改父进程或其他同级子进程中的数据。

这意味着没有必要序列化数据帧以将其传递给孩子。您所需要的只是对数据框的引用。对于Process,您可以直接传递引用

p = multiprocessing.Process(target=worker_fctn, args=(my_dataframe,))
p.start()
p.join()

如果您使用Queue 或其他工具(例如Pool),则数据可能会被序列化。您可以使用工作人员已知但实际上并未传递给工作人员的全局变量来解决该问题。

剩下的是返回数据。它仅在子级中,仍需要序列化才能返回给父级。

【讨论】:

  • 这是不正确的。 Python 多处理仍然使用序列化/pickle 将数据传递给子进程。所以在这种情况下,数据帧仍然会被 PICKLED。更多详情:stackoverflow.com/questions/52600240/…
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2017-07-25
  • 2010-11-18
  • 2011-06-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-02-02
相关资源
最近更新 更多