如何用xarray的apply_ufunc结合dask加速逐日均温最大值计算?
使用Dask并行计算每日最高气温(替代groupby)
我明白你想把现有的groupby('time.dayofyear').max()换成Dask并行化实现来提速,而且刚接触Dask觉得官方文档不够接地气。下面我一步步给你讲怎么用apply_ufunc实现,还有一些apply_ufunc的核心实用指引。
先看适配你的场景的代码实现
你的数据集是3小时分辨率,2008年(闰年)刚好2928个时间步(366×8),用apply_ufunc的核心是按日期分组后并行求最大值。不过先给你推荐一个更简洁且高效的替代方案——用xarray的resample结合Dask,它内部已经做了并行优化,代码更易读:
import xarray as xr # 读取数据,保持time维度按天拆分的chunk设置(8个时间步刚好是一天) TAS = xr.open_dataset(INFILE_template.format(YR, YR), chunks={'time': 8}) # 按天重采样并计算最大值,自动利用Dask并行 DAYMAX = TAS.tas.resample(time="D").max() # 保存结果 DAYMAX.to_netcdf(OUTFILE_template.format(YR, YR))
如果你确实想学习apply_ufunc的用法,下面是用它实现分组求每日最高温的代码:
import xarray as xr import dask.array as da TAS = xr.open_dataset(INFILE_template.format(YR, YR), chunks={'time': 8}) # 生成日期分组标签(把时间转成"YYYY-MM-DD"格式的字符串) date_groups = TAS.time.dt.strftime("%Y-%m-%d") # 定义Dask原生的分组求最大函数 def daily_max(arr, group_labels): return da.groupby_reduction(arr, group_labels, reduction=da.max, axis=0) # 用apply_ufunc调用并行计算 DAYMAX = xr.apply_ufunc( daily_max, TAS.tas, date_groups, input_core_dims=[["time"], ["time"]], # 指定要处理的核心维度 output_core_dims=[["date"]], # 输出的新维度 vectorize=True, # 自动对lat/lon维度做广播处理 dask="parallelized", # 启用Dask并行 output_dtypes=[TAS.tas.dtype] ) # 给date维度添加坐标 DAYMAX = DAYMAX.assign_coords(date=sorted(date_groups.unique())) DAYMAX.to_netcdf(OUTFILE_template.format(YR, YR))
apply_ufunc核心实用指引(接地气版)
如果你想深入掌握这个工具,这几个要点一定要记牢:
- 核心维度(core_dims):必须明确指定输入/输出中需要被函数处理的维度(比如例子中的
time),剩下的维度(如lat/lon)会自动并行广播处理。 - Dask模式:设置
dask="parallelized"会让xarray把任务拆分成Dask任务块,适配多核或集群环境;如果你的函数本身支持Dask数组,可设置dask="allowed"。 - 向量化(vectorize):当函数只处理核心维度、需要对其他维度广播时,开启这个参数可以避免手动写循环。
- 分组聚合:Dask原生的
da.groupby_reduction是高效分组计算的首选,支持max、min、mean等绝大多数常用聚合操作。
补充:原groupby效率问题的小说明
你原代码的groupby('time.dayofyear').max()其实也会调用Dask,但resample的逻辑更贴合"按时间间隔分组"的需求,尤其是当处理跨年份数据时,dayofyear会把不同年份的同一天归为一组,而resample是严格按日历日期分组,更符合你"每日最高温"的需求。
内容的提问来源于stack exchange,提问作者Shrad
相关产品推荐
相关产品推荐

