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).
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 × outer_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.
where an outer chunk is sample_chunk × ∏inner_shape × itemsize and a stored chunk is
sample_chunk × ∏inner_chunk × itemsize.
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.
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 staying nonzero 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 17 x pinned = 544.0 MiB, 100 lent, 0 allocated
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.
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.
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 × outer_chunk_bytes≤ what you have). 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 and the
scatter memcpy. 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.