Skip to content

pipeline: throughput starves at low array count (parallelism limited by #arrays, not pool size) #154

Description

@nclack

Symptom

damacy throughput (and the decode stage's output rate) collapses when a scenario
has only a few distinct arrays, even though there are many chunks/files available
and the IO pool is large. The pipeline appears unable to extract enough
independent work from a small number of arrays to keep IO + GPU decode saturated.

Evidence

From a cold-cache survey on the fast /mnt/main0 mount (L40), comparing two
real scenarios with the same chunk geometry (32 MiB uint16 decode chunks) but
very different array counts:

arrays throughput decode stage GB/s_out decode load
512 1669 samp/s 51.3 GB/s 99%
2 38 samp/s 3.2 GB/s 100%

Same GPU, same chunk type — the decode stage shows ~100% wall occupancy in both,
but its output rate is ~16x lower with 2 arrays. "Decode-bound" at 3.2 GB/s
means the decode stage is occupied yet starved/serialized, not saturated.

The cleanest single-variable evidence: on the same 2-array scenario, with the
same patches (a CPU tensorstore benchmark was made to reproduce damacy's
xorshift sampling bit-for-bit), a 32-thread CPU reader gets 203 samp/s vs
damacy's 38 — 5.3x faster. tensorstore happily issues 32 concurrent reads
into 2 arrays; damacy does not keep that many independent operations in flight.
So this is a damacy pipeline-parallelism limit, not a property of the data or
chunking. (The 2-array dataset has ~576 distinct chunk files, so per-file
parallelism is available — it just isn't being used.)

Hypothesis (needs confirming)

The read-ahead / decode parallelism keys off the number of distinct arrays (or
shards) present in the lookahead window. With few arrays, in-flight depth can't
grow regardless of n_io_threads / metadata_io_concurrency, so IO and GPU
decode idle behind a serialized per-array dependency. Worth checking against the
counters (metadata_backend_read_max_active, in-flight read depth, distinct
shards per wave) whether depth tracks array count rather than the configured
pool size.

Suggested repro / investigation

The synthetic scenario machinery (bench/scenario.py + gen_dataset.py) can
isolate this cleanly: hold chunk_shape/shard_shape/dtype/codec and the patch
fixed, and sweep n_zarrs (e.g. 1, 2, 4, 16, 64, 512). If throughput and decode
GB/s_out scale with n_zarrs and plateau, that confirms the pipeline derives its
parallelism from array count rather than from available chunks/files. From there:
can the planner/dispatcher interleave more in-flight reads/decodes within a single
array (across its chunks/shards) so that a 1-2 array dataset still saturates the
IO pool and the GPU decode stream?

Why it matters

Several training datasets have very few arrays (e.g. timelapse sets with 2-11
large arrays). On those, damacy currently loses to a CPU reader despite GPU
decode, purely from this starvation — so it's worth fixing independently of the
data-format and backend decisions.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions