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)