【发布时间】: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