【问题标题】:With xarray, how to parallelize 1D operations on a multidimensional Dataset?使用 xarray,如何在多维数据集上并行化一维操作?
【发布时间】:2020-03-06 19:42:57
【问题描述】:

我有一个 4D xarray 数据集。我想在特定维度(此处为时间)上的两个变量之间进行线性回归,并将回归参数保留在 3D 数组(其余维度)中。 我设法通过使用此串行代码获得了我想要的结果,但它相当慢:

# add empty arrays to store results of the regression
res_shape = tuple(v for k,v in ds[x].sizes.items() if k != 'year')
res_dims = tuple(k for k,v in ds[x].sizes.items() if k != 'year')
ds[sl] = (res_dims, np.empty(res_shape, dtype='float32'))
ds[inter] = (res_dims, np.empty(res_shape, dtype='float32'))
# Iterate in kept dimensions
for lat in ds.coords['latitude']:
    for lon in ds.coords['longitude']:
        for duration in ds.coords['duration']:
            locator = {'longitude':lon, 'latitude':lat, 'duration':duration}
            sel = ds.loc[locator]
            res = scipy.stats.linregress(sel[x], sel[y])
            ds[sl].loc[locator] = res.slope
            ds[inter].loc[locator] = res.intercept

我怎样才能加速和并行化这个操作?

我知道apply_ufunc 可能是一个选项(并且可以与 dask 并行),但我没有设法正确设置参数。

以下问题相关但没有答案:

编辑 2:将之前的编辑移至答案

【问题讨论】:

  • 纬度和经度不是要一起考虑的吗?给两个点(x1,y1)和(x2,y2),为什么要迭代(x1,y1),(x1,y2),(x2,y1),(x2,y2)?
  • @JohnZwinck 两个点可能在同一纬度,但经度不同。
  • 是的。您不应该迭代 (x,y) 对,而不是所有 x 与所有 y 的笛卡尔积吗?

标签: python dask python-xarray


【解决方案1】:

可以通过传递vectorize=True 使用apply_ufunc()scipy.stats.linregress(和其他非ufunc)应用于xarray 数据集,如下所示:

# return a tuple of DataArrays
res = xr.apply_ufunc(scipy.stats.linregress, ds[x], ds[y],
        input_core_dims=[['year'], ['year']],
        output_core_dims=[[], [], [], [], []],
        vectorize=True)
# add the data to the existing dataset
for arr_name, arr in zip(array_names, res):
    ds[arr_name] = arr

尽管仍然是串行的,apply_ufunc 在这种特定情况下比循环实现快大约 36 倍。

但是,使用 dask 的并行化仍然没有实现多个输出,例如来自 scipy.stats.linregress 的输出:

NotImplementedError:来自 apply_ufunc 的多个输出尚不支持 dask='parallelized'

【讨论】:

  • 我正在尝试使用 sp.stats.ranksums 但没有运气。我收到关于我修复的维度的 ValueError(您的“年份”,实际上对我来说也是“年份”),因为这两个数据集的年份值不同。
【解决方案2】:

LCT 之前的回答涵盖了这里应该说的大部分内容,然而我认为可以将dask='parallelized' 与多个输出结合起来,就像你从scipy.stats.linregress 获得的那样。

这里的技巧是将多个输出堆叠到一个数组中然后输出,您还必须使用output_core_dims kwarg 来指定来自apply_ufunc() 调用的 DataArray 输出现在将有一个额外的维度:

def new_linregress(x, y):
    # Wrapper around scipy linregress to use in apply_ufunc
    slope, intercept, r_value, p_value, std_err = stats.linregress(x, y)
    return np.array([slope, intercept, r_value, p_value, std_err])
# return a new DataArray
stats = xr.apply_ufunc(new_linregress, ds[x], ds[y],
                       input_core_dims=[['year'], ['year']],
                       output_core_dims=[["parameter"]],
                       vectorize=True,
                       dask="parallelized",
                       output_dtypes=['float64'],
                       output_sizes={"parameter": 5},
                      )

注意此方法目前仅适用于 dask='parallelized',前提是您有 dask<2.0,但如果您有类似 dask='allowed' 的其他内容,它似乎适用于多个输出。请查看此Github issue 了解更多详情。

希望对你有帮助!

编辑:我被告知只要您有xarray>=0.15.0dask<2.0 问题就已得到纠正!所以现在可以使用dask='parallelized' 来加快速度。 :)

【讨论】:

  • 这个答案太棒了!试图在时间维度上做一个类似的线性回归问题,并绞尽脑汁看着xarray.pydata.org/en/v0.15.1/examples/…。将所有内容放入一个 numpy 数组以避免“NotImplementedError”的技巧非常好,只需在 :D 之后将它们拆分出来
猜你喜欢
  • 2017-06-08
  • 2021-12-17
  • 2021-10-26
  • 2020-04-17
  • 2020-12-06
  • 1970-01-01
  • 2020-10-11
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多