【问题标题】:Why is this multiprocessing code slower than serial? How would you do this better?为什么这个多处理代码比串行代码慢?你会如何做得更好?
【发布时间】:2021-09-22 21:23:06
【问题描述】:

我显然在 multiprocessing 上做错了,但我不确定是什么——我希望在这个任务上看到一些加速,但是在分叉进程中运行这个测试函数所花费的时间是 2比它在主要过程中花费的时间多几个数量级。这是一项不平凡的任务,所以我不认为这是工作负载太小而无法从多处理中受益的情况,如this question 和基本上所有其他关于multiprocessing 的SO 问题。而且我知道启动新进程会产生开销,但我的函数会返回执行实际计算所花费的时间,我认为这会在完成分叉开销之后发生。

我查看了一堆文档和示例,使用 mapmap_asyncapplyapply_async 而不是 imap_unordered 尝试了此代码的版本,但我得到了可比的结果案例。在这一点上我完全被迷惑了......任何帮助理解我做错了什么都会很棒,修改后的代码 sn-p 提供了一个如何通过并行化此任务来获得性能提升的示例。谢谢!

import time

t_start = time.time()
from multiprocessing import Pool
from numpy.random import rand, randint
import numpy.linalg as la
import numpy as np


def f(sp):
    (M, i) = sp
    t0 = time.time()
    M = (M @ M.T) / 1000 + np.eye(M.shape[0])
    M_inv = la.inv(M)
    t_elapsed = time.time() - t0
    return i, M.shape[0], la.det(M_inv), t_elapsed


randmat = lambda m: rand(m, m)

N = 20
n_m = 1500
specs = list(zip([randmat(n_m) for _ in range(N)], range(N)))

t0 = time.time()
for result in [f(sp) for sp in specs]:
    print(result)

print(f"\n--- serial time: {time.time()-t0}; total elapsed: {time.time()-t_start }\n")

t0 = time.time()
with Pool(processes=10) as pool:
    multiple_results = pool.imap_unordered(f, specs)
    for result in multiple_results:
        print(result)

print(f"\n--- parallel time: {time.time()-t0}\n")

输出:

(0, 1500, 2.613708465497732e-76, 0.17858004570007324)
(1, 1500, 2.3314319199405457e-76, 0.18518280982971191)
(2, 1500, 2.4510533542449015e-76, 0.18424344062805176)
...(snip)...
(17, 1500, 2.0972534465354807e-76, 0.18465876579284668)
(18, 1500, 2.4890185099760677e-76, 0.18526124954223633)
(19, 1500, 3.0716539033944427e-76, 0.17455506324768066)

--- serial time: 5.365333557128906; total elapsed: 5.747828006744385

(0, 1500, 2.613708465497732e-76, 9.31627368927002)
(1, 1500, 2.3314319199405457e-76, 9.709473848342896)
(5, 1500, 2.6716145027956763e-76, 10.101540327072144)
...(snip)...
(19, 1500, 3.0716539033944427e-76, 10.48097825050354)
(18, 1500, 2.4890185099760677e-76, 10.82164478302002)
(17, 1500, 2.0972534465354807e-76, 10.97563886642456)

--- parallel time: 40.98197340965271

(系统信息:Mint 20.1、AMD Ryzen 5 2600(6 核、12 线程)、Python 3.8)

更新

我相信@Paul 下面的回答可能是正确的。进一步研究这一点,我提出了一个更加病态的测试用例来解决@Charles Duffy 对序列化成本高昂的担忧——在下面的例子中,我只向每个进程发送一个 100 元素的微小 numpy 向量,而内部每个函数调用的时间从 ~0.03 秒到大约 100 秒!!差了三个数量级以上!疯子!我只能想象在multiprocessing 无法处理的后台发生了某种灾难性的 CPU 访问争用。

...但这似乎也是一个与multiprocessingnumpy 之间的交互有关的问题,因为我尝试了ray,我得到了我期望的那种性能提升并行化。

新结果 tl;dr

  • 串行时间:4.73s
  • ray时间:0.112s
  • multiprocessing 时间:217.48s (!!!)

代码 v2

import time

t_start = time.time()

import multiprocessing as mp
from numpy.random import randn
import numpy.linalg as la
import numpy as np
import ray


num_vecs = 20
vec_size = 100
inputs = [(randn(vec_size, 1), i, t_start) for i in range(num_vecs)]


def f(input):
    (v, i, t_start) = input
    t0 = time.time()
    det_sum = 0
    M = (v @ v.T) + np.diag(v[:, 0])
    for _ in range(50):
        M = M @ (M.T)
        M = M @ (la.inv(M + np.eye(M.shape[0])) / 2)
        det_sum += la.det(M)
    t_inner = time.time() - t0
    t_since_start = time.time() - t_start
    return i, det_sum, t_inner, t_since_start


def print_result(r):
    print(
        f"id: {r[0]:2}, det_sum: {r[1]:.3e}, inner time: {r[2]:.4f}, time since start: {r[3]:.4f}"
    )


t0 = time.time()
for result in [f(sp) for sp in inputs]:
    print_result(result)
