使用xarray + dask的内存错误 - 使用groupby或apply_ufunc?

Tho*_*las 6 python out-of-memory dask python-xarray pandas-groupby

我使用xarray作为分析流体湍流数据的工作流程的基础,但我无法正确利用dask来限制笔记本电脑上的内存使用量.

我有一个n带维度的数据阵列('t', 'x', 'z'),我在z维度上将其拆分为5个块:

<xarray.DataArray 'n' (t: 801, x: 960, z: 512)>
dask.array<shape=(801, 960, 512), dtype=float32, chunksize=(801, 960, 5)>
Coordinates:
* t              (t) int64 200 201 202 203 204 205 206 207 208 209 210 211 ...
* x              (x) int64 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 ...
* z              (z) int64 0 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 ...  
Run Code Online (Sandbox Code Playgroud)

我想计算n超过t的均方根波动,并返回带维数的简化数据阵列('x', 'z').我想利用dask一次只在一个块上执行此操作,因为我的笔记本电脑上只有几GB RAM可用.

我写了一个通用的ufunc来计算一个dask数组的均方根:

def rms_gufunc(x, axis):
    """Generalized ufunc to calculate root mean squared value of a 1d dask array."""
    squares = np.square(x)
    return np.sqrt(squares.mean(axis=-1))
Run Code Online (Sandbox Code Playgroud)

但现在我不确定应用这个的最佳方法是什么.据我所知,我可以使用(1)xarray.apply_ufunc或(2)groupby.reduce.

1.使用apply_ufunc:

我可以使用xarray.apply_ufunc来应用这个函数:

def rms(da, dim=None):
    """
    Reduces a dataarray by calculating the root mean square along dimension dim.
    """

    if dim is None:
        raise ValueError('Must supply a dimension along which to calculate rms')

    return xr.apply_ufunc(rms_gufunc, da,
                          input_core_dims=[[dim]],
                          dask='parallelized', output_dtypes=[da.dtype])

n_rms = rms(data['n'], dim='t')
n_rms.load()  # Trigger computation
Run Code Online (Sandbox Code Playgroud)

这似乎有效,但似乎比必要的更复杂?

2.使用groupby.reduce:

xarray文档似乎暗示这是一个"拆分 - 应用 - 组合"操作,我应该可以通过类似的东西来完成

n_rms = data['n'].groupby('z').reduce(rms_gufunc, dim='t')
Run Code Online (Sandbox Code Playgroud)

然而这导致了一个MemoryError,我很确定这不是我想用groupby步骤实现的目标.我应该用来将groupby_bins数据分成我制作的块z吗?

我想知道a)如果我apply_ufunc正确使用,和b)我将如何使用相同的东西groupby.

R z*_* zu 2

由于它是一个 3D 数组,我假设以下情况。

\n\n

水分子在 xz 平面上的速度 (960 \xce\xbcm x 512 \xce\xbcm) 随时间 (801 fs) 变化。\n求 xz 平面上每个元素在整个时间内的速度均方根函数。

\n\n

numpy 代码将是:

\n\n
xz_plane_rmsf = (data ** 2).mean(axis=0)\n
Run Code Online (Sandbox Code Playgroud)\n\n

其中 data 是一个 3D numpy 数组,带有shape=(801, 960, 512). \n 的第 0、1、2 维分别data代表时间、x 坐标和 z 坐标。的每个元素data代表水分子在时间 t 和坐标 x 和 z 处的平均速度。

\n\n

Dask 数组的等效代码是:

\n\n
# Make lazy array\nxz_plane_rmsf = (data ** 2).mean(axis=0)\n# Evaluate the array\nxz_plane_rmsf = xz_plane_rmsf.compute()\n
Run Code Online (Sandbox Code Playgroud)\n\n

其中data是 3D Dask 数组。

\n\n

唯一剩下的问题是将 xarray 转换为 Dask 数组。\n我不使用 xarray,但看起来它已经是 Dask 数组:

\n\n
<xarray.DataArray 'n' (t: 801, x: 960, z: 512)>\ndask.array<shape=(801, 960, 512), dtype=float32, chunksize=(801, 960, 5)>\n
Run Code Online (Sandbox Code Playgroud)\n