Skip to content

Tuning: batch, shuffle window, concurrency, memory

This page is practical guidance for setting up InSituDataset on your own store. For why the engine behaves this way — the read plan, the pool, the prefetch pipeline — see Architecture; for the measurements behind the numbers here, see Benchmarks.

Read Two operating points first if you are serving inference rather than training — several defaults below are chosen for training and are the wrong call for a latency-sensitive single pass.

The mental model

Hold these two sentences and the rest follows:

  • A sample is one slice of the sample axis — a timestep, an observation, a model state, a microscopy Z-plane, whatever your rows are (it can be any single physical axis, not just axis 0) — spanning the whole inner extent of its chunk.
  • A batch draws batch_size shuffled samples from a window of block_chunks sample-axis chunks that the loader keeps decoded in memory at once.

So the loader reads each chunk once, holds a rolling window of them, and serves shuffled batches out of that window — concurrency fills the window, the window bounds memory.

Two operating points

The same knobs point in different directions depending on what you are optimizing, and the difference is large enough that guidance for one is actively wrong for the other.

  • Training optimizes steady-state throughput per GB of RAM. It runs for hours, so the first batch costs nothing, and it re-reads the data every epoch, so caching and shuffle quality matter.
  • Inference / scoring optimizes time to first batch. It is single-pass (no shuffle to protect, nothing to cache across epochs) and often short-lived, so a cold start you could ignore in training is the whole latency budget.

max_inflight is where they diverge, and not in the way its name suggests. Sweeping it over real WeatherBench2 ERA5 on a GPU loop:

max_inflight samples/s stall peak RSS cold TTFB warm TTFB
1 1332.1 3.7% 1095 MB 3462 ms 87 ms
4 1330.3 3.6% 1121 MB 656 ms 85 ms
16 1326.9 3.7% 1198 MB 347 ms 88 ms
32 (default) 1332.1 3.4% 1246 MB 317 ms 80 ms

Steady-state throughput does not move at all — 1332.1 at either end, stall flat — while cold TTFB moves 11×. On a loop with enough compute per byte, read-ahead is a cold-start dial, not a throughput dial. Dropping it to 1 buys ~150 MB (12%) for free in training, and costs an inference service more than three seconds on its first request.

The corollary for the other two axes: block_chunks is a training knob (it buys shuffle quality, and only .train shuffles — eval views are deterministic), and the cross-epoch cache is a training knob (single-pass scoring never reads a chunk twice).

block_chunks guidance on this page is under revision

This page presents block_chunks as a shuffle-and-residency knob — the sizing advice below, and the per-regime rows at the end, all follow from that. It also affects throughput, by a mechanism the page does not yet describe: a batch may draw from any chunk in its shuffle-block, so the whole block must be co-resident to gather, and on a store whose read latency has a long tail one slow read can hold the concurrency window open for the rest of the block. Wider blocks amortise that; on a pass already saturated downstream of admission they only cost residency.

Which way it lands is therefore a property of the store and the box, not of the loader, and the guidance here is not yet conditioned on it. Until it is, treat the block_chunks sizing advice below as necessary but not sufficient — correct about RAM and shuffle quality, silent about read-latency tails — and measure both ends of your range rather than taking a default from this page.

The knobs you set

All of these are InSituDataset(...) arguments except the last, which is fixed when the store is written.

plain name argument what it does default
batch size batch_size samples per batch 32
shuffle window block_chunks outer chunks held resident + shuffled across at once 16, raised if needed to hold one batch
reads in flight max_inflight concurrent stored-chunk GETs — sets cold start; sets throughput only while IO-bound 32
batch queue prefetch_depth assembled batches queued ahead of your training step 2
cache cache_budget_bytes, cache_dir decoded data retained across epochs (decode-once) off
stored-chunk size inner_chunks (write time) the fetch unit — how each chunk is split for IO —

batch_size is the ordinary ML knob and barely touches IO. block_chunks trades shuffle quality against RAM. max_inflight saturates the network while you are IO-bound — once compute is the constraint it stops buying throughput and buys cold start instead, which is why the table above matters more than its name suggests. The cache is how repeated epochs (or repeated scoring passes) skip re-reading. inner_chunks is the one decision made when the data is written, and it sets how cheap concurrency can be.

