【发布时间】: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