【发布时间】:2021-08-25 22:11:20
【问题描述】:
我有一个非常大的数组(这里有 200 万个单元格),并且想为数组中的每个单元格执行一个工作流。这是我的测试代码:
import numpy as np
import dask
from dask.distributed import Client, LocalCluster
import dask.bag as db
# invoke 8 workers
cluster = LocalCluster(n_workers=8)
client = Client(cluster)
# test workflow to be applied to each cell. The real case is much more complex than this.
def g(x):
np.sqrt(np.abs(x)) ** np.log(np.abs(x))
# test array with 2,000,000 cells. Values are normally distributed.
test_array = np.random.randn(2000000)
然后我使用串行和并行计算来执行这个工作流。
%%time
# serial computation
results_serial = np.zeros((2000000, 1))
for i in range(len(test_array)):
results_serial[i] = g(test_array[i])
这需要大约 11 秒才能在我的机器上运行。但是对于使用dask.bag的并行计算:
%%time
# parallel computation
b = db.from_sequence(test_array, npartitions=24)
b = b.map(g)
results_parallel = b.compute()
在我的机器上运行大约需要 90 秒,这比串行计算要慢得多。我想知道为什么我们会看到这种情况,以及使用 dask.bag 或 Dask 的其他模块加速并行案例的建议解决方案是什么?
以下是此代码的笔记本版本的链接,其中包含更多 cmets: https://github.com/whyjz/dask-playground/blob/main/dask-test.ipynb
【问题讨论】:
标签: dask dask-distributed