【发布时间】:2021-09-22 21:23:06
【问题描述】:
我显然在 multiprocessing 上做错了,但我不确定是什么——我希望在这个任务上看到一些加速,但是在分叉进程中运行这个测试函数所花费的时间是 2比它在主要过程中花费的时间多几个数量级。这是一项不平凡的任务,所以我不认为这是工作负载太小而无法从多处理中受益的情况,如this question 和基本上所有其他关于multiprocessing 的SO 问题。而且我知道启动新进程会产生开销,但我的函数会返回执行实际计算所花费的时间,我认为这会在完成分叉开销之后发生。
我查看了一堆文档和示例,使用 map、map_async、apply 和 apply_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 访问争用。
...但这似乎也是一个与multiprocessing 和numpy 之间的交互有关的问题,因为我尝试了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 建议
numpy和multiprocessing之间可能存在不良互动,有什么见解吗?
标签: python numpy python-multiprocessing