【发布时间】:2015-12-16 15:49:39
【问题描述】:
下面的结果是我有一个令人尴尬的并行 for 循环,我正在尝试线程化。解释这个问题有点冗长,但尽管冗长,我认为这应该是一个相当微不足道的问题,多处理模块旨在轻松解决。
我有一个包含 k 个不同函数的大长度 N 数组,以及一个长度为 N 的 abcissa 数组。感谢@senderle 在Efficient algorithm for evaluating a 1-d array of functions on a same-length 1d numpy array 中描述的巧妙解决方案,我有一个基于numpy 的快速算法,我可以使用它来评估abcissa 处的函数以返回一个长度为N 的纵坐标数组:
def apply_indexed_fast(abcissa, func_indices, func_table):
""" Returns the output of an array of functions evaluated at a set of input points
if the indices of the table storing the required functions are known.
Parameters
----------
func_table : array_like
Length k array of function objects
abcissa : array_like
Length Npts array of points at which to evaluate the functions.
func_indices : array_like
Length Npts array providing the indices to use to choose which function
operates on each abcissa element. Thus func_indices is an array of integers
ranging between 0 and k-1.
Returns
-------
out : array_like
Length Npts array giving the evaluation of the appropriate function on each
abcissa element.
"""
func_argsort = func_indices.argsort()
func_ranges = list(np.searchsorted(func_indices[func_argsort], range(len(func_table))))
func_ranges.append(None)
out = np.zeros_like(abcissa)
for i in range(len(func_table)):
f = func_table[i]
start = func_ranges[i]
end = func_ranges[i+1]
ix = func_argsort[start:end]
out[ix] = f(abcissa[ix])
return out
我现在要做的是使用多处理来并行化该函数内的 for 循环。在描述我的方法之前,为了清楚起见,我将简要概述@senderle 开发的算法是如何工作的。如果你能阅读上面的代码并立即理解,则跳过下一段文字。
首先我们找到对输入func_indices进行排序的索引数组,我们用它来定义长度为k的func_ranges整数数组。 func_ranges 的整数条目控制应用于输入 abcissa 的适当子数组的函数,其工作原理如下。令 f 为输入 func_table 中的第 i 个函数。那么我们应该应用函数 f 的输入 abcissa 切片是 slice(func_ranges[i], func_ranges[i+1])。因此,一旦计算出 func_ranges,我们就可以在输入 func_table 上运行一个简单的 for 循环,并连续将每个函数对象应用于适当的切片,填充我们的输出数组。有关此算法的最小示例,请参见下面的代码。
def trivial_functional(i):
def f(x):
return i*x
return f
k = 250
func_table = np.array([trivial_functional(j) for j in range(k)])
Npts = 1e6
abcissa = np.random.random(Npts)
func_indices = np.random.random_integers(0,len(func_table)-1,Npts)
result = apply_indexed_fast(abcissa, func_indices, func_table)
所以我现在的目标是使用多处理来并行化这个计算。我认为这会很简单,使用我通常的技巧来实现令人尴尬的并行 for 循环。但是我在下面的尝试引发了一个我不明白的异常。
from multiprocessing import Pool, cpu_count
def apply_indexed_parallelized(abcissa, func_indices, func_table):
func_argsort = func_indices.argsort()
func_ranges = list(np.searchsorted(func_indices[func_argsort], range(len(func_table))))
func_ranges.append(None)
out = np.zeros_like(abcissa)
num_cores = cpu_count()
pool = Pool(num_cores)
def apply_funci(i):
f = func_table[i]
start = func_ranges[i]
end = func_ranges[i+1]
ix = func_argsort[start:end]
out[ix] = f(abcissa[ix])
pool.map(apply_funci, range(len(func_table)))
pool.close()
return out
result = apply_indexed_parallelized(abcissa, func_indices, func_table)
PicklingError: Can't pickle <type 'function'>: attribute lookup __builtin__.function failed
我在 SO 的其他地方看到过这个:Multiprocessing: How to use Pool.map on a function defined in a class?。我一一尝试了那里提出的每种方法;在所有情况下,我都会收到“打开的文件过多”错误,因为线程从未关闭,或者适应的算法只是挂起。这似乎应该有一个简单的解决方案,因为这只不过是线程化一个令人尴尬的并行 for 循环。
【问题讨论】:
-
如果我没记错的话,您的问题是由于使用了一个无法从
__main__命名空间访问的函数(即将一个函数传递给在另一个函数内部定义的池,或者在主范围之外) )。 -
是的,我之前遇到过这个建议,但我还没有看到任何人实施解决方案。我的第一个想法是,这是 python 强制执行的一个极其严格的条件。你能看到如何调整这个问题以符合这个要求吗?我不知道该怎么做,因为 func_table 和 func_indices 必须绑定在我的函数的命名空间内。
-
一种解决方案是传入一个类的实例。但是,您有一个更深层次的问题。您将 numpy 数组视为共享内存。相反,将制作一个副本并独立地传递给每个进程。原始数组不会被修改。
-
是的,你是对的,这是一个更深层次的问题。这是我之前第一次尝试以这种方式并行化某些东西,所以我以前没有遇到过这种情况。您知道如何以允许线程共享访问 numpy 数组的方式使用多处理模块吗?或者你能给我指出一个我可以阅读的资源吗?
-
其他关于 numpy 和 multiporocessing 的问题参考
sharedmem, github.com/sturlamolden/sharedmem-numpy。但它只是节省内存,而不是时间。
标签: python performance numpy parallel-processing scientific-computing