open_mfdataset - different behavior with dask.distributed.LocalCluster
Nobody has claimed this yet.
- Dominant language
- Python
- Stars
- 4.2k
- Forks
- 1.4k
- Avg merge
- 2d 15h
- Merged PRs (30d)
- 14
Description
Big fan of Xarray! Not that familiar with submitting tickets like this, so my apologies for rule breaking. Also, if this belongs over in the dask project, I can move there.
dask 2.6.0
numpy 1.17.3
xarray 0.14.1
netCDF4 1.5.3
I am attempting to use open_mfdataset on nc files I've generated through dask/xarray after initializing the dask LocalCluster. I've found that I am able to compute successfully when I don't run the distributed cluster. But if I do, I get a variety of issues. I've got a synthetic data generating example here. Running the soundspeed.compute() will sometimes succeed, and will sometimes cause worker restarts resulting in hdf errors and no return.
I was thinking it was something with serialization, i've seen other tickets with similar issues, but I don't see how it applies to my test case.
Example code:
import numpy as np
import xarray as xr
import os
from dask.distributed import Client
cl = Client()
outpth = r'D:\dasktest\data_dir\EM2040\converted\test'
mint = 0
maxt = 1000
for i in range(100):
times = np.arange(mint, maxt)
beams = np.arange(250)
sectors=['40107_0_260000', '40107_1_320000', '40107_2_290000']
soundspeed = np.random.randn(1000,3,250)
ds = xr.Dataset({'soundspeed': (('time','sectors','beams'), soundspeed)},
{'time': times, 'sectors': sectors, 'beams':beams},)
ds.to_netcdf(os.path.join(outpth, 'test{}.nc'.format(i)), mode='w')
mint = maxt
maxt += 1000
fils = [os.path.join(outpth, x) for x in os.listdir(outpth) if os.path.splitext(x)[1] == '.nc']
tst = xr.open_mfdataset(fils, concat_dim='time', combine='nested')
tst.soundspeed.compute()
I've found that running this example with <10 files reduces the number of errors I'm getting dramatically. I've tried this on different machines in different domain environments just to be sure.
I really just want to make sure I'm not making a silly mistake somewhere. Appreciate the help.
My last run on actual data:
>>> ra.soundspeed.compute()
distributed.nanny - WARNING - Restarting worker
distributed.nanny - WARNING - Restarting worker
distributed.nanny - WARNING - Restarting worker
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001F83F1E2360>, key=BasicIndexer((slice(None, None, None), slice(None, None, None)))))), (slice(0, 1719, None), slice(0, 3, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
Traceback (most recent call last):
File "<stdin>", line 1, in <module>
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\dataarray.py", line 837, in compute
return new.load(**kwargs)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\dataarray.py", line 811, in load
ds = self._to_temp_dataset().load(**kwargs)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\dataset.py", line 649, in load
evaluated_data = da.compute(*lazy_data.values(), **kwargs)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\dask\base.py", line 436, in compute
results = schedule(dsk, keys, **kwargs)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\distributed\client.py", line 2545, in get
results = self.gather(packed, asynchronous=asynchronous, direct=direct)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\distributed\client.py", line 1845, in gather
asynchronous=asynchronous,
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\distributed\client.py", line 762, in sync
self.loop, func, *args, callback_timeout=callback_timeout, **kwargs
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\distributed\utils.py", line 333, in sync
raise exc.with_traceback(tb)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\distributed\utils.py", line 317, in f
result[0] = yield future
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\tornado\gen.py", line 735, in run
value = future.result()
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\distributed\client.py", line 1701, in _gather
raise exception.with_traceback(traceback)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\dask\array\core.py", line 106, in getter
c = np.asarray(c)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\numpy\core\_asarray.py", line 85, in asarray
return array(a, dtype, copy=False, order=order)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\indexing.py", line 481, in __array__
return np.asarray(self.array, dtype=dtype)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\numpy\core\_asarray.py", line 85, in asarray
return array(a, dtype, copy=False, order=order)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\indexing.py", line 643, in __array__
return np.asarray(self.array, dtype=dtype)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\numpy\core\_asarray.py", line 85, in asarray
return array(a, dtype, copy=False, order=order)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\indexing.py", line 547, in __array__
return np.asarray(array[self.key], dtype=None)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\backends\netCDF4_.py", line 72, in __getitem__
key, self.shape, indexing.IndexingSupport.OUTER, self._getitem
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\indexing.py", line 827, in explicit_indexing_adapter
result = raw_indexing_method(raw_key.tuple)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\backends\netCDF4_.py", line 83, in _getitem
original_array = self.get_array(needs_lock=False)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\backends\netCDF4_.py", line 62, in get_array
ds = self.datastore._acquire(needs_lock)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\backends\netCDF4_.py", line 360, in _acquire
with self._manager.acquire_context(needs_lock) as root:
File "C:\PydroXL_19\envs\dasktest\lib\contextlib.py", line 81, in __enter__
return next(self.gen)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\backends\file_manager.py", line 186, in acquire_context
file, cached = self._acquire_with_cache_info(needs_lock)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\backends\file_manager.py", line 204, in _acquire_with_cache_info
file = self._opener(*self._args, **kwargs)
File "netCDF4\_netCDF4.pyx", line 2321, in netCDF4._netCDF4.Dataset.__init__
File "netCDF4\_netCDF4.pyx", line 1885, in netCDF4._netCDF4._ensure_nc_success
OSError: [Errno -101] NetCDF: HDF error: b'D:\\dasktest\\data_dir\\EM2040\\converted\\rangeangle_20.nc'
My last run on the synthetic data set generated above:
>>> tst.soundspeed.compute()
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FC8AA20>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FCB82D0>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FCB8240>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FCB81F8>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FCB81B0>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FCB8360>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FCB83A8>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FCB8510>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FCB8750>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FCB8990>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FCB8BD0>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FCB8E10>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FC9D090>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FC9D2D0>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FC9D510>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FC9D750>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FC9DC18>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FC9DBD0>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FC9DCA8>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
Traceback (most recent call last):
File "<stdin>", line 1, in <module>
distributed.worker - WARNING - Compute Failed
Function: getter
args: (ImplicitToExplicitIndexingAdapter(array=CopyOnWriteArray(array=LazilyOuterIndexedArray(array=<xarray.backends.netCDF4_.NetCDF4ArrayWrapper object at 0x000001BB5FC9DD38>, key=BasicIndexer((slice(None, None, None), slice(None, None, None), slice(None, None, None)))))), (slice(0, 1000, None), slice(0, 3, None), slice(0, 250, None)))
kwargs: {}
Exception: OSError(-101, 'NetCDF: HDF error')
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\dataarray.py", line 837, in compute
return new.load(**kwargs)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\dataarray.py", line 811, in load
ds = self._to_temp_dataset().load(**kwargs)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\dataset.py", line 649, in load
evaluated_data = da.compute(*lazy_data.values(), **kwargs)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\dask\base.py", line 436, in compute
results = schedule(dsk, keys, **kwargs)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\distributed\client.py", line 2545, in get
results = self.gather(packed, asynchronous=asynchronous, direct=direct)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\distributed\client.py", line 1845, in gather
asynchronous=asynchronous,
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\distributed\client.py", line 762, in sync
self.loop, func, *args, callback_timeout=callback_timeout, **kwargs
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\distributed\utils.py", line 333, in sync
raise exc.with_traceback(tb)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\distributed\utils.py", line 317, in f
result[0] = yield future
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\tornado\gen.py", line 735, in run
value = future.result()
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\distributed\client.py", line 1701, in _gather
raise exception.with_traceback(traceback)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\dask\array\core.py", line 106, in getter
c = np.asarray(c)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\numpy\core\_asarray.py", line 85, in asarray
return array(a, dtype, copy=False, order=order)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\indexing.py", line 481, in __array__
return np.asarray(self.array, dtype=dtype)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\numpy\core\_asarray.py", line 85, in asarray
return array(a, dtype, copy=False, order=order)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\indexing.py", line 643, in __array__
return np.asarray(self.array, dtype=dtype)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\numpy\core\_asarray.py", line 85, in asarray
return array(a, dtype, copy=False, order=order)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\indexing.py", line 547, in __array__
return np.asarray(array[self.key], dtype=None)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\backends\netCDF4_.py", line 72, in __getitem__
key, self.shape, indexing.IndexingSupport.OUTER, self._getitem
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\core\indexing.py", line 827, in explicit_indexing_adapter
result = raw_indexing_method(raw_key.tuple)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\backends\netCDF4_.py", line 83, in _getitem
original_array = self.get_array(needs_lock=False)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\backends\netCDF4_.py", line 62, in get_array
ds = self.datastore._acquire(needs_lock)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\backends\netCDF4_.py", line 360, in _acquire
with self._manager.acquire_context(needs_lock) as root:
File "C:\PydroXL_19\envs\dasktest\lib\contextlib.py", line 81, in __enter__
return next(self.gen)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\backends\file_manager.py", line 186, in acquire_context
file, cached = self._acquire_with_cache_info(needs_lock)
File "C:\PydroXL_19\envs\dasktest\lib\site-packages\xarray\backends\file_manager.py", line 204, in _acquire_with_cache_info
file = self._opener(*self._args, **kwargs)
File "netCDF4\_netCDF4.pyx", line 2321, in netCDF4._netCDF4.Dataset.__init__
File "netCDF4\_netCDF4.pyx", line 1885, in netCDF4._netCDF4._ensure_nc_success
OSError: [Errno -101] NetCDF: HDF error: b'D:\\dasktest\\data_dir\\EM2040\\converted\\test\\test4.nc'
Contributor guide
First steps
- Read the whole issue, then the project's contributing guide.
- Comment on the issue to say you are picking it up — it saves two people doing the same work.
- Fork the repository and make your change on a branch.
- Open a pull request that references the issue number.
Research direction
Start with the synthetic example using Client(), xr.open_mfdataset(..., combine='nested'), and tst.soundspeed.compute(), then follow the traceback through xarray's netCDF4 backend and file manager. Compare behavior with and without the distributed client and with fewer files. Done means the reported HDF errors and worker restarts are reproducible, explained, and covered by a regression test or a documented limitation.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- numpy, python
- Domain
- data, distributed-systems
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 25/100