【问题标题】:Large numpy arrays in shared memory for multiprocessing: Is something wrong with this approach?共享内存中用于多处理的大型 numpy 数组:这种方法有问题吗?
【发布时间】:2018-03-30 10:02:04
【问题描述】:

多处理是一个很棒的工具,但使用大内存块并不是那么直接。您可以在每个进程中加载​​块并将结果转储到磁盘上,但有时您需要将结果存储在内存中。最重要的是,使用花哨的 numpy 功能。

我已经阅读/谷歌了很多,并想出了一些答案:

Use numpy array in shared memory for multiprocessing

Share Large, Read-Only Numpy Array Between Multiprocessing Processes

Python multiprocessing global numpy arrays

How do I pass large numpy arrays between python subprocesses without saving to disk?

等等等等等等。

它们都有缺点: 不那么主流的库(sharedmem);全局存储变量;不太容易阅读代码、管道等。

我的目标是在我的工作人员中无缝使用 numpy,而不用担心转换和其他东西。

经过多次试验,我想出了this。它适用于我的 ubuntu 16、python 3.6、16GB、8 核机器。与以前的方法相比,我做了很多“捷径”。没有全局共享状态,没有需要在 worker 内部转换为 numpy 的纯内存指针,作为进程参数传递的大型 numpy 数组等。

Pastebin link above,但我会在这里放几个sn-ps。

一些进口:

import numpy as np
import multiprocessing as mp
import multiprocessing.sharedctypes
import ctypes

分配一些共享内存并将其包装到一个 numpy 数组中:

def create_np_shared_array(shape, dtype, ctype)
     . . . . 
    shared_mem_chunck = mp.sharedctypes.RawArray(ctype, size)
    numpy_array_view = np.frombuffer(shared_mem_chunck, dtype).reshape(shape)
    return numpy_array_view

创建共享数组并在其中放入一些东西

src = np.random.rand(*SHAPE).astype(np.float32)
src_shared = create_np_shared_array(SHAPE,np.float32,ctypes.c_float)
dst_shared = create_np_shared_array(SHAPE,np.float32,ctypes.c_float)
src_shared[:] = src[:]  # Some numpy ops accept an 'out' array where to store the results

产生进程:

p = mp.Process(target=lengthly_operation,args=(src_shared, dst_shared, k, k + STEP))
p.start()
p.join()

以下是一些结果(完整参考请参见 pastebin 代码):

Serial version: allocate mem 2.3741257190704346 exec: 17.092209577560425 total: 19.46633529663086 Succes: True
Parallel with trivial np: allocate mem 2.4535582065582275 spawn  process: 0.00015354156494140625 exec: 3.4581971168518066 total: 5.911908864974976 Succes: False
Parallel with shared mem np: allocate mem 4.535916328430176 (pure alloc:4.014216661453247 copy: 0.5216996669769287) spawn process: 0.00015664100646972656 exec: 3.6783478260040283 total: 8.214420795440674 Succes: True

我还做了一个cProfile(为什么在分配共享内存时要多花 2 秒?)并意识到有一些对tempfile.py{method 'write' of '_io.BufferedWriter' objects} 的调用。

问题

  • 我做错了吗?
  • (大)阵列是否来回腌制而我没有获得任何加快速度?请注意,第二次运行(使用常规 np 数组未通过正确性测试)
  • 有没有办法进一步改善时序、代码清晰度等? (针对多处理范例)

备注

  • 我不能使用进程池,因为 mem 必须在 fork 处继承,而不是作为参数发送。

【问题讨论】:

  • 如果我在磁盘上有 numpy 友好格式的数据,我会考虑的一件事是对该文件执行 mmap 并将它们包装到一个 numpy 数组中。在多个进程中对同一文件执行mmap(设置MAP_SHARED | PROT_READ)应该会导致数据最多存储在RAM中一次。

标签: python numpy multiprocessing


【解决方案1】:

共享数组的分配速度很慢,因为显然它是先写入磁盘的,所以可以通过 mmap 进行共享。有关参考,请参阅 heap.pysharedctypes.py。 这就是为什么tempfile.py 出现在分析器中的原因。我认为这种方法的优点是在崩溃的情况下会清理共享内存,而 POSIX 共享内存无法保证这一点。

感谢 fork,您的代码不会发生酸洗,正如您所说,内存是继承的。第二次运行不起作用的原因是因为不允许子进程写入父进程的内存。相反,私有页面是动态分配的,只有在子进程结束时才会被丢弃。

我只有一个建议:你不必自己指定ctype,可以通过np.ctypeslib._typecodes从numpy dtype中找出正确的类型。或者只是对所有内容使用c_byte,并使用 dtype itemsize 来计算缓冲区的大小,无论如何它都会被 numpy 强制转换。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-10-20
    • 2011-12-15
    • 2013-07-21
    • 2021-03-19
    • 1970-01-01
    • 1970-01-01
    • 2014-08-16
    相关资源
    最近更新 更多