【问题标题】:Why can't I use join() before closing pool in python multiprocessing为什么我不能在 python 多处理中关闭池之前使用 join()
【发布时间】:2020-09-27 10:43:08
【问题描述】:

我有一个类,它有一个方法可以进行一些并行计算并且经常被调用。因此,我希望我的池在类的构造函数中初始化一次,而不是在每次调用此方法时都创建一个新池。在这种方法中,我想使用 apply_async() 为所有工作进程启动一个任务,然后等待(阻塞)并聚合每个任务的结果。我的代码如下所示:

class Foo:
     def __init__(self, ...):
         # ...
         self.pool = mp.Pool(mp.cpu_count())

     def do_parallel_calculations(self, ...):
         for _ in range(mp.cpu_count()):
              self.pool.apply_async(calc_func, args=(...), callback=aggregate_result)
         
         # wait for results to be aggregated to a global var by the callback
         self.pool.join()  # <-- ValueError: Pool is still running
         
         # do something with the aggregated result of all worker processes

但是,当我运行它时,我在 self.pool.join() 中收到一个错误,上面写着:“ValueError:Pool is still running”。现在,在我看到的所有示例中,self.pool.close() 在 self.pool.join() 之前被调用,我认为这就是我收到此错误的原因,但我不想在那里关闭我的池下次调用此方法!我不能不使用 self.pool.join() 因为我需要一种方法来等待所有进程完成,并且我不想浪费地手动旋转,例如使用“while not global_flag: pass”。

我可以做些什么来实现我想要做的事情?为什么不让多处理让我加入一个仍然开放的池?这似乎是一件非常合理的事情。

【问题讨论】:

  • 您无法加入正在运行的池。所以,你关闭并加入它,然后创建一个新的。这是多处理的代价之一。或者你一开始就不使用异步操作,它有不同的价格标签。

标签: python python-multiprocessing


【解决方案1】:

让我们用一个真实的例子来具体说明:

import multiprocessing as mp


def calc_func(x):
    return x * x


class Foo:
    def __init__(self):
        self.pool = mp.Pool(mp.cpu_count())

    def do_parallel_calculations(self, values):
        results = []
        for value in values:
            results.append(self.pool.apply_async(calc_func, args=(value,)))
        for result in results:
            print(result.get())

if __name__ == '__main__':
    foo = Foo()
    foo.do_parallel_calculations([1,2,3])

【讨论】:

  • 我猜您在我开发答案时或多或少地实现了相同的解决方案。
【解决方案2】:

我想我设法通过在 apply_async() 返回的 AsyncResult 对象上调用 get() 来做到这一点。于是代码变成了:

def do_parallel_calculations(self, ...):
     results = []
     for _ in range(mp.cpu_count()):
          results.append(self.pool.apply_async(calc_func, args=(...)))
     aggregated_result = 0
     for result in results:
          aggregated_result += result.get()

其中 calc_func() 返回单个任务结果,不需要回调和全局变量。

这并不理想,因为我以任意顺序等待它们,而不是按照它们实际完成的顺序(最有效的方法是减少结果),但由于我只有 4 个内核,因此几乎不应该很明显。

【讨论】:

  • 您正在按照提交的顺序等待他们。但是,只要您在继续之前等待所有这些都完成,即使假设第一个提交的是最后一个完成,这也不会比等待完成顺序慢得多:获得第一个结果后,结果为其他任务将可用,通过按完成顺序获取结果可以节省的多余时间只是重复调用 result.get()(立即返回)所花费的时间,加上汇总结果的时间。
猜你喜欢
  • 2011-06-04
  • 2017-11-19
  • 2013-08-13
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-08-09
  • 1970-01-01
  • 2020-12-25
相关资源
最近更新 更多