print(f"\n--- serial time: {time.time()-t0}; total elapsed: {time.time()-t_start }\n")

ray.init(num_cpus=10)
g = ray.remote(f)
t0 = time.time()
results = ray.get([g.remote(s) for s in inputs])
for result in results:
    print(result)
print(f"\n--- parallel time: {time.time()-t0}\n")

t0 = time.time()
with mp.Pool(processes=10) as pool:
    multiple_results = pool.imap_unordered(f, inputs)
    for result in multiple_results:
        print(result)
print(f"\n--- parallel time: {time.time()-t0}\n")

输出 v2

id:  0, det_sum: 1.427e-133, inner time: 2.8998, time since start: 3.1620
id:  1, det_sum: 3.294e-118, inner time: 0.3816, time since start: 3.5436
id:  2, det_sum: 2.729e-114, inner time: 0.0569, time since start: 3.6005
...(snip)...
id: 17, det_sum: 2.372e-104, inner time: 0.0344, time since start: 4.8887
id: 18, det_sum: 3.523e-116, inner time: 0.0509, time since start: 4.9396
id: 19, det_sum: 9.242e-101, inner time: 0.0549, time since start: 4.9945

--- serial time: 4.734628677368164; total elapsed: 4.996868848800659

id:  0, det_sum: 1.427e-133, inner time: 0.0436, time since start: 6.1446
id:  1, det_sum: 3.294e-118, inner time: 0.0465, time since start: 6.1541
id:  2, det_sum: 2.729e-114, inner time: 0.0436, time since start: 6.1517
...(snip)...
id: 17, det_sum: 2.372e-104, inner time: 0.0438, time since start: 6.2027
id: 18, det_sum: 3.523e-116, inner time: 0.0394, time since start: 6.1995
id: 19, det_sum: 9.242e-101, inner time: 0.0413, time since start: 6.2032

--- parallel time: 0.1118767261505127

id:  0, det_sum: 1.427e-133, inner time: 101.0206, time since start: 107.2395
id:  2, det_sum: 2.729e-114, inner time: 102.6551, time since start: 108.8744
id:  5, det_sum: 2.063e-111, inner time: 104.2321, time since start: 110.4516
...(snip)...
id: 18, det_sum: 3.523e-116, inner time: 102.0273, time since start: 223.5556
id: 16, det_sum: 5.887e-99, inner time: 102.9106, time since start: 223.5907
id: 19, det_sum: 9.242e-101, inner time: 101.1289, time since start: 223.6742

--- parallel time: 217.47953820228577

【问题讨论】:

  • multiprocessing 涉及大量开销。除非您分析工作负载并确定序列化/反序列化将比实际计算便宜,否则期望它使事情变得更快是不合理的。
  • 启动一个新进程并不是那么昂贵——fork() 是一个便宜的系统调用。 将数据输入和输出该过程...这很昂贵。
  • 谢谢@CharlesDuffy 没看到。
  • 无论如何,我将从采样分析器开始,以获取有关挂钟时间的非推测性数字。如果它序列化/反序列化,那会像大拇指一样突出。如果是其他问题……好吧,我们也会知道,您将能够更好地区分新问题和现有问题。
  • 感谢@CharlesDuffy,我添加了一个新的测试用例,我无法想象它是序列化绑定(只发送一个 50 elt numpy 数组),它的性能非常糟糕。 @Paul 建议 numpymultiprocessing 之间可能存在不良互动,有什么见解吗?

标签: python numpy python-multiprocessing


【解决方案1】:

我相信您的numpy 很可能已经在单进程模型中利用了您的多核架构。比如来自here

但现在许多架构都有一个 BLAS,它也利用了多核机器。如果您的 numpy/scipy 是使用其中之一编译的,那么 dot() 将在您不执行任何操作的情况下并行计算(如果这更快)。其他矩阵运算也是如此,例如求逆、奇异值分解、行列式等。

你可以检查一下:

>>> import numpy as np
>>> np.show_config()

而且,作为一个简单的测试,如果您增加矩阵的大小并直接运行它,您是否看到numpy 正在使用您的多个内核?例如,运行时观看top

>>> n_m = 20000
>>> M = np.random.rand(n_m, n_m)
>>> M = (M @ M.T) / 1000 + np.eye(M.shape[0])

这可能足够慢,以至于您可以在一个进程中查看它是否已经在使用多个内核。

您可以想象,如果它已经在这样做,那么将其拆分为不同的进程纯粹会增加开销,因此会更慢。

【讨论】:

  • 谢谢@Paul,我认为你是对的——我用更多的测试用例来解决这个问题,我用一个更糟糕的例子编辑了原始问题。然而,ray 模块似乎能够以某种方式提供性能提升而不会遇到这个问题。我不太了解这些内部结构是如何工作的,无法理解为什么multiprocessing 会给出如此糟糕的结果而ray 工作正常,但至少我有前进的方向......
猜你喜欢
  • 1970-01-01
  • 2021-01-24
  • 1970-01-01
  • 2012-11-26
  • 1970-01-01
  • 2013-08-08
  • 1970-01-01
  • 1970-01-01
  • 2014-08-31
相关资源
最近更新 更多