Python多处理模块中apply()和apply_async()之间的区别

Pie*_*esi 5 python parallel-processing multiprocessing

我目前有一段代码产生多个进程,如下所示:

pool = Pool(processes=None)
results = [pool.apply(f, args=(arg1, arg2, arg3)) for arg3 in arg_list]
Run Code Online (Sandbox Code Playgroud)

我的想法是,这将使用之后可用的所有核心将工作划分为核心processes=None.但是,多处理模块docs中Pool.apply()方法的文档如下:

相当于apply()内置函数.它会一直阻塞,直到结果准备就绪,因此apply_async()更适合并行执行工作.此外,func仅在池中的一个工作程序中执行.

第一个问题: 我不清楚这一点.如何apply在工人之间分配工作,以及与什么方式不同apply_async?如果任务分布在工人之间,那怎么可能func只在其中一个工人中执行?

我的猜测:我的猜测是apply,在我当前的实现中,给一个具有一组参数的工作者提供一个任务,然后等待该工作完成,然后将下一组参数提供给另一个工作者.通过这种方式,我将工作发送到不同的流程,但没有发生并行性.这似乎是这种情况,因为apply事实上只是:

def apply(self, func, args=(), kwds={}):
    '''
    Equivalent of `func(*args, **kwds)`.
    Pool must be running.
    '''
    return self.apply_async(func, args, kwds).get()
Run Code Online (Sandbox Code Playgroud)

第二个问题:我也想更好地理解为什么在出台的文件,部分16.6.1.5.('使用工人池'),他们说,即使是apply_async像这样 的建筑[pool.apply_async(os.getpid, ()) for i in range(4)] 可能会使用更多的工艺,但它不确定它会.是什么决定是否使用多个流程?

Ili*_*ija 10

您已指向 Python2.7 文档,因此我将根据 Python2.7 多处理实现来回答我的问题。它可能在 Python3.X 上有所不同,但应该不会有太大不同。

apply和之间的区别apply_async

当您查看它们在下面实际实现时,这两者之间的区别实际上是自我描述的。在这里,我将复制/粘贴multiprocessing/pool.py机器人功能的代码。

def apply(self, func, args=(), kwds={}):
    '''
    Equivalent of `apply()` builtin
    '''
    assert self._state == RUN
    return self.apply_async(func, args, kwds).get()
Run Code Online (Sandbox Code Playgroud)

如您所见,apply实际上是在调用,apply_async但就在返回结果之前,get被调用。这基本上会apply_async阻塞直到返回结果。

def apply_async(self, func, args=(), kwds={}, callback=None):
    '''
    Asynchronous equivalent of `apply()` builtin
    '''
    assert self._state == RUN
    result = ApplyResult(self._cache, callback)
    self._taskqueue.put(([(result._job, None, func, args, kwds)], None))
    return result
Run Code Online (Sandbox Code Playgroud)

apply_async将任务放入任务队列中,并返回一个handle已提交的任务。有了它,handle您可以分别调用get或wait获取结果或等待任务完成。任务完成后,它返回的内容作为参数传递给callback函数。

例子:

from multiprocessing import Pool
from time import sleep


def callback(a):
    print a


def worker(i, n):
    print 'Entering worker ', i
    sleep(n)
    print 'Exiting worker'
    return 'worker_response'


if __name__ == '__main__':
    pool = Pool(4)
    a = [pool.apply_async(worker, (i, 4), callback=callback) for i in range(8)]
    for i in a:
        i.wait()
Run Code Online (Sandbox Code Playgroud)

结果:

Entering worker  0
Entering worker  1
Entering worker  2
Entering worker  3
Exiting worker
Exiting worker
Exiting worker
Exiting worker
Entering worker  4
Entering worker  5
worker_response
Entering worker  6
worker_response
Entering worker  7
worker_response
worker_response
Exiting worker
Exiting worker
Exiting worker
Exiting worker
worker_response
worker_response
worker_response
worker_response
Run Code Online (Sandbox Code Playgroud)

注意,使用时apply_async必须等待结果或等待任务完成。如果你不评论我例子的最后两行,你的脚本会在你运行后立即完成。

为什么apply_async可能会使用更多的进程

关于如何apply描述和工作,我理解这一点。由于apply运行任务将其发送到 a 中的可用进程Pool,因此apply_async将任务添加到队列中,然后任务队列线程将它们发送到 中的可用进程Pool。这就是为什么当您使用apply_async.

为了更好地理解作者试图传达的想法,我多次浏览了这一部分。让我们在这里检查一下:

# evaluate "os.getpid()" asynchronously
res = pool.apply_async(os.getpid, ()) # runs in *only* one process
print res.get(timeout=1)              # prints the PID of that process

# launching multiple evaluations asynchronously *may* use more processes
multiple_results = [pool.apply_async(os.getpid, ()) for i in range(4)]
print [res.get(timeout=1) for res in multiple_results]
Run Code Online (Sandbox Code Playgroud)

如果我们试图通过查看前一个示例来理解最后一个示例,那么当您多次连续调用apply_async它时,当然可能会同时运行更多的调用。这可能取决于当时Pool使用了多少进程。这就是为什么他们说可能。