【问题标题】:Output Queue of a Python multiprocessing is providing more results than expectedPython 多处理的输出队列提供的结果比预期的要多
【发布时间】:2014-02-07 19:00:29
【问题描述】:

根据以下代码,我希望结果列表的长度与多进程提供的项目范围之一相同:

import multiprocessing as mp

def worker(working_queue, output_queue):
    while True:
        if working_queue.empty() is True:
            break #this is supposed to end the process.
        else:
            picked = working_queue.get()
            if picked % 2 == 0: 
                output_queue.put(picked)
            else:
                working_queue.put(picked+1)
    return

if __name__ == '__main__':
    static_input = xrange(100)    
    working_q = mp.Queue()
    output_q = mp.Queue()
    for i in static_input:
        working_q.put(i)
    processes = [mp.Process(target=worker,args=(working_q, output_q)) for i in range(mp.cpu_count())]
    for proc in processes:
        proc.start()
    for proc in processes:
        proc.join()
    results_bank = []
    while True:
        if output_q.empty() is True:
            break
        else:
            results_bank.append(output_q.get())
    print len(results_bank) # length of this list should be equal to static_input, which is the range used to populate the input queue. In other words, this tells whether all the items placed for processing were actually processed.
    results_bank.sort()
    print results_bank

有人知道如何让这段代码正常运行吗?

【问题讨论】:

  • 顺便说一句,如果您能帮助我了解是什么让多处理 python 代码对操作系统平台不敏感,我将不胜感激。如果在 Windows 7 或 MacOS 中运行,上述代码的行为会有所不同;在前者中,控制台无响应,而在后者中,结果中的项目重复。

标签: python queue multiprocessing


【解决方案1】:

这段代码永远不会停止:

每个工作人员从队列中获取一个项目,只要它不为空:

picked = working_queue.get()

并为每一个得到一个新的:

working_queue.put(picked+1)

因此,队列永远不会为空,除非进程之间的时间恰好在某个进程调用empty() 时队列为空。因为队列长度最初是 100 并且您拥有与 cpu_count() 一样多的进程,如果这在任何实际系统上停止,我会感到惊讶。

执行稍加修改的代码证明我错了,它确实会在某个时候停止,这让我很惊讶。用一个进程执行代码似乎有一个错误,因为一段时间后进程冻结但不返回。对于多个进程,结果会有所不同。

在循环迭代中添加一个较短的睡眠周期会使代码的行为符合我的预期和上面的解释。 Queue.putQueue.getQueue.empty 之间似乎存在一些时间问题,尽管它们应该是线程安全的。删除 empty 测试也会得到预期的结果(不会卡在空队列中)。

找到不同行为的原因。放入队列的对象不会立即刷新。因此empty 可能会返回False,尽管队列中有项目等待刷新。

来自documentation

注意:当一个对象被放入队列时,该对象被腌制并且一个 后台线程稍后将腌制数据刷新到底层 管道。这会产生一些有点令人惊讶的后果,但是 不应该造成任何实际困难——如果他们真的打扰的话 然后,您可以改为使用由管理器创建的队列。

  1. 将对象放入空队列后,队列的 empty() 方法返回 False 并且 get_nowait() 可以在不引发 Queue.Empty 的情况下返回。

  2. 如果多个进程正在对对象进行排队,则对象可能会在另一端无序接收。但是,由同一进程排队的对象将始终按预期顺序排列。

【讨论】:

  • 感谢您的回答。我在两个空测试中都包含了缺失的 else,但它没有用。我试图删除空测试,但该过程仍然卡住。我使用空测试来确保流程在某个时候结束。否则,进程在向空队列询问数据时会卡住。我认为 .emtpy() 方法的不可靠性只是最后几轮进程的问题,当队列中的项目很少时,CPU“错误地”将队列读取为空,而其他 CPU 只是在排队更多项目.在这种情况下,单个 CPU 应该能够完成工作。
  • 是否可以在全局变量上使用 .lock,将其长度与原始排队数据的数量进行比较,作为在进程陷入空队列之前结束进程的解决方案?
  • @Yag 一个简单的解决方法是向get() 添加一个小超时,这样当队列不再交付项目而不是检查empty 时会引发异常。当然,似乎也不能保证它一直有效。否则我猜你需要使用共享状态。虽然我不知道可能有更好的解决方案。
  • 感谢您的 cmets。超时加上异常确实很有用。由于并行设计,维护空测试很方便,因为超时异常可能是临时空队列的结果,该队列刚刚被重新填充。我试图解决这个with a manager, but a normal queue 终于为我工作了。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2016-11-13
  • 1970-01-01
  • 2015-01-12
  • 2012-02-23
  • 1970-01-01
  • 2015-01-10
  • 1970-01-01
相关资源
最近更新 更多