diff --git a/pointCollection/io_utils.py b/pointCollection/io_utils.py index 6aa9a65..3b13267 100644 --- a/pointCollection/io_utils.py +++ b/pointCollection/io_utils.py @@ -37,6 +37,15 @@ # 6.6 MiB in 256 KiB ones. DEFAULT_REMOTE_BLOCK_SIZE = 256 * 1024 +# Ceiling on the bytes a windowed remote read keeps cached per open file. +# fsspec's default 'readahead' cache holds ONE block, and reading a range out +# of a chunked HDF5 file walks through every field's chunks in turn, so a +# block is evicted and fetched again many times: an ATL11 tile read pulled +# 1.4-3.3x the whole granules' size (IS 1.71 GB where 0.26 GB was needed). +# A 'blockcache' keeps the blocks already fetched (LRU); only blocks actually +# read are held, so memory is what the read touches, up to this ceiling. +DEFAULT_REMOTE_CACHE_BYTES = 1024**3 + # pc.indexedH5 is the class that reads and writes this format, so 'indexedH5' # is its canonical name. geoIndex files written before that spelling was # settled on, and calling code following the geoIndex file_type convention, @@ -302,10 +311,14 @@ def open_remote(filename, mode='rb', fs=None, block_size=None, daac='NSIDC'): get_s3fs(daac=daac). block_size : int or NoneType, default None Bytes fetched per range request. None leaves the filesystem's own - default (5 MiB for s3fs) in place; see DEFAULT_REMOTE_BLOCK_SIZE for + default in place (50 MiB for s3fs 2026.7; fsspec's generic default is + 5 MiB), which suits a whole-file read; see DEFAULT_REMOTE_BLOCK_SIZE for why a windowed read wants a smaller one. Passed per file rather than to the session, so it applies to a caller-supplied fs too -- including an earthaccess DAAC session, whose constructor takes no such argument. + A read with a block_size also gets an fsspec 'blockcache' (see + DEFAULT_REMOTE_CACHE_BYTES), so a block is fetched once however often + the read comes back to it. daac : str or NoneType, default 'NSIDC' DAAC whose credentials are needed, if fs is None. None selects the default AWS credential chain; see get_s3fs(). @@ -318,11 +331,18 @@ def open_remote(filename, mode='rb', fs=None, block_size=None, daac='NSIDC'): fs = get_s3fs(daac=daac) if block_size is None: return fs.open(filename, mode) - try: - return fs.open(filename, mode, block_size=block_size) - except TypeError: - # a filesystem (or a stand-in) whose open() takes no block_size - return fs.open(filename, mode) + attempts = [{'block_size': block_size}] + if 'r' in mode: + attempts.insert(0, {'block_size': block_size, 'cache_type': 'blockcache', + 'cache_options': {'maxblocks': max(1, DEFAULT_REMOTE_CACHE_BYTES // block_size)}}) + for kwargs in attempts: + try: + return fs.open(filename, mode, **kwargs) + except TypeError: + # a filesystem (or a stand-in) whose open() takes no cache_type, + # or no block_size: fall back to what it does take + pass + return fs.open(filename, mode) def as_gdal_path(filename): """ diff --git a/tests/test_remote_cache.py b/tests/test_remote_cache.py new file mode 100644 index 0000000..5ae2c49 --- /dev/null +++ b/tests/test_remote_cache.py @@ -0,0 +1,101 @@ +""" +open_remote's block cache: a windowed read does not refetch blocks. + +No network access is required. CountingFS is a real fsspec filesystem whose +files are local bytes, and it counts the bytes each range request fetches, so +fsspec's own caches -- the thing under test -- run exactly as they do over S3. + +The HDF5 file has many compressed fields in several groups, like an ATL11 +granule. Reading ONE window of every field (what a geoIndex query does, one +merged range per granule and beam pair) with fsspec's default one-block +'readahead' cache fetches some blocks again and again: on real ATL11 a +~1,400-point range read pulled ~30 MB, against ~7 MB with the block cache, +and a whole IS tile 1.71 GB against 0.26 GB (2026-09-25). This file shows the +same effect at a smaller ratio. +""" +import h5py +import numpy as np +import pytest +from fsspec.spec import AbstractBufferedFile, AbstractFileSystem + +import pointCollection as pc + +N = 100_000 +NAMES = [f'group{g}/field{i}' for g in range(4) for i in range(10)] +WINDOW = slice(50_000, 51_400) + + +class CountingFile(AbstractBufferedFile): + def _fetch_range(self, start, end): + data = self.fs.blobs[self.path][start:end] + self.fs.fetched += len(data) + return data + + +class CountingFS(AbstractFileSystem): + protocol = 'counting' + # fsspec reuses a filesystem instance built with the same arguments; each + # test needs its own byte count + cachable = False + + def __init__(self, blobs): + super().__init__() + self.blobs = blobs + self.fetched = 0 + + def _open(self, path, mode='rb', block_size=None, cache_type='readahead', + cache_options=None, **kwargs): + return CountingFile(self, path, mode, block_size=block_size or 5 * 2**20, + cache_type=cache_type, cache_options=cache_options, + size=len(self.blobs[path])) + + +@pytest.fixture(scope='module') +def granule(tmp_path_factory): + path = tmp_path_factory.mktemp('remote_cache') / 'granule.h5' + rng = np.random.default_rng(0) + with h5py.File(path, 'w') as h5f: + # create every field first and fill them afterwards, so the chunks + # are laid out apart from the metadata, as in a written granule + dsets = [h5f.create_dataset(name, shape=(N,), dtype='f8', chunks=(10_000,), + compression='gzip') for name in NAMES] + for ds in dsets: + ds[:] = rng.normal(size=N) + return path.read_bytes() + + +def read_window(fd): + with h5py.File(fd, 'r') as h5f: + return {name: h5f[name][WINDOW] for name in NAMES} + + +def test_windowed_read_does_not_refetch(granule): + fs = CountingFS({'granule.h5': granule}) + with pc.io_utils.open_remote('granule.h5', fs=fs, + block_size=pc.io_utils.DEFAULT_REMOTE_BLOCK_SIZE) as fd: + cached = read_window(fd) + + # the same read on the old one-block cache + old = CountingFS({'granule.h5': granule}) + with old.open('granule.h5', 'rb', block_size=pc.io_utils.DEFAULT_REMOTE_BLOCK_SIZE, + cache_type='readahead') as fd: + plain = read_window(fd) + + assert old.fetched > 1.5 * fs.fetched + for name in NAMES: + assert np.array_equal(cached[name], plain[name]) + + +def test_block_cache_requested_only_with_a_block_size(): + seen = [] + + class RecordingFS: + def open(self, path, mode='rb', **kwargs): + seen.append(kwargs) + + pc.io_utils.open_remote('s3://b/k.h5', fs=RecordingFS(), block_size=256 * 1024) + pc.io_utils.open_remote('s3://b/k.h5', fs=RecordingFS()) + assert seen[0]['cache_type'] == 'blockcache' + assert (seen[0]['cache_options']['maxblocks'] * 256 * 1024 + == pc.io_utils.DEFAULT_REMOTE_CACHE_BYTES) + assert seen[1] == {} # a whole-file read keeps the filesystem's defaults