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
32 changes: 26 additions & 6 deletions pointCollection/io_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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().
Expand All @@ -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):
"""
Expand Down
101 changes: 101 additions & 0 deletions tests/test_remote_cache.py
Original file line number Diff line number Diff line change
@@ -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
Loading