【问题标题】:Mutliprocessing in Python in a for loop and passing multiple ArgumentsPython 中的 for 循环中的多处理并传递多个参数
【发布时间】:2018-09-21 14:10:42
【问题描述】:

我正在使用 python 脚本进行大量计算。由于它受 CPU 限制,我通常使用线程模块的方法没有产生任何性能改进。

我现在尝试使用多处理而不是多线程来更好地使用我的 CPU 并加快冗长的计算。

我在 stackoverflow 上找到了一些示例代码,但我没有让脚本接受多个参数。有人可以帮我解决这个问题吗?在我很确定我使用 Pool.map 错误之前,我从未使用过这些模块。 - 任何帮助表示赞赏。也欢迎使用其他方式来完成多处理。

from multiprocessing import Pool

def calculation(foo, bar, foobar, baz):
    # Do a lot of calculations based on the variables
    # Later the result is written to a file.
    result = foo * bar * foobar * baz
    print(result)

if __name__ == '__main__':
    for foo in range(3):
        for bar in range(5):
            for baz in range(4):
                for foobar in range(10):

                    Pool.map(calculation, foo, bar, foobar, baz)
                    Pool.close()
                    Pool.join()

【问题讨论】:

    标签: python multithreading python-3.x python-multiprocessing


    【解决方案1】:

    正如您所怀疑的那样,您使用 map 的方式不止一种。

    • map 的要点是在可迭代的所有元素上调用函数。就像内置的 map 函数一样,但是是并行的。如果您想排队单个呼叫,只需使用apply_async

    • 对于您特别询问的问题:map 采用单参数函数。如果您想传递多个参数,您可以修改或包装您的函数以采用单个元组而不是多个参数(我将在最后展示),或者只使用starmap。或者,如果你想使用apply_async,它需要一个包含多个参数的函数,但你将apply_async 传递给一个参数元组,而不是单独的参数。

    • 您需要在Pool 实例上调用map,而不是Pool 类。您要做的类似于尝试从文件类型中read,而不是从特定打开的文件中读取。
    • 您尝试在每次迭代后关闭并加入Pool。在完成所有这些之前,您不想这样做,否则您的代码将等待第一个完成,然后为第二个引发异常。

    因此,可行的最小更改是:

    if __name__ == '__main__':
        pool = Pool()
        for foo in range(3):
            for bar in range(5):
                for baz in range(4):
                    for foobar in range(10):
                        pool.apply_async(calculation, (foo, bar, foobar, baz))
        pool.close()
        pool.join()
    

    请注意,我将所有内容都保存在 if __name__ == '__main__': 块中,包括新的 Pool() 构造函数。我不会在后面的示例中展示这一点,但对于所有示例来说都是必需的,原因在文档的 Programming guidelines 部分中进行了说明。1


    如果您想使用 map 函数之一,则需要一个可迭代的完整参数,如下所示:

    pool = Pool()
    args = ((foo, bar, foobar, baz) 
            for foo in range(3) 
            for bar in range(5) 
            for baz in range(4) 
            for foobar in range(10))
    pool.starmap(calculation, args)
    pool.close()
    pool.join()
    

    或者,更简单地说:

    pool = Pool()
    pool.starmap(calculate, itertools.product(range(3), range(5), range(4), range(10)))
    pool.close()
    pool.join()
    

    假设您没有使用旧版本的 Python,您可以通过在 with 语句中使用 Pool 来进一步简化它:

    with Pool() as pool:
        pool.starmap(calculate, 
                     itertools.product(range(3), range(5), range(4), range(10)))
    

    使用mapstarmap 的一个问题是,它会做额外的工作来确保按顺序返回结果。但是您只是返回None 并忽略它,那为什么会这样呢?

    使用apply_async 没有这个问题。

    您也可以将map 替换为imap_unordered,但没有istarmap_unordered,因此您需要将函数包装为不需要starmap

    def starcalculate(args):
        return calculate(*args)
    
    with Pool() as pool:
        pool.imap_unordered(starcalculate,
                            itertools.product(range(3), range(5), range(4), range(10)))
    

    1。如果您使用spawnforkserver 启动方法——并且spawn 是Windows 上的默认值——每个子进程都相当于import 对你的模块执行操作。因此,所有不受__main__ 保护的顶级代码都将在每个孩子中运行。该模块试图保护您免受由此带来的一些最糟糕的后果(例如,不是用指数爆炸的孩子创建新孩子来对您的计算机进行分叉轰炸,而是经常遇到异常),但它不能使代码真正工作.

    【讨论】:

    • 非常感谢您快速而详细的回答。我终于接受了你最后一个使用 pool.imap_unordered 的建议,但偶然发现了一个小问题。似乎整个脚本,甚至在计算函数之外(定义变量等发生在那里)都被执行。我必须在哪里放置这样的代码('startup'),这样它只会在开始时执行一次。
    • 在几乎所有多处理脚本中你仍然需要__name__ == '__main__' 保护,即使我没有在我的单行示例中展示它。那是问题吗? (如果是这样,您认为我需要更清楚地说明并在答案中解释吗?)
    • 也许一个简短的解释对像我这样的其他人来说会很好,幸运的是,如果你错过了抛出的错误,那么它很容易理解。我现在遇到的问题是,我在开始时读取了一个包含大量“只读”数据的文件,这是每个进程都需要的。如何将这些数据传递给每个进程?全局变量显然不起作用。后来我想将该数据输出到文件中,但也不知道该怎么做。 (是否要保存以在函数末尾附加到文件?)我应该快速更新问题还是会弄乱你的答案?
    • @Marco 这是一个新问题。事实上,他们两个。很有可能在 StackOverflow 上已经为每个问题提供了一个很好的答案,因此请先搜索(但要注意仅适用于 Linux 的答案)。但如果没有,请创建新问题,每个问题都有一个适当的minimal reproducible example。首先阅读multiprocessing 文档的编程指南部分——我认为它不会回答你所有的问题,但它至少会给你一些背景知识。此外,在这里探索concurrent.futures.ProcessPoolExecutor(请参阅文档中的示例)是否会让您的生活更轻松。
    • 漂亮的答案!
    猜你喜欢
    • 1970-01-01
    • 2020-04-20
    • 2020-09-25
    • 2019-06-27
    • 2018-06-15
    • 2016-03-31
    • 2021-09-17
    • 1970-01-01
    • 2018-08-06
    相关资源
    最近更新 更多