Merged
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
12 changes: 12 additions & 0 deletions docs/release.rst
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,18 @@ This release od Zarr Python is is the first release of Zarr to not supporting Py
See `this link <https://github.com/zarr-developers/zarr-python/milestone/11?closed=1>` for the full list of closed and
merged PR tagged with the 2.6 milestone.

* Add ability to partially read and decompress arrays, see :issue:`667`. It is
only available to chunks stored using fs-spec and using bloc as a compressor.

For certain analysis case when only a small portion of chunks is needed it can
be advantageous to only access and decompress part of the chunks. Doing
partial read and decompression add high latency to many of the operation so
should be used only when the subset of the data is small compared to the full
chunks and is stored contiguously (that is to say either last dimensions for C
layout, firsts for F). Pass ``partial_decompress=True`` as argument when
creating an ``Array``, or when using ``open_array``. No option exists yet to
apply partial read and decompress on a per-operation basis.

2.5.0
-----

Expand Down
2 changes: 1 addition & 1 deletion requirements_dev_minimal.txt
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
# library requirements
asciitree==0.3.3
fasteners==0.15
numcodecs==0.6.4
numcodecs==0.7.2
msgpack-python==0.5.6
setuptools-scm==3.3.3
# test requirements
Expand Down
152 changes: 128 additions & 24 deletions zarr/core.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,22 +11,41 @@

from zarr.attrs import Attributes
from zarr.codecs import AsType, get_codec
from zarr.errors import ArrayNotFoundError, ReadOnlyError
from zarr.indexing import (BasicIndexer, CoordinateIndexer, MaskIndexer,
OIndex, OrthogonalIndexer, VIndex, check_fields,
check_no_multi_fields, ensure_tuple,
err_too_many_indices, is_contiguous_selection,
is_scalar, pop_fields)
from zarr.errors import ArrayNotFoundError, ReadOnlyError, ArrayIndexError
from zarr.indexing import (
BasicIndexer,
CoordinateIndexer,
MaskIndexer,
OIndex,
OrthogonalIndexer,
VIndex,
PartialChunkIterator,
check_fields,
check_no_multi_fields,
ensure_tuple,
err_too_many_indices,
is_contiguous_selection,
is_scalar,
pop_fields,
)
from zarr.meta import decode_array_metadata, encode_array_metadata
from zarr.storage import array_meta_key, attrs_key, getsize, listdir
from zarr.util import (InfoReporter, check_array_shape, human_readable_size,
is_total_slice, nolock, normalize_chunks,
normalize_resize_args, normalize_shape,
normalize_storage_path)
from zarr.util import (
InfoReporter,
check_array_shape,
human_readable_size,
is_total_slice,
nolock,
normalize_chunks,
normalize_resize_args,
normalize_shape,
normalize_storage_path,
PartialReadBuffer,
)


# noinspection PyUnresolvedReferences
class Array(object):
class Array:
"""Instantiate an array from an initialized store.

Parameters
Expand All@@ -51,6 +70,12 @@ class Array(object):
If True (default), user attributes will be cached for attribute read
operations. If False, user attributes are reloaded from the store prior
to all attribute read operations.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Attributes
----------
Expand DownExpand Up@@ -102,8 +127,17 @@ class Array(object):

"""

def __init__(self, store, path=None, read_only=False, chunk_store=None,
synchronizer=None, cache_metadata=True, cache_attrs=True):
def __init__(
self,
store,
path=None,
read_only=False,
chunk_store=None,
synchronizer=None,
cache_metadata=True,
cache_attrs=True,
partial_decompress=False,
):
# N.B., expect at this point store is fully initialized with all
# configuration metadata fully specified and normalized

