【问题标题】:Multiprocessing and Queues多处理和队列
【发布时间】:2017-08-08 14:42:48
【问题描述】:

`此代码尝试使用队列将任务提供给多个工作进程。

我想计算不同进程数量和不同数据处理方法之间的速度差异。

但是输出并没有像我想象的那样。

from multiprocessing import Process, Queue
import time
result = []

base = 2

data = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 23, 45, 76, 4567, 65423, 45, 4, 3, 21]

# create queue for new tasks
new_tasks = Queue(maxsize=0)

# put tasks in queue
print('Putting tasks in Queue')
for i in data:
    new_tasks.put(i)

# worker function definition
def f(q, p_num):
    print('Starting process: {}'.format(p_num))
    while not q.empty():
        # mimic some process being done
        time.sleep(0.05)
        print(q.get(), p_num)
    print('Finished', p_num)

print('initiating processes')
processes = []
for i in range(0, 2):
    if __name__ == '__main__':
        print('Creating process {}'.format(i))
        p = Process(target=f, args=(new_tasks, i))
        processes.append(p)
#record start time
start = time.time()

# start process
for p in processes:
    p.start()

# wait for processes to finish processes
for p in processes:
    p.join()

#record end time
end = time.time()

# print time result
print('Time taken: {}'.format(end-start))

我预计会这样:

Putting tasks in Queue
initiating processes
Creating process 0
Creating process 1
Starting process: 1
Starting process: 0
1 1
2 0
3 1
4 0
5 1
6 0
7 1
8 0
9 1
10 0
11 1
23 0
45 1
76 0
4567 1
65423 0
45 1
4 0
3 1
21 0
Finished 1
Finished 0
Time taken: <some-time>

但是我实际上得到了这个:

Putting tasks in Queue
initiating processes
Creating process 0
Creating process 1
Time taken: 0.01000523567199707
Putting tasks in Queue
Putting tasks in Queue
initiating processes
Time taken: 0.0
Starting process: 1
initiating processes
Time taken: 0.0
Starting process: 0
1 1
2 0
3 1
4 0
5 1
6 0
7 1
8 0
9 1
10 0
11 1
23 0
45 1
76 0
4567 1
65423 0
45 1
4 0
3 1
21 0
Finished 0

似乎有两个主要问题,我不确定它们的相关性如何:

  1. 打印语句,例如: Putting tasks in Queue initiating processes Time taken: 0.0 通过代码系统地重复 - 我说系统地重复,因为它们每次都准确地重复。

  2. 第二个进程永远不会结束,它永远不会识别队列是空的,因此无法退出

【问题讨论】:

  • 我听起来你的代码格式有问题:你应该只有一个 Time taken:... 打印输出。
  • 另外你不应该轮询q.empty(),因为一个贪婪的线程可能会窃取最后一个项目,而让所有其他线程等待永远不会出现的项目。您应该使用的是队列结束标记。每个线程一个。
  • 否则这是个好问题。您在编写代码和收集输出方面付出了一些努力并显示了您期望发生的事情。
  • @quamrana 你在哪里对,轮询 q.empty() 是 2 的解决方案
  • 你可以是 PEP8 兼容的,但是一个额外的或缺少的缩进可以完全改变一个 python 程序。

标签: python queue multiprocessing


【解决方案1】:

1) 我无法重现此内容。

2) 看下面的代码:

while not q.empty():
    time.sleep(0.05)
    print(q.get(), p_num)

每一行都可以由任何进程以任何顺序运行。现在考虑q 有一个项目和两个进程A 和B。现在考虑以下执行顺序:

# A runs
while not q.empty():
    time.sleep(0.05)

# B runs
while not q.empty():
    time.sleep(0.05)

# A runs
print(q.get(), p_num)  # Removes and prints the last element of q

# B runs
print(q.get(), p_num)  # q is now empty so q.get() blocks forever

交换time.sleep 和q.get 的顺序可以消除我所有运行中的阻塞,但仍有可能有多个进程进入循环而只剩下一个项目。

解决此问题的方法是使用非阻塞 get 调用并捕获 queue.Empty 异常:

import queue

while True:
    time.sleep(0.05)
    try:
        print(q.get(False), p_num)
    except queue.Empty:
        break

【讨论】:

  • 很好的答案,很好的解释,对我的问题的另一半有什么想法吗?
  • 不,查看您的代码,我看不到任何方式可以多次打印该行。也许尝试处理问题中的代码,看看您正在运行的代码是否有任何区别
【解决方案2】:

你的工作线程应该是这样的:

def f(q, p_num):
    print('Starting process: {}'.format(p_num))
    while True:
        value = q.get()
        if value is None:
            break
        # mimic some process being done
        time.sleep(0.05)
        print(value, p_num)
    print('Finished', p_num)

并且队列应该在真实数据之后填充标记:

for i in data:
    new_tasks.put(i)
for _ in range(num_of_threads):
    new_tasks.put(None)

【讨论】:

  • 您已选择使用if 然后break 而不是try - except。这只是速度问题吗?
  • 另外,我的队列应该填充哪些“标记”?我在队列或多处理或多线程的文档中找不到它们的提及? (一个链接会很可爱,我不希望你在 cmets 上写一篇文章)
  • @quarana 项目肯定需要在队列中才能正确通过进程而不共享或“锁定”?
  • 我使用None 作为标记。您需要传递尽可能多的这些线程,因为您有线程。即使一个线程异常贪婪并消耗了所有排队的项目,它也会找到第一个None并退出。然后所有其余的惰性线程,当他们调用get() 时,每个都会找到一个None 并在他们到达那里时退出。 None 标记可以在真实数据进入后的任何时候被推送到队列中,但在您希望线程退出之前。
  • 这可能是@ikkuh 在查找队列末尾方面有更好的答案。
猜你喜欢
  • 1970-01-01
  • 2016-04-18
  • 2020-05-18
  • 2016-05-02
  • 1970-01-01
  • 2015-10-11
  • 2010-10-29
  • 1970-01-01
相关资源
最近更新 更多