【问题标题】:Optimizing dask.distributed scheduling for data reduction优化 dask.distributed 调度以减少数据
【发布时间】:2020-05-28 05:00:03
【问题描述】:

我有一个关于 dask.distributed 中任务的调度/执行顺序的问题,以应对大型原始数据集的强数据缩减。

我们将 dask.distributed 用于从电影帧中提取信息的代码。它的具体应用是结晶学,但通常的步骤是:

  1. 将存储为 HDF5 文件中的 3D 数组(或其中一些连接的文件)中的电影帧读取到 dask 数组中。这显然是 I/O 繁重的工作
  2. 将这些帧分组到通常包含 10 个移动静止图像的连续子堆栈中,其中的帧被聚合(求和或平均),从而生成单个 2D 图像。
  3. 对 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


    【解决方案1】:

    首先,你所做的事情本身并没有什么愚蠢的!

    一般来说,Dask 试图减少它所持有的临时对象的数量,并且它还通过可并行性(图的宽度和工作人员的数量)来平衡这一点。调度很复杂,Dask 还使用了另一种优化,它将任务融合在一起,使它们更加优化。如果有很多小块,您可能会遇到问题:https://docs.dask.org/en/latest/array-best-practices.html?highlight=chunk%20size#select-a-good-chunk-size

    Dask 确实有许多optimization configurations,我建议在考虑其他块大小后使用它们。我还鼓励您通读following issue,因为围绕调度配置进行了健康的讨论。

    最后,您可能会考虑额外的工作人员 memory configuration,因为您可能希望更严格地控​​制每个工作人员应该使用多少内存

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-06-12
      • 1970-01-01
      • 1970-01-01
      • 2011-01-22
      • 2017-02-15
      • 1970-01-01
      • 2011-10-29
      • 1970-01-01
      相关资源
      最近更新 更多