【问题标题】:Parallel tasks on subsets of a dask array wrapped in an xarray Dataset包装在 xarray 数据集中的 dask 数组子集上的并行任务
【发布时间】:2020-07-13 10:49:49
【问题描述】:

我有一个存储为 zarr 的大型 xarray.Dataset。我想对其执行一些自定义操作,这些操作不能仅通过使用 Dask 集群将自动处理的类似 numpy 的函数来完成。因此,我将数据集划分为小的子集,并为每个子集向我的 Dask 集群提交一个格式为

的任务
def my_task(zarr_path, subset_index):
    ds = xarray.open_zarr(zarr_path)  # this returns an xarray.Dataset containing a dask.array
    sel = ds.sel(partition_index)
    sel  = sel.load()  # I want to get the data into memory
    # then do my custom operations
    ...

但是,我注意到这会创建一个“任务中的任务”:当工作人员收到“my_task”时,它会依次将任务提交到集群以加载数据集的相关部分。为了避免这种情况并确保在工作人员中执行完整的任务,我提交的是任务:

def my_task_2(zarr_path, subset_index):
    with dask.config.set(scheduler="threading"):
        my_task(zarr_path, subset_index)

这是最好的方法吗?这种情况的最佳做法是什么?

【问题讨论】:

    标签: dask python-xarray


    【解决方案1】:

    通常使用apply_ufunc 或map_blocks 等方法在Xarray 数据集中的块之间并行应用函数。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2022-11-10
      • 1970-01-01
      • 2020-10-28
      • 1970-01-01
      • 2021-05-19
      • 1970-01-01
      • 2021-08-17
      相关资源
      最近更新 更多