Dask发送29.88 MiB大型任务图导致计算启动缓慢的问题排查与优化咨询
我最近在用Dask的时候碰到个麻烦:计算启动要花大概5分钟,还收到了关于大型任务图的警告。完整的警告内容如下:
UserWarning: Sending large graph of size 29.88 MiB.
This may cause some slowdown.
Consider loading the data with Dask directly
or using futures or delayed objects to embed the data into the graph without repetition.
See also https://docs.dask.org/en/stable/best-practices.html#load-data-with-dask for more information.
下面是我的复现代码,我用了xee(xarray的一个扩展)从Google Earth Engine拉取数据到xarray结构里,分块大小是xee认为的“安全”值——不会超过Earth Engine的查询请求限制:
import ee import json import xarray as xr from dask.distributed import performance_report import dask # 虽然已经在Dask worker上认证了Google Earth Engine,但本地机器也需要认证! with open(json_key, 'r') as file: data = json.load(file) credentials = ee.ServiceAccountCredentials(data["client_email"], json_key) ee.Initialize(credentials = credentials, opt_url='https://earthengine-highvolume.googleapis.com') WSDemo = ee.FeatureCollection("projects/robust-raster/assets/boundaries/WSDemoSHP_Albers") California = ee.FeatureCollection("projects/robust-raster/assets/boundaries/California") # ic = ee.ImageCollection('LANDSAT/LC08/C02/T1_L2').filterDate('2014-01-01', '2014-12-31') ic = ee.ImageCollection('LANDSAT/LC08/C02/T1_L2').filterDate('2020-05-01', '2020-08-31') xarray_data = xr.open_dataset(ic, engine='ee', crs="EPSG:3310", scale=30, geometry=WSDemo.geometry()) xarray_data = xarray_data.chunk({"time": 48, "X": 512, "Y": 256}) def compute_ndvi(df): # 执行计算 df['ndvi'] = (df['SR_B5'] - df['SR_B4']) / (df['SR_B5'] + df['SR_B4']) return df def user_function_wrapper(ds, user_func, *args, **kwargs): df_input = ds.to_dataframe().reset_index() df_output = user_func(df_input, *args, **kwargs) df_output = df_output.set_index(list(ds.dims)) ds_output = df_output.to_xarray() return ds_output test = xr.map_blocks(user_function_wrapper, xarray_data, args=(compute_ndvi,)) # 生成单分块运行的Dask报告 with performance_report(filename="dask_report.html"): test.compute()
你可以忽略那些把分块转成DataFrame的代码——我正在做一个功能,让用户能编写兼容pandas DataFrame的自定义函数。我现在的疑问是:29.88 MiB的任务图真的算大到会导致这么慢的启动吗?接下来我打算研究futures或者delayed对象的用法,但先想问问,我的代码里有没有明显的错误,才触发了这个警告和启动缓慢的问题?
备注:内容来源于stack exchange,提问作者Adriano Matos

