【问题标题】:Easy parallelization of numpy.apply_along_axis()?numpy.apply_along_axis() 的简单并行化?
【发布时间】:2018-01-13 13:36:18
【问题描述】:

如何通过numpy.apply_along_axis() 将函数应用于 NumPy 数组的元素,以便利用多核?这似乎是一件很自然的事情,在对正在应用的函数的所有调用都是独立的常见情况下。

在我的特殊情况下——如果这很重要——应用轴是轴 0:np.apply_along_axis(func, axis=0, arr=param_grid)np 是 NumPy)。

我快速浏览了 Numba,但我似乎无法获得这种并行化,循环如下:

@numba.jit(parallel=True)
result = np.empty(shape=params.shape[1:])
for index in np.ndindex(*result.shape)):  # All the indices of params[0,...]
    result[index] = func(params[(slice(None),) + index])  # Applying func along axis 0

显然还有一个NumPy 中的编译选项用于通过 OpenMP 进行并行化,但它似乎无法通过 MacPorts 访问。

也可以考虑将数组分割成几块,然后使用线程(以避免复制数据)并在每块上并行应用函数。这比我正在寻找的更复杂(如果全局解释器锁没有被充分释放,可能无法正常工作)。

如果能够以简单的方式使用多个内核来执行简单的可并行任务,例如将函数应用于数组的所有元素(这基本上是这里需要的,函数 @ 987654327@ 采用一维参数数组。

【问题讨论】:

  • apply_along_axis 是纯 Python 代码,执行您所展示的操作,除了它将感兴趣的轴移到末尾,其余部分执行 ndindex(arr.shape[:-1])。替代方案已在stackoverflow.com/questions/45067268/… 等帖子中讨论过
  • 由于可以将 n-d 问题重新塑造为 2d(您的兴趣轴加上其余部分),因此基本问题是 1d 列表理解。遍历行。另一个 SO 问题:stackoverflow.com/questions/44239498/…
  • 我希望这些 StackOverflow 问题包含一个我可以使用的使用多个内核的解决方案!现在,我不确定 Python 列表理解如何成功地比 np.apply_along_axis() 更快,但至少可以通过探索 np.apply_along_axis() 的简单替代方案来使单核版本更快……
  • 各种 SO 已经研究过在数组的行上使用 multiprocessing.pool.map

标签: python arrays performance numpy parallel-processing


【解决方案1】:

好的,我解决了:一个想法是使用标准的multiprocessing 模块并将原始数组分成几块(以限制与工作人员的通信开销)。这可以相对容易地完成,如下所示:

import multiprocessing

import numpy as np

def parallel_apply_along_axis(func1d, axis, arr, *args, **kwargs):
    """
    Like numpy.apply_along_axis(), but takes advantage of multiple
    cores.
    """        
    # Effective axis where apply_along_axis() will be applied by each
    # worker (any non-zero axis number would work, so as to allow the use
    # of `np.array_split()`, which is only done on axis 0):
    effective_axis = 1 if axis == 0 else axis
    if effective_axis != axis:
        arr = arr.swapaxes(axis, effective_axis)

    # Chunks for the mapping (only a few chunks):
    chunks = [(func1d, effective_axis, sub_arr, args, kwargs)
              for sub_arr in np.array_split(arr, multiprocessing.cpu_count())]

    pool = multiprocessing.Pool()
    individual_results = pool.map(unpacking_apply_along_axis, chunks)
    # Freeing the workers:
    pool.close()
    pool.join()

    return np.concatenate(individual_results)

Pool.map() 中应用的函数unpacking_apply_along_axis() 应该是独立的(以便子进程可以导入它),并且只是处理Pool.map() 仅接受一个参数这一事实的薄包装器:

def unpacking_apply_along_axis((func1d, axis, arr, args, kwargs)):
    """
    Like numpy.apply_along_axis(), but with arguments in a tuple
    instead.

    This function is useful with multiprocessing.Pool().map(): (1)
    map() only handles functions that take a single argument, and (2)
    this function can generally be imported from a module, as required
    by map().
    """
    return np.apply_along_axis(func1d, axis, arr, *args, **kwargs)

(在 Python 3 中,这应该写成

def unpacking_apply_along_axis(all_args):
    (func1d, axis, arr, args, kwargs) = all_args

因为argument unpacking was removed)。

在我的特殊情况下,这导致 2 个具有超线程的内核的速度提高了 2 倍。接近 4 倍的因子会更好,但速度已经很不错了,只需几行代码,对于具有更多内核的机器(这很常见)应该更好。也许有一种方法可以避免数据复制和使用共享内存(可能通过multiprocessing module 本身)?

【讨论】:

  • 感谢您在上面回答中的代码,这对我来说很有希望,因为我也面临着同样的问题,一个令人尴尬的可并行化问题,而np.apply_along_axis() 只能让我到目前为止。我要补充的另一个想法是 dask 可以帮助并行化,据说无论如何我还没有设法让它工作,因为它需要使用 dask 数组而不是 numpy 数组,并且并非所有 numpy 功能都可以从 dask API。见dask.org
  • Dask 确实是一个不错的尝试工具。 Vaex 可能是另一个 (github.com/maartenbreddels/vaex)。
  • 在使用您在上面提供的代码的修改版本之后,我能够重新设计我的多处理方法(我之前使用了一种不同类型的循环,它没有利用拆分和连接) ,这样做给了我类似的性能,但内存占用大大减少。再次感谢,这对我帮助很大。
  • 我在 unpacking_apply_along_axis 的函数定义上遇到语法错误,可能与使用双 ((。使用 Python 3 时出现语法错误,在 Python 2 中它似乎有效。我不能'不知道究竟是为什么。对此有任何帮助吗?
  • 确实如此。我刚刚在答案中添加了 Python 3 案例。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-09-30
  • 2012-04-04
  • 2021-07-17
  • 1970-01-01
相关资源
最近更新 更多