【问题标题】:Multiprocessing Queue - child processes gets stuck sometimes and does not reap多处理队列 - 子进程有时会卡住并且不会收获
【发布时间】:2019-06-28 04:10:54
【问题描述】:

首先,如果标题有点奇怪,我深表歉意,但我真的想不出如何将我面临的问题放在一行中。

所以我有以下代码

import time
from multiprocessing import Process, current_process, Manager
from multiprocessing import JoinableQueue as Queue

# from threading import Thread, current_thread
# from queue import Queue


def checker(q):
    count = 0
    while True:
        if not q.empty():
            data = q.get()
            # print(f'{data} fetched by {current_process().name}')
            # print(f'{data} fetched by {current_thread().name}')
            q.task_done()
            count += 1
        else:
            print('Queue is empty now')
            print(current_process().name, '-----', count)
            # print(current_thread().name, '-----', count)


if __name__ == '__main__':
    t = time.time()
    # m = Manager()
    q = Queue()
    # with open("/tmp/c.txt") as ifile:
    #     for line in ifile:
    #         q.put((line.strip()))
    for i in range(1000):
        q.put(i)
    time.sleep(0.1)
    procs = []
    for _ in range(2):
        p = Process(target=checker, args=(q,), daemon=True)
        # p = Thread(target=checker, args=(q,))
        p.start()
        procs.append(p)
    q.join()
    for p in procs:
        p.join()

示例输出

1:当进程刚刚挂起时

Queue is empty now
Process-2 ----- 501
output hangs at this point

2:当一切正常时。

Queue is empty now
Process-1 ----- 515
Queue is empty now
Process-2 ----- 485

Process finished with exit code 0

现在这种行为是间歇性的,有时会发生,但并非总是如此。

我也尝试使用Manager.Queue() 来代替multiprocessing.Queue(),但没有成功,并且都出现了同样的问题。

我用multiprocessingmultithreading 对此进行了测试,我得到了完全相同的行为,与multithreading 相比,与multiprocessing 相比,这种行为的发生率要低得多。

所以我认为我在概念上遗漏了一些东西或做错了,但我现在无法抓住它,因为我在这方面花费了太多时间,现在我的脑海中没有看到可能非常基本的东西。

感谢您的帮助。

【问题讨论】:

  • 显然如果在队列q为空之前调用join(),可能会出现死锁:stackoverflow.com/questions/31665328/…
  • 即使我在加入进程之前添加了q.join(),它仍然无法按预期工作。
  • 可能是因为multiprocessing.Queue 没有任何名为join() 的方法。来自文档:Queue implements all the methods of queue.Queue except for task_done() and join(). 第一个“队列”是指您在上面的示例中使用的multiprocessing.Queue
  • 我试过JoinabaleQueue
  • @TuanDT 供您参考刚刚更新了问题。

标签: python multiprocessing


【解决方案1】:

我相信您在 checker 方法中存在竞争条件。您检查队列是否为空,然后在单独 步骤中将下一个任务出列。将这两种操作分开而不互斥或锁定通常不是一个好主意,因为队列的状态可能会在检查和弹出之间发生变化。它可能是非空的,但另一个进程可能会在通过检查的进程能够这样做之前将等待的工作出列。

但是,我通常更喜欢沟通而不是锁定;它更不容易出错,并使一个人的意图更清晰。在这种情况下,我会向工作进程发送一个标记值(例如None),以指示所有工作都已完成。然后每个工作人员将下一个对象(始终是线程安全的)出列,如果对象是None,则子进程退出。

下面的示例代码是您的程序的简化版本,应该可以在没有竞争的情况下工作:

def checker(q):
    while True:
        data = q.get()
        if data is None:
            print(f'process f{current_process().name} ending')
            return
        else:
            pass # do work

if __name__ == '__main__':
    q = Queue()
    for i in range(1000):
        q.put(i)
    procs = []
    for _ in range(2):
        q.put(None) # Sentinel value
        p = Process(target=checker, args=(q,), daemon=True)
        p.start()
        procs.append(p)
    for proc in procs:
        proc.join()

【讨论】:

  • 你能解释一下你所说的队列是线程安全的,因为这意味着我们在线程函数中使用时不必在队列周围加锁。此外,我确实在检查 q 是否不为空的同一步骤中从队列中获取值。那么单独的步骤是怎样的呢?
  • 线程安全 意味着一个操作可以在多个线程(或进程)中使用而没有竞争条件。所以不需要显式锁定。你确实有不同的步骤:你首先调用if not q.empty(),然后在下一行调用q.get()。在这两行之间,另一个进程可能会从队列中检索下一项,这意味着该进程会“认为”队列中有数据,但实际上在尝试检索时它是空的。
  • 我不明白,因为如果队列线程安全,那么 q.emptyq.get 如何导致竞争条件,因为这些是队列上的方法。
  • 队列上的单个方法是线程安全的,但这并不意味着它们上的任何操作序列都是线程安全的。后一种想法不可能奏效,因为我们可以将不同顺序的操作放在一起,对吧?队列如何“知道”某些操作集构成线程安全序列?必须有某种锁,上面写着“在我解锁之前,所有操作都是线程安全的”。但根据定义,这是一种互斥机制,与队列本身是分开的。
  • 是的,没错。 q.empty() 可以返回True,但随后的q.get() 调用仍然会阻塞,因为队列已被它们之间的另一个进程清空。这是教科书的比赛条件。您可以添加互斥,但添加哨兵很可能会更快、更清晰、更不容易出错,就像我在示例中所做的那样。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-08-18
  • 1970-01-01
  • 2015-09-19
  • 2017-09-07
  • 1970-01-01
  • 2012-03-07
相关资源
最近更新 更多