【问题标题】:Pool, queue, hang池,队列,挂起
【发布时间】:2017-07-25 19:39:34
【问题描述】:

我想使用队列来保存结果,因为我希望消费者(串行而非并行)在工作人员产生结果时处理工作人员的结果。

现在,我想知道为什么下面的程序会挂起。

import multiprocessing as mp
import time
import numpy as np
def worker(arg):
    time.sleep(0.2)
    q, arr = arg 
    q.put(arr[0])

p = mp.Pool(4)
x = np.array([4,4])
q = mp.Queue()

for i in range(4):
    x[0] = i 
    #worker((q,x))
    p.apply_async(worker, args=((q, x),)) 

print("done_apply")
time.sleep(0.2)
for i in range(4):
    print(q.get())

【问题讨论】:

  • 我不确定我是否明白你在问什么。您显示的代码是否由于某种死锁而无法正常工作,或者它是否可以正常工作并且您正试图针对某些潜在的未来问题对其进行强化?
  • 它挂起。我找到了一个解决方案,但它使用管理器,并且需要复制输入。对不起,有问题的错字:'when' -> 为什么。

标签: python parallel-processing queue multiprocessing pool


【解决方案1】:

Queue 对象无法共享。通过找到这个answer,我首先得出了与OP相同的结论。

不幸的是,这段代码中还有其他问题(这并不能使它与链接的答案完全相同)

  • worker(arg) 应该是 worker(*arg) 以使 args 解包工作。没有它,我的进程也会被锁定(我承认我不知道为什么。它应该抛出异常,但我猜多处理和异常不能很好地协同工作)
  • 将相同的x 传递给工作人员会得到相同的结果(apply 有效,但apply_async 无效

另一件事:为了使代码可移植,将主代码包装为if __name__ == "__main__":,在 Windows 上是必需的,因为进程生成的不同

为我输出 0,3,2,1 的完全固定代码:

import multiprocessing as mp
import time
import numpy as np
def worker(*arg):  # there are 2 arguments to "worker"
#def worker(q, arr):  # is probably even better
    time.sleep(0.2)
    q, arr = arg
    q.put(arr[0])

if __name__ == "__main__":
    p = mp.Pool(4)

    m = mp.Manager()  # use a manager, Queue objects cannot be shared
    q = m.Queue()

    for i in range(4):
        x = np.array([4,4])  # create array each time (or make a copy)
        x[0] = i
        p.apply_async(worker, args=(q, x))

    print("done_apply")
    time.sleep(0.2)
    for i in range(4):
        print(q.get())

【讨论】:

  • 是的,也适用于您要修改的所有其他共享对象(列表、字典)。
  • 每个worker返回一个结果,队列持有它。对于每个结果,除了创建结果的进程之外,任何进程都不会修改结果。所以..这意味着不需要经理?
  • docs.python.org/3/library/…:“在进行并发编程时,通常最好尽量避免使用共享状态。在使用多个进程时尤其如此”。在您的示例中,它在没有经理的情况下锁定。所以我想这是必要的,文档建议你这样做,所以我不会绕过这个建议......问题不在于结果,而是Queue对象的共享。
  • 多处理的日志示例使用没有管理器的队列:docs.python.org/3/howto/logging-cookbook.html
【解决方案2】:

将 apply_async 更改为 apply 会给出错误消息:

"Queue objects should only be shared between processes through inheritance"

解决方案:

import multiprocessing as mp
import time
import numpy as np
def worker(arg):
    time.sleep(0.2)
    q, arr = arg
    q.put(arr[0])

p = mp.Pool(4)
x = np.array([4,4])
m = mp.Manager()
q = m.Queue()

for i in range(4):
    x[0] = i
    #worker((q,x))
    p.apply_async(worker, args=((q, x),))

print("done_apply")
time.sleep(0.2)
for i in range(4):
    print(q.get())

结果:

done_apply
3
3
3
3

显然,我需要手动制作 numpy 数组的副本,因为所需的结果应该是 0、1、2、3 的任意顺序,而不是 3、3、3、3。

【讨论】:

  • 我还听说manager是由主进程管理的。这让经理变慢了。
【解决方案3】:

我认为您选择使用multiprocessing.Pool 和您自己的queue 是您遇到的主要问题的根源。使用池会预先创建子进程,然后将作业分配给这些子进程。但是由于您不能(轻松)将queue 传递给已经存在的进程,因此这不适合您的问题。

相反,您应该摆脱自己的队列并使用池中内置的队列来获取由worker 编辑的值return,或者完全废弃池并使用multiprocessing.Process 启动新进程为您必须完成的每项任务。

我还要注意,您的代码在修改 x 数组的主线程和在旧值发送到工作进程之前序列化旧值的线程之间的主进程中存在竞争条件。很多时候,您最终可能会发送同一数组的多个副本(带有最终值),而不是您想要的几个不同的值。

这是一个快速且未经测试的版本,可以丢弃队列:

def worker(arr):
    time.sleep(0.2)
    return arr[0]

if __name__ == "__main__":
    p = mp.Pool(4)
    results = p.map(worker, [np.array([i, 4]) for i in range(4)])
    p.join()
    for result in results:
        print(result)

这是一个删除 Pool 并保留队列的版本:

def worker(q, arr): 
    time.sleep(0.2)
    q.put(arr[0])

if __name__ == "__main__":
    q = m.Queue()
    processes = []

    for i in range(4):
        p = mp.Process(target=worker, args=(q, np.array([i, 4])))
        p.start()
        processes.append(p)

    for i in range(4):
        print(q.get())

    for p in processes:
        p.join()

请注意,在上一个版本中,在我们尝试join 进程之前,我们get 来自队列的结果可能很重要(尽管如果我们只处理四个值,则可能不是)。如果队列被填满,如果我们按照其他顺序执行,可能会发生死锁。工作进程可能会被阻止尝试写入队列,而主进程被阻止等待工作进程退出。

【讨论】:

  • 嗯。我现在正在使用池和队列而不是进程。如果有 10 个任务和 4 个进程(每个 cpu 核心 1 个),使用 pool 只会创建 4 个进程。使用 Process 需要创建 10 个进程。有没有重用流程?
  • 队列用于在工作进程的生产者和主进程的消费者之间传递结果。每 30 个周期,程序从硬盘读取并启动 120 个异步计算任务。之后,程序运行消费者,一旦其中一个工作人员将结果放入队列,消费者就会从队列中获取 120 个结果。消费者完成后,整个事情会重复。然后在最后,程序调用pool.close、pool.join,然后consumer处理所有剩余的结果。工人任务使用的时间比消费者 + 硬盘多得多(30 倍)。
  • 是否有任何简单的方法或“内置在池中的队列”允许消费者在工人仍在处理任务时从工人那里获取结果?到目前为止,我在启动消费者之前手动计算了 120 个任务。应该有一个更简单的方法......我可能应该问一个新问题。
  • 我不完全确定我明白你在问什么,但如果你的意思是主进程如何消耗池中的一些结果,你可以使用pool.imap(或@987654336 @) 以获取与进一步计算并行运行的结果的迭代器。我建议reading the docs。
猜你喜欢
  • 1970-01-01
  • 2013-02-24
  • 1970-01-01
  • 1970-01-01
  • 2014-10-06
  • 1970-01-01
  • 2020-04-11
  • 2021-03-01
  • 1970-01-01
相关资源
最近更新 更多