【发布时间】:2019-07-05 10:08:25
【问题描述】:
为了加速某项任务,我将 Process 子类化以创建一个将处理来自样本的数据的工作人员。一些管理类将为其提供数据并读取输出(使用两个Queue 实例)。对于异步操作,我使用put_nowait 和get_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