The memory model

Peak memory is the sum of three independently-bounded pieces — none grows with epoch length or dataset size:

  • Shuffle window: block_chunks × resident_chunk_bytes — the decoded chunks held resident.
  • Reads in flight: max_inflight × stored_chunk_bytes — the fetch pipeline.
  • Batch queue: prefetch_depth × batch_bytes — assembled batches awaiting the consumer.

InSituDataset.print_summary() evaluates all of this for your geometry and configuration, before any read happens — including the ragged multiplier below, the run length gather will get, and an estimated peak. Reach for it rather than doing the arithmetic by hand:

ds = InSituDataset(store, manifest, batch_size=32, block_chunks=16)
ds.print_summary()              # ds.print_summary(iterations=2) if you will zip two views

where a stored chunk is sample_chunk × ∏inner_chunk × itemsize, and a resident chunk is the stored chunks that compose it, held whole: n_stored_chunks × stored_chunk_bytes. That equals the logical sample_chunk × ∏inner_shape × itemsize exactly when the chunk grid divides the array evenly, and exceeds it when it does not — see Ragged chunk grids below.

The point to internalize: raising concurrency costs stored-chunk-sized memory, not outer-chunk-sized — but only when the data is inner (spatially) chunked. If each outer chunk is a single stored chunk, those two sizes are equal and concurrency gets expensive (the "fat, single inner" regime below).

If you set cache_budget_bytes above the working set, residency rises to that budget on purpose — that extra memory is the cross-epoch cache. Point cache_dir at local NVMe to spill it to disk instead of RAM.

Several iterations at once multiply the budget

The three bounds above describe one iteration. One InSituDataset owns one chunk pool, and every active iteration shares it — zip(ds.train, ds.val), or two DataLoaders over the same dataset. That is supported, and chunks a windowed read pulls across a split boundary are decoded once and reused by both. But each iteration holds its own references to the chunks it is working on, so residency is the sum, not the maximum:

cache_budget_bytes  >=  n_concurrent_iterations x (block_chunks x resident_chunk_bytes)

The auto-sized default is deliberately computed for one iteration, and stays that way. The engine cannot know how many iterations you intend to run, so any automatic multiplier would be a guess that silently costs memory for the single-iteration case — which is almost every case. Running several is the explicit choice, so sizing for it is yours too.

