【问题标题】:Process finishes but cannot be joined?进程完成但无法加入?
【发布时间】:2019-07-05 10:08:25
【问题描述】:

为了加速某项任务,我将 Process 子类化以创建一个将处理来自样本的数据的工作人员。一些管理类将为其提供数据并读取输出(使用两个Queue 实例)。对于异步操作,我使用put_nowaitget_nowait。最后,我向我的进程发送了一个特殊的退出代码,它打破了它的内部循环。然而……它永远不会发生。这是一个最小的可重现示例:

import multiprocessing as mp

class Worker(mp.Process):
  def __init__(self, in_queue, out_queue):
    super(Worker, self).__init__()
    self.input_queue = in_queue
    self.output_queue = out_queue

  def run(self):
    while True:
      received = self.input_queue.get(block=True)
      if received is None:
        break
      self.output_queue.put_nowait(received)
    print("\tWORKER DEAD")


class Processor():
  def __init__(self):
    # prepare
    in_queue = mp.Queue()
    out_queue = mp.Queue()
    worker = Worker(in_queue, out_queue)
    # get to work
    worker.start()
    in_queue.put_nowait(list(range(10**5))) # XXX
    # clean up
    print("NOTIFYING")
    in_queue.put_nowait(None)
    #out_queue.get() # XXX
    print("JOINING")
    worker.join()

Processor()

这段代码永远不会完成,像这样永久挂起:

NOTIFYING
JOINING
    WORKER DEAD

为什么?

我用XXX 标记了两行。在第一个中,如果我发送的数据较少(例如,10**4),一切都会正常完成(进程按预期加入)。同样在第二个中,如果我在通知工人完成后get()。我知道我遗漏了一些东西,但 documentation 中似乎没有任何相关内容。

【问题讨论】:

    标签: python-3.x queue python-multiprocessing


    【解决方案1】:

    文档提到

    当一个对象被放入队列时,该对象被腌制,然后后台线程将腌制的数据刷新到底层管道。这会产生一些后果 [...] 将对象放入空队列后,在队列的 empty() 方法返回 False 和 get_nowait() 可以在不引发 queue.Empty 的情况下返回之前可能会有一个无限小的延迟。

    https://docs.python.org/3.7/library/multiprocessing.html#pipes-and-queues

    还有那个

    无论何时使用队列,您都需要确保所有已放入队列的项目最终都会在进程加入之前被移除。否则,您无法确定将项目放入队列的进程将终止。

    https://docs.python.org/3.7/library/multiprocessing.html#multiprocessing-programming

    这意味着您描述的行为可能是由工作人员中的self.output_queue.put_nowait(received) 与处理器__init__ 中的worker.join() 加入工作人员之间的竞争条件引起的。如果加入比将其送入队列更快,那么一切都很好。如果太慢,队列中有一个项目,worker不会加入。

    取消注释主进程中的out_queue.get() 将清空队列,从而允许加入。但是,如果队列已经为空,则返回队列很重要,因此使用超时可能是尝试等待竞速条件结束的一种选择,例如 out_qeue.get(timeout=10)

    保护主例程可能也很重要,尤其是对于 Windows (python multiprocessing on windows, if __name__ == "__main__")

    【讨论】:

    • 您要求我使用JoinableQueue 而不仅仅是Queue。我试过了,但我得到了ValueError: task_done() called too many times。所以我应该理解一些关于队列/管道的事情。此外,即使您的建议有效,我仍然会问:这背后的原因是什么? “任务已完成”是什么意思,我为什么需要这样做?如果这很重要,那么为什么会有没有这种能力的队列类呢?
    • 对于加入队列的困惑,我们深表歉意。我的印象是多处理使用与队列模块相同的队列,但显然没有。
    • 也许这有助于更好地理解 mp 队列:docs.python.org/3.7/library/…:“如果子进程已将项目放入队列 [...],则该进程将不会终止,直到所有缓冲项目都已已冲入管道。”
    • 另外:请记住,将项目放入队列的进程将在终止之前等待,直到所有缓冲的项目都由“馈送”线程馈送到底层管道。 docs.python.org/3.7/library/…
    • 您好,感谢您的反馈。我目前正在路上,但会尽快将其包装成正确的答案。
    猜你喜欢
    • 2019-10-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2010-09-07
    • 2018-11-08
    • 1970-01-01
    相关资源
    最近更新 更多