【问题标题】:Dask delayed / dask array no responseDask延迟/ dask数组无响应
【发布时间】:2023-04-04 05:10:01
【问题描述】:

我有一个分布式 dask 集群设置,我用它来加载和转换一堆数据。像魅力一样工作。

我想用它做一些并行处理。这是我的功能

el = 5000
n_using = 26
n_across= 6

mat = np.random.random((el,n_using,n_across))
idx = np.tril_indices(n_across*2, -n_across)

def get_vals(c1, m, el, idx):
    m1 = m[c1,:,:]
    corr_vals = np.zeros((el, (n_across//2)*(n_across+1)))
    for c2 in range(c1+1, el):
        corr = np.corrcoef(m1.T, m[c2,:,:].T)
        corr_vals[c2] = corr[idx]
        
    return corr_vals

lazy_get_val = dask.delayed(get_vals, pure=True)

这是我正在尝试做的单处理器版本:

arrays = [get_vals(c1, mat, el, idx) for c1 in range(el)]
all_corr = np.stack(arrays, axis=0)

工作正常,但需要几个小时。 这是我在 dask 中的做法:

lazy_list = [lazy_get_val(c1, mat, el, idx) for c1 in range(el)]
arrays = [da.from_delayed(lazy_item, dtype=float, shape=(el, 21)) for lazy_item in lazy_list]
all_corr = da.stack(arrays, axis=0)

即使它运行all_corr[1].compute(),它也只是坐在那里不响应。当我中断内核时,它似乎卡在/distributed/utils.py:

~/.../lib/python3.6/site-packages/distributed/utils.py in sync(loop, func, *args, **kwargs)

    249     else:
    250         while not e.is_set():
--> 251             e.wait(10)
    252     if error[0]:
    253         six.reraise(*error[0])

对调试有什么建议吗?


其他:

  • 如果我使用较小的 mat (el=1000) 运行它,它运行良好。
  • 如果我创建el = 5000,它就会挂起。
  • 如果我中断内核并使用el = 1000 再次运行它,它就会挂起。

【问题讨论】:

  • 我可以问你一个mcve 吗?你很接近,但我不知道这里应该放什么垫子。理想情况下,潜在的问题提问者可以复制粘贴以重现您的问题的简单内容将是理想的。
  • 对不起。我已经更新了它。再次感谢您的调查。
  • @MRocklin - 对此有什么想法吗?我切换到使用 joblib,它工作得很好,但我宁愿使用 dask,因为我已经设置了一个非常漂亮的集群。

标签: python dask dask-distributed dask-delayed


【解决方案1】:

在示例中添加导入后,我运行了一些东西,构建图表时速度非常慢。这可以通过避免将 numpy 数组直接放在延迟调用中来改进,如下所示:

# mat = np.random.random((el,n_using,n_across))
# idx = np.tril_indices(n_across*2, -n_across)
mat = dask.delayed(np.random.random)((el,n_using,n_across))
idx = dask.delayed(np.tril_indices)(n_across*2, -n_across)

或者通过将 pure=True 关键字删除到 dask.delayed (当您设置 pure=True 时,它​​必须对所有输入的内容进行哈希处理以获得它们的唯一键,您这样做了 5000 次)。我通过使用 IPython 中的 %snakeviz 魔法分析您的代码发现了这一点。

然后我跑了all_corr[1].compute(),这很好。然后我跑了all_corr.compute(),它似乎会进展到完成,但不是很快。我怀疑您的任务太小以至于开销太大,或者您的代码在 Python for 循环中花费了太多时间,因此遇到了 GIL 问题。不确定是哪个。

我建议尝试的下一件事是使用 dask.distributed 调度程序,它可以更好地处理 GIL 问题并加剧开销问题。了解其执行情况可能有助于隔离问题。

【讨论】:

  • 嗨,马特 - 感谢您的帮助。从延迟的随机调用中制作matidx 确实可以解决它们的问题,但不能解决我的“真正”问题。我生成了matidx 作为随机矩阵来创建一个mcve,但它们来自我的另一个来源。我在 dask 调度程序中尝试了它,它似乎运行良好,认为它需要 20 分钟(与使用 joblib.Parallel 的 4 分钟相比)。任务可能太小了——我会坚持使用 joblib。再次感谢您的帮助。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-12-25
  • 2022-12-11
  • 2017-07-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多