【问题标题】:Avoiding deadlocks due to queue overflow with multiprocessing.JoinableQueue使用 multiprocessing.JoinableQueue 避免由于队列溢出导致的死锁
【发布时间】:2019-05-01 09:48:57
【问题描述】:

假设我们有一个multiprocessing.Pool,其中工作线程共享一个multiprocessing.JoinableQueue,使工作项出队并可能使更多工作入队:

def worker_main(queue):
    while True:
        work = queue.get()
        for new_work in process(work):
            queue.put(new_work)
        queue.task_done()

当队列填满时,queue.put() 将阻塞。只要至少有一个进程从队列中读取queue.get(),它就会释放队列中的空间以解除对写入者的阻塞。但所有进程都可能同时在 queue.put() 处阻塞。

有没有办法避免像这样被卡住?

【问题讨论】:

    标签: python python-multiprocessing


    【解决方案1】:

    根据process(work) 创建更多项目的频率,除了无限最大大小的队列之外,可能没有其他解决方案。

    简而言之,您的队列必须足够大,以容纳您随时可以拥有的整个积压工作项。


    由于queue is implemented with semaphores,可能确实有a hard size limit of SEM_VALUE_MAX其中in MacOS is 32767。因此,如果这还不够,您将需要对该实现进行子类化或使用put(block=False) 并处理queue.Full(例如,将多余的项目放在其他地方)。

    或者,查看one of the 3rd-party implementations of distributed work item queue for Python

    【讨论】:

    • 我没有指定最大队列大小,但似乎有一个隐含的(在 macOS 上)。
    • 我看不到在上面的例子中如何处理 Queue.Full。我不能失去工作。我可以把它放在一边,但当队列耗尽时,主进程将终止。
    • 我会接受这一点,因为您找到了根本原因,即 macOS 上的 SEM_VALUE_MAX 低得离谱,尽管不幸的是没有好的解决方案。 (我转而使用 Redis,这将比共享内存慢得多,唉。)查看 xnu 源,似乎没有理由设置这个限制,所以我呼吁苹果提高它。跨度>
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-08-02
    • 1970-01-01
    • 1970-01-01
    • 2014-06-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多