【问题标题】:Python multiprocessing is taking much longer than single processingPython 多处理比单处理花费更长的时间
【发布时间】:2013-11-09 06:05:24
【问题描述】:

我正在依次对 3 个不同的 numpy 2D 数组执行一些大型计算。阵列很大,每个 25000x25000。每次计算都需要大量时间,因此我决定在服务器上的 3 个 CPU 内核上并行运行其中的 3 个。我遵循标准的多处理指南并创建 2 个进程和一个工作函数。两个计算通过 2 个进程运行,第三个计算在本地运行,没有单独的进程。我将巨大的数组作为进程的参数传递,例如:

p1 = Process(target = Worker, args = (queue1, array1, ...)) # Some other params also going

p2 = Process(target = Worker, args = (queue2, array2, ...)) # Some other params also going

Worker 函数将两个 numpy 向量(一维数组)发送回一个附加在队列中的列表中,例如:

queue.put([v1, v2])

我没有使用multiprocessing.pool

但令人惊讶的是我没有得到加速,它实际上运行速度慢了 3 倍。传递大型数组需要时间吗?我无法弄清楚发生了什么。我应该使用共享内存对象而不是传递数组吗?

如果有人能提供帮助,我将不胜感激。

谢谢。

【问题讨论】:

  • 作为一种基准测试,尝试计算腌制/解封阵列所需的时间。正因为您正在连续调度工作人员,您必须等待整个 pickle>unpickle 循环在第一个数组上完成,然后才能开始。您也可能会淹没您的 i/o 流或输出队列。考虑尝试使用进程池来执行此操作,并对每个工作迭代可能的最小和最简单的数据集合进行操作,这样您的工作人员就可以在您完成填充输入缓冲区之前开始工作,而不必提取尽可能多的数据在开始之前结束。
  • 子进程不共享内存,因此要将参数从父进程传递给子进程,它们必须经过序列化>反序列化过程。 mutliprocessing 将参数传递给子进程的方式是在父进程中腌制参数,然后使用pickle 模块在子进程中解开它们。处理小而简单的数据通常没什么大不了的,但是由于您使用的是巨大的 numpy 数组,因此您的数据既不小也不简单。基本上,如果您可以将必须​​传递给工作人员的数据减少为少量更原始的类型,这可能会有所帮助。
  • @Saullo Castro 谢谢。我不知道幕后这个泡菜+解酒。我是python编程的新手,现在我明白了。我现在试试 np.memmap。如果有任何问题,我会在这里发布。不过非常感谢。
  • 一旦你开始在巨大的数组上做一些事情,你的内存布局就变得非常重要。在您调用的方法中执行 += 之类的操作可以加快速度,只是因为内存分配比实际执行您想要执行的操作需要更多时间。
  • @Saullo Castro 只是一个简单的问题,我查看了工人池。这很简单。我只想知道使用工人池是否比手动创建多个子流程更可取?我知道我有 3 个矩阵或 np.arrays 所以我需要 3 个子进程。那么,如果我不使用池并像我一样手动创建进程,会有什么问题吗?在这两种情况下都可能存在pickle-unpickle 和内部内存分配和映射问题。如果你能在这方面提供帮助。或者如果其他人可以提供帮助。

标签: python arrays numpy process multiprocessing


【解决方案1】:

这是使用np.memmapPool 的示例。看到您可以定义进程和工作人员的数量。在这种情况下,您无法控制队列,这可以使用multiprocessing.Queue 来实现:

from multiprocessing import Pool

import numpy as np

def mysum(array_file_name, col1, col2, shape):
    a = np.memmap(array_file_name, shape=shape, mode='r+')
    a[:, col1:col2] = np.random.random((shape[0], col2-col1))
    ans = a[:, col1:col2].sum()
    del a
    return ans

if __name__ == '__main__':
    nop = 1000 # number_of_processes
    now = 3 # number of workers
    p = Pool(now)
    array_file_name = 'test.array'
    shape = (250000, 250000)
    a = np.memmap(array_file_name, shape=shape, mode='w+')
    del a
    cols = [[shape[1]*i/nop, shape[1]*(i+1)/nop] for i in range(nop)]
    results = []
    for c1, c2 in cols:
        r = p.apply_async(mysum, args=(array_file_name, c1, c2, shape))
        results.append(r)
    p.close()
    p.join()

    final_result = sum([r.get() for r in results])
    print final_result

如果可能,您可以使用共享内存并行处理获得更好的性能。请参阅此相关问题:

【讨论】:

  • @Saulllo,抱歉回复晚了,感谢提供代码。我试过了,没问题,但需要一点改变。 results.append(r) 不会将结果附加到结果中。我必须为它指定回调。像 r = pool.map_async(print_num, tasks, callback = results.append) 然后它工作。但无论如何,这是一个很好的例子。 np.memmap 是个好主意。我的问题还没有解决,它给出了像 Exception Type: UnpickleableError Exception Value: Cannot pickle objects 我必须以某种方式解决的错误,但它不适用于数组,我现在不发送它们超过进程。
  • 错误的详细调用栈::: response = middleware_method(request, response); request.session.save(); session_data = self.encode(self._get_session(no_load=must_create)); pickled = pickle.dumps(session_dict, pickle.HIGHEST_PROTOCOL)
【解决方案2】:

我的问题似乎已解决。我在里面使用了一个 django 模块,我在里面调用了 multiprocessing.pool.map_async。我的工作函数是类本身内部的一个函数。这就是问题所在。多进程不能在另一个进程中调用同一类的函数,因为子进程不共享内存。所以在子进程内部没有类的活动实例。可能这就是它没有被调用的原因。据我了解。我从类中删除了该函数,并将其放在同一个文件中,但在类之外,就在类定义开始之前。有效。我也得到了适度的加速。还有一件事是面临同样问题的人请不要读取大型数组并在进程之间传递。酸洗和解酸会花费很多时间,而且您不会加快速度而是减慢速度。尝试读取子进程内部的数组。

如果可能,请使用 numpy.memmap 数组,它们非常快。

【讨论】:

  • 抱歉重新打开一个老问题,但是如果你将池化函数拉出类,它还能调用其他类函数还是需要删除这些函数?我实际上遇到了和你一样的问题。
  • 我认为它不能调用其他类函数。我所做的就是将该函数保留在课堂之外,使用管道将一些数据推送到其中,并在完成后从中获取结果。我将它用作完成一项工作的实用程序功能。不与他人交流。
猜你喜欢
  • 1970-01-01
  • 2018-08-18
  • 1970-01-01
  • 2017-01-24
  • 2021-02-14
  • 2016-03-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多