【问题标题】:parallelized algorithm for evaluating a 1-d array of functions on a same-length 1d numpy array用于在相同长度的一维 numpy 数组上评估一维函数数组的并行算法
【发布时间】: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


【解决方案1】:

警告/警告:

您可能不想将multiprocessing 应用于此问题。你会发现对大型数组的操作相对简单,问题将是内存绑定numpy。瓶颈是将数据从 RAM 移动到 CPU 缓存。 CPU 缺乏数据,因此在问题上投入更多的 CPU 并没有多大帮助。此外,您当前的方法将腌制并为输入序列中的每个项目制作整个数组的副本,这会增加大量开销。

numpy + multiprocessing 在很多情况下非常有效,但您需要确保处理的是 CPU 密集型问题。理想情况下,这是一个 CPU 密集型问题,输入和输出相对较小,以减轻酸洗输入和输出的开销。对于numpy 最常用于解决的许多问题,情况并非如此。


您当前方法的两个问题

回答你的问题:

您的直接错误是由于传入了一个无法从全局范围访问的函数(即在函数内部定义的函数)。

但是,您还有另一个问题。您将 numpy 数组视为可以由每个进程修改的共享内存。相反,当使用multiprocessing 时,原始数组将被腌制(有效地制作副本)并独立地传递给每个进程。原始数组永远不会被修改。


避免PicklingError

作为重现错误的最小示例,请考虑以下几点:

import multiprocessing

def apply_parallel(input_sequence):
    def func(x):
        pass
    pool = multiprocessing.Pool()
    pool.map(func, input_sequence)
    pool.close()

foo = range(100)
apply_parallel(foo)

这将导致:

PicklingError: Can't pickle <type 'function'>: attribute lookup 
               __builtin__.function failed

当然,在这个简单的示例中,我们可以简单地将函数定义移回__main__ 命名空间。但是,在你的例子中,你需要它来引用传入的数据。让我们看一个更接近你正在做的例子:

import numpy as np
import multiprocessing

def parallel_rolling_mean(data, window):
    data = np.pad(data, window, mode='edge')
    ind = np.arange(len(data)) + window

    def func(i):
        return data[i-window:i+window+1].mean()

    pool = multiprocessing.Pool()
    result = pool.map(func, ind)
    pool.close()
    return result

foo = np.random.rand(20).cumsum()
result = parallel_rolling_mean(foo, 10)

有多种方法可以处理这个问题,但常见的方法是:

import numpy as np
import multiprocessing

class RollingMean(object):
    def __init__(self, data, window):
        self.data = np.pad(data, window, mode='edge')
        self.window = window

    def __call__(self, i):
        start = i - self.window
        stop = i + self.window + 1
        return self.data[start:stop].mean()

def parallel_rolling_mean(data, window):
    func = RollingMean(data, window)
    ind = np.arange(len(data)) + window

    pool = multiprocessing.Pool()
    result = pool.map(func, ind)
    pool.close()
    return result

foo = np.random.rand(20).cumsum()
result = parallel_rolling_mean(foo, 10)

太棒了!有效!


但是,如果您将其扩展到大型阵列,您很快就会发现它要么运行得非常慢(您可以通过在 pool.map 调用中增加 chunksize 来加快速度)或者您将快速运行内存不足(一旦你增加了chunksize)。

multiprocessing 腌制输入,以便可以将其传递给单独且独立的 python 进程。这意味着您正在为您操作的每个 i 复制整个 数组。

我们稍后会回到这一点......


multiprocessing 进程之间不共享内存

multiprocessing 模块通过挑选输入并将它们传递给独立进程来工作。这意味着,如果您在一个进程中修改某些内容,其他进程将看不到修改。

不过multiprocessing也提供primitives that live in shared memory,可以被子进程访问和修改。有一个few different waysadapting numpy arrays 使用共享内存multiprocessing.Array。但是,我建议首先避免使用这些(如果您不熟悉,请阅读 false sharing)。在某些情况下它非常有用,但通常是为了节省内存,而不是提高性能。

因此,最好在单个进程中对大型数组进行所有修改(这对于一般 IO 来说也是非常有用的模式)。它不一定是“主要”过程,但这样想是最容易的。

例如,假设我们想让parallel_rolling_mean 函数接受一个输出数组来存储内容。一个有用的模式类似于以下内容。注意迭代器的使用和只在主进程中修改输出:

import numpy as np
import multiprocessing

def parallel_rolling_mean(data, window, output):
    def windows(data, window):
        padded = np.pad(data, window, mode='edge')
        for i in xrange(len(data)):
            yield padded[i:i + 2*window + 1]

    pool = multiprocessing.Pool()
    results = pool.imap(np.mean, windows(data, window))
    for i, result in enumerate(results):
        output[i] = result
    pool.close()

foo = np.random.rand(20000000).cumsum()
output = np.zeros_like(foo)
parallel_rolling_mean(foo, 10, output)
print output

希望这个例子有助于澄清一些事情。


chunksize 和性能

关于性能的一个简短说明:如果我们扩大它,它会很快变得非常缓慢。如果您查看系统监视器(例如top/htop),您可能会注意到您的内核大部分时间都处于空闲状态。

默认情况下,主进程为每个进程腌制每个输入并立即将其传入,然后等待它们完成腌制下一个输入。在许多情况下,这意味着主进程工作,然后在工作进程忙时处于空闲状态,然后在主进程正在酸洗下一个输入时,工作进程处于空闲状态。

关键是增加chunksize参数。这将导致pool.imap 为每个进程“预先腌制”指定数量的输入。基本上,主线程可以保持忙碌的酸洗输入,而工作进程可以保持忙碌的处理。缺点是您正在使用更多内存。如果每个输入都占用大量 RAM,这可能是个坏主意。但是,如果没有,这可以显着加快速度。

举个简单的例子:

import numpy as np
import multiprocessing

def parallel_rolling_mean(data, window, output):
    def windows(data, window):
        padded = np.pad(data, window, mode='edge')
        for i in xrange(len(data)):
            yield padded[i:i + 2*window + 1]

    pool = multiprocessing.Pool()
    results = pool.imap(np.mean, windows(data, window), chunksize=1000)
    for i, result in enumerate(results):
        output[i] = result
    pool.close()

foo = np.random.rand(2000000).cumsum()
output = np.zeros_like(foo)
parallel_rolling_mean(foo, 10, output)
print output

使用chunksize=1000,处理一个200万个元素的数组需要21秒:

python ~/parallel_rolling_mean.py  83.53s user 1.12s system 401% cpu 21.087 total

但是使用chunksize=1(默认值)大约需要八倍的时间(2 分 41 秒)。

python ~/parallel_rolling_mean.py  358.26s user 53.40s system 246% cpu 2:47.09 total

其实用默认的chunksize,其实比单进程实现同样的事情要差很多,只需要45秒:

python ~/sequential_rolling_mean.py  45.11s user 0.06s system 99% cpu 45.187 total

【讨论】:

  • 哇,非常有启发性的答案,乔。我已经读了好几遍,并不断从中学习新东西。非常感谢您抽出宝贵时间撰写如此周到的答案。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-08-14
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多