【问题标题】:Why queue is not multiprocesing safe?为什么队列不是多处理安全的?
【发布时间】:2015-02-23 23:50:38
【问题描述】:

我创建了一些多处理代码 - 检测问题非常简单,但我发现了一些问题 - 队列没有同步更新。

# coding=utf-8
import multiprocessing

def do_work(input_queue, output_queue):
  print multiprocessing.current_process().name
  input_queue.put(1)
  while not input_queue.empty():
    output_queue.put(input_queue.get() + 1)

def main():
  input_queue = multiprocessing.Queue()
  output_queue = multiprocessing.Queue()
  for i in range(8):
    input_queue.put(i)

  processes = []
  for i in range(2):
    process = multiprocessing.Process(name = str(i),
                                      target = do_work,
                                      args = (input_queue,
                                              output_queue), )
    processes.append(process)
    process.start()
  for process in processes:
    process.join()
  results = []
  while not output_queue.empty():
    results.append(output_queue.get())
  print len(results), results

if __name__ == '__main__':
  main()

有时结果是 - 看起来不错:

process 0
process 1
10 [2, 1, 3, 4, 6, 5, 8, 7, 2, 2]

但有时结果会有所不同,例如值 1 未在进程开始时放置:

process 0
process 1
9 [1, 2, 3, 4, 5, 6, 7, 8, 2]

看起来打印没有问题,因为它是在主线程中完成的,但队列不支持进程间锁定。你能提出一些建议吗?

【问题讨论】:

  • 这不会使数据结构线程不安全;这是异步行为,这是我完全期望发生的。
  • Makoto 10 值 = 10 个结果未知顺序但不是 9 个结果 - 再次查看我简化的演示文稿。
  • 你从哪里得到 10 个?
  • @Makoto print len(results), results 一次是 10,一次是 9 - 再次阅读问题。我不知道为什么?
  • 我认为你需要锁

标签: python python-multiprocessing


【解决方案1】:

您的代码略有改动:

# coding=utf-8
import multiprocessing

def do_work(input_queue, output_queue, lock):
  with lock:
    input_queue.put(1)
    print input_queue.empty(), input_queue.qsize()
    while not input_queue.empty():
      output_queue.put(input_queue.get() + 1)

def main():
  input_queue = multiprocessing.Queue()
  output_queue = multiprocessing.Queue()
  lock = multiprocessing.Lock()
  for i in range(8):
    input_queue.put(i)

  processes = []
  for i in range(2):
    process = multiprocessing.Process(name = str(i),
                                      target = do_work,
                                      args = (input_queue,
                                              output_queue, lock), )
    processes.append(process)
    process.start()
  for process in processes:
    process.join()
  results = []
  while not output_queue.empty():
    results.append(output_queue.get())
  print len(results), results

if __name__ == '__main__':
  main()

请注意,现在整个进程处于锁定状态,因此不可能出现竞争条件,并且它还会打印输入队列的大小以及它是否为空。现在这是其中一次运行的输出:

False 9
True 1
9 [1, 2, 3, 4, 5, 6, 7, 8, 2]

注意第二个进程如何说队列是空的,但同时只有一个元素。原因在文档中:

empty() 如果队列为空,则返回 True,否则返回 False。因为 多线程/多处理语义,这是不可靠

要修复它,您可以将条件 while not input_queue.empty() 替换为 while input_queue.qsize() > 0。当你这样做时,你会看到你的代码挂起。这是有道理的,因为您首先检查队列的大小,然后尝试将其弹出。考虑以下场景:队列中有一个元素,两个线程都看到,并尝试弹出。一个成功,另一个现在尝试从一个空队列中弹出,然后阻塞。要解决这个问题,请尝试执行非阻塞弹出,如果失败则重试:

# coding=utf-8
import multiprocessing
import Queue

def do_work(input_queue, output_queue):
  input_queue.put(1)
  while input_queue.qsize() > 0:
    try:
      output_queue.put(input_queue.get(False) + 1)
    except Queue.Empty:
      pass

def main():
  input_queue = multiprocessing.Queue()
  output_queue = multiprocessing.Queue()
  for i in range(8):
    input_queue.put(i)

  processes = []
  for i in range(2):
    process = multiprocessing.Process(name = str(i),
                                      target = do_work,
                                      args = (input_queue,
                                              output_queue) )
    processes.append(process)
    process.start()
  for process in processes:
    process.join()
  results = []
  while True:
    try:
      results.append(output_queue.get(False))
    except Queue.Empty:
      break
  print len(results), results

if __name__ == '__main__':
  main()

【讨论】:

  • 锁定是对的,但这并不能解释为什么重复值会出现在队列中。
  • @Makoto,因为他插入 1 三次,然后重新插入从输入队列到输出队列的所有内容。重复很有意义
  • @makato 这是错误的代码它将给出 10 或 9 个结果 - 已测试!
  • 似乎这个答案一针见血。不知道为什么我没有注意到do_work 方法中额外的put 调用...
  • @Ishamael Lock 将使代码成为单线程,因此它不是解决方案(整个工作人员都自行锁定)。这段代码可以给出 9 或 10 个结果,所以同样的错误。文档说空不可靠,但什么是可靠的——在其他地方它并没有说不可靠——很难学习。
猜你喜欢
  • 1970-01-01
  • 2019-06-04
  • 2015-01-11
  • 1970-01-01
  • 2015-01-04
  • 2012-12-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多