Expand All@@ -118,6 +152,7 @@ def __init__(self, store, path=None, read_only=False, chunk_store=None,
self._synchronizer = synchronizer
self._cache_metadata = cache_metadata
self._is_view = False
self._partial_decompress = partial_decompress

# initialize metadata
self._load_metadata()
Expand DownExpand Up@@ -1580,8 +1615,17 @@ def _set_selection(self, indexer, value, fields=None):
self._chunk_setitems(lchunk_coords, lchunk_selection, chunk_values,
fields=fields)

def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
out_is_ndarray, fields, out_selection):
def _process_chunk(
self,
out,
cdata,
chunk_selection,
drop_axes,
out_is_ndarray,
fields,
out_selection,
partial_read_decode=False,
):
"""Take binary data from storage and fill output array"""
if (out_is_ndarray and
not fields and
Expand All@@ -1604,8 +1648,9 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
# optimization: we want the whole chunk, and the destination is
# contiguous, so we can decompress directly from the chunk
# into the destination array

if self._compressor:
if isinstance(cdata, PartialReadBuffer):
cdata = cdata.read_full()
self._compressor.decode(cdata, dest)
else:
chunk = ensure_ndarray(cdata).view(self._dtype)
Expand All@@ -1614,6 +1659,33 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
return

# decode chunk
try:
if partial_read_decode:
cdata.prepare_chunk()
# size of chunk
tmp = np.empty(self._chunks, dtype=self.dtype)
index_selection = PartialChunkIterator(chunk_selection, self.chunks)
for start, nitems, partial_out_selection in index_selection:
expected_shape = [
len(
range(*partial_out_selection[i].indices(self.chunks[0] + 1))
)
if i < len(partial_out_selection)
else dim
for i, dim in enumerate(self.chunks)
]
cdata.read_part(start, nitems)
chunk_partial = self._decode_chunk(
cdata.buff,
start=start,
nitems=nitems,
expected_shape=expected_shape,
)
tmp[partial_out_selection] = chunk_partial
out[out_selection] = tmp[chunk_selection]
return
except ArrayIndexError:
cdata = cdata.read_full()
chunk = self._decode_chunk(cdata)

# select data from chunk
Expand DownExpand Up@@ -1688,11 +1760,36 @@ def _chunk_getitems(self, lchunk_coords, lchunk_selection, out, lout_selection,
out_is_ndarray = False

ckeys = [self._chunk_key(ch) for ch in lchunk_coords]
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
if (
self._partial_decompress
and self._compressor
and self._compressor.codec_id == "blosc"
and hasattr(self._compressor, "decode_partial")
and not fields
and self.dtype != object
and hasattr(self.chunk_store, "getitems")
):
partial_read_decode = True
cdatas = {
ckey: PartialReadBuffer(ckey, self.chunk_store)
for ckey in ckeys
if ckey in self.chunk_store
}
else:
partial_read_decode = False
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
for ckey, chunk_select, out_select in zip(ckeys, lchunk_selection, lout_selection):
if ckey in cdatas:
self._process_chunk(out, cdatas[ckey], chunk_select, drop_axes,
out_is_ndarray, fields, out_select)
self._process_chunk(
out,
cdatas[ckey],
chunk_select,
drop_axes,
out_is_ndarray,
fields,
out_select,
partial_read_decode=partial_read_decode,
)
else:
# check exception type
if self._fill_value is not None:
Expand All@@ -1706,7 +1803,8 @@ def _chunk_setitems(self, lchunk_coords, lchunk_selection, values, fields=None):
ckeys = [self._chunk_key(co) for co in lchunk_coords]
cdatas = [self._process_for_setitem(key, sel, val, fields=fields)
for key, sel, val in zip(ckeys, lchunk_selection, values)]
self.chunk_store.setitems({k: v for k, v in zip(ckeys, cdatas)})
values = {k: v for k, v in zip(ckeys, cdatas)}
self.chunk_store.setitems(values)

def _chunk_setitem(self, chunk_coords, chunk_selection, value, fields=None):
"""Replace part or whole of a chunk.
Expand DownExpand Up@@ -1800,11 +1898,17 @@ def _process_for_setitem(self, ckey, chunk_selection, value, fields=None):
def _chunk_key(self, chunk_coords):
return self._key_prefix + '.'.join(map(str, chunk_coords))

def _decode_chunk(self, cdata):

def _decode_chunk(self, cdata, start=None, nitems=None, expected_shape=None):
# decompress
if self._compressor:
chunk = self._compressor.decode(cdata)
# only decode requested items
if (
all([x is not None for x in [start, nitems]])
and self._compressor.codec_id == "blosc"
) and hasattr(self._compressor, "decode_partial"):
chunk = self._compressor.decode_partial(cdata, start, nitems)
else:
chunk = self._compressor.decode(cdata)
else:
chunk = cdata

Expand All@@ -1829,7 +1933,7 @@ def _decode_chunk(self, cdata):

# ensure correct chunk shape
chunk = chunk.reshape(-1, order='A')
chunk = chunk.reshape(self._chunks, order=self._order)
chunk = chunk.reshape(expected_shape or self._chunks, order=self._order)

return chunk

Expand Down
31 changes: 26 additions & 5 deletions zarr/creation.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -362,11 +362,26 @@ def array(data, **kwargs):
return z


def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
compressor='default', fill_value=0, order='C', synchronizer=None,
filters=None, cache_metadata=True, cache_attrs=True, path=None,
object_codec=None, chunk_store=None, storage_options=None,
**kwargs):
def open_array(
store=None,
mode="a",
shape=None,
chunks=True,
dtype=None,
compressor="default",
fill_value=0,
order="C",
synchronizer=None,
filters=None,
cache_metadata=True,
cache_attrs=True,
path=None,
object_codec=None,
chunk_store=None,
storage_options=None,
partial_decompress=False,
**kwargs
):
"""Open an array using file-mode-like semantics.

Parameters
Expand DownExpand Up@@ -415,6 +430,12 @@ def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
storage_options : dict
If using an fsspec URL to create the store, these will be passed to
the backend implementation. Ignored otherwise.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Returns
-------
Expand Down
4 changes: 4 additions & 0 deletions zarr/errors.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,6 +15,10 @@ def __init__(self, *args):
super().__init__(self._msg.format(*args))


class ArrayIndexError(IndexError):
pass


class _BaseZarrIndexError(IndexError):
_msg = ""

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
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
12 changes: 12 additions & 0 deletions docs/release.rst
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,18 @@ This release od Zarr Python is is the first release of Zarr to not supporting Py
See `this link <https://github.com/zarr-developers/zarr-python/milestone/11?closed=1>` for the full list of closed and
merged PR tagged with the 2.6 milestone.

* Add ability to partially read and decompress arrays, see :issue:`667`. It is
only available to chunks stored using fs-spec and using bloc as a compressor.

For certain analysis case when only a small portion of chunks is needed it can
be advantageous to only access and decompress part of the chunks. Doing
partial read and decompression add high latency to many of the operation so
should be used only when the subset of the data is small compared to the full
chunks and is stored contiguously (that is to say either last dimensions for C
layout, firsts for F). Pass ``partial_decompress=True`` as argument when
creating an ``Array``, or when using ``open_array``. No option exists yet to
apply partial read and decompress on a per-operation basis.

2.5.0
-----

Expand Down
2 changes: 1 addition & 1 deletion requirements_dev_minimal.txt
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
# library requirements
asciitree==0.3.3
fasteners==0.15
numcodecs==0.6.4
numcodecs==0.7.2
msgpack-python==0.5.6
setuptools-scm==3.3.3
# test requirements
Expand Down
152 changes: 128 additions & 24 deletions zarr/core.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,22 +11,41 @@

from zarr.attrs import Attributes
from zarr.codecs import AsType, get_codec
from zarr.errors import ArrayNotFoundError, ReadOnlyError
from zarr.indexing import (BasicIndexer, CoordinateIndexer, MaskIndexer,
OIndex, OrthogonalIndexer, VIndex, check_fields,
check_no_multi_fields, ensure_tuple,
err_too_many_indices, is_contiguous_selection,
is_scalar, pop_fields)
from zarr.errors import ArrayNotFoundError, ReadOnlyError, ArrayIndexError
from zarr.indexing import (
BasicIndexer,
CoordinateIndexer,
MaskIndexer,
OIndex,
OrthogonalIndexer,
VIndex,
PartialChunkIterator,
check_fields,
check_no_multi_fields,
ensure_tuple,
err_too_many_indices,
is_contiguous_selection,
is_scalar,
pop_fields,
)
from zarr.meta import decode_array_metadata, encode_array_metadata
from zarr.storage import array_meta_key, attrs_key, getsize, listdir
from zarr.util import (InfoReporter, check_array_shape, human_readable_size,
is_total_slice, nolock, normalize_chunks,
normalize_resize_args, normalize_shape,
normalize_storage_path)
from zarr.util import (
InfoReporter,
check_array_shape,
human_readable_size,
is_total_slice,
nolock,
normalize_chunks,
normalize_resize_args,
normalize_shape,
normalize_storage_path,
PartialReadBuffer,
)


# noinspection PyUnresolvedReferences
class Array(object):
class Array:
"""Instantiate an array from an initialized store.

Parameters
Expand All@@ -51,6 +70,12 @@ class Array(object):
If True (default), user attributes will be cached for attribute read
operations. If False, user attributes are reloaded from the store prior
to all attribute read operations.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Attributes
----------
Expand DownExpand Up@@ -102,8 +127,17 @@ class Array(object):

"""

def __init__(self, store, path=None, read_only=False, chunk_store=None,
synchronizer=None, cache_metadata=True, cache_attrs=True):
def __init__(
self,
store,
path=None,
read_only=False,
chunk_store=None,
synchronizer=None,
cache_metadata=True,
cache_attrs=True,
partial_decompress=False,
):
# N.B., expect at this point store is fully initialized with all
# configuration metadata fully specified and normalized

Expand All@@ -118,6 +152,7 @@ def __init__(self, store, path=None, read_only=False, chunk_store=None,
self._synchronizer = synchronizer
self._cache_metadata = cache_metadata
self._is_view = False
self._partial_decompress = partial_decompress

# initialize metadata
self._load_metadata()
Expand DownExpand Up@@ -1580,8 +1615,17 @@ def _set_selection(self, indexer, value, fields=None):
self._chunk_setitems(lchunk_coords, lchunk_selection, chunk_values,
fields=fields)

def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
out_is_ndarray, fields, out_selection):
def _process_chunk(
self,
out,
cdata,
chunk_selection,
drop_axes,
out_is_ndarray,
fields,
out_selection,
partial_read_decode=False,
):
"""Take binary data from storage and fill output array"""
if (out_is_ndarray and
not fields and
Expand All@@ -1604,8 +1648,9 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
# optimization: we want the whole chunk, and the destination is
# contiguous, so we can decompress directly from the chunk
# into the destination array

if self._compressor:
if isinstance(cdata, PartialReadBuffer):
cdata = cdata.read_full()
self._compressor.decode(cdata, dest)
else:
chunk = ensure_ndarray(cdata).view(self._dtype)
Expand All@@ -1614,6 +1659,33 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
return

# decode chunk
try:
if partial_read_decode:
cdata.prepare_chunk()
# size of chunk
tmp = np.empty(self._chunks, dtype=self.dtype)
index_selection = PartialChunkIterator(chunk_selection, self.chunks)
for start, nitems, partial_out_selection in index_selection:
expected_shape = [
len(
range(*partial_out_selection[i].indices(self.chunks[0] + 1))
)
if i < len(partial_out_selection)
else dim
for i, dim in enumerate(self.chunks)
]
cdata.read_part(start, nitems)
chunk_partial = self._decode_chunk(
cdata.buff,
start=start,
nitems=nitems,
expected_shape=expected_shape,
)
tmp[partial_out_selection] = chunk_partial
out[out_selection] = tmp[chunk_selection]
return
except ArrayIndexError:
cdata = cdata.read_full()
chunk = self._decode_chunk(cdata)

# select data from chunk
Expand DownExpand Up@@ -1688,11 +1760,36 @@ def _chunk_getitems(self, lchunk_coords, lchunk_selection, out, lout_selection,
out_is_ndarray = False

ckeys = [self._chunk_key(ch) for ch in lchunk_coords]
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
if (
self._partial_decompress
and self._compressor
and self._compressor.codec_id == "blosc"
and hasattr(self._compressor, "decode_partial")
and not fields
and self.dtype != object
and hasattr(self.chunk_store, "getitems")
):
partial_read_decode = True
cdatas = {
ckey: PartialReadBuffer(ckey, self.chunk_store)
for ckey in ckeys
if ckey in self.chunk_store
}
else:
partial_read_decode = False
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
for ckey, chunk_select, out_select in zip(ckeys, lchunk_selection, lout_selection):
if ckey in cdatas:
self._process_chunk(out, cdatas[ckey], chunk_select, drop_axes,
out_is_ndarray, fields, out_select)
self._process_chunk(
out,
cdatas[ckey],
chunk_select,
drop_axes,
out_is_ndarray,
fields,
out_select,
partial_read_decode=partial_read_decode,
)
else:
# check exception type
if self._fill_value is not None:
Expand All@@ -1706,7 +1803,8 @@ def _chunk_setitems(self, lchunk_coords, lchunk_selection, values, fields=None):
ckeys = [self._chunk_key(co) for co in lchunk_coords]
cdatas = [self._process_for_setitem(key, sel, val, fields=fields)
for key, sel, val in zip(ckeys, lchunk_selection, values)]
self.chunk_store.setitems({k: v for k, v in zip(ckeys, cdatas)})
values = {k: v for k, v in zip(ckeys, cdatas)}
self.chunk_store.setitems(values)

def _chunk_setitem(self, chunk_coords, chunk_selection, value, fields=None):
"""Replace part or whole of a chunk.
Expand DownExpand Up@@ -1800,11 +1898,17 @@ def _process_for_setitem(self, ckey, chunk_selection, value, fields=None):
def _chunk_key(self, chunk_coords):
return self._key_prefix + '.'.join(map(str, chunk_coords))

def _decode_chunk(self, cdata):

def _decode_chunk(self, cdata, start=None, nitems=None, expected_shape=None):
# decompress
if self._compressor:
chunk = self._compressor.decode(cdata)
# only decode requested items
if (
all([x is not None for x in [start, nitems]])
and self._compressor.codec_id == "blosc"
) and hasattr(self._compressor, "decode_partial"):
chunk = self._compressor.decode_partial(cdata, start, nitems)
else:
chunk = self._compressor.decode(cdata)
else:
chunk = cdata

Expand All@@ -1829,7 +1933,7 @@ def _decode_chunk(self, cdata):

# ensure correct chunk shape
chunk = chunk.reshape(-1, order='A')
chunk = chunk.reshape(self._chunks, order=self._order)
chunk = chunk.reshape(expected_shape or self._chunks, order=self._order)

return chunk

Expand Down
31 changes: 26 additions & 5 deletions zarr/creation.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -362,11 +362,26 @@ def array(data, **kwargs):
return z


def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
compressor='default', fill_value=0, order='C', synchronizer=None,
filters=None, cache_metadata=True, cache_attrs=True, path=None,
object_codec=None, chunk_store=None, storage_options=None,
**kwargs):
def open_array(
store=None,
mode="a",
shape=None,
chunks=True,
dtype=None,
compressor="default",
fill_value=0,
order="C",
synchronizer=None,
filters=None,
cache_metadata=True,
cache_attrs=True,
path=None,
object_codec=None,
chunk_store=None,
storage_options=None,
partial_decompress=False,
**kwargs
):
"""Open an array using file-mode-like semantics.

Parameters
Expand DownExpand Up@@ -415,6 +430,12 @@ def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
storage_options : dict
If using an fsspec URL to create the store, these will be passed to
the backend implementation. Ignored otherwise.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Returns
-------
Expand Down
4 changes: 4 additions & 0 deletions zarr/errors.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,6 +15,10 @@ def __init__(self, *args):
super().__init__(self._msg.format(*args))


class ArrayIndexError(IndexError):
pass


class _BaseZarrIndexError(IndexError):
_msg = ""

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
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
12 changes: 12 additions & 0 deletions docs/release.rst
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,18 @@ This release od Zarr Python is is the first release of Zarr to not supporting Py
See `this link <https://github.com/zarr-developers/zarr-python/milestone/11?closed=1>` for the full list of closed and
merged PR tagged with the 2.6 milestone.

* Add ability to partially read and decompress arrays, see :issue:`667`. It is
only available to chunks stored using fs-spec and using bloc as a compressor.

For certain analysis case when only a small portion of chunks is needed it can
be advantageous to only access and decompress part of the chunks. Doing
partial read and decompression add high latency to many of the operation so
should be used only when the subset of the data is small compared to the full
chunks and is stored contiguously (that is to say either last dimensions for C
layout, firsts for F). Pass ``partial_decompress=True`` as argument when
creating an ``Array``, or when using ``open_array``. No option exists yet to
apply partial read and decompress on a per-operation basis.

2.5.0
-----

Expand Down
2 changes: 1 addition & 1 deletion requirements_dev_minimal.txt
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
# library requirements
asciitree==0.3.3
fasteners==0.15
numcodecs==0.6.4
numcodecs==0.7.2
msgpack-python==0.5.6
setuptools-scm==3.3.3
# test requirements
Expand Down
152 changes: 128 additions & 24 deletions zarr/core.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,22 +11,41 @@

from zarr.attrs import Attributes
from zarr.codecs import AsType, get_codec
from zarr.errors import ArrayNotFoundError, ReadOnlyError
from zarr.indexing import (BasicIndexer, CoordinateIndexer, MaskIndexer,
OIndex, OrthogonalIndexer, VIndex, check_fields,
check_no_multi_fields, ensure_tuple,
err_too_many_indices, is_contiguous_selection,
is_scalar, pop_fields)
from zarr.errors import ArrayNotFoundError, ReadOnlyError, ArrayIndexError
from zarr.indexing import (
BasicIndexer,
CoordinateIndexer,
MaskIndexer,
OIndex,
OrthogonalIndexer,
VIndex,
PartialChunkIterator,
check_fields,
check_no_multi_fields,
ensure_tuple,
err_too_many_indices,
is_contiguous_selection,
is_scalar,
pop_fields,
)
from zarr.meta import decode_array_metadata, encode_array_metadata
from zarr.storage import array_meta_key, attrs_key, getsize, listdir
from zarr.util import (InfoReporter, check_array_shape, human_readable_size,
is_total_slice, nolock, normalize_chunks,
normalize_resize_args, normalize_shape,
normalize_storage_path)
from zarr.util import (
InfoReporter,
check_array_shape,
human_readable_size,
is_total_slice,
nolock,
normalize_chunks,
normalize_resize_args,
normalize_shape,
normalize_storage_path,
PartialReadBuffer,
)


# noinspection PyUnresolvedReferences
class Array(object):
class Array:
"""Instantiate an array from an initialized store.

Parameters
Expand All@@ -51,6 +70,12 @@ class Array(object):
If True (default), user attributes will be cached for attribute read
operations. If False, user attributes are reloaded from the store prior
to all attribute read operations.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Attributes
----------
Expand DownExpand Up@@ -102,8 +127,17 @@ class Array(object):

"""

def __init__(self, store, path=None, read_only=False, chunk_store=None,
synchronizer=None, cache_metadata=True, cache_attrs=True):
def __init__(
self,
store,
path=None,
read_only=False,
chunk_store=None,
synchronizer=None,
cache_metadata=True,
cache_attrs=True,
partial_decompress=False,
):
# N.B., expect at this point store is fully initialized with all
# configuration metadata fully specified and normalized

Expand All@@ -118,6 +152,7 @@ def __init__(self, store, path=None, read_only=False, chunk_store=None,
self._synchronizer = synchronizer
self._cache_metadata = cache_metadata
self._is_view = False
self._partial_decompress = partial_decompress

# initialize metadata
self._load_metadata()
Expand DownExpand Up@@ -1580,8 +1615,17 @@ def _set_selection(self, indexer, value, fields=None):
self._chunk_setitems(lchunk_coords, lchunk_selection, chunk_values,
fields=fields)

def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
out_is_ndarray, fields, out_selection):
def _process_chunk(
self,
out,
cdata,
chunk_selection,
drop_axes,
out_is_ndarray,
fields,
out_selection,
partial_read_decode=False,
):
"""Take binary data from storage and fill output array"""
if (out_is_ndarray and
not fields and
Expand All@@ -1604,8 +1648,9 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
# optimization: we want the whole chunk, and the destination is
# contiguous, so we can decompress directly from the chunk
# into the destination array

if self._compressor:
if isinstance(cdata, PartialReadBuffer):
cdata = cdata.read_full()
self._compressor.decode(cdata, dest)
else:
chunk = ensure_ndarray(cdata).view(self._dtype)
Expand All@@ -1614,6 +1659,33 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
return

# decode chunk
try:
if partial_read_decode:
cdata.prepare_chunk()
# size of chunk
tmp = np.empty(self._chunks, dtype=self.dtype)
index_selection = PartialChunkIterator(chunk_selection, self.chunks)
for start, nitems, partial_out_selection in index_selection:
expected_shape = [
len(
range(*partial_out_selection[i].indices(self.chunks[0] + 1))
)
if i < len(partial_out_selection)
else dim
for i, dim in enumerate(self.chunks)
]
cdata.read_part(start, nitems)
chunk_partial = self._decode_chunk(
cdata.buff,
start=start,
nitems=nitems,
expected_shape=expected_shape,
)
tmp[partial_out_selection] = chunk_partial
out[out_selection] = tmp[chunk_selection]
return
except ArrayIndexError:
cdata = cdata.read_full()
chunk = self._decode_chunk(cdata)

# select data from chunk
Expand DownExpand Up@@ -1688,11 +1760,36 @@ def _chunk_getitems(self, lchunk_coords, lchunk_selection, out, lout_selection,
out_is_ndarray = False

ckeys = [self._chunk_key(ch) for ch in lchunk_coords]
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
if (
self._partial_decompress
and self._compressor
and self._compressor.codec_id == "blosc"
and hasattr(self._compressor, "decode_partial")
and not fields
and self.dtype != object
and hasattr(self.chunk_store, "getitems")
):
partial_read_decode = True
cdatas = {
ckey: PartialReadBuffer(ckey, self.chunk_store)
for ckey in ckeys
if ckey in self.chunk_store
}
else:
partial_read_decode = False
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
for ckey, chunk_select, out_select in zip(ckeys, lchunk_selection, lout_selection):
if ckey in cdatas:
self._process_chunk(out, cdatas[ckey], chunk_select, drop_axes,
out_is_ndarray, fields, out_select)
self._process_chunk(
out,
cdatas[ckey],
chunk_select,
drop_axes,
out_is_ndarray,
fields,
out_select,
partial_read_decode=partial_read_decode,
)
else:
# check exception type
if self._fill_value is not None:
Expand All@@ -1706,7 +1803,8 @@ def _chunk_setitems(self, lchunk_coords, lchunk_selection, values, fields=None):
ckeys = [self._chunk_key(co) for co in lchunk_coords]
cdatas = [self._process_for_setitem(key, sel, val, fields=fields)
for key, sel, val in zip(ckeys, lchunk_selection, values)]
self.chunk_store.setitems({k: v for k, v in zip(ckeys, cdatas)})
values = {k: v for k, v in zip(ckeys, cdatas)}
self.chunk_store.setitems(values)

def _chunk_setitem(self, chunk_coords, chunk_selection, value, fields=None):
"""Replace part or whole of a chunk.
Expand DownExpand Up@@ -1800,11 +1898,17 @@ def _process_for_setitem(self, ckey, chunk_selection, value, fields=None):
def _chunk_key(self, chunk_coords):
return self._key_prefix + '.'.join(map(str, chunk_coords))

def _decode_chunk(self, cdata):

def _decode_chunk(self, cdata, start=None, nitems=None, expected_shape=None):
# decompress
if self._compressor:
chunk = self._compressor.decode(cdata)
# only decode requested items
if (
all([x is not None for x in [start, nitems]])
and self._compressor.codec_id == "blosc"
) and hasattr(self._compressor, "decode_partial"):
chunk = self._compressor.decode_partial(cdata, start, nitems)
else:
chunk = self._compressor.decode(cdata)
else:
chunk = cdata

Expand All@@ -1829,7 +1933,7 @@ def _decode_chunk(self, cdata):

# ensure correct chunk shape
chunk = chunk.reshape(-1, order='A')
chunk = chunk.reshape(self._chunks, order=self._order)
chunk = chunk.reshape(expected_shape or self._chunks, order=self._order)

return chunk

Expand Down
31 changes: 26 additions & 5 deletions zarr/creation.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -362,11 +362,26 @@ def array(data, **kwargs):
return z


def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
compressor='default', fill_value=0, order='C', synchronizer=None,
filters=None, cache_metadata=True, cache_attrs=True, path=None,
object_codec=None, chunk_store=None, storage_options=None,
**kwargs):
def open_array(
store=None,
mode="a",
shape=None,
chunks=True,
dtype=None,
compressor="default",
fill_value=0,
order="C",
synchronizer=None,
filters=None,
cache_metadata=True,
cache_attrs=True,
path=None,
object_codec=None,
chunk_store=None,
storage_options=None,
partial_decompress=False,
**kwargs
):
"""Open an array using file-mode-like semantics.

Parameters
Expand DownExpand Up@@ -415,6 +430,12 @@ def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
storage_options : dict
If using an fsspec URL to create the store, these will be passed to
the backend implementation. Ignored otherwise.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Returns
-------
Expand Down
4 changes: 4 additions & 0 deletions zarr/errors.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,6 +15,10 @@ def __init__(self, *args):
super().__init__(self._msg.format(*args))


class ArrayIndexError(IndexError):
pass


class _BaseZarrIndexError(IndexError):
_msg = ""

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
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
12 changes: 12 additions & 0 deletions docs/release.rst
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,18 @@ This release od Zarr Python is is the first release of Zarr to not supporting Py
See `this link <https://github.com/zarr-developers/zarr-python/milestone/11?closed=1>` for the full list of closed and
merged PR tagged with the 2.6 milestone.

* Add ability to partially read and decompress arrays, see :issue:`667`. It is
only available to chunks stored using fs-spec and using bloc as a compressor.

For certain analysis case when only a small portion of chunks is needed it can
be advantageous to only access and decompress part of the chunks. Doing
partial read and decompression add high latency to many of the operation so
should be used only when the subset of the data is small compared to the full
chunks and is stored contiguously (that is to say either last dimensions for C
layout, firsts for F). Pass ``partial_decompress=True`` as argument when
creating an ``Array``, or when using ``open_array``. No option exists yet to
apply partial read and decompress on a per-operation basis.

2.5.0
-----

Expand Down
2 changes: 1 addition & 1 deletion requirements_dev_minimal.txt
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
# library requirements
asciitree==0.3.3
fasteners==0.15
numcodecs==0.6.4
numcodecs==0.7.2
msgpack-python==0.5.6
setuptools-scm==3.3.3
# test requirements
Expand Down
152 changes: 128 additions & 24 deletions zarr/core.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,22 +11,41 @@

from zarr.attrs import Attributes
from zarr.codecs import AsType, get_codec
from zarr.errors import ArrayNotFoundError, ReadOnlyError
from zarr.indexing import (BasicIndexer, CoordinateIndexer, MaskIndexer,
OIndex, OrthogonalIndexer, VIndex, check_fields,
check_no_multi_fields, ensure_tuple,
err_too_many_indices, is_contiguous_selection,
is_scalar, pop_fields)
from zarr.errors import ArrayNotFoundError, ReadOnlyError, ArrayIndexError
from zarr.indexing import (
BasicIndexer,
CoordinateIndexer,
MaskIndexer,
OIndex,
OrthogonalIndexer,
VIndex,
PartialChunkIterator,
check_fields,
check_no_multi_fields,
ensure_tuple,
err_too_many_indices,
is_contiguous_selection,
is_scalar,
pop_fields,
)
from zarr.meta import decode_array_metadata, encode_array_metadata
from zarr.storage import array_meta_key, attrs_key, getsize, listdir
from zarr.util import (InfoReporter, check_array_shape, human_readable_size,
is_total_slice, nolock, normalize_chunks,
normalize_resize_args, normalize_shape,
normalize_storage_path)
from zarr.util import (
InfoReporter,
check_array_shape,
human_readable_size,
is_total_slice,
nolock,
normalize_chunks,
normalize_resize_args,
normalize_shape,
normalize_storage_path,
PartialReadBuffer,
)


# noinspection PyUnresolvedReferences
class Array(object):
class Array:
"""Instantiate an array from an initialized store.

Parameters
Expand All@@ -51,6 +70,12 @@ class Array(object):
If True (default), user attributes will be cached for attribute read
operations. If False, user attributes are reloaded from the store prior
to all attribute read operations.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Attributes
----------
Expand DownExpand Up@@ -102,8 +127,17 @@ class Array(object):

"""

def __init__(self, store, path=None, read_only=False, chunk_store=None,
synchronizer=None, cache_metadata=True, cache_attrs=True):
def __init__(
self,
store,
path=None,
read_only=False,
chunk_store=None,
synchronizer=None,
cache_metadata=True,
cache_attrs=True,
partial_decompress=False,
):
# N.B., expect at this point store is fully initialized with all
# configuration metadata fully specified and normalized

Expand All@@ -118,6 +152,7 @@ def __init__(self, store, path=None, read_only=False, chunk_store=None,
self._synchronizer = synchronizer
self._cache_metadata = cache_metadata
self._is_view = False
self._partial_decompress = partial_decompress

# initialize metadata
self._load_metadata()
Expand DownExpand Up@@ -1580,8 +1615,17 @@ def _set_selection(self, indexer, value, fields=None):
self._chunk_setitems(lchunk_coords, lchunk_selection, chunk_values,
fields=fields)

def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
out_is_ndarray, fields, out_selection):
def _process_chunk(
self,
out,
cdata,
chunk_selection,
drop_axes,
out_is_ndarray,
fields,
out_selection,
partial_read_decode=False,
):
"""Take binary data from storage and fill output array"""
if (out_is_ndarray and
not fields and
Expand All@@ -1604,8 +1648,9 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
# optimization: we want the whole chunk, and the destination is
# contiguous, so we can decompress directly from the chunk
# into the destination array

if self._compressor:
if isinstance(cdata, PartialReadBuffer):
cdata = cdata.read_full()
self._compressor.decode(cdata, dest)
else:
chunk = ensure_ndarray(cdata).view(self._dtype)
Expand All@@ -1614,6 +1659,33 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
return

# decode chunk
try:
if partial_read_decode:
cdata.prepare_chunk()
# size of chunk
tmp = np.empty(self._chunks, dtype=self.dtype)
index_selection = PartialChunkIterator(chunk_selection, self.chunks)
for start, nitems, partial_out_selection in index_selection:
expected_shape = [
len(
range(*partial_out_selection[i].indices(self.chunks[0] + 1))
)
if i < len(partial_out_selection)
else dim
for i, dim in enumerate(self.chunks)
]
cdata.read_part(start, nitems)
chunk_partial = self._decode_chunk(
cdata.buff,
start=start,
nitems=nitems,
expected_shape=expected_shape,
)
tmp[partial_out_selection] = chunk_partial
out[out_selection] = tmp[chunk_selection]
return
except ArrayIndexError:
cdata = cdata.read_full()
chunk = self._decode_chunk(cdata)

# select data from chunk
Expand DownExpand Up@@ -1688,11 +1760,36 @@ def _chunk_getitems(self, lchunk_coords, lchunk_selection, out, lout_selection,
out_is_ndarray = False

ckeys = [self._chunk_key(ch) for ch in lchunk_coords]
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
if (
self._partial_decompress
and self._compressor
and self._compressor.codec_id == "blosc"
and hasattr(self._compressor, "decode_partial")
and not fields
and self.dtype != object
and hasattr(self.chunk_store, "getitems")
):
partial_read_decode = True
cdatas = {
ckey: PartialReadBuffer(ckey, self.chunk_store)
for ckey in ckeys
if ckey in self.chunk_store
}
else:
partial_read_decode = False
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
for ckey, chunk_select, out_select in zip(ckeys, lchunk_selection, lout_selection):
if ckey in cdatas:
self._process_chunk(out, cdatas[ckey], chunk_select, drop_axes,
out_is_ndarray, fields, out_select)
self._process_chunk(
out,
cdatas[ckey],
chunk_select,
drop_axes,
out_is_ndarray,
fields,
out_select,
partial_read_decode=partial_read_decode,
)
else:
# check exception type
if self._fill_value is not None:
Expand All@@ -1706,7 +1803,8 @@ def _chunk_setitems(self, lchunk_coords, lchunk_selection, values, fields=None):
ckeys = [self._chunk_key(co) for co in lchunk_coords]
cdatas = [self._process_for_setitem(key, sel, val, fields=fields)
for key, sel, val in zip(ckeys, lchunk_selection, values)]
self.chunk_store.setitems({k: v for k, v in zip(ckeys, cdatas)})
values = {k: v for k, v in zip(ckeys, cdatas)}
self.chunk_store.setitems(values)

def _chunk_setitem(self, chunk_coords, chunk_selection, value, fields=None):
"""Replace part or whole of a chunk.
Expand DownExpand Up@@ -1800,11 +1898,17 @@ def _process_for_setitem(self, ckey, chunk_selection, value, fields=None):
def _chunk_key(self, chunk_coords):
return self._key_prefix + '.'.join(map(str, chunk_coords))

def _decode_chunk(self, cdata):

def _decode_chunk(self, cdata, start=None, nitems=None, expected_shape=None):
# decompress
if self._compressor:
chunk = self._compressor.decode(cdata)
# only decode requested items
if (
all([x is not None for x in [start, nitems]])
and self._compressor.codec_id == "blosc"
) and hasattr(self._compressor, "decode_partial"):
chunk = self._compressor.decode_partial(cdata, start, nitems)
else:
chunk = self._compressor.decode(cdata)
else:
chunk = cdata

Expand All@@ -1829,7 +1933,7 @@ def _decode_chunk(self, cdata):

# ensure correct chunk shape
chunk = chunk.reshape(-1, order='A')
chunk = chunk.reshape(self._chunks, order=self._order)
chunk = chunk.reshape(expected_shape or self._chunks, order=self._order)

return chunk

Expand Down
31 changes: 26 additions & 5 deletions zarr/creation.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -362,11 +362,26 @@ def array(data, **kwargs):
return z


def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
compressor='default', fill_value=0, order='C', synchronizer=None,
filters=None, cache_metadata=True, cache_attrs=True, path=None,
object_codec=None, chunk_store=None, storage_options=None,
**kwargs):
def open_array(
store=None,
mode="a",
shape=None,
chunks=True,
dtype=None,
compressor="default",
fill_value=0,
order="C",
synchronizer=None,
filters=None,
cache_metadata=True,
cache_attrs=True,
path=None,
object_codec=None,
chunk_store=None,
storage_options=None,
partial_decompress=False,
**kwargs
):
"""Open an array using file-mode-like semantics.

Parameters
Expand DownExpand Up@@ -415,6 +430,12 @@ def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
storage_options : dict
If using an fsspec URL to create the store, these will be passed to
the backend implementation. Ignored otherwise.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Returns
-------
Expand Down
4 changes: 4 additions & 0 deletions zarr/errors.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,6 +15,10 @@ def __init__(self, *args):
super().__init__(self._msg.format(*args))


class ArrayIndexError(IndexError):
pass


class _BaseZarrIndexError(IndexError):
_msg = ""

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
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
12 changes: 12 additions & 0 deletions docs/release.rst
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,18 @@ This release od Zarr Python is is the first release of Zarr to not supporting Py
See `this link <https://github.com/zarr-developers/zarr-python/milestone/11?closed=1>` for the full list of closed and
merged PR tagged with the 2.6 milestone.

* Add ability to partially read and decompress arrays, see :issue:`667`. It is
only available to chunks stored using fs-spec and using bloc as a compressor.

For certain analysis case when only a small portion of chunks is needed it can
be advantageous to only access and decompress part of the chunks. Doing
partial read and decompression add high latency to many of the operation so
should be used only when the subset of the data is small compared to the full
chunks and is stored contiguously (that is to say either last dimensions for C
layout, firsts for F). Pass ``partial_decompress=True`` as argument when
creating an ``Array``, or when using ``open_array``. No option exists yet to
apply partial read and decompress on a per-operation basis.

2.5.0
-----

Expand Down
2 changes: 1 addition & 1 deletion requirements_dev_minimal.txt
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
# library requirements
asciitree==0.3.3
fasteners==0.15
numcodecs==0.6.4
numcodecs==0.7.2
msgpack-python==0.5.6
setuptools-scm==3.3.3
# test requirements
Expand Down
152 changes: 128 additions & 24 deletions zarr/core.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,22 +11,41 @@

from zarr.attrs import Attributes
from zarr.codecs import AsType, get_codec
from zarr.errors import ArrayNotFoundError, ReadOnlyError
from zarr.indexing import (BasicIndexer, CoordinateIndexer, MaskIndexer,
OIndex, OrthogonalIndexer, VIndex, check_fields,
check_no_multi_fields, ensure_tuple,
err_too_many_indices, is_contiguous_selection,
is_scalar, pop_fields)
from zarr.errors import ArrayNotFoundError, ReadOnlyError, ArrayIndexError
from zarr.indexing import (
BasicIndexer,
CoordinateIndexer,
MaskIndexer,
OIndex,
OrthogonalIndexer,
VIndex,
PartialChunkIterator,
check_fields,
check_no_multi_fields,
ensure_tuple,
err_too_many_indices,
is_contiguous_selection,
is_scalar,
pop_fields,
)
from zarr.meta import decode_array_metadata, encode_array_metadata
from zarr.storage import array_meta_key, attrs_key, getsize, listdir
from zarr.util import (InfoReporter, check_array_shape, human_readable_size,
is_total_slice, nolock, normalize_chunks,
normalize_resize_args, normalize_shape,
normalize_storage_path)
from zarr.util import (
InfoReporter,
check_array_shape,
human_readable_size,
is_total_slice,
nolock,
normalize_chunks,
normalize_resize_args,
normalize_shape,
normalize_storage_path,
PartialReadBuffer,
)


# noinspection PyUnresolvedReferences
class Array(object):
class Array:
"""Instantiate an array from an initialized store.

Parameters
Expand All@@ -51,6 +70,12 @@ class Array(object):
If True (default), user attributes will be cached for attribute read
operations. If False, user attributes are reloaded from the store prior
to all attribute read operations.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Attributes
----------
Expand DownExpand Up@@ -102,8 +127,17 @@ class Array(object):

"""

def __init__(self, store, path=None, read_only=False, chunk_store=None,
synchronizer=None, cache_metadata=True, cache_attrs=True):
def __init__(
self,
store,
path=None,
read_only=False,
chunk_store=None,
synchronizer=None,
cache_metadata=True,
cache_attrs=True,
partial_decompress=False,
):
# N.B., expect at this point store is fully initialized with all
# configuration metadata fully specified and normalized

Expand All@@ -118,6 +152,7 @@ def __init__(self, store, path=None, read_only=False, chunk_store=None,
self._synchronizer = synchronizer
self._cache_metadata = cache_metadata
self._is_view = False
self._partial_decompress = partial_decompress

# initialize metadata
self._load_metadata()
Expand DownExpand Up@@ -1580,8 +1615,17 @@ def _set_selection(self, indexer, value, fields=None):
self._chunk_setitems(lchunk_coords, lchunk_selection, chunk_values,
fields=fields)

def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
out_is_ndarray, fields, out_selection):
def _process_chunk(
self,
out,
cdata,
chunk_selection,
drop_axes,
out_is_ndarray,
fields,
out_selection,
partial_read_decode=False,
):
"""Take binary data from storage and fill output array"""
if (out_is_ndarray and
not fields and
Expand All@@ -1604,8 +1648,9 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
# optimization: we want the whole chunk, and the destination is
# contiguous, so we can decompress directly from the chunk
# into the destination array

if self._compressor:
if isinstance(cdata, PartialReadBuffer):
cdata = cdata.read_full()
self._compressor.decode(cdata, dest)
else:
chunk = ensure_ndarray(cdata).view(self._dtype)
Expand All@@ -1614,6 +1659,33 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
return

# decode chunk
try:
if partial_read_decode:
cdata.prepare_chunk()
# size of chunk
tmp = np.empty(self._chunks, dtype=self.dtype)
index_selection = PartialChunkIterator(chunk_selection, self.chunks)
for start, nitems, partial_out_selection in index_selection:
expected_shape = [
len(
range(*partial_out_selection[i].indices(self.chunks[0] + 1))
)
if i < len(partial_out_selection)
else dim
for i, dim in enumerate(self.chunks)
]
cdata.read_part(start, nitems)
chunk_partial = self._decode_chunk(
cdata.buff,
start=start,
nitems=nitems,
expected_shape=expected_shape,
)
tmp[partial_out_selection] = chunk_partial
out[out_selection] = tmp[chunk_selection]
return
except ArrayIndexError:
cdata = cdata.read_full()
chunk = self._decode_chunk(cdata)

# select data from chunk
Expand DownExpand Up@@ -1688,11 +1760,36 @@ def _chunk_getitems(self, lchunk_coords, lchunk_selection, out, lout_selection,
out_is_ndarray = False

ckeys = [self._chunk_key(ch) for ch in lchunk_coords]
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
if (
self._partial_decompress
and self._compressor
and self._compressor.codec_id == "blosc"
and hasattr(self._compressor, "decode_partial")
and not fields
and self.dtype != object
and hasattr(self.chunk_store, "getitems")
):
partial_read_decode = True
cdatas = {
ckey: PartialReadBuffer(ckey, self.chunk_store)
for ckey in ckeys
if ckey in self.chunk_store
}
else:
partial_read_decode = False
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
for ckey, chunk_select, out_select in zip(ckeys, lchunk_selection, lout_selection):
if ckey in cdatas:
self._process_chunk(out, cdatas[ckey], chunk_select, drop_axes,
out_is_ndarray, fields, out_select)
self._process_chunk(
out,
cdatas[ckey],
chunk_select,
drop_axes,
out_is_ndarray,
fields,
out_select,
partial_read_decode=partial_read_decode,
)
else:
# check exception type
if self._fill_value is not None:
Expand All@@ -1706,7 +1803,8 @@ def _chunk_setitems(self, lchunk_coords, lchunk_selection, values, fields=None):
ckeys = [self._chunk_key(co) for co in lchunk_coords]
cdatas = [self._process_for_setitem(key, sel, val, fields=fields)
for key, sel, val in zip(ckeys, lchunk_selection, values)]
self.chunk_store.setitems({k: v for k, v in zip(ckeys, cdatas)})
values = {k: v for k, v in zip(ckeys, cdatas)}
self.chunk_store.setitems(values)

def _chunk_setitem(self, chunk_coords, chunk_selection, value, fields=None):
"""Replace part or whole of a chunk.
Expand DownExpand Up@@ -1800,11 +1898,17 @@ def _process_for_setitem(self, ckey, chunk_selection, value, fields=None):
def _chunk_key(self, chunk_coords):
return self._key_prefix + '.'.join(map(str, chunk_coords))

def _decode_chunk(self, cdata):

def _decode_chunk(self, cdata, start=None, nitems=None, expected_shape=None):
# decompress
if self._compressor:
chunk = self._compressor.decode(cdata)
# only decode requested items
if (
all([x is not None for x in [start, nitems]])
and self._compressor.codec_id == "blosc"
) and hasattr(self._compressor, "decode_partial"):
chunk = self._compressor.decode_partial(cdata, start, nitems)
else:
chunk = self._compressor.decode(cdata)
else:
chunk = cdata

Expand All@@ -1829,7 +1933,7 @@ def _decode_chunk(self, cdata):

# ensure correct chunk shape
chunk = chunk.reshape(-1, order='A')
chunk = chunk.reshape(self._chunks, order=self._order)
chunk = chunk.reshape(expected_shape or self._chunks, order=self._order)

return chunk

Expand Down
31 changes: 26 additions & 5 deletions zarr/creation.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -362,11 +362,26 @@ def array(data, **kwargs):
return z


def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
compressor='default', fill_value=0, order='C', synchronizer=None,
filters=None, cache_metadata=True, cache_attrs=True, path=None,
object_codec=None, chunk_store=None, storage_options=None,
**kwargs):
def open_array(
store=None,
mode="a",
shape=None,
chunks=True,
dtype=None,
compressor="default",
fill_value=0,
order="C",
synchronizer=None,
filters=None,
cache_metadata=True,
cache_attrs=True,
path=None,
object_codec=None,
chunk_store=None,
storage_options=None,
partial_decompress=False,
**kwargs
):
"""Open an array using file-mode-like semantics.

Parameters
Expand DownExpand Up@@ -415,6 +430,12 @@ def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
storage_options : dict
If using an fsspec URL to create the store, these will be passed to
the backend implementation. Ignored otherwise.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Returns
-------
Expand Down
4 changes: 4 additions & 0 deletions zarr/errors.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,6 +15,10 @@ def __init__(self, *args):
super().__init__(self._msg.format(*args))


class ArrayIndexError(IndexError):
pass


class _BaseZarrIndexError(IndexError):
_msg = ""

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
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
12 changes: 12 additions & 0 deletions docs/release.rst
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,18 @@ This release od Zarr Python is is the first release of Zarr to not supporting Py
See `this link <https://github.com/zarr-developers/zarr-python/milestone/11?closed=1>` for the full list of closed and
merged PR tagged with the 2.6 milestone.

* Add ability to partially read and decompress arrays, see :issue:`667`. It is
only available to chunks stored using fs-spec and using bloc as a compressor.

For certain analysis case when only a small portion of chunks is needed it can
be advantageous to only access and decompress part of the chunks. Doing
partial read and decompression add high latency to many of the operation so
should be used only when the subset of the data is small compared to the full
chunks and is stored contiguously (that is to say either last dimensions for C
layout, firsts for F). Pass ``partial_decompress=True`` as argument when
creating an ``Array``, or when using ``open_array``. No option exists yet to
apply partial read and decompress on a per-operation basis.

2.5.0
-----

Expand Down
2 changes: 1 addition & 1 deletion requirements_dev_minimal.txt
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
# library requirements
asciitree==0.3.3
fasteners==0.15
numcodecs==0.6.4
numcodecs==0.7.2
msgpack-python==0.5.6
setuptools-scm==3.3.3
# test requirements
Expand Down
152 changes: 128 additions & 24 deletions zarr/core.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,22 +11,41 @@

from zarr.attrs import Attributes
from zarr.codecs import AsType, get_codec
from zarr.errors import ArrayNotFoundError, ReadOnlyError
from zarr.indexing import (BasicIndexer, CoordinateIndexer, MaskIndexer,
OIndex, OrthogonalIndexer, VIndex, check_fields,
check_no_multi_fields, ensure_tuple,
err_too_many_indices, is_contiguous_selection,
is_scalar, pop_fields)
from zarr.errors import ArrayNotFoundError, ReadOnlyError, ArrayIndexError
from zarr.indexing import (
BasicIndexer,
CoordinateIndexer,
MaskIndexer,
OIndex,
OrthogonalIndexer,
VIndex,
PartialChunkIterator,
check_fields,
check_no_multi_fields,
ensure_tuple,
err_too_many_indices,
is_contiguous_selection,
is_scalar,
pop_fields,
)
from zarr.meta import decode_array_metadata, encode_array_metadata
from zarr.storage import array_meta_key, attrs_key, getsize, listdir
from zarr.util import (InfoReporter, check_array_shape, human_readable_size,
is_total_slice, nolock, normalize_chunks,
normalize_resize_args, normalize_shape,
normalize_storage_path)
from zarr.util import (
InfoReporter,
check_array_shape,
human_readable_size,
is_total_slice,
nolock,
normalize_chunks,
normalize_resize_args,
normalize_shape,
normalize_storage_path,
PartialReadBuffer,
)


# noinspection PyUnresolvedReferences
class Array(object):
class Array:
"""Instantiate an array from an initialized store.

Parameters
Expand All@@ -51,6 +70,12 @@ class Array(object):
If True (default), user attributes will be cached for attribute read
operations. If False, user attributes are reloaded from the store prior
to all attribute read operations.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Attributes
----------
Expand DownExpand Up@@ -102,8 +127,17 @@ class Array(object):

"""

def __init__(self, store, path=None, read_only=False, chunk_store=None,
synchronizer=None, cache_metadata=True, cache_attrs=True):
def __init__(
self,
store,
path=None,
read_only=False,
chunk_store=None,
synchronizer=None,
cache_metadata=True,
cache_attrs=True,
partial_decompress=False,
):
# N.B., expect at this point store is fully initialized with all
# configuration metadata fully specified and normalized

Expand All@@ -118,6 +152,7 @@ def __init__(self, store, path=None, read_only=False, chunk_store=None,
self._synchronizer = synchronizer
self._cache_metadata = cache_metadata
self._is_view = False
self._partial_decompress = partial_decompress

# initialize metadata
self._load_metadata()
Expand DownExpand Up@@ -1580,8 +1615,17 @@ def _set_selection(self, indexer, value, fields=None):
self._chunk_setitems(lchunk_coords, lchunk_selection, chunk_values,
fields=fields)

def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
out_is_ndarray, fields, out_selection):
def _process_chunk(
self,
out,
cdata,
chunk_selection,
drop_axes,
out_is_ndarray,
fields,
out_selection,
partial_read_decode=False,
):
"""Take binary data from storage and fill output array"""
if (out_is_ndarray and
not fields and
Expand All@@ -1604,8 +1648,9 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
# optimization: we want the whole chunk, and the destination is
# contiguous, so we can decompress directly from the chunk
# into the destination array

if self._compressor:
if isinstance(cdata, PartialReadBuffer):
cdata = cdata.read_full()
self._compressor.decode(cdata, dest)
else:
chunk = ensure_ndarray(cdata).view(self._dtype)
Expand All@@ -1614,6 +1659,33 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
return

# decode chunk
try:
if partial_read_decode:
cdata.prepare_chunk()
# size of chunk
tmp = np.empty(self._chunks, dtype=self.dtype)
index_selection = PartialChunkIterator(chunk_selection, self.chunks)
for start, nitems, partial_out_selection in index_selection:
expected_shape = [
len(
range(*partial_out_selection[i].indices(self.chunks[0] + 1))
)
if i < len(partial_out_selection)
else dim
for i, dim in enumerate(self.chunks)
]
cdata.read_part(start, nitems)
chunk_partial = self._decode_chunk(
cdata.buff,
start=start,
nitems=nitems,
expected_shape=expected_shape,
)
tmp[partial_out_selection] = chunk_partial
out[out_selection] = tmp[chunk_selection]
return
except ArrayIndexError:
cdata = cdata.read_full()
chunk = self._decode_chunk(cdata)

# select data from chunk
Expand DownExpand Up@@ -1688,11 +1760,36 @@ def _chunk_getitems(self, lchunk_coords, lchunk_selection, out, lout_selection,
out_is_ndarray = False

ckeys = [self._chunk_key(ch) for ch in lchunk_coords]
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
if (
self._partial_decompress
and self._compressor
and self._compressor.codec_id == "blosc"
and hasattr(self._compressor, "decode_partial")
and not fields
and self.dtype != object
and hasattr(self.chunk_store, "getitems")
):
partial_read_decode = True
cdatas = {
ckey: PartialReadBuffer(ckey, self.chunk_store)
for ckey in ckeys
if ckey in self.chunk_store
}
else:
partial_read_decode = False
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
for ckey, chunk_select, out_select in zip(ckeys, lchunk_selection, lout_selection):
if ckey in cdatas:
self._process_chunk(out, cdatas[ckey], chunk_select, drop_axes,
out_is_ndarray, fields, out_select)
self._process_chunk(
out,
cdatas[ckey],
chunk_select,
drop_axes,
out_is_ndarray,
fields,
out_select,
partial_read_decode=partial_read_decode,
)
else:
# check exception type
if self._fill_value is not None:
Expand All@@ -1706,7 +1803,8 @@ def _chunk_setitems(self, lchunk_coords, lchunk_selection, values, fields=None):
ckeys = [self._chunk_key(co) for co in lchunk_coords]
cdatas = [self._process_for_setitem(key, sel, val, fields=fields)
for key, sel, val in zip(ckeys, lchunk_selection, values)]
self.chunk_store.setitems({k: v for k, v in zip(ckeys, cdatas)})
values = {k: v for k, v in zip(ckeys, cdatas)}
self.chunk_store.setitems(values)

def _chunk_setitem(self, chunk_coords, chunk_selection, value, fields=None):
"""Replace part or whole of a chunk.
Expand DownExpand Up@@ -1800,11 +1898,17 @@ def _process_for_setitem(self, ckey, chunk_selection, value, fields=None):
def _chunk_key(self, chunk_coords):
return self._key_prefix + '.'.join(map(str, chunk_coords))

def _decode_chunk(self, cdata):

def _decode_chunk(self, cdata, start=None, nitems=None, expected_shape=None):
# decompress
if self._compressor:
chunk = self._compressor.decode(cdata)
# only decode requested items
if (
all([x is not None for x in [start, nitems]])
and self._compressor.codec_id == "blosc"
) and hasattr(self._compressor, "decode_partial"):
chunk = self._compressor.decode_partial(cdata, start, nitems)
else:
chunk = self._compressor.decode(cdata)
else:
chunk = cdata

Expand All@@ -1829,7 +1933,7 @@ def _decode_chunk(self, cdata):

# ensure correct chunk shape
chunk = chunk.reshape(-1, order='A')
chunk = chunk.reshape(self._chunks, order=self._order)
chunk = chunk.reshape(expected_shape or self._chunks, order=self._order)

return chunk

Expand Down
31 changes: 26 additions & 5 deletions zarr/creation.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -362,11 +362,26 @@ def array(data, **kwargs):
return z


def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
compressor='default', fill_value=0, order='C', synchronizer=None,
filters=None, cache_metadata=True, cache_attrs=True, path=None,
object_codec=None, chunk_store=None, storage_options=None,
**kwargs):
def open_array(
store=None,
mode="a",
shape=None,
chunks=True,
dtype=None,
compressor="default",
fill_value=0,
order="C",
synchronizer=None,
filters=None,
cache_metadata=True,
cache_attrs=True,
path=None,
object_codec=None,
chunk_store=None,
storage_options=None,
partial_decompress=False,
**kwargs
):
"""Open an array using file-mode-like semantics.

Parameters
Expand DownExpand Up@@ -415,6 +430,12 @@ def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
storage_options : dict
If using an fsspec URL to create the store, these will be passed to
the backend implementation. Ignored otherwise.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Returns
-------
Expand Down
4 changes: 4 additions & 0 deletions zarr/errors.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,6 +15,10 @@ def __init__(self, *args):
super().__init__(self._msg.format(*args))


class ArrayIndexError(IndexError):
pass


class _BaseZarrIndexError(IndexError):
_msg = ""

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
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
12 changes: 12 additions & 0 deletions docs/release.rst
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,18 @@ This release od Zarr Python is is the first release of Zarr to not supporting Py
See `this link <https://github.com/zarr-developers/zarr-python/milestone/11?closed=1>` for the full list of closed and
merged PR tagged with the 2.6 milestone.

* Add ability to partially read and decompress arrays, see :issue:`667`. It is
only available to chunks stored using fs-spec and using bloc as a compressor.

For certain analysis case when only a small portion of chunks is needed it can
be advantageous to only access and decompress part of the chunks. Doing
partial read and decompression add high latency to many of the operation so
should be used only when the subset of the data is small compared to the full
chunks and is stored contiguously (that is to say either last dimensions for C
layout, firsts for F). Pass ``partial_decompress=True`` as argument when
creating an ``Array``, or when using ``open_array``. No option exists yet to
apply partial read and decompress on a per-operation basis.

2.5.0
-----

Expand Down
2 changes: 1 addition & 1 deletion requirements_dev_minimal.txt
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
# library requirements
asciitree==0.3.3
fasteners==0.15
numcodecs==0.6.4
numcodecs==0.7.2
msgpack-python==0.5.6
setuptools-scm==3.3.3
# test requirements
Expand Down
152 changes: 128 additions & 24 deletions zarr/core.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,22 +11,41 @@

from zarr.attrs import Attributes
from zarr.codecs import AsType, get_codec
from zarr.errors import ArrayNotFoundError, ReadOnlyError
from zarr.indexing import (BasicIndexer, CoordinateIndexer, MaskIndexer,
OIndex, OrthogonalIndexer, VIndex, check_fields,
check_no_multi_fields, ensure_tuple,
err_too_many_indices, is_contiguous_selection,
is_scalar, pop_fields)
from zarr.errors import ArrayNotFoundError, ReadOnlyError, ArrayIndexError
from zarr.indexing import (
BasicIndexer,
CoordinateIndexer,
MaskIndexer,
OIndex,
OrthogonalIndexer,
VIndex,
PartialChunkIterator,
check_fields,
check_no_multi_fields,
ensure_tuple,
err_too_many_indices,
is_contiguous_selection,
is_scalar,
pop_fields,
)
from zarr.meta import decode_array_metadata, encode_array_metadata
from zarr.storage import array_meta_key, attrs_key, getsize, listdir
from zarr.util import (InfoReporter, check_array_shape, human_readable_size,
is_total_slice, nolock, normalize_chunks,
normalize_resize_args, normalize_shape,
normalize_storage_path)
from zarr.util import (
InfoReporter,
check_array_shape,
human_readable_size,
is_total_slice,
nolock,
normalize_chunks,
normalize_resize_args,
normalize_shape,
normalize_storage_path,
PartialReadBuffer,
)


# noinspection PyUnresolvedReferences
class Array(object):
class Array:
"""Instantiate an array from an initialized store.

Parameters
Expand All@@ -51,6 +70,12 @@ class Array(object):
If True (default), user attributes will be cached for attribute read
operations. If False, user attributes are reloaded from the store prior
to all attribute read operations.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Attributes
----------
Expand DownExpand Up@@ -102,8 +127,17 @@ class Array(object):

"""

def __init__(self, store, path=None, read_only=False, chunk_store=None,
synchronizer=None, cache_metadata=True, cache_attrs=True):
def __init__(
self,
store,
path=None,
read_only=False,
chunk_store=None,
synchronizer=None,
cache_metadata=True,
cache_attrs=True,
partial_decompress=False,
):
# N.B., expect at this point store is fully initialized with all
# configuration metadata fully specified and normalized

Expand All@@ -118,6 +152,7 @@ def __init__(self, store, path=None, read_only=False, chunk_store=None,
self._synchronizer = synchronizer
self._cache_metadata = cache_metadata
self._is_view = False
self._partial_decompress = partial_decompress

# initialize metadata
self._load_metadata()
Expand DownExpand Up@@ -1580,8 +1615,17 @@ def _set_selection(self, indexer, value, fields=None):
self._chunk_setitems(lchunk_coords, lchunk_selection, chunk_values,
fields=fields)

def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
out_is_ndarray, fields, out_selection):
def _process_chunk(
self,
out,
cdata,
chunk_selection,
drop_axes,
out_is_ndarray,
fields,
out_selection,
partial_read_decode=False,
):
"""Take binary data from storage and fill output array"""
if (out_is_ndarray and
not fields and
Expand All@@ -1604,8 +1648,9 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
# optimization: we want the whole chunk, and the destination is
# contiguous, so we can decompress directly from the chunk
# into the destination array

if self._compressor:
if isinstance(cdata, PartialReadBuffer):
cdata = cdata.read_full()
self._compressor.decode(cdata, dest)
else:
chunk = ensure_ndarray(cdata).view(self._dtype)
Expand All@@ -1614,6 +1659,33 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
return

# decode chunk
try:
if partial_read_decode:
cdata.prepare_chunk()
# size of chunk
tmp = np.empty(self._chunks, dtype=self.dtype)
index_selection = PartialChunkIterator(chunk_selection, self.chunks)
for start, nitems, partial_out_selection in index_selection:
expected_shape = [
len(
range(*partial_out_selection[i].indices(self.chunks[0] + 1))
)
if i < len(partial_out_selection)
else dim
for i, dim in enumerate(self.chunks)
]
cdata.read_part(start, nitems)
chunk_partial = self._decode_chunk(
cdata.buff,
start=start,
nitems=nitems,
expected_shape=expected_shape,
)
tmp[partial_out_selection] = chunk_partial
out[out_selection] = tmp[chunk_selection]
return
except ArrayIndexError:
cdata = cdata.read_full()
chunk = self._decode_chunk(cdata)

# select data from chunk
Expand DownExpand Up@@ -1688,11 +1760,36 @@ def _chunk_getitems(self, lchunk_coords, lchunk_selection, out, lout_selection,
out_is_ndarray = False

ckeys = [self._chunk_key(ch) for ch in lchunk_coords]
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
if (
self._partial_decompress
and self._compressor
and self._compressor.codec_id == "blosc"
and hasattr(self._compressor, "decode_partial")
and not fields
and self.dtype != object
and hasattr(self.chunk_store, "getitems")
):
partial_read_decode = True
cdatas = {
ckey: PartialReadBuffer(ckey, self.chunk_store)
for ckey in ckeys
if ckey in self.chunk_store
}
else:
partial_read_decode = False
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
for ckey, chunk_select, out_select in zip(ckeys, lchunk_selection, lout_selection):
if ckey in cdatas:
self._process_chunk(out, cdatas[ckey], chunk_select, drop_axes,
out_is_ndarray, fields, out_select)
self._process_chunk(
out,
cdatas[ckey],
chunk_select,
drop_axes,
out_is_ndarray,
fields,
out_select,
partial_read_decode=partial_read_decode,
)
else:
# check exception type
if self._fill_value is not None:
Expand All@@ -1706,7 +1803,8 @@ def _chunk_setitems(self, lchunk_coords, lchunk_selection, values, fields=None):
ckeys = [self._chunk_key(co) for co in lchunk_coords]
cdatas = [self._process_for_setitem(key, sel, val, fields=fields)
for key, sel, val in zip(ckeys, lchunk_selection, values)]
self.chunk_store.setitems({k: v for k, v in zip(ckeys, cdatas)})
values = {k: v for k, v in zip(ckeys, cdatas)}
self.chunk_store.setitems(values)

def _chunk_setitem(self, chunk_coords, chunk_selection, value, fields=None):
"""Replace part or whole of a chunk.
Expand DownExpand Up@@ -1800,11 +1898,17 @@ def _process_for_setitem(self, ckey, chunk_selection, value, fields=None):
def _chunk_key(self, chunk_coords):
return self._key_prefix + '.'.join(map(str, chunk_coords))

def _decode_chunk(self, cdata):

def _decode_chunk(self, cdata, start=None, nitems=None, expected_shape=None):
# decompress
if self._compressor:
chunk = self._compressor.decode(cdata)
# only decode requested items
if (
all([x is not None for x in [start, nitems]])
and self._compressor.codec_id == "blosc"
) and hasattr(self._compressor, "decode_partial"):
chunk = self._compressor.decode_partial(cdata, start, nitems)
else:
chunk = self._compressor.decode(cdata)
else:
chunk = cdata

Expand All@@ -1829,7 +1933,7 @@ def _decode_chunk(self, cdata):

# ensure correct chunk shape
chunk = chunk.reshape(-1, order='A')
chunk = chunk.reshape(self._chunks, order=self._order)
chunk = chunk.reshape(expected_shape or self._chunks, order=self._order)

return chunk

Expand Down
31 changes: 26 additions & 5 deletions zarr/creation.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -362,11 +362,26 @@ def array(data, **kwargs):
return z


def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
compressor='default', fill_value=0, order='C', synchronizer=None,
filters=None, cache_metadata=True, cache_attrs=True, path=None,
object_codec=None, chunk_store=None, storage_options=None,
**kwargs):
def open_array(
store=None,
mode="a",
shape=None,
chunks=True,
dtype=None,
compressor="default",
fill_value=0,
order="C",
synchronizer=None,
filters=None,
cache_metadata=True,
cache_attrs=True,
path=None,
object_codec=None,
chunk_store=None,
storage_options=None,
partial_decompress=False,
**kwargs
):
"""Open an array using file-mode-like semantics.

Parameters
Expand DownExpand Up@@ -415,6 +430,12 @@ def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
storage_options : dict
If using an fsspec URL to create the store, these will be passed to
the backend implementation. Ignored otherwise.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Returns
-------
Expand Down
4 changes: 4 additions & 0 deletions zarr/errors.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,6 +15,10 @@ def __init__(self, *args):
super().__init__(self._msg.format(*args))


class ArrayIndexError(IndexError):
pass


class _BaseZarrIndexError(IndexError):
_msg = ""

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
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
12 changes: 12 additions & 0 deletions docs/release.rst
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,18 @@ This release od Zarr Python is is the first release of Zarr to not supporting Py
See `this link <https://github.com/zarr-developers/zarr-python/milestone/11?closed=1>` for the full list of closed and
merged PR tagged with the 2.6 milestone.

* Add ability to partially read and decompress arrays, see :issue:`667`. It is
only available to chunks stored using fs-spec and using bloc as a compressor.

For certain analysis case when only a small portion of chunks is needed it can
be advantageous to only access and decompress part of the chunks. Doing
partial read and decompression add high latency to many of the operation so
should be used only when the subset of the data is small compared to the full
chunks and is stored contiguously (that is to say either last dimensions for C
layout, firsts for F). Pass ``partial_decompress=True`` as argument when
creating an ``Array``, or when using ``open_array``. No option exists yet to
apply partial read and decompress on a per-operation basis.

2.5.0
-----

Expand Down
2 changes: 1 addition & 1 deletion requirements_dev_minimal.txt
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
# library requirements
asciitree==0.3.3
fasteners==0.15
numcodecs==0.6.4
numcodecs==0.7.2
msgpack-python==0.5.6
setuptools-scm==3.3.3
# test requirements
Expand Down
152 changes: 128 additions & 24 deletions zarr/core.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,22 +11,41 @@

from zarr.attrs import Attributes
from zarr.codecs import AsType, get_codec
from zarr.errors import ArrayNotFoundError, ReadOnlyError
from zarr.indexing import (BasicIndexer, CoordinateIndexer, MaskIndexer,
OIndex, OrthogonalIndexer, VIndex, check_fields,
check_no_multi_fields, ensure_tuple,
err_too_many_indices, is_contiguous_selection,
is_scalar, pop_fields)
from zarr.errors import ArrayNotFoundError, ReadOnlyError, ArrayIndexError
from zarr.indexing import (
BasicIndexer,
CoordinateIndexer,
MaskIndexer,
OIndex,
OrthogonalIndexer,
VIndex,
PartialChunkIterator,
check_fields,
check_no_multi_fields,
ensure_tuple,
err_too_many_indices,
is_contiguous_selection,
is_scalar,
pop_fields,
)
from zarr.meta import decode_array_metadata, encode_array_metadata
from zarr.storage import array_meta_key, attrs_key, getsize, listdir
from zarr.util import (InfoReporter, check_array_shape, human_readable_size,
is_total_slice, nolock, normalize_chunks,
normalize_resize_args, normalize_shape,
normalize_storage_path)
from zarr.util import (
InfoReporter,
check_array_shape,
human_readable_size,
is_total_slice,
nolock,
normalize_chunks,
normalize_resize_args,
normalize_shape,
normalize_storage_path,
PartialReadBuffer,
)


# noinspection PyUnresolvedReferences
class Array(object):
class Array:
"""Instantiate an array from an initialized store.

Parameters
Expand All@@ -51,6 +70,12 @@ class Array(object):
If True (default), user attributes will be cached for attribute read
operations. If False, user attributes are reloaded from the store prior
to all attribute read operations.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Attributes
----------
Expand DownExpand Up@@ -102,8 +127,17 @@ class Array(object):

"""

def __init__(self, store, path=None, read_only=False, chunk_store=None,
synchronizer=None, cache_metadata=True, cache_attrs=True):
def __init__(
self,
store,
path=None,
read_only=False,
chunk_store=None,
synchronizer=None,
cache_metadata=True,
cache_attrs=True,
partial_decompress=False,
):
# N.B., expect at this point store is fully initialized with all
# configuration metadata fully specified and normalized

Expand All@@ -118,6 +152,7 @@ def __init__(self, store, path=None, read_only=False, chunk_store=None,
self._synchronizer = synchronizer
self._cache_metadata = cache_metadata
self._is_view = False
self._partial_decompress = partial_decompress

# initialize metadata
self._load_metadata()
Expand DownExpand Up@@ -1580,8 +1615,17 @@ def _set_selection(self, indexer, value, fields=None):
self._chunk_setitems(lchunk_coords, lchunk_selection, chunk_values,
fields=fields)

def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
out_is_ndarray, fields, out_selection):
def _process_chunk(
self,
out,
cdata,
chunk_selection,
drop_axes,
out_is_ndarray,
fields,
out_selection,
partial_read_decode=False,
):
"""Take binary data from storage and fill output array"""
if (out_is_ndarray and
not fields and
Expand All@@ -1604,8 +1648,9 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
# optimization: we want the whole chunk, and the destination is
# contiguous, so we can decompress directly from the chunk
# into the destination array

if self._compressor:
if isinstance(cdata, PartialReadBuffer):
cdata = cdata.read_full()
self._compressor.decode(cdata, dest)
else:
chunk = ensure_ndarray(cdata).view(self._dtype)
Expand All@@ -1614,6 +1659,33 @@ def _process_chunk(self, out, cdata, chunk_selection, drop_axes,
return

# decode chunk
try:
if partial_read_decode:
cdata.prepare_chunk()
# size of chunk
tmp = np.empty(self._chunks, dtype=self.dtype)
index_selection = PartialChunkIterator(chunk_selection, self.chunks)
for start, nitems, partial_out_selection in index_selection:
expected_shape = [
len(
range(*partial_out_selection[i].indices(self.chunks[0] + 1))
)
if i < len(partial_out_selection)
else dim
for i, dim in enumerate(self.chunks)
]
cdata.read_part(start, nitems)
chunk_partial = self._decode_chunk(
cdata.buff,
start=start,
nitems=nitems,
expected_shape=expected_shape,
)
tmp[partial_out_selection] = chunk_partial
out[out_selection] = tmp[chunk_selection]
return
except ArrayIndexError:
cdata = cdata.read_full()
chunk = self._decode_chunk(cdata)

# select data from chunk
Expand DownExpand Up@@ -1688,11 +1760,36 @@ def _chunk_getitems(self, lchunk_coords, lchunk_selection, out, lout_selection,
out_is_ndarray = False

ckeys = [self._chunk_key(ch) for ch in lchunk_coords]
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
if (
self._partial_decompress
and self._compressor
and self._compressor.codec_id == "blosc"
and hasattr(self._compressor, "decode_partial")
and not fields
and self.dtype != object
and hasattr(self.chunk_store, "getitems")
):
partial_read_decode = True
cdatas = {
ckey: PartialReadBuffer(ckey, self.chunk_store)
for ckey in ckeys
if ckey in self.chunk_store
}
else:
partial_read_decode = False
cdatas = self.chunk_store.getitems(ckeys, on_error="omit")
for ckey, chunk_select, out_select in zip(ckeys, lchunk_selection, lout_selection):
if ckey in cdatas:
self._process_chunk(out, cdatas[ckey], chunk_select, drop_axes,
out_is_ndarray, fields, out_select)
self._process_chunk(
out,
cdatas[ckey],
chunk_select,
drop_axes,
out_is_ndarray,
fields,
out_select,
partial_read_decode=partial_read_decode,
)
else:
# check exception type
if self._fill_value is not None:
Expand All@@ -1706,7 +1803,8 @@ def _chunk_setitems(self, lchunk_coords, lchunk_selection, values, fields=None):
ckeys = [self._chunk_key(co) for co in lchunk_coords]
cdatas = [self._process_for_setitem(key, sel, val, fields=fields)
for key, sel, val in zip(ckeys, lchunk_selection, values)]
self.chunk_store.setitems({k: v for k, v in zip(ckeys, cdatas)})
values = {k: v for k, v in zip(ckeys, cdatas)}
self.chunk_store.setitems(values)

def _chunk_setitem(self, chunk_coords, chunk_selection, value, fields=None):
"""Replace part or whole of a chunk.
Expand DownExpand Up@@ -1800,11 +1898,17 @@ def _process_for_setitem(self, ckey, chunk_selection, value, fields=None):
def _chunk_key(self, chunk_coords):
return self._key_prefix + '.'.join(map(str, chunk_coords))

def _decode_chunk(self, cdata):

def _decode_chunk(self, cdata, start=None, nitems=None, expected_shape=None):
# decompress
if self._compressor:
chunk = self._compressor.decode(cdata)
# only decode requested items
if (
all([x is not None for x in [start, nitems]])
and self._compressor.codec_id == "blosc"
) and hasattr(self._compressor, "decode_partial"):
chunk = self._compressor.decode_partial(cdata, start, nitems)
else:
chunk = self._compressor.decode(cdata)
else:
chunk = cdata

Expand All@@ -1829,7 +1933,7 @@ def _decode_chunk(self, cdata):

# ensure correct chunk shape
chunk = chunk.reshape(-1, order='A')
chunk = chunk.reshape(self._chunks, order=self._order)
chunk = chunk.reshape(expected_shape or self._chunks, order=self._order)

return chunk

Expand Down
31 changes: 26 additions & 5 deletions zarr/creation.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -362,11 +362,26 @@ def array(data, **kwargs):
return z


def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
compressor='default', fill_value=0, order='C', synchronizer=None,
filters=None, cache_metadata=True, cache_attrs=True, path=None,
object_codec=None, chunk_store=None, storage_options=None,
**kwargs):
def open_array(
store=None,
mode="a",
shape=None,
chunks=True,
dtype=None,
compressor="default",
fill_value=0,
order="C",
synchronizer=None,
filters=None,
cache_metadata=True,
cache_attrs=True,
path=None,
object_codec=None,
chunk_store=None,
storage_options=None,
partial_decompress=False,
**kwargs
):
"""Open an array using file-mode-like semantics.

Parameters
Expand DownExpand Up@@ -415,6 +430,12 @@ def open_array(store=None, mode='a', shape=None, chunks=True, dtype=None,
storage_options : dict
If using an fsspec URL to create the store, these will be passed to
the backend implementation. Ignored otherwise.
partial_decompress : bool, optional
If True and while the chunk_store is a FSStore and the compresion used
is Blosc, when getting data from the array chunks will be partially
read and decompressed when possible.

.. versionadded:: 2.7

Returns
-------
Expand Down
4 changes: 4 additions & 0 deletions zarr/errors.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -15,6 +15,10 @@ def __init__(self, *args):
super().__init__(self._msg.format(*args))


class ArrayIndexError(IndexError):
pass


class _BaseZarrIndexError(IndexError):
_msg = ""

Expand Down
Loading