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_sizeshuffled samples from a window ofblock_chunkssample-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:
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:
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:
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¶
- At write time, pick
inner_chunksso 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. - Start with the defaults (
max_inflight=32,block_chunks=16). 32 reads in flight saturates in-region S3 in most cases. - Size
block_chunksto 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. - Tune
max_inflightby 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:
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.
- 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_bytesto hold the split (andcache_diron NVMe to spill); epoch 0 warms it and later epochs read decode-once. - Raise
block_chunksas far as RAM allows — it is your shuffle quality. - If RAM is tight,
max_inflightis 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_inflighthigh. 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_chunkssmall — 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_chunksmake 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.