【发布时间】:2020-05-28 05:00:03
【问题描述】:
我有一个关于 dask.distributed 中任务的调度/执行顺序的问题,以应对大型原始数据集的强数据缩减。
我们将 dask.distributed 用于从电影帧中提取信息的代码。它的具体应用是结晶学,但通常的步骤是:
- 将存储为 HDF5 文件中的 3D 数组(或其中一些连接的文件)中的电影帧读取到 dask 数组中。这显然是 I/O 繁重的工作
- 将这些帧分组到通常包含 10 个移动静止图像的连续子堆栈中,其中的帧被聚合(求和或平均),从而生成单个 2D 图像。
- 对 2D 图像(例如某些特征的位置)运行几个计算量大的分析函数,返回一个结果字典,与电影本身相比,它可以忽略不计。
我们通过在第 1 步和第 2 步中使用 dask.array API 来实现这一点(后者使用 map_blocks,块/块大小为一个或几个聚合子堆栈),然后将数组块转换为dask.delayed 对象(使用to_delayed),它们被传递给执行实际数据缩减功能的函数(步骤 3)。我们注意在步骤 2 中正确对齐 HDF5 数组的块、dask 计算和聚合范围,以便每个最终延迟对象(tasks 的元素)的任务图非常干净。下面是示例代码:
def sum_sub_stacks(mov):
# aggregation function
sub_stk = []
for k in range(mov.shape[0]//10):
sub_stk.append(mov[k*10:k*10+10,...].sum(axis=0, keepdims=True))
return np.concatenate(sub_stk)
def get_info(mov):
# reduction function
results = []
for frame in mov:
results.append({
'sum': frame.sum(),
'variance': frame.var()
# ...actually much more complex/expensive stuff
})
return results
# connect to dask.distributed scheduler
client = Client(address='127.0.0.1:8786')
# 1: get the movie
fh = h5py.File('movie_stack.h5')
movie = da.from_array(fh['/entry/data/raw_counts'], chunks=(100,-1,-1))
# 2: sum sub-stacks within movie
movie_aggregated = movie.map_blocks(sum_sub_stacks,
chunks=(10,) + movie.chunks[1:],
dtype=movie.dtype)
# 3: create and run reduction tasks
tasks = [delayed(get_info)(chk)
for chk in movie_aggregated.to_delayed().ravel()]
info = client.compute(tasks, sync=True)
理想的操作调度显然是让每个工作人员在单个块上执行 1-2-3 序列,然后继续下一个块,这将保持 I/O 负载恒定、CPU 最大化和内存不足.
相反,首先所有工作人员都试图从文件中读取尽可能多的块(第 1 步),这会造成 I/O 瓶颈并迅速耗尽工作人员内存,从而导致本地驱动器崩溃。通常,在某些时候,worker 最终会移动到第 2/3 步,这会快速释放内存并正确使用所有 CPU,但在其他情况下,worker 会以不协调的方式被杀死,或者整个计算会停止。中间情况也会发生,幸存的工人仅在一段时间内表现得合理。
是否有任何方法可以提示调度程序以如上所述的首选顺序处理任务,或者是否有其他方法可以改善调度行为?还是这种代码/做事方式本身就很愚蠢?
【问题讨论】:
标签: python dask dask-distributed