【问题标题】:How to successfully utilize Queue.join() with multiprocessing?如何通过多处理成功利用 Queue.join()?
【发布时间】: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 Trueq.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


    【解决方案1】:

    队列和池

    一次运行一个...单独的问题。

    实际上,这是同一个问题的一部分。这一切都意味着你是不是 使用由Pool 管理的多处理池。你现在做的是 把你所有的任务放在一个队列里,把它们直接放回去,然后 使用池一次处理一个,该池一次只能获得一个任务 时间。这两种范式是互斥的:如果你想使用一个池来 为你做事,你不需要排队;如果您需要处理队列 你自己,你可能不想使用pool

    multiprocessing.Pool 和伴随的方法产生正确数量的工人 进程,将你的函数序列化给它们,然后在内部设置一个队列 并处理发送任务和获取结果。这比做起来容易得多 手动,通常是正确的处理方式:

    当你使用 pool 时,你会做这样的事情:

    results = pool.map(compute_primes, [(0,100_000) for _ in range(8)])
    

    在所有池完成之前会为您阻塞,或者:

    results = pool.map_async(compute_primes, [(0, 100_000) for _ in range(8)])
    results.wait() # wait
    

    除非您计划在结果输入时对其进行处理,在这种情况下您不需要 完全使用results.wait()

    for _ in range(8):
        result = results.get()
        do_stuff(result)
    

    您确实使用 pool.join()pool.close() 只是为了确保 pool 已关闭 优雅地下降,但这与获得结果无关。

    你的例子

    您的第一个示例有效,因为您这样做:

    • 将任务放入队列中
    • 一一取出处理
    • 加入一个空队列 -> 立即离开

    你的第二个例子失败了,因为你这样做了:

    • 将任务放入队列中
    • 完成一项任务
    • 等待队列为空或完成 -> 无限期阻塞

    在这种情况下,您根本不需要排队。

    手动使用队列

    除此之外:你从哪里得到你的Queuemultiprocessing.Queue 不是 可加入;你需要multiprocessing.JoinableQueuethreading.Queue应该 不能与multiprocessing 一起使用。 queue.Queue,同样,不应使用 使用 `multiprocessing.

    什么时候使用任务队列?当您只想应用一堆 一堆函数的参数。也许您想使用自定义类。 也许你想做一些有趣的事情。也许你想做一些事情 一种类型的论点,但如果论点属于某种类型,则另一种说法, 并且代码以这种方式组织得更好。在这些情况下,子类化Process(或 Thread 用于多线程)你自己可能会更清晰。都不是 似乎适用于这种情况。

    join 与队列一起使用

    .join() 用于 task 队列。它阻塞,直到队列中的每个任务都有 已标记为完成。当您想要卸载某些处理时,这很方便 到一堆进程,但在你做任何事情之前等待它们。那么你 通常做这样的事情:

    tasks = JoinableQueue()
    for t in qs:
        tasks.put(t)
    start_multiprocessing() # dummy fn
    tasks.join() # wait for everything to be done
    

    但是在这种情况下,您不想这样做,或者不想这样做。

    【讨论】:

    • 来自 queue.Queue,将更新答案。
    • @NikolaKapralov 我已经编辑了答案,以明确目前出了什么问题。
    • 感谢您的回答。所以,简而言之,Pool 一开始就是错误的做法。这也很好地回答了我的第二个问题,即为什么将 queue.Queue 替换为 python 列表时,它似乎可以完美地工作。如果我理解正确, Pool 用于自动化任务,我的工作人员将自动从列表中获取项目,我想 Processes 将与 queue.Queue 一起使用,其中每个进程都将手动处理一个 Queue 项目。我说的对吗?
    • 我已经重构了答案,以明确说明您根本不想使用队列
    【解决方案2】:

    我不希望为 Pool 构造函数指定参数,除非出于某种原因,我需要很少的并发进程。通过在没有参数的情况下构造 Pool,潜在并发进程的数量将因计算机的 CPU 架构而异。以下是我将如何执行您的任务(假设我完全了解您的用例):

    from multiprocessing import Pool
    
    
    def genPrime():  # prime number generator
        D = {}
        q = 2
        while True:
            if q not in D:
                yield q
                D[q * q] = [q]
            else:
                for p in D[q]:
                    D.setdefault(p + q, []).append(p)
                del D[q]
            q += 1
    
    
    def compute_primes(n):
        g = genPrime()
        return [next(g) for _ in range(n)]
    
    
    NCOMPUTATIONS = 8
    NPRIMES = 30_000
    
    
    def main():
        with Pool() as pool:
            ar = []
            for _ in range(NCOMPUTATIONS):
                ar.append(pool.apply_async(compute_primes, [NPRIMES]))
            for _ar in ar:
                result = _ar.get() # waits for process to terminate and get its return value
                assert len(result) == NPRIMES
    
    
    if __name__ == '__main__':
        main()
    

    [请注意我不是genPrime函数的作者]

    【讨论】:

    • 感谢您的解释!实际上,我有一个版本的代码,它的编写方式与您的相同,并且可以正常工作。但是给我留下的印象是 Pool 工作人员从 ar 中选择任务,我认为这是一个 Queue,现在我明白事实并非如此。因为我认为这是一个队列,所以我决定尝试以“正确的方式”来做,经过一番搜索,得出的结论是 queue.Queue 与 get() 和 join() 是要走的路。感谢您提及 Pool 的 CPU 计数。
    • 队列对于多线程非常有用,因此对于多处理也是如此。查看 multiprocessing.Value() 以了解在子进程之间共享数据的机制。它“在幕后”使用共享内存。
    猜你喜欢
    • 2016-06-30
    • 2011-10-20
    • 2021-09-25
    • 2021-10-22
    • 2020-05-20
    • 2011-06-17
    • 2020-02-05
    • 2017-07-28
    • 2020-06-11
    相关资源
    最近更新 更多