【问题标题】:Why doesn't Python process with input and output queues not join once it is done?为什么 Python 处理完成后不加入输入和输出队列?
【发布时间】:2017-07-27 08:21:21
【问题描述】:

这个使用 multiprocessing 的简单 Python3 程序似乎无法按预期工作。

所有输入进程共享一个输入队列,它们从中消费数据。他们都共享一个输出队列,一旦完成,他们就会在其中写入结果。我发现这个程序挂在进程join()。这是为什么呢?

#!/usr/bin/env python3

import multiprocessing

def worker_func(in_q, out_q):
    print("A worker has started")    
    w_results = {}
    while not in_q.empty():
        v = in_q.get()
        w_results[v] = v
    out_q.put(w_results)
    print("A worker has finished")

def main():

    # Input queue to share among processes
    fpaths = [str(i) for i in range(10000)]
    in_q = multiprocessing.Queue()
    for fpath in fpaths:
        in_q.put(fpath)

    # Create processes and start them
    N_PROC = 2
    out_q = multiprocessing.Queue()
    workers = []
    for _ in range(N_PROC):
        w = multiprocessing.Process(target=worker_func, args=(in_q, out_q,))
        w.start()
        workers.append(w)
    print("Done adding workers")

    # Wait for processes to finish
    for w in workers:
        w.join()
    print("Done join of workers")

    # Collate worker results
    out_results = {}
    while not out_q.empty():
        out_results.update(out_q.get())

if __name__ == "__main__":
    main()

N_PROC = 2:

$ python3 test.py
Done adding workers
A worker has started
A worker has started
A worker has finished
<---- I do not get "A worker has finished" from second worker
<---- I do not get "Done join of workers"

即使使用单个子进程N_PROC = 1,它也不起作用:

$ python3 test.py
Done adding workers
A worker has started
A worker has finished
<---- I do not get "Done join of workers"

如果我尝试使用 1000 个项目的较小输入队列,一切正常。

我知道一些旧的 StackOverflow 问题说队列有限制。为什么 Python3 文档中没有记录这一点?

我可以使用什么替代解决方案?我想使用多处理(不是线程),在 N 个进程之间拆分输入。一旦他们的共享输入队列为空,我希望每个进程都收集其结果(可以是像 dict 这样的大/复杂数据结构)并将其返回给父进程。如何做到这一点?

【问题讨论】:

    标签: python queue multiprocessing


    【解决方案1】:

    这是由您的设计引起的经典错误。当工作人员终止时,他们会因为无法将所有数据放入out_q 而停止,从而使您的程序死锁。这与队列下的管道缓冲区的大小有关。

    当您使用multiprocessing.Queue 时,您应该在尝试加入馈送进程之前将其清空,以确保Process 不会因为等待所有对象放入Queue 而停止。因此,在加入流程之前拨打您的out_q.get 电话应该可以解决您的问题:。您可以使用哨兵模式来检测计算的结束。

    #!/usr/bin/env python3
    
    import multiprocessing
    from multiprocessing.queues import Empty
    
    def worker_func(in_q, out_q):
        print("A worker has started")    
        w_results = {}
        while not in_q.empty():
            try:
                v = in_q.get(timeout=1)
                w_results[v] = v
            except Empty:
                pass
        out_q.put(w_results)
        out_q.put(None)
        print("A worker has finished")
    
    def main():
    
        # Input queue to share among processes
        fpaths = [str(i) for i in range(10000)]
        in_q = multiprocessing.Queue()
        for fpath in fpaths:
            in_q.put(fpath)
    
        # Create processes and start them
        N_PROC = 2
        out_q = multiprocessing.Queue()
        workers = []
        for _ in range(N_PROC):
            w = multiprocessing.Process(target=worker_func, args=(in_q, out_q,))
            w.start()
            workers.append(w)
        print("Done adding workers")
    
        # Collate worker results
        out_results = {}
        n_proc_end = 0
        while not n_proc_end == N_PROC:
            res = out_q.get()
            if res is None:
                n_proc_end += 1
            else:
                out_results.update(res)
    
        # Wait for processes to finish
        for w in workers:
            w.join()
        print("Done join of workers")
    
    if __name__ == "__main__":
        main()
    

    另外,请注意您的代码中包含竞争条件。队列in_q 可以在您检查not in_q.empty()get 之间清空。您应该使用非阻塞获取来确保您不会死锁,等待空队列。

    最后,您尝试实现类似于multiprocessing.Pool 的东西,它以更健壮的方式处理这种通信。您还可以查看concurrent.futures API,它更加健壮,从某种意义上说,设计得更好。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-12-27
      • 1970-01-01
      • 1970-01-01
      • 2015-05-17
      • 1970-01-01
      • 1970-01-01
      • 2019-02-15
      相关资源
      最近更新 更多