Run two without raising the budget and the loader stops with residency budget exhausted: ..., naming how many iterations are sharing the pool. It cannot free a slot, because every resident chunk is legitimately referenced by one of them. Raise cache_budget_bytes (or lower block_chunks, which shrinks each iteration's share) — not max_inflight, which is a concurrency dial and is not what is binding here.

This got stricter once pin accounting became owner-scoped

Before that fix, starting a second iteration silently released the first one's references. That freed budget by accident and hid the requirement — and the cost was that the first iteration's in-use chunks could be evicted mid-gather, delivering plausible wrong data. A budget that appeared to work for zip(ds.train, ds.val) was relying on that bug. Sizing for the sum is the honest requirement.

The batch queue is a high-water mark

Batch buffers are pooled and held for the dataset's lifetime, so that third bound is the largest number of batches ever simultaneously live — not the steady state. A transient burst that holds many at once (list(ds.val) to materialize a split, gradient accumulation over N steps, an exported tensor kept past its batch) permanently raises the resident floor to that burst's size.

It cannot raise peak memory. The floor is exactly what was simultaneously live at the peak, which fresh-allocation-per-batch also held at that instant; pooling only declines to give it back afterwards. If the burst fit, the floor fits.

close() is the only release, and it drops the cross-epoch chunk cache and store session with it — so size the budget for your burst rather than planning to reclaim between phases. The per-epoch log line reports the pool directly, and allocated continuing to climb after the first epoch is the signal that batches are not coming back:

epoch 3 (train): chunks 151/151 hit (100%), peak resident 51;
                 batch buffers (pool) 17 x pinned = 544.0 MiB, 400 lent, 17 allocated

The chunk figures are that pass's own; the buffer figures are the pool's running totals, marked (pool) because one buffer pool serves every iteration open on the dataset. Read allocated as a number that settles — 17 after four epochs of 100 batches each is a pool that converged after the first — rather than one that returns to zero each epoch.

Under pin_host_buffers / as_torch(..., device=...) the floor is page-locked, which the kernel cannot reclaim — so it is bounded separately at RAM/8 by default. Past that the loader hands out ordinary pageable memory and warns once, rather than raising or stalling; raise pin_budget_bytes if you meant to exceed it.

Sharing a cache_dir between processes

One writer at a time, or any number of readers — never both. The pool takes an advisory lock on the directory for its lifetime, so the second opener is refused at construction rather than corrupting the first:

already open second opener
writer writer refused
writer reader (readonly_cache=True) refused
reader writer refused
reader reader allowed

This applies whenever cache_dir is set — with or without persist=True — because the two write the same filenames. Without it, a second process re-admitting a chunk the first has mapped truncates the file underneath it, and the reader carries on with right-shape, right-dtype, wrong numbers.

The workload it is shaped around — one job warms a cache, several score against it — is spelled readonly_cache=True:

# Run once: warm the cache.
warm = InSituDataset(store, manifest, cache_dir="/data/era5-cache", persist=True)
for batch in warm.all: ...
warm.close()

# Then, concurrently, as many readers as you like.
ds = InSituDataset(store, manifest, cache_dir="/data/era5-cache", readonly_cache=True)

A read-only opener writes nothing — no chunk files, no log entries — and a cache miss raises, naming the array and chunk. That is what makes it a contract (this cache is complete for what I am about to read) rather than a performance hint; the usual cause of a miss is a different split, sample_range or transform set than the run that warmed it. reset_stale_cache=True is rejected in this mode. Excluding a reader from a live writer follows from the same promise: a cache still being warmed is not complete, so it is better refused at construction than failed forty batches in.

Invalidation — reset_stale_cache deleting entries, or a run re-admitting chunks it evicted — is therefore something no reader is ever present for. Underneath that, a chunk file is written to a temp name and renamed into place, so a reader holding a mapping keeps its own inode and keeps reading real data whether that file is replaced or deleted. That second layer is what stands where the lock cannot be taken (below): the worst a bypassed lock costs an active reader is a future open — a miss — never wrong numbers.

If you hit the lock

The error names the holder's PID, host and start time. That is a hint read from the lockfile; the lock itself is authoritative, so confirm with:

fuser -v /data/era5-cache/.insitu.lock
lsof /data/era5-cache/.insitu.lock

Do not delete the lockfile. It is the first thing people try and it is actively harmful: the lock is held by the kernel against an open file description, not by the file, so deleting it releases nothing — and it makes the next two processes lock different inodes, reintroducing exactly the corruption the check prevents.

There is no such thing as a stale lock, and so no cleanup procedure. The kernel releases it when the process dies — SIGKILL, OOM and spot preemption included. If you are seeing the error, that process is alive.

What a transform edit invalidates

Cache identity is per array: each manifest entry carries the hash of the chunk_transforms scoped to its array (see applies). Editing a transform scoped to 2m_temperature therefore raises naming ['2m_temperature'], and reset_stale_cache=True deletes just that array's chunk files — 10m_u_component_of_wind and the rest keep theirs, because their bytes did not change.

Two cases that are deliberately not stale:

  • Adding a variable. An array with no entries yet is a cold start for that array. You can widen a configuration without a wipe.
  • Reading a subset. An array this run does not open is left alone entirely — entries and files. Two configurations can share one cache_dir, each warming its own arrays. This is what makes an ablation / feature-importance sweep cheap: drop a variable from the run and its cached chunks are neither validated nor deleted, so the run that puts it back reads it warm. Dropping a variable is not an edit to it.

Narrowing a scope is different, and it is stale: a variable that leaves applies([...], scale) while the run still reads it has cached chunks holding scaled bytes that the new configuration must not serve, so its files are deleted and re-decoded. Only that one variable — the transform is hashed without its scope, so the variables that stayed hash identically.

Put the two together and the rule is: a variable is invalidated when the pipeline over it changes, not when your configuration stops mentioning it.

Re-fitted statistics may not invalidate

Without cloudpickle (insitubatch[cache]) the fingerprint falls back to hashing a transform's source plus its repr, and numpy summarizes any array over 1000 elements in repr. A StandardScaler holding per-gridpoint mean/std therefore hashes the same after a re-fit, and the cache reopens as a hit carrying the old normalization. Install the extra, or pass StandardScaler(..., cache_key=...) with a version you bump on every fit.

A load that drops entries rewrites the log (to a temp name, then a rename), so a reset is not re-read and re-rejected on every subsequent open. That is safe only because one writer holds the lock above.

Work out what one chunk costs before setting a budget

cache_budget_bytes below what one pass needs co-resident raises at construction, naming the floor, rather than being quietly raised to it. Ask the geometry what one chunk costs, then check that against the memory you actually have:

geom = open_geometries(store, variables=["temperature_2m"])["temperature_2m"]
geom.chunk_bytes                      # bytes one sample-axis chunk occupies once resident
2 * block_chunks * geom.chunk_bytes   # the floor: current block + one read-ahead
                                      # (3 blocks when a windowed shuffled pass
                                      #  releases its spill -- see below)

Compare that with free -h (the available column, not free) and leave room for the model and the framework's own allocator. Across several variables the floor is the sum, and a chunk_transform can make the transformed output the binding term instead of the stored tiles — so for anything but a single-variable run, ask describe(), which reports the number the engine will actually use rather than one you assembled by hand.

A windowed split can hold chunks and still draw nothing

A windowed view reads anchor + offset, so anchors whose window runs off the array are dropped. A split lying entirely in that tail keeps its chunks and yields no batches — the same symptom as a split that rounded to zero chunks, and a different cause. describe() reports both numbers, and print_summary() shows them side by side whenever the views are windowed:

splits (chunks)  train 5 (232 drawable)  val 1 (0 drawable)  test 0 (0 drawable)

A split with chunks but nothing drawable is a split-yields-nothing warning, naming the offsets responsible. A split you asked to be empty — fractions=(1.0, 0.0, 0.0) — is not a finding and is not warned about.

Two iterations at once are refused, not left to stall

Every active iteration holds its own references, so N of them need N floors. The budget is sized for one, because the engine cannot know how many you intend to run. Starting a second one the budget cannot hold now raises there, naming the number to pass:

cache_budget_bytes=49152 cannot hold 2 concurrent iterations: each needs 49152 bytes
co-resident, so 2 need 98304. ... Pass cache_budget_bytes=98304 or more, or iterate the
splits one after another rather than together.

zip(ds.train, ds.val) and two DataLoaders both count as two. Iterating the splits one after another — the ordinary epoch — never has two live at once, however many passes it runs.

The arithmetic is complete when the second iteration starts, so it is answered there rather than left to starve mid-epoch: whether an under-sized pool reaches the stall at all depends on how the two interleave, so the same configuration would otherwise succeed and fail on consecutive runs.

Windowed and shuffled: hold the split, or release and re-read

A windowed view reads anchor + offset, and shuffle permutes chunk order, so a chunk one block reads may be needed again many blocks later. Held from first use to last, that makes the resident set close to the whole train split rather than two blocks — the cost of decode-once when several views share one array.

Passing persist=True (with a cache_dir) changes the trade automatically: each block's chunks are handed back as that block drains, and a chunk a later block needs is admitted again from the on-disk cache. Residency falls to the three-block floor. Measured on a 256-chunk split with a three-chunk lead, the floor goes from 259 chunks to 12, and from 316 to 18 for leads {0, 24, 240} — with the same samples delivered.

It is a rule rather than a knob because it is only a good trade when the second admission is local. persist=True is what makes an evicted slot revivable; with cache_dir alone the backing is unlinked on eviction and the re-read pays the fetch and the decode again. Put the cache on NVMe.

The per-epoch summary reports the churn, because a hit rate cannot show it — a re-read served from the cache counts as a hit:

epoch 0 (train): chunks 246/502 hit (49%), peak resident 12,
                 249 re-read (spill released per block, revived from the cache), 479 evicted

Rising re-reads at a steady hit rate mean the budget is churning rather than holding, and the floor is the number to raise.

Do not compute this as sample_chunk_size × prod(inner_shape) × itemsize. Residency is the array's stored tiles, kept whole: a grid that does not divide the array evenly still stores full-size edge chunks and the loader does not clip them. On NOAA GFS analysis — 1440 time steps per chunk over a 721×1440 field, tiled 400×400 — that product gives 5.98 GB while the real cost is 7.37 GB, eight whole tiles. Sizing a budget from the smaller number under-provisions by 19%, and the pool starves mid-epoch, which fails in the shape of a hang rather than an error.

The difference is not a rounding detail on an archive that chunks the sample axis deeply. One chunk of that store is 7.37 GB, so block_chunks=4 asks for 59 GB — and a box that cannot supply it should say so before the first read, not be OOM-killed halfway through an epoch. Passing no budget still sizes itself automatically.

Put cache_dir on local NVMe, not NFS

The cache is an mmap tier, so this is what it is built for. Over a network filesystem it is slow and — more to the point — unarbitrated: flock may be emulated per client, so two processes on different hosts can each believe they hold the write lock. The loader warns when it can detect one (Linux, via the mount table), but detection is not a fix. A network cache_dir and a platform with no POSIX locking (Windows) are the two configurations where two writers can still corrupt each other, and both warn saying so.

Shuffle quality

block_chunks is also the shuffle-quality knob. Each batch is drawn from the samples in the current window — block_chunks × samples-per-chunk of them — so set the window comfortably larger than batch_size; otherwise a batch is just one or two chunks' worth of correlated samples. A window smaller than a batch is not merely poor shuffle, it is unschedulable (the batch would need more blocks resident than the residency floor holds), so the loader raises block_chunks to ceil(batch_size / samples-per-chunk) when you ask for less and logs that it did. On finely chunked stores — one sample per chunk, like per-frame or per-spectrum archives — that floor is what the default 16 becomes. The chunks are re-permuted every epoch, so even a modest window converges toward a full-dataset shuffle over many epochs — the regime training actually runs in. shuffle_quality scores an emitted order 0–1 (1 ≈ global) if you want to measure it; Architecture explains why the block-local shuffle converges.

This section is training-only: only .train shuffles — .val, .test and .all are deterministic — so a scoring pass has no shuffle quality to protect.

Which stage is actually the bottleneck?

The knobs above only help if you turn the right one, and the symptom that sends people to the wrong one is always the same: the batch queue is empty. Slow storage, a saturated decode pool, and a residency budget too small to admit the next chunk all present that way, and they want opposite fixes. ds.last_pass separates them.

for batch in ds.train:
    ...

st = ds.last_pass
print(st.limiting_stage)                    # e.g. "residency"
print(st.times.admission_parked_s)          # the evidence behind it

Or read it from the per-epoch log line (logging.getLogger("insitubatch").setLevel(logging.INFO)), which ends in limited by: <stage> -- <what to do>.

Try it: a misconfigured loader that diagnoses itself

This runs as written — the store is one of the public benchmark stores, no credentials and no build step. era5_c1 is the chunk-per-sample (GRIB-like) end of the family: 6000 full 721x1440 fields, one per chunk.

Keep the sleep. It stands in for a training step, and without one the question is degenerate: a for batch in ds.train: pass loop takes batches as fast as they appear, so the queue is empty on every sample however fast the loader is. The report says so if you drop it, but then all it can tell you is throughput.

import logging
import time

from insitubatch import InSituDataset, obstore_store, open_geometries, split_by_chunk

logging.basicConfig(level=logging.INFO, format="%(message)s")

store = obstore_store(
    "gs://insitubatch-bench-insitubatch/era5_c1.zarr", skip_signature=True
)
geoms = open_geometries(store)
manifest = split_by_chunk(geoms["t2m"], fractions=(0.8, 0.1, 0.1))

# max_inflight=1 is the mistake. One read at a time against an object store.
ds = InSituDataset(store, manifest, batch_size=16, block_chunks=8, max_inflight=1)

for i, batch in enumerate(ds.train):
    time.sleep(0.15)  # stand in for a training step
    if i == 9:
        break
epoch 0 (train): 10 batches in 15.3s | queue 0/2 (empty 100%, fed 0%)
  | inflight peak 1/1 | resident 32 chunks, 127 MiB of 127 MiB | consumer 1.35s
  | wait fetch 14.06s, parked 15.18s | cpu decode 1.05s, gather 0.32s
  | limited by: store -- ... the loader waited on the store. Raise max_inflight ...
    Admission also parked on the residency budget (50%), but that is most likely
    backpressure from this stage rather than a budget that is too small.

Two things to read. The verdict is store, and the parked total — nearly as large as the fetch total — is called out as backpressure rather than a cause: with one read in flight, chunks are fetched slowly and therefore sit resident longer, so the budget stays full (127 MiB of 127 MiB). Raising cache_budget_bytes would fix none of it.

Now raise max_inflight, and keep raising it:

max_inflight wall for 10 batches verdict
1 15.3 s store (+ residency backpressure)
8 1.8 s store
16 1.1 s store, permits saturated
32 1.6 s store, permits saturated

16 is the knee, and 32 is worse. Once every permit is in use the advice changes to say so: "Raise max_inflight and re-measure — but if throughput stops improving, the store or the network is the floor and this is the honest answer rather than a knob left unturned." That is the truth for this run: it reads GCS cross-cloud over the public internet. Run it in-region and this is where it stops.

Note wait fetch rises as the pass gets faster (14.06s at max_inflight=1, 15.86s at 16). That is the summing rule, not a contradiction: 16 tiles now wait concurrently, so the total is task-seconds against a 1.1s wall. Read the ratios between stages, never any stage against the clock.

The rows are in precedence order: a loop body that does essentially nothing is answered first, because "the loader kept up, so your training step is the constraint" is only an answer where a training step exists.

what the report shows limiting stage what to do
consumer_s is a few % of wall the dominant producer stage, reported not diagnosed your loop has no real training step, so read this as a throughput measurement rather than a starved pipeline
batch queue fed (fed_frac high) consumer nothing — the loader kept up and your training step is the constraint. This is the goal
queue empty, fetch_wait_s dominant store raise max_inflight; check the store is in-region and on the fast backend
as above, and inflight_peak == max_inflight store same, but re-measure: if throughput stops improving the network is the floor, not a knob left unturned
queue empty, decode_s / assemble_s dominant decode raise decode_threads; check the chunk_transform is vectorized numpy that releases the GIL
queue empty, admission_parked_s dominant residency raise cache_budget_bytes, or lower batch_size / block_chunks / concurrent iterations
queue empty, gather_s / batch_transform_s dominant gather check the gather run length in describe(), and any batch_transform

residency is only named when nothing else explains the pressure. Parking is backpressure: whatever is slow downstream holds chunks pinned until the budget fills, so a parked total that merely tracks another stage is a symptom. When a runner-up is large enough to explain it, the report names that stage instead and says why.

The residency row is the one worth knowing exists. A budget-starved loader is otherwise indistinguishable from slow storage — it is exactly how the pre-#39 deadlock presented — and admission_parked_s is the only counter that tells them apart. It is the non-terminal neighbour of the state the loader raises residency budget exhausted on: same cause, caught before it becomes provably fatal.

limiting_stage returns "unknown" rather than guessing when the evidence does not separate the candidates — no samples, starvation too rare to matter, or no stage owning enough of the accounted time. A confident wrong verdict costs more than an admission.

Two clocks, on purpose

Waiting is measured with perf_counter (waiting is the quantity); in-thread cost is measured with thread_time (CPU actually burned). A wall clock around a thread hop measures GIL wait, not work — that is how the scatter memcpy once read as 51% of the hot path when its real share is 7.7–10.1%.

So fetch_wait_s is summed across the tiles in flight and will exceed wall time. Read the stages against each other, never against the clock.

The recipe

  1. At write time, pick inner_chunks so a stored chunk is ~10–50 MB: small enough that many reads in flight stay cheap, large enough that per-request overhead doesn't dominate.
  2. Start with the defaults (max_inflight=32, block_chunks=16). 32 reads in flight saturates in-region S3 in most cases.
  3. Size block_chunks to your RAM budget (block_chunks × resident_chunk_bytes ≤ what you have — on a ragged grid that is more than the outer chunk's logical size). Which direction to push it from there is the branch below.
  4. Tune max_inflight by the metric your operating point cares about — raise it until decoded MB/s stops climbing, and check TTFB, because past the IO-bound region only TTFB keeps responding. From the repo you can measure the knee directly:
python -m bench.probe_decode --url <store> --concurrency 1,4,8,16,32

If throughput is flat across that range you are compute-bound, and the knob is now buying cold start alone — see the table above before turning it down to reclaim memory.

  1. Sanity-check concurrency cost (max_inflight × stored_chunk_bytes). If it's large, your stored chunks are too big — chunk the inner dims (step 1).

Then branch on what you are running:

For multi-epoch training

  • Set cache_budget_bytes to hold the split (and cache_dir on NVMe to spill); epoch 0 warms it and later epochs read decode-once.
  • Raise block_chunks as far as RAM allows — it is your shuffle quality.
  • If RAM is tight, max_inflight is the cheapest thing to give up: measured, dropping it to 1 cost nothing in steady-state throughput and returned ~12% of peak RSS. You pay the cold start once, at the start of a run that lasts hours.

For inference / single-pass scoring

  • Leave max_inflight high. It is the cold-start knob, worth ~11× on time to first batch, and the ~150 MB it costs is the cheapest latency you will ever buy. This is the one place not to economize.
  • Keep block_chunks small — there is no shuffle to protect (eval views are deterministic), so the window only needs to hold the read plan.
  • Skip the cross-epoch cache unless you score the same data more than once; a single pass never reads a chunk twice, so the budget is pure overhead.
  • If the first batch still dominates, the remaining lever is at write time: smaller inner_chunks make the first window's reads finish sooner (step 1).

Regimes

regime shape guidance
GRIB (chunk=1) 1 sample/chunk, single inner concurrency follows block_chunks; a worker loader is competitive here single-pass (nothing to amortize), but the cross-epoch cache wins repeated passes
moderate ~8–40 samples/chunk, single inner the common case; max_inflight ≈ 32. insitu's edge grows with sample_chunk (each chunk is read once, not re-decoded per sample)
fat, single inner huge outer chunk, single inner the stored chunk is the outer chunk, so concurrency costs full-chunk memory. Rechunk the inner dims, or shrink sample_chunk
fat, spatial huge outer chunk, inner grid the sweet spot: small stored chunks make high max_inflight cheap; keep block_chunks small for low residency

Advanced: decode threads

decode_threads (on SchedulerConfig) sizes the pool that runs codec decode. It defaults to auto (min(32, cpu+4)) and rarely needs changing; on a busy box ~8 can beat auto by avoiding oversubscription. It is only reachable when you drive a Scheduler directly — InSituDataset uses the auto default. There is no separate inner-fan-out cap: the inner grid is dialed by max_inflight, which fetches at stored-chunk granularity.

The decode pool is process-wide. It is created once and outlives every scheduler, so decode_threads takes effect on the first dataset built in the process; a later, different value is ignored and logs a warning. That is deliberate rather than a limitation: how many threads can usefully decode at once is a property of the machine — cores, memory bandwidth — not of the dataset being read, and two datasets in one process should not run two pools competing for the same cores. If you want a non-default size, set it on the first dataset you create.

Ragged chunk grids cost more than their arrays

A zarr chunk grid that does not divide the array evenly still stores full-size chunks at the edges. An ERA5-shaped array of 721 latitudes chunked at 180 occupies five stored chunks — 900 rows of storage for 721 rows of data.

The pool holds stored chunks whole, padding included, because the buffer unit is the stored chunk (one unit, one shape). So resident bytes per chunk are n_tiles × prod(chunk_shape) × itemsize, which can exceed the logical chunk:

geometry logical chunk actually resident ratio
720×1440 @ 45×90 (divides evenly) 63.3 MiB 63.3 MiB 1.000×
2048×2048 @ 256×256 (divides evenly) 16.0 MiB 16.0 MiB 1.000×
721×1440 @ 180×360 3.96 MiB 4.94 MiB 1.248×
as above, short final outer chunk 1.98 MiB 3.96 MiB 1.997×

The automatic budget accounts for this, and so does resident_bytes — the number you are told is the number you are using. What it means for you is that a ragged grid needs a proportionally larger cache_budget_bytes if you set one by hand. If your grid divides evenly, this section costs you nothing.

Your chunk_transform is unaffected: it receives the logical chunk, clipped, as a real contiguous array. It never sees padding, never masks an edge, and never needs to know the chunk grid.