【发布时间】:2021-11-28 10:57:58
【问题描述】:
我正在尝试学习 Python 中的多处理库,但我无法让我的代码与 queue.Queue 一起使用。简而言之,我不知道在我的代码中将queue.Queue.join() 方法放在哪里。它是在while循环中还是在它之外?如果它超出了while循环,我应该写while q.not_empty吗?当文档明确提到使用join() 时,我为什么要使用q.not_empty?
这是我的代码。我期待我的 4 个核心同时返回由我的函数计算的素数数量,每个核心 2 次,总共 8 次计算。主要的计算功能没有问题。
import queue
def main():
q = queue.Queue()
[q.put((compute_primes, (1, 30000))) for _ in range(8)]
with multiprocessing.Pool(processes=4) as pool:
while q.not_empty:
result = q.get()
function = pool.apply_async(result[0], args=(result[1][0], result[1][1]))
function.get()
q.join()
使用上面的代码,如果队列为空,我会跳出循环。但这应该是不现实的,为什么我之后还需要q.join()?
使用下面的代码,我无法跳出循环。更改为while True 和q.join() 的位置
def main():
q = queue.Queue()
[q.put((compute_primes, (1, 30000))) for _ in range(8)]
with multiprocessing.Pool(processes=4) as pool:
while True:
result = q.get()
function = pool.apply_async(result[0], args=(result[1][0], result[1][1]))
function.get()
q.join()
我应该把q.join放在哪里?
附:这段代码也不能有效地并行化任务,它本质上是一个一个地计算函数,我不明白为什么,但这是一个不同的问题。
附: 2
主函数代码
def compute_primes(start, end):
start_time = time.time()
primes = []
for number in range(start, end + 1):
flag = True
for i in range(2, number):
if (number % i) == 0:
flag = False
break
if flag:
primes.append(number)
end_time = time.time()
print(f"Time taken: {end_time - start_time}\n"
f"Amount primes: {len(primes)}")
return primes
【问题讨论】:
标签: python parallel-processing multiprocessing queue