Tre*_*ith 7 python multiprocessing dask python-xarray netcdf4
我对 xarray 相当陌生,我目前正在尝试利用它来对某些 NetCDF 进行子集化。我在共享服务器上运行它,想知道如何最好地限制 xarray 使用的处理能力,以便它与其他人很好地合作。我已经阅读了 dask 和 xarray 文档,但我似乎不清楚如何设置 cpus/线程的上限。以下是空间子集的示例:
import glob
import os
import xarray as xr
from multiprocessing.pool import ThreadPool
import dask
wd = os.getcwd()
test_data = os.path.join(wd, 'test_data')
lat_bnds = (43, 50)
lon_bnds = (-67, -80)
output = 'test_data_subset'
def subset_nc(ncfile, lat_bnds, lon_bnds, output):
if not glob.os.path.exists(output):
glob.os.makedirs(output)
outfile = os.path.join(output, os.path.basename(ncfile).replace('.nc', '_subset.nc'))
with dask.config.set(scheduler='threads', pool=ThreadPool(5)):
ds = xr.open_dataset(ncfile, decode_times=False)
ds_sub = ds.where(
(ds.lon >= min(lon_bnds)) & (ds.lon <= max(lon_bnds)) & (ds.lat >= min(lat_bnds)) & (ds.lat <= max(lat_bnds)),
drop=True)
comp = dict(zlib=True, complevel=5)
encoding = {var: comp for var in ds.data_vars}
ds_sub.to_netcdf(outfile, format='NETCDF4', encoding=encoding)
list_files = glob.glob(os.path.join(test_data, '*'))
print(list_files)
for i in list_files:
subset_nc(i, lat_bnds, lon_bnds, output)
Run Code Online (Sandbox Code Playgroud)
我已经通过移动ThreadPool配置尝试了一些变化,但我仍然看到服务器中的活动过多top(> 3000%cpu活动)。我不确定问题出在哪里。
这实际上在这里解决了https://github.com/pydata/xarray/issues/2417#issuecomment-460298993(GitHub
问题似乎来自这里的提问者)
建议的解决方案@jhamman(再次是非常好的配置文件匹配)是将环境变量设置OMP_NUM_THREADS为所需的线程数(建议大约是所需核心数的 2 倍,大概是为了利用现代 Intel/AMD 线程技术)