Skip to content

Commit 36925bd

Browse files
authored
Avoid closed/locked HDF5 handles when lazily loading OMX data
1 parent 6688081 commit 36925bd

2 files changed

Lines changed: 57 additions & 5 deletions

File tree

‎sharrow/dataset.py‎

Lines changed: 25 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -366,7 +366,7 @@ def from_omx(
366366

367367
arrays = {}
368368
filename = omx_file_name(omx)
369-
if filename is not None:
369+
if _is_reopenable(filename):
370370
# fast path: parallel chunk decoding via h5py
371371
import concurrent.futures
372372

@@ -442,6 +442,24 @@ def _fast_load_omx_array(filename, name):
442442
return omx_reader.read_dataset(f["data"][name])
443443

444444

445+
def _is_reopenable(filename) -> bool:
446+
"""Check whether a file can be independently opened for reading with h5py.
447+
448+
Reopening can fail if the file name is unknown (e.g. an in-memory file),
449+
or if the file is already open elsewhere in a mode that locks it.
450+
"""
451+
if filename is None:
452+
return False
453+
import h5py
454+
455+
try:
456+
with h5py.File(filename, "r"):
457+
pass
458+
except Exception: # noqa: BLE001
459+
return False
460+
return True
461+
462+
445463
def from_omx_3d(
446464
omx: openmatrix.File | str | Iterable[openmatrix.File | str],
447465
index_names=("otaz", "dtaz", "time_period"),
@@ -538,18 +556,20 @@ def from_omx_3d(
538556
omx_data_map[k] = n
539557

540558
omx_filenames = [omx_file_name(i) for i in omx]
559+
omx_reopenable = [_is_reopenable(i) for i in omx_filenames]
541560

542561
import dask.array
543562

544563
def _lazy_omx_array(k):
545564
# Build a lazy dask array for one matrix table. When the source
546-
# file is on disk, defer to the parallel chunk-decoding reader;
547-
# otherwise fall back to wrapping the open file handle's node.
565+
# file can be independently reopened, defer to the parallel
566+
# chunk-decoding reader; otherwise read the data eagerly, as the
567+
# open file handle may be closed before the dask graph is computed.
548568
n = omx_data_map[k]
549569
filename = omx_filenames[n]
550570
node = omx_data[n][k]
551-
if filename is None:
552-
return dask.array.from_array(node)
571+
if not omx_reopenable[n]:
572+
return dask.array.from_array(np.asarray(node[:]))
553573
return dask.array.from_delayed(
554574
dask.delayed(_fast_load_omx_array)(filename, k),
555575
shape=tuple(node.shape),

‎sharrow/tests/test_datasets.py‎

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -451,3 +451,35 @@ def test_from_omx_compressed_blosc():
451451
with h5py.File(f, mode="r") as back:
452452
ds = sh.dataset.from_omx(back, indexes="taz")
453453
np.testing.assert_array_equal(ds["DIST"].values, arr)
454+
455+
456+
def test_from_omx_3d_to_zarr():
457+
"""Lazily loaded 3d skims remain readable when writing to zarr."""
458+
matrices = _random_matrices()
459+
with tempfile.TemporaryDirectory() as tempdir:
460+
f = Path(tempdir).joinpath("skims.omx")
461+
_write_compressed_omx(f, matrices)
462+
skims = sh.dataset.from_omx_3d(str(f), time_periods=["AM", "PM"])
463+
zarr_path = Path(tempdir).joinpath("skims.zarr")
464+
skims[["TIME"]].to_zarr(zarr_path, mode="w")
465+
back = xr.open_zarr(zarr_path)
466+
np.testing.assert_array_equal(
467+
back["TIME"].sel(time_period="AM").values, matrices["TIME__AM"]
468+
)
469+
470+
471+
def test_from_omx_3d_writable_handle():
472+
"""A file handle open for writing does not block lazy loading."""
473+
matrices = _random_matrices()
474+
with tempfile.TemporaryDirectory() as tempdir:
475+
f = Path(tempdir).joinpath("skims.omx")
476+
_write_compressed_omx(f, matrices)
477+
with openmatrix.open_file(f, mode="a") as back:
478+
skims = sh.dataset.from_omx_3d(
479+
back, time_periods=["AM", "PM"], max_float_precision=64
480+
)
481+
computed = skims.compute()
482+
np.testing.assert_array_equal(computed["DIST"].values, matrices["DIST"])
483+
np.testing.assert_array_equal(
484+
computed["TIME"].sel(time_period="PM").values, matrices["TIME__PM"]
485+
)

0 commit comments

Comments
 (0)