为什么多处理比单核慢?使用 joblib 或 dask 会有所不同吗?

Pyt*_*ous 5 python numba joblib python-multiprocessing dask

问题

我正在尝试优化一些适合所谓的令人尴尬的并行计算的计算,但我发现使用 python 的multiprocessing包实际上会减慢速度。

我的问题是:我做错了什么,或者并行化实际上减慢了速度是否有内在原因?是因为我在使用 numba 吗?像 joblib 或 dak 这样的其他软件包会有很大的不同吗?

有很多类似的问题,其中的答案总是开销比节省的时间要多,但所有这些问题往往都围绕着非常简单的函数展开,而我原以为带有嵌套循环的东西更适合并行化. 我也没有找到 joblib、multiprocessing 和 dask 之间的比较。

我的功能

我有一个函数,它将一维 numpy 数组作为形状 n 的参数,并输出一个形状为 (nxt) 的 numpy 数组,其中每一行都是独立的,即输出的第 0 行仅取决于输入的第 0 项,项目 1 上的第 1 行,等等。像这样:

在此处输入图片说明

底层计算使用numba进行了优化,可将速度提高多个数量级。

玩具示例 - 结果

我无法分享确切的代码,所以我想出了一个玩具示例。中定义的计算my_fun_numba实际上无关紧要,它只是一些非常平庸的数字运算以保持 CPU 忙碌。

对于玩具示例,我的 PC 上的结果是这些,它们与我使用实际代码得到的结果非常相似。

如您所见,将输入数组拆分为不同的块并将它们中的每一个发送到 multiprocessing.pool 实际上比仅在单个核心上使用 numba 会减慢速度。 在此处输入图片说明

我试过的

我在装饰器中尝试了cache和nogil选项的各种组合numba.jit,但差异很小。

我已经用 PyCharm 分析了代码(不是 timeit.Timer 部分,只是一次运行),如果我正确理解输出,似乎大部分时间都花在等待 pool 上。

按时间排序:

分析 - 按时间排序

按自己时间排序:

分析 - 按自己的时间排序

玩具示例 - 代码

import numpy as np
import pandas as pd
import multiprocessing
from multiprocessing import Pool
import numba
import timeit


@numba.jit(nopython = True, nogil = True, cache = True)
def my_fun_numba(x):
    dim2 = 10
    out = np.empty((len(x), dim2))
    n = len(x)
    for r in range(n):   
        for c in range(dim2):
            out[r,c] = np.cos(x[r]) ** 2 + np.sin(x[r]) ** 2
    return out

def my_fun_non_numba(x):
    dim2 = 10
    out = np.empty((len(x), dim2))
    n = len(x)
    for r in range(n):   
        for c in range(dim2):
            out[r,c] = np.cos(x[r]) ** 2 + np.sin(x[r]) ** 2
    return out


def my_func_parallel(inp, func, cpus = None):
    if cpus == None:
        cpus = max(1, multiprocessing.cpu_count() - 1)
    else:
        cpus = cpus
        
    inp_split = np.array_split(inp,cpus)
    pool = Pool(cpus)    
    out = np.vstack(pool.map(func, inp_split) )    
    pool.close()
    pool.join()
    return out

if __name__ == "__main__":
    inputs = np.array([100,10e3,1e6] ).astype(int)
    res = pd.DataFrame(index = inputs, columns =['no paral, no numba','no paral, numba','numba 6 cores','numba 12 cores'])
    
    r = 3
    n = 1

    
    for i in inputs:
        my_arg = np.arange(0,i)

        
        res.loc[i, 'no paral, no numba'] = min(
            timeit.Timer("my_fun_non_numba(my_arg)", globals=globals()).repeat(repeat=r, number=n)
            )
        
        res.loc[i, 'no paral, numba'] = min(
            timeit.Timer("my_fun_numba(my_arg)", globals=globals()).repeat(repeat=r, number=n)
            )
        
        res.loc[i, 'numba 6 cores'] = min(
            timeit.Timer("my_func_parallel(my_arg, my_fun_numba, cpus = 6)", globals=globals()).repeat(repeat=r, number=n)
            )
        
        res.loc[i, 'numba 12 cores'] = min(
            timeit.Timer("my_func_parallel(my_arg, my_fun_numba, cpus = 12)", globals=globals()).repeat(repeat=r, number=n)
            )
  
Run Code Online (Sandbox Code Playgroud)