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.
由于它是一个 3D 数组,我假设以下情况。
\n\n水分子在 xz 平面上的速度 (960 \xce\xbcm x 512 \xce\xbcm) 随时间 (801 fs) 变化。\n求 xz 平面上每个元素在整个时间内的速度均方根函数。
\n\nnumpy 代码将是:
\n\nxz_plane_rmsf = (data ** 2).mean(axis=0)\nRun 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 处的平均速度。
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()\nRun Code Online (Sandbox Code Playgroud)\n\n其中data是 3D Dask 数组。
唯一剩下的问题是将 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)>\nRun Code Online (Sandbox Code Playgroud)\n