【问题标题】:Python: How to parallelize functions used in nested for loops with many function inputs?Python:如何并行化用于具有许多函数输入的嵌套 for 循环中的函数?
【发布时间】:2019-02-02 11:10:42
【问题描述】:

我希望加快一个缓慢运行的循环,但不要认为我会用最好的方法来解决这个问题。我想并行化一些运行我编写的函数的代码,并且在尝试弄清楚如何在使用 python 的 multiprocessing 模块时准确地制定输入参数时遇到了一些麻烦。我拥有的代码基本上是以下形式:

a = some_value
b = some_value
c = some_value
for i in range(1,101):
    for j in range(1,101):
        b = np.array([i*0.001,j*0.001]).reshape((2,1))
        (A,B,C,D) = function(a,b,c,d)

所以我的函数本身有多种参数,但对于这个特殊用途,我只需要改变一个变量(这是一个包含两个值的数组)并创建一个值网格。此外,所有其他输入都是整数。我熟悉通过以下示例代码使用工人池并行化此类循环的非常简单的示例:

pool = mp.Pool(processes=4)
input_parameters = *list of iterables for multiprocessing*
result = pool.map(paramest.parameter_estimate_ND, input_parameters)

迭代列表是使用itertools 模块创建的。由于我只更改了函数的一个输入变量,并且在构造此类输入参数时遇到麻烦之前声明了所有其他变量。所以我真正想要的是使用multiprocessing同时运行不同的输入来加快for循环的执行。

然后我的问题是,如何构造使用multiprocessing 来并行化在函数上运行的代码,同时只更改特定变量的输入?

我是否以最好的方式处理这个问题?有没有更好的方法来做这样的事情?

谢谢!

【问题讨论】:

  • 你是在做数学运算还是在使用其他python工具?因为如果你只是在做数学,有更好的工具来加速你的 for 循环
  • 我不确定您所说的其他 python '工具' 是什么意思,但我正在并行化的函数使用一些输入数据并执行一些数学运算(使用 numpy)。本质上是使用数据来测试一些条件并创建许多不同的输出数组。
  • 我的意思是我认为你不需要使用进程,因为已经有像 Numba 这样的库可以优化你的 for 循环并并行化它,只要你只是在内部for循环中的函数

标签: python multiprocessing nested-loops python-multiprocessing python-multithreading


【解决方案1】:

通常,您只需要担心嵌套循环的内部循环的并行化。假设对 function 的每次调用都足够繁重,值得作为一项任务运行,那么一次将 100 个调用放入池中应该绰绰有余。


那么,如何并行化该内部循环?

只要把它变成一个函数:

def wrapper(a, c, d, i, j):
    b = np.array([i*0.001,j*0.001]).reshape((2,1))
    return function(a,b,c,d)

现在:

for i in range(1,101):
    pfunc = partial(function, a, c, d, i)
    ABCDs = pool.map(pfunc, range(1, 101))

或者,您甚至可以在 i 循环内定义包装函数,而不是创建部分函数:

for i in range(1,101):
    def wrapper(j):
        b = np.array([i*0.001,j*0.001]).reshape((2,1))
        return function(a,b,c,d)
    ABCDs = pool.map(wrapper, range(1, 101))

如果您在通过池的队列传递闭包变量时遇到问题,那很容易;您实际上不需要捕获变量,只需要捕获值,所以:

for i in range(1,101):
    def wrapper(j, *, a=a, c=c, d=d, i=i):
        b = np.array([i*0.001,j*0.001]).reshape((2,1))
        return function(a,b,c,d)
    ABCDs = pool.map(wrapper, range(1, 101))

如果事实证明仅j 不够并行,您可以轻松地将其更改为映射到(i, j)

def wrapper(i, j, *, a=a, b=b, c=c, d=d):
    b = np.array([i*0.001,j*0.001]).reshape((2,1))
    return function(a,b,c,d)
for i in range(1,101):
    ABCDs = pool.map(wrapper, itertools.product(range(1, 101), range(1, 101)))

ABCDs 将是 A, B, C, D 值的可迭代,所以很可能,无论您想对 A, B, C, D 做什么,都只是:

    for A, B, C, D in ABCDs:
        # whatever

【讨论】:

  • 谢谢!这非常有效!一开始我遇到了麻烦,因为我在定义包装函数之前输入了pool = mp.Pool(processes=4)。只需在放置该行代码之前定义包装器即可修复所有问题并完全按照我的意愿工作。
猜你喜欢
  • 2018-09-30
  • 1970-01-01
  • 2016-05-07
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-04-18
  • 2021-07-01
  • 1970-01-01
相关资源
最近更新 更多