From cd3c4f0e60217c1cd6992f880da8707dcbe584b6 Mon Sep 17 00:00:00 2001 From: Ben Dichter Date: Wed, 30 Sep 2026 01:25:29 -0400 Subject: [PATCH 1/2] feat: read only the needed rows of uncompressed chunks BytesCodec now implements partial decoding. An uncompressed chunk is stored in C order, so each row along its first axis is a contiguous run of bytes; for a selection that does not touch every row, the rows from the first to the last one it touches are fetched with a single range request and the selection is applied to them. The pipeline already uses partial decoding when the array-to-bytes codec supports it and there are no array-to-array or bytes-to-bytes codecs, so this applies to uncompressed arrays only. Reading 100 rows of a single-chunk 20 MB array now fetches the bytes of those rows instead of the whole chunk. Co-Authored-By: Claude Opus 5.5 --- src/zarr/codecs/bytes.py | 83 +++++++++++++++++++++++++++++- tests/test_codecs/test_bytes.py | 90 +++++++++++++++++++++++++++++++++ 2 files changed, 171 insertions(+), 2 deletions(-) diff --git a/src/zarr/codecs/bytes.py b/src/zarr/codecs/bytes.py index 622161c458..8120db6785 100644 --- a/src/zarr/codecs/bytes.py +++ b/src/zarr/codecs/bytes.py @@ -5,7 +5,10 @@ from dataclasses import dataclass, replace from typing import TYPE_CHECKING, ClassVar, Final, Literal -from zarr.abc.codec import ArrayBytesCodec +import numpy as np + +from zarr.abc.codec import ArrayBytesCodec, ArrayBytesCodecPartialDecodeMixin +from zarr.abc.store import RangeByteRequest from zarr.codecs._deprecated_enum import _coerce_enum_input, _DeprecatedStrEnumMeta from zarr.core.common import JSON, parse_named_configuration from zarr.core.dtype.common import HasEndianness @@ -14,8 +17,10 @@ if TYPE_CHECKING: from typing import Self + from zarr.abc.store import ByteGetter from zarr.core.array_spec import ArraySpec from zarr.core.buffer import Buffer, NDBuffer + from zarr.core.indexing import SelectorTuple EndianLiteral = Literal["little", "big"] @@ -40,7 +45,7 @@ def _parse_endian(data: object) -> EndianLiteral: @dataclass(frozen=True) -class BytesCodec(ArrayBytesCodec): +class BytesCodec(ArrayBytesCodec, ArrayBytesCodecPartialDecodeMixin): """bytes codec""" is_fixed_size = True @@ -137,6 +142,38 @@ async def _decode_single( ) -> NDBuffer: return self._decode_sync(chunk_bytes, chunk_spec) + async def _decode_partial_single( + self, + byte_getter: ByteGetter, + selection: SelectorTuple, + chunk_spec: ArraySpec, + ) -> NDBuffer | None: + """Read only the part of an uncompressed chunk that a selection needs. + + The chunk is stored in C order, so each row along its first axis is a + contiguous run of bytes. The rows from the first to the last one the + selection touches are fetched with a single range request, and the + selection is applied to them. A selection that touches every row reads + the whole chunk, as before. + """ + rows = _selected_rows(selection, chunk_spec.shape) + if rows is None or rows == (0, chunk_spec.shape[0]): + chunk_bytes = await byte_getter.get(prototype=chunk_spec.prototype) + if chunk_bytes is None: + return None + return self._decode_sync(chunk_bytes, chunk_spec)[selection] + first, stop = rows + row_items = int(np.prod(chunk_spec.shape[1:])) + row_bytes = chunk_spec.dtype.to_native_dtype().itemsize * row_items + chunk_bytes = await byte_getter.get( + prototype=chunk_spec.prototype, + byte_range=RangeByteRequest(first * row_bytes, stop * row_bytes), + ) + if chunk_bytes is None: + return None + rows_spec = replace(chunk_spec, shape=(stop - first, *chunk_spec.shape[1:])) + return self._decode_sync(chunk_bytes, rows_spec)[_shift_rows(selection, first)] + def _encode_sync( self, chunk_array: NDBuffer, @@ -165,3 +202,45 @@ async def _encode_single( def compute_encoded_size(self, input_byte_length: int, _chunk_spec: ArraySpec) -> int: return input_byte_length + + +def _selected_rows(selection: SelectorTuple, shape: tuple[int, ...]) -> tuple[int, int] | None: + """The first row along axis 0 that a selection touches and one past the last. + + Returns None when the rows cannot be determined, in which case the whole + chunk is read. + """ + if len(shape) == 0 or not isinstance(selection, tuple) or len(selection) == 0: + return None + first_axis = selection[0] + if isinstance(first_axis, slice): + rows = range(*first_axis.indices(shape[0])) + if len(rows) == 0: + return None + return min(rows[0], rows[-1]), max(rows[0], rows[-1]) + 1 + if isinstance(first_axis, int | np.integer): + row = int(first_axis) % shape[0] + return row, row + 1 + if isinstance(first_axis, np.ndarray): + indices = np.nonzero(first_axis)[0] if first_axis.dtype == bool else first_axis + if indices.size == 0: + return None + return int(indices.min()) % shape[0], int(indices.max()) % shape[0] + 1 + return None + + +def _shift_rows(selection: SelectorTuple, first: int) -> SelectorTuple: + """The selection relative to the rows fetched, which start at row `first`.""" + first_axis = selection[0] + if isinstance(first_axis, slice): + start, stop, step = first_axis.start, first_axis.stop, first_axis.step + shifted: object = slice( + None if start is None else start - first, + None if stop is None else stop - first, + step, + ) + elif isinstance(first_axis, np.ndarray) and first_axis.dtype == bool: + shifted = first_axis[first:] + else: + shifted = first_axis - first + return (shifted, *selection[1:]) diff --git a/tests/test_codecs/test_bytes.py b/tests/test_codecs/test_bytes.py index ead778f526..925436d583 100644 --- a/tests/test_codecs/test_bytes.py +++ b/tests/test_codecs/test_bytes.py @@ -338,3 +338,93 @@ def test_bytes_codec_evolve_structured_single_byte_fields_clears_endian() -> Non spec = _make_array_spec(dtype) evolved = codec.evolve_from_array_spec(spec) assert evolved.endian is None + + +class _RangeLoggingStore(zarr.storage.MemoryStore): + """A memory store that records the byte ranges it serves for chunk keys.""" + + def __init__(self) -> None: + super().__init__() + self.reads: list[tuple[str, Any, int]] = [] + + async def get(self, key: str, prototype: Any = None, byte_range: Any = None) -> Any: + buf = await super().get(key, prototype, byte_range) + if buf is not None and not key.endswith("zarr.json"): + self.reads.append((key, byte_range, len(buf))) + return buf + + +def _single_chunk_array( + data: np.ndarray[Any, Any], **kwargs: Any +) -> tuple[Any, _RangeLoggingStore]: + store = _RangeLoggingStore() + arr = zarr.create_array( + store, shape=data.shape, dtype=data.dtype, chunks=data.shape, fill_value=0, **kwargs + ) + arr[...] = data + store.reads.clear() + return arr, store + + +@pytest.mark.parametrize( + ("selection", "rows_read"), + [ + (np.s_[500:600, 3], 100), + (np.s_[500:600], 100), + (np.s_[777, 5], 1), + (np.s_[-1], 1), + (np.s_[10:20:3, ::2], 10), # rows 10 to 19 are fetched + (np.s_[...], 10_000), + ], +) +def test_uncompressed_partial_read(selection: Any, rows_read: int) -> None: + """Reading part of an uncompressed chunk fetches only the rows it touches.""" + data = np.arange(100_000, dtype="int16").reshape(10_000, 10) + arr, store = _single_chunk_array(data, compressors=None) + np.testing.assert_array_equal(arr[selection], data[selection]) + assert sum(n for *_, n in store.reads) == rows_read * 10 * 2 + + +@pytest.mark.parametrize("dtype", [">u2", " None: + """Partial reads return the same values as numpy for several dtypes, byte orders, and shapes.""" + shape = {1: (1000,), 2: (100, 7), 3: (20, 5, 3)}[ndim] + data = (np.arange(int(np.prod(shape))) % 200).astype(dtype).reshape(shape) + arr, _ = _single_chunk_array(data, compressors=None) + for selection in [np.s_[3:9], np.s_[4], np.s_[-3:], np.s_[1:15:4]]: + np.testing.assert_array_equal(arr[selection], data[selection]) + np.testing.assert_array_equal(arr.oindex[[1, 5, 2]], data[[1, 5, 2]]) + coords = tuple(np.array([0, 6, 3]) % n for n in shape) + np.testing.assert_array_equal(arr.vindex[coords], data[coords]) + + +def test_uncompressed_partial_read_across_chunks() -> None: + """A selection spanning several chunks reads only the touched rows of each.""" + data = np.arange(40_000, dtype="int32").reshape(4000, 10) + store = _RangeLoggingStore() + arr = zarr.create_array( + store, shape=data.shape, dtype="int32", chunks=(1000, 10), compressors=None + ) + arr[...] = data + store.reads.clear() + np.testing.assert_array_equal(arr[990:1010, 2], data[990:1010, 2]) + assert sorted((key, n) for key, _, n in store.reads) == [("c/0/0", 400), ("c/1/0", 400)] + + +def test_compressed_chunks_are_read_whole() -> None: + """With a compressor the chunk bytes cannot be split, so the whole chunk is fetched.""" + data = np.arange(100_000, dtype="int16").reshape(10_000, 10) + arr, store = _single_chunk_array(data) # default compressor + np.testing.assert_array_equal(arr[500:600, 3], data[500:600, 3]) + ((_, byte_range, _),) = store.reads + assert byte_range is None + + +def test_uncompressed_partial_read_missing_chunk() -> None: + """An unwritten chunk reads as the fill value without fetching anything.""" + store = _RangeLoggingStore() + arr = zarr.create_array( + store, shape=(100, 4), dtype="int16", chunks=(100, 4), compressors=None, fill_value=7 + ) + np.testing.assert_array_equal(arr[10:20], np.full((10, 4), 7, dtype="int16")) From 1222e64bed404166732ee439a3ba9a915bc9b37c Mon Sep 17 00:00:00 2001 From: Ben Dichter Date: Wed, 30 Sep 2026 08:30:37 -0400 Subject: [PATCH 2/2] Fix boolean masks and stores that ignore ranges in BytesCodec partial reads A boolean mask on the first axis, as from arr.oindex[mask], was trimmed to start at the first selected row but not to end at the last, so its length no longer matched the rows fetched and numpy raised an IndexError. The row window and the shifted selection are now computed together in _row_window, which trims the mask at both ends. This also gives mypy the narrowing it needs. FsspecStore over HTTP returns the whole object when a server ignores the Range header. The partial read then failed to reshape a chunk that the full read would have decoded, so a read that worked before this branch raised an error. A response longer than the requested range is now sliced to it. Co-Authored-By: Claude Opus 5.5 --- src/zarr/codecs/bytes.py | 67 ++++++++++++++++----------------- tests/test_codecs/test_bytes.py | 25 +++++++++++- 2 files changed, 57 insertions(+), 35 deletions(-) diff --git a/src/zarr/codecs/bytes.py b/src/zarr/codecs/bytes.py index 8120db6785..5a57179e1e 100644 --- a/src/zarr/codecs/bytes.py +++ b/src/zarr/codecs/bytes.py @@ -20,7 +20,7 @@ from zarr.abc.store import ByteGetter from zarr.core.array_spec import ArraySpec from zarr.core.buffer import Buffer, NDBuffer - from zarr.core.indexing import SelectorTuple + from zarr.core.indexing import Selector, SelectorTuple EndianLiteral = Literal["little", "big"] @@ -156,13 +156,13 @@ async def _decode_partial_single( selection is applied to them. A selection that touches every row reads the whole chunk, as before. """ - rows = _selected_rows(selection, chunk_spec.shape) - if rows is None or rows == (0, chunk_spec.shape[0]): + window = _row_window(selection, chunk_spec.shape) + if window is None or window[:2] == (0, chunk_spec.shape[0]): chunk_bytes = await byte_getter.get(prototype=chunk_spec.prototype) if chunk_bytes is None: return None return self._decode_sync(chunk_bytes, chunk_spec)[selection] - first, stop = rows + first, stop, rows_selection = window row_items = int(np.prod(chunk_spec.shape[1:])) row_bytes = chunk_spec.dtype.to_native_dtype().itemsize * row_items chunk_bytes = await byte_getter.get( @@ -171,8 +171,12 @@ async def _decode_partial_single( ) if chunk_bytes is None: return None + if len(chunk_bytes) > (stop - first) * row_bytes: + # The store sent the whole chunk, as an HTTP server that ignores + # the Range header does. + chunk_bytes = chunk_bytes[first * row_bytes : stop * row_bytes] rows_spec = replace(chunk_spec, shape=(stop - first, *chunk_spec.shape[1:])) - return self._decode_sync(chunk_bytes, rows_spec)[_shift_rows(selection, first)] + return self._decode_sync(chunk_bytes, rows_spec)[rows_selection] def _encode_sync( self, @@ -204,43 +208,38 @@ def compute_encoded_size(self, input_byte_length: int, _chunk_spec: ArraySpec) - return input_byte_length -def _selected_rows(selection: SelectorTuple, shape: tuple[int, ...]) -> tuple[int, int] | None: - """The first row along axis 0 that a selection touches and one past the last. +def _row_window( + selection: SelectorTuple, shape: tuple[int, ...] +) -> tuple[int, int, SelectorTuple] | None: + """The rows along axis 0 that a selection touches, and the selection relative to them. - Returns None when the rows cannot be determined, in which case the whole - chunk is read. + Returns the first row, one past the last row, and the selection shifted so + that it indexes an array holding only those rows. Returns None when the rows + cannot be determined, in which case the whole chunk is read. """ if len(shape) == 0 or not isinstance(selection, tuple) or len(selection) == 0: return None first_axis = selection[0] + shifted: Selector if isinstance(first_axis, slice): rows = range(*first_axis.indices(shape[0])) - if len(rows) == 0: + if len(rows) == 0 or rows.step < 0: return None - return min(rows[0], rows[-1]), max(rows[0], rows[-1]) + 1 - if isinstance(first_axis, int | np.integer): - row = int(first_axis) % shape[0] - return row, row + 1 - if isinstance(first_axis, np.ndarray): - indices = np.nonzero(first_axis)[0] if first_axis.dtype == bool else first_axis + first, stop = rows[0], rows[-1] + 1 + shifted = slice(0, stop - first, rows.step) + elif isinstance(first_axis, int | np.integer): + first = int(first_axis) % shape[0] + stop = first + 1 + shifted = 0 + elif isinstance(first_axis, np.ndarray): + if first_axis.dtype == bool: + indices = np.nonzero(first_axis)[0] + else: + indices = first_axis % shape[0] if indices.size == 0: return None - return int(indices.min()) % shape[0], int(indices.max()) % shape[0] + 1 - return None - - -def _shift_rows(selection: SelectorTuple, first: int) -> SelectorTuple: - """The selection relative to the rows fetched, which start at row `first`.""" - first_axis = selection[0] - if isinstance(first_axis, slice): - start, stop, step = first_axis.start, first_axis.stop, first_axis.step - shifted: object = slice( - None if start is None else start - first, - None if stop is None else stop - first, - step, - ) - elif isinstance(first_axis, np.ndarray) and first_axis.dtype == bool: - shifted = first_axis[first:] + first, stop = int(indices.min()), int(indices.max()) + 1 + shifted = first_axis[first:stop] if first_axis.dtype == bool else indices - first else: - shifted = first_axis - first - return (shifted, *selection[1:]) + return None + return first, stop, (shifted, *selection[1:]) diff --git a/tests/test_codecs/test_bytes.py b/tests/test_codecs/test_bytes.py index 925436d583..bee634dde4 100644 --- a/tests/test_codecs/test_bytes.py +++ b/tests/test_codecs/test_bytes.py @@ -392,9 +392,13 @@ def test_uncompressed_partial_read_values(dtype: str, ndim: int) -> None: shape = {1: (1000,), 2: (100, 7), 3: (20, 5, 3)}[ndim] data = (np.arange(int(np.prod(shape))) % 200).astype(dtype).reshape(shape) arr, _ = _single_chunk_array(data, compressors=None) - for selection in [np.s_[3:9], np.s_[4], np.s_[-3:], np.s_[1:15:4]]: + selections: list[Any] = [np.s_[3:9], np.s_[4], np.s_[-3:], np.s_[1:15:4]] + for selection in selections: np.testing.assert_array_equal(arr[selection], data[selection]) np.testing.assert_array_equal(arr.oindex[[1, 5, 2]], data[[1, 5, 2]]) + mask = np.zeros(shape[0], dtype=bool) + mask[[2, 5, 11]] = True + np.testing.assert_array_equal(arr.oindex[mask], data[mask]) coords = tuple(np.array([0, 6, 3]) % n for n in shape) np.testing.assert_array_equal(arr.vindex[coords], data[coords]) @@ -428,3 +432,22 @@ def test_uncompressed_partial_read_missing_chunk() -> None: store, shape=(100, 4), dtype="int16", chunks=(100, 4), compressors=None, fill_value=7 ) np.testing.assert_array_equal(arr[10:20], np.full((10, 4), 7, dtype="int16")) + + +class _IgnoresRangeStore(zarr.storage.MemoryStore): + """A memory store that sends the whole value for every read, as an HTTP + server that ignores the Range header does.""" + + async def get(self, key: str, prototype: Any = None, byte_range: Any = None) -> Any: + return await super().get(key, prototype) + + +def test_uncompressed_partial_read_store_ignores_range() -> None: + data = np.arange(400, dtype="