Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 44 additions & 1 deletion echopype/tests/echodata/test_echodata.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,10 @@
import fsspec
from pathlib import Path
import shutil
import warnings

from xarray import DataTree
from zarr.errors import GroupNotFoundError
from zarr.errors import GroupNotFoundError, ZarrUserWarning

import echopype
from echopype.echodata import EchoData
Expand Down Expand Up @@ -850,3 +851,45 @@ def test_echodata_delete(caplog, ek60_path):

# Check that it doesn't exist
assert not os.path.exists(temp_zarr_path)


@pytest.mark.unit
def test_to_zarr_no_consolidated_metadata_warning(test_path):
"""
This test checks that we explicitly use `consolidated=False`, otherwise
xarray will attempt to write consolidated metadata but Zarr v3 does not
support this yet.
TODO remove this test once Zarr v3 supports consolidated metadata.
"""
with warnings.catch_warnings(record=True) as caught_warnings:
warnings.simplefilter("always", ZarrUserWarning)
ed = open_raw(
test_path["EK60"] / "ncei-wcsd/SH1701/TEST-D20170114-T202932.raw",
sonar_model="EK60",
)
ed.to_zarr(overwrite=True)
assert not any(
isinstance(w.message, ZarrUserWarning)
and "Consolidated metadata is currently not part in the Zarr format 3 specification"
in str(w.message)
for w in caught_warnings
)


@pytest.mark.unit
def test_to_zarr_no_deprecated_chunk_warning(test_path):
"""
This test checks that dask rechunking receives dictionary and not tuple.
"""
with warnings.catch_warnings(record=True) as caught_warnings:
warnings.simplefilter("always")
ed = open_raw(
test_path["EK60"] / "ncei-wcsd/SH1701/TEST-D20170114-T202932.raw",
sonar_model="EK60",
)
ed.to_zarr(overwrite=True)
assert not any(
isinstance(w.message, FutureWarning)
and "Supplying chunks as dimension-order tuples" in str(w.message)
for w in caught_warnings
)
17 changes: 11 additions & 6 deletions echopype/utils/coding.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@

import numpy as np
import xarray as xr
from dask.array import Array as DaskArray
from dask.array.core import auto_chunks
from dask.utils import parse_bytes
from xarray import coding
Expand Down Expand Up @@ -217,17 +218,21 @@ def set_zarr_encodings(
# tolerance then no need to rechunk
if chunk_diff < chunk_size_tolerance:
rechunk = False
chunks = existing_chunks
dask_formatted_chunks = existing_chunks

if rechunk:
# Use dask auto chunk to determine the optimal chunk
# spread for optimal chunk size
chunks = _get_dask_auto_chunk(val, chunk_size=chunk_size)
# Dask expects chunk dictionary but Zarr expects list-like iterable of
# values in encoding
chunks = [*chunks.values()]
dask_formatted_chunks = _get_dask_auto_chunk(val, chunk_size=chunk_size)

encoding[name]["chunks"] = chunks
# Dask expects chunk dictionary but Zarr expects list-like iterable of
# values in encoding
zarr_formatted_chunks = [*dask_formatted_chunks.values()]

encoding[name]["chunks"] = zarr_formatted_chunks
# Add temporary dask formatted_chunks to encoding for use in dask rechunking
if isinstance(ds[name].data, DaskArray):
encoding[name]["dask_formatted_chunks"] = dask_formatted_chunks
if PREFERRED_CHUNKS in encoding[name]:
# Remove 'preferred_chunks', use chunks only instead
encoding[name].pop(PREFERRED_CHUNKS)
Expand Down
17 changes: 14 additions & 3 deletions echopype/utils/io.py
Original file line number Diff line number Diff line change
Expand Up @@ -72,11 +72,22 @@ def save_file(ds, path, mode, engine, group=None, compression_settings=None, **k
if engine == "netcdf4":
ds.to_netcdf(path=path, mode=mode, group=group, encoding=encoding, engine=engine, **kwargs)
elif engine == "zarr":
# Ensure that encoding and chunks match
# Rechunk Dask array according to dask formatted chunk encoding. The zarr chunks are based
# on the dask formatted chunks, so we avoid any potential chunking mismatch between them.
# If we do not do this, we will need to call `ds.to_zrar(...,align_chunks=True,...)`, which
# also "rechunks the Dask array to align with Zarr chunks before writing"
# (https://docs.xarray.dev/en/latest/generated/xarray.Dataset.to_zarr.html).
for var, enc in encoding.items():
if isinstance(ds[var].data, DaskArray):
ds[var] = ds[var].chunk(enc.get("chunks", {}))
ds.to_zarr(store=path.root, mode=mode, group=group, encoding=encoding, **kwargs)
ds[var] = ds[var].chunk(enc.get("dask_formatted_chunks", {}))
# # Remove dask_formatted_chunks from encoding
encoding[var].pop("dask_formatted_chunks")
# Set consolidated to False to avoid warning about consolidated metadata not being supported
# by Zarr 3
# TODO remove consolidated=False when Zarr 3 supports consolidated metadata
ds.to_zarr(
store=path.root, mode=mode, group=group, consolidated=False, encoding=encoding, **kwargs
)
else:
raise ValueError(f"{engine} is not a supported save format")

Expand Down
Loading