【问题标题】:How to use multiprocessing.Queue.get method?如何使用 multiprocessing.Queue.get 方法?
【发布时间】:2019-04-07 12:29:08
【问题描述】:

下面的代码将三个数字放在一个队列中。然后它尝试从队列中取回号码。但它永远不会。如何从队列中获取数据?

import multiprocessing

queue = multiprocessing.Queue()

for i in range(3):
    queue.put(i)

while not queue.empty():
    print queue.get()

【问题讨论】:

    标签: python queue multiprocessing python-multiprocessing


    【解决方案1】:

    在使用get之前检查queue

    import multiprocessing
    
    queue = multiprocessing.Queue()
    
    for i in range(3):
        queue.put(i)
    
    while not queue.empty():
        if not queue.empty():
            print queue.get()
    

    【讨论】:

    • 在问题情况下“if”语句不会被执行,对吧?
    【解决方案2】:

    在阅读@Martijn Pieters 之后,我最初删除了这个答案,因为他更详细和更早地描述了“为什么这不起作用”。然后 我意识到,OP 示例中的用例不太适合

    的规范冠冕堂皇的标题

    “如何使用 multiprocessing.Queue.get 方法”。

    那不是因为有 演示不涉及子进程,但是因为在实际应用程序中,几乎没有一个队列是预先填充的,并且只在之后读取,但是读取 并且写入发生在中间的等待时间之间。 Martijn 展示的扩展演示代码在通常情况下不起作用,因为当排队跟不上读取速度时,while 循环会过早中断。所以这里是重新加载的答案,它能够处理通常的交错提要和读取场景:


    不要依赖 queue.empty 检查同步。

    在将对象放入空队列后,队列的 empty() 方法返回 False 和 get_nowait() 可以在不引发 queue.Empty 的情况下返回之前可能会有一个无限小的延迟。 ...

    empty()

    如果队列为空,则返回 True,否则返回 False。由于多线程/多处理语义,这是不可靠的。 docs

    使用队列中的for msg in iter(queue.get, sentinel):.get(),通过传递一个标记值来跳出循环...iter(callable, sentinel)?

    from multiprocessing import Queue
    
    SENTINEL = None
    
    if __name__ == '__main__':
    
        queue = Queue()
    
        for i in [*range(3), SENTINEL]:
            queue.put(i)
    
        for msg in iter(queue.get, SENTINEL):
            print(msg)
    

    ...如果您需要非阻塞解决方案,请使用get_nowait() 并处理可能的queue.Empty 异常。

    from multiprocessing import Queue
    from queue import Empty
    import time
    
    SENTINEL = None
    
    if __name__ == '__main__':
    
        queue = Queue()
    
        for i in [*range(3), SENTINEL]:
            queue.put(i)
    
        while True:
            try:
                msg = queue.get_nowait()
                if msg == SENTINEL:
                    break
                print(msg)
            except Empty:
                # do other stuff
                time.sleep(0.1)
    

    如果只有一个进程且该进程中只有一个线程正在读取队列,也可以将最后一个代码 sn-p 交换为:

    while True:
        if not queue.empty():  # this is not an atomic operation ...
            msg = queue.get()  # ... thread could be interrupted in between
            if msg == SENTINEL:
                break
            print(msg)
        else:
            # do other stuff
            time.sleep(0.1)
    

    由于线程可以在检查if not queue.empty()queue.get() 之间删除GIL,因此这不适用于进程中的多线程队列读取。如果多个进程正在从队列中读取,这同样适用。

    不过,对于单一生产者/单一消费者场景,使用 multiprocessing.Pipe 而不是 multiprocessing.Queue 就足够了,而且性能更高。

    【讨论】:

      【解决方案3】:

      您的代码确实有效,有时

      那是因为队列不是瞬间不是空的。该实现涉及到支持多个进程之间的通信,因此涉及到线程和管道,这导致empty 状态的持续时间比您的代码允许的时间长。

      Pipes and Queues section中的说明:

      当一个对象被放入队列时,该对象被腌制,然后后台线程将腌制数据刷新到底层管道。这会产生一些令人惊讶的后果,但不会造成任何实际困难 - 如果它们真的困扰您,那么您可以改用由经理创建的队列。

      1. 将对象放入空队列后在队列的empty() 方法返回False 之前可能会有一个无限小的延迟 [...]

      (我的粗体强调)

      如果您添加一个循环来检查是否为空首先然后您的代码可以工作:

      queue = multiprocessing.Queue()
      
      for i in range(3):
          queue.put(i)
      
      while queue.empty():
          print 'queue is still empty'
      
      while not queue.empty():
          print queue.get()
      

      当您运行上述代码时,大多数情况下'queue is still empty' 会出现一次。有时根本不出现,有时会打印两次。

      【讨论】:

        猜你喜欢
        • 2018-04-15
        • 2023-03-04
        • 2012-12-17
        • 2013-09-19
        • 2012-03-10
        • 2012-08-31
        • 2019-08-17
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多