【问题标题】:How can you code a nested concurrency in python?如何在 python 中编写嵌套并发?
【发布时间】:2019-05-28 05:04:30
【问题描述】:

我的代码有以下方案:

class A():
    def evaluate(self):
        b = B()
        for i in range(30):
            b.run()

class B():
    def run(self):
        pass

if __name__ == '__main__':
    a = A()
    for i in range(10):
        a.evaluate()

我想要有两个级别的并发,第一个是evaluate 方法,第二个是run 方法(嵌套并发)。问题是如何使用Pool class of the multiprocessing 模块引入这种并发性?我应该明确传递核心数量吗?该解决方案不应创建大于multiprocessing.cpu_count() 数量的进程。

注意:假设核心数大于 10。

编辑: 我看到很多 cmets 说由于 GIL,python 没有真正的并发性,这对于 python 多线程是正确的,但对于多处理,这不是正确的外观here,我也已经计时了这个@987654323 @doed,结果表明它可以比顺序执行更快。

【问题讨论】:

  • 您将创建嵌套进程。管理器应负责维护最大进程数。
  • 如果你的程序是计算密集型的,你可能不会从 Python 的并发中获得任何好处。
  • 根据 Pool 的文档,如果没有明确传递,Pool 将创建等于 cpu_count() 的工人数量。因此,根据我的理解@juanpa.arrivillaga,如果我两次创建了一个 Pool 的对象,它可能会创建太多的工作人员。
  • 为什么要创建两次 Pool 对象?
  • 对于每个并发级别,我必须创建 Pool 对象。实际上这是我的问题之一,我不知道如何在不创建两个池的情况下对其进行编码? @juanpa.arrivillaga

标签: python python-3.x concurrency nested python-multiprocessing


【解决方案1】:

您的评论涉及一个可能的解决方案。为了获得“嵌套”并发,您可以拥有 2 个单独的池。这将导致“平面”结构程序而不是嵌套程序。此外,它将 A 与 B 分离,A 现在对 b 一无所知,它只是发布到一个通用队列。下面的示例使用单个进程来说明如何连接通过异步队列进行通信的并发工作人员,但它可以很容易地替换为池:

import multiprocessing as mp


class A():
    def __init__(self, in_q, out_q):
      self.in_q = in_q
      self.out_q = out_q

    def evaluate(self):
        """
        Reads from input does work and process output
        """
        while True:
          job = self.in_q.get()
          for i in range(30):
            self.out_q.put(i)

class B():
    def __init__(self, in_q):
      self.in_q = in_q

    def run(self):
        """
        Loop over queue and process items, optionally configure
        with another queue to "sink" the processing pipeline
        """
        while True:
           job = self.in_q.get()

if __name__ == '__main__':
    # create the queues to wire up our concurrent worker pools
    A_q = mp.Queue()
    AB_q = mp.Queue()

    a = A(in_q=A_q, out_q=AB_q)
    b = B(in_q=AB_q)

    p = mp.Process(target=a.evaluate)
    p.start()

    p2 = mp.Process(target=b.run)
    p2.start()

    for i in range(10):
        A_q.put(i)

    p.join()
    p2.join()

这是 golang 中常见的模式。

【讨论】:

    猜你喜欢
    • 2020-03-17
    • 2018-03-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-08-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多