【问题标题】:Is it possible to share a numpy array that's not empty between processes?是否可以在进程之间共享一个不为空的 numpy 数组?
【发布时间】:2021-05-31 00:48:41
【问题描述】:

我以为SharedMemory 会保留目标数组的值,但当我实际尝试时,似乎没有。

from multiprocessing import Process, Semaphore, shared_memory
import numpy as np
import time
 
dtype_eV = np.dtype({ 'names':['idx', 'value', 'size'], \
                    'formats':['int32', 'float64', 'float64'] })
 
 
def worker_writer(id, number, a, shm):
    exst_shm = shared_memory.SharedMemory(name=shm)
    b = np.ndarray(a.shape, dtype=a.dtype, buffer=exst_shm.buf)
 
    for i in range(5):
        time.sleep(0.5)
        b['idx'][i] = i
 
def worker_reader(id, number, a, shm):
    exst_shm = shared_memory.SharedMemory(name=shm)
    b = np.ndarray(a.shape, dtype=a.dtype, buffer=exst_shm.buf)
 
    for i in range(5):
        time.sleep(1)
        print(b['idx'][i], b['value'][i])
 
 
if __name__ == "__main__":
    a = np.zeros(5, dtype=dtype_eV)
    a['value'] = 100
    shm = shared_memory.SharedMemory(create=True, size=a.nbytes)  
    c = np.ndarray(a.shape, dtype=a.dtype, buffer=shm.buf)
    th1 = Process(target=worker_writer, args=(1, 50000000, a, shm.name))
    th2 = Process(target=worker_reader, args=(2, 50000000, a, shm.name))
 
    th1.start()
    th2.start()
    th1.join()
    th2.join()

'''
result:
0 0.0
1 0.0
2 0.0
3 0.0
4 0.0
'''

在上面的代码中,两个进程可以共享一个数组(a)并访问它。但是共享之前给出的值(a['value'] = 100)丢失了。是自然的还是分享后有什么方法可以保持价值?

【问题讨论】:

  • "但是共享之前给出的值(a['value'] = 100) 丢失了。" - 什么?为什么会出现这个值?您将其存储在a,而不是共享内存中。
  • 我不明白a 的意义是什么。很明显,您希望它做一些有用的事情,但您从未读取过存储在那里的任何内容 - 您读取的只是元数据,例如其形状、dtype 和缓冲区大小。
  • @user2357112-supports-monica 谢谢你的回答。我想我误解了这个概念。我认为必须有一个原始数组,并通过 SharedMemory 在进程之间共享。
  • @SergeBallesta: 错误 - 事实上,buffer 确实 指定了要使用的数组的内存。 (有关演示,请参阅 here。)不过,该内存不是写入 100 的内存。
  • 您应该只使用a 来计算nbytes,那么它不是很有用。请改用c。您永远不会将 a 中包含的任何值传递给 shm 构造函数,那么它如何知道 ['value'] = 100

标签: python multiprocessing


【解决方案1】:

这是一个如何使用 numpy 使用 shared_memory 的示例。它是从我的其他几个answers 粘贴在一起的,但是shared_memory 有几个陷阱需要牢记:

  • 当您从 shm 对象创建 numpy ndarray 时,它不会阻止 shm 被垃圾回收。不幸的副作用是下次尝试访问数组时,会出现段错误。从 another question 我创建了一个快速的 ndarray 子类,仅将 shm 作为属性附加,因此引用会保留,并且不会被 GC。
  • 另一个陷阱是,在 Windows 上,操作系统会跟踪何时删除内存,而不是授予您删除内存的权限。这意味着即使您不调用 unlink,如果没有对该特定内存段的活动引用(由名称给出),内存也会被删除。解决这个问题的方法是确保在主进程上保持一个 shm 打开,该主进程的寿命比所有子进程都长。在末尾调用 close 和 unlink 会保持对末尾的引用,并确保在其他平台上不会泄漏内存。
import numpy as np
import multiprocessing as mp
from multiprocessing.shared_memory import SharedMemory

class SHMArray(np.ndarray): #copied from https://numpy.org/doc/stable/user/basics.subclassing.html#slightly-more-realistic-example-attribute-added-to-existing-array
    '''an ndarray subclass that holds on to a ref of shm so it doesn't get garbage collected too early.'''
    def __new__(cls, input_array, shm=None):
        obj = np.asarray(input_array).view(cls)
        obj.shm = shm
        return obj

    def __array_finalize__(self, obj):
        if obj is None: return
        self.shm = getattr(obj, 'shm', None)

def child_func(name, shape, dtype):
    shm = SharedMemory(name=name)
    arr = SHMArray(np.ndarray(shape, buffer=shm.buf, dtype=dtype), shm)
    arr[:] += 5
    shm.close() #be sure to cleanup your shm's locally when they're not needed (referring to arr after this will segfault)

if __name__ == "__main__":
    shape = (10,) # 1d array 10 elements long
    dtype = 'f4' # 32 bit floats
    dummy_array = np.ndarray(shape, dtype=dtype) #dumy array to calculate nbytes
    shm = SharedMemory(create=True, size=dummy_array.nbytes)
    arr = np.ndarray(shape, buffer=shm.buf, dtype=dtype) #create the real arr backed by the shm
    arr[:] = 0
    print(arr) #should print arr full of 0's
    p1 = mp.Process(target=child_func, args=(shm.name, shape, dtype))
    p1.start()
    p1.join()
    print(arr) #should print arr full of 5's
    shm.close() #be sure to cleanup your shm's
    shm.unlink() #call unlink when the actual memory can be deleted

【讨论】:

  • 感谢您的回答和其他参考!我从他们身上学到了很多!
  • 嗨@Aaron,我有一个关于在Windows 笔记本电脑here 上运行multiprocessing.Pool 的问题。我希望你能花一些时间来检查一下这个问题。非常感谢您的帮助!
猜你喜欢
  • 2021-05-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-04-17
  • 2015-03-17
  • 1970-01-01
  • 2019-04-05
  • 2012-12-15
相关资源
最近更新 更多