【发布时间】: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