【问题标题】:Python multiprocessing: how to limit the number of waiting processes?Python多处理:如何限制等待进程的数量?
【发布时间】:2012-06-18 20:21:45
【问题描述】:

当使用 Pool.apply_async 运行大量任务(带有大参数)时,进程被分配并进入等待状态,等待进程的数量没有限制。最终会吃掉所有内存,如下例所示:

import multiprocessing
import numpy as np

def f(a,b):
    return np.linalg.solve(a,b)

def test():

    p = multiprocessing.Pool()
    for _ in range(1000):
        p.apply_async(f, (np.random.rand(1000,1000),np.random.rand(1000)))
    p.close()
    p.join()

if __name__ == '__main__':
    test()

我正在寻找一种限制等待队列的方法,即只有有限数量的等待进程,并且 Pool.apply_async 在等待队列已满时被阻塞。

【问题讨论】:

    标签: python multiprocessing pool


    【解决方案1】:

    这是最佳答案的猴子修补替代方案:

    import queue
    from multiprocessing.pool import ThreadPool as Pool
    
    
    class PatchedQueue():
      """
      Wrap stdlib queue and return a Queue(maxsize=...)
      when queue.SimpleQueue is accessed
      """
    
      def __init__(self, simple_queue_max_size=5000):
        self.simple_max = simple_queue_max_size  
    
      def __getattr__(self, attr):
        if attr == "SimpleQueue":
          return lambda: queue.Queue(maxsize=self.simple_max)
        return getattr(queue, attr)
    
    
    class BoundedPool(Pool):
      # Override queue in this scope to use the patcher above
      queue = PatchedQueue()
    
    pool = BoundedPool()
    pool.apply_async(print, ("something",))
    

    这在 Python 3.8 中按预期工作,其中多处理池使用 queue.SimpleQueue 设置任务队列。听起来multiprocessing.Pool 的实现可能自 2.7 以来发生了变化

    【讨论】:

    • 我没有测试过ThreadPool的,但是如果我修改成from multiprocessing.pool import Pool就不行了(限制没变,好像SimpleQueue没变到Queue)。知道如何解决这个问题吗?
    【解决方案2】:

    等待pool._taskqueue 是否超过所需大小:

    import multiprocessing
    import time
    
    import numpy as np
    
    
    def f(a,b):
        return np.linalg.solve(a,b)
    
    def test(max_apply_size=100):
        p = multiprocessing.Pool()
        for _ in range(1000):
            p.apply_async(f, (np.random.rand(1000,1000),np.random.rand(1000)))
    
            while p._taskqueue.qsize() > max_apply_size:
                time.sleep(1)
    
        p.close()
        p.join()
    
    if __name__ == '__main__':
        test()
    

    【讨论】:

    • 只想补充一点,我发现这是解决多处理内存问题的最简单方法。我使用了 max_apply_size = 10 这对我的问题很有效,这是一个缓慢的文件转换。正如@ecatmur 所建议的那样使用信号量似乎是一个更强大的解决方案,但对于简单的脚本来说可能是矫枉过正。
    • TaylorMonacelli 您的编辑被拒绝说明了 SO 上的 mod 问题。您的编辑修复了一个错误。 @greg-449 是一个“驱动模式”,只批准了 15% 的编辑,并给出了一个荒谬的拒绝理由。
    【解决方案3】:

    在这种情况下,您可以使用 maxsize 参数添加显式队列并使用 queue.put() 而不是 pool.apply_async()。然后工作进程可以:

    for a, b in iter(queue.get, sentinel):
        # process it
    

    如果您想将内存中创建的输入参数/结果的数量限制为大约活动工作进程的数量,那么您可以使用pool.imap*() 方法:

    #!/usr/bin/env python
    import multiprocessing
    import numpy as np
    
    def f(a_b):
        return np.linalg.solve(*a_b)
    
    def main():
        args = ((np.random.rand(1000,1000), np.random.rand(1000))
                for _ in range(1000))
        p = multiprocessing.Pool()
        for result in p.imap_unordered(f, args, chunksize=1):
            pass
        p.close()
        p.join()
    
    if __name__ == '__main__':
        main()
    

    【讨论】:

    • 使用imap 没有区别。输入队列仍然是无限的,使用这个解决方案最终会吃掉所有的内存。
    • @Radim:答案中的imap 代码即使你给它一个无限生成器也可以工作。
    • 不幸的是,Python 2 中没有(没有查看 py3 中的代码)。有关一些解决方法,请参阅this SO answer
    【解决方案4】:

    multiprocessing.Pool 有一个_taskqueue 类型为multiprocessing.Queue 的成员,它接受一个可选的maxsize 参数;不幸的是,它在没有maxsize 参数集的情况下构造它。

    我建议使用multiprocessing.Pool.__init__ 的复制粘贴将multiprocessing.Pool 子类化,从而将maxsize 传递给_taskqueue 构造函数。

    猴子修补对象(池或队列)也可以,但你必须猴子修补 pool._taskqueue._maxsizepool._taskqueue._sem 所以它会很脆弱:

    pool._taskqueue._maxsize = maxsize
    pool._taskqueue._sem = BoundedSemaphore(maxsize)
    

    【讨论】:

    • 我使用的是 Python 2.7.3,_taskqueue 的类型是 Queue.Queue。这意味着它是一个简单的队列,而不是 multiprocessing.Queue。子类化 multiprocessing.Pool 和覆盖 init 工作正常,但猴子修补对象没有按预期工作。但是,这是我正在寻找的 hack,谢谢。
    猜你喜欢
    • 2014-06-07
    • 2013-12-01
    • 1970-01-01
    • 2011-04-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-09-19
    • 2021-12-05
    相关资源
    最近更新 更多