Skip to content

fix: bound in-flight memory at high thread counts (stacked on #723) - #726

Draft
dougnukem wants to merge 6 commits into
OpenGene:masterfrom
dougnukem:arch/pack-size-memory-cap-723
Draft

dougnukem wants to merge 6 commits into
OpenGene:masterfrom
dougnukem:arch/pack-size-memory-cap-723

Conversation

@dougnukem

@dougnukem dougnukem commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

Stacked on #723: this branch is built on top of
fix/backpressure-scales-with-threads-721, so the diff below includes #723's
commits until that one merges -- the only new commit here is the last one
("scale pack size down at high thread counts..."). Opening against master
directly (rather than a fork-internal PR against the unmerged branch) so
this is visible without waiting on #723 first.

Summary

packInMemLimit(threads) (added in #723) can't safely drop below 2*threads
without reintroducing the #721 deadlock, so it has no slack left to also cap
total buffered memory -- that limit already grows linearly with thread count
on purpose. This scales the pack size down instead once thread count
exceeds the point where packInMemLimit's own floor stops mattering, so
packInMemLimit(threads) * packSize(threads) (roughly the total buffered
read count) stays close to constant instead of growing linearly with
threads, with a floor (100 reads/pack) so packs don't shrink far enough for
per-pack overhead to dominate. See the comment on packSize() in
src/common.h for the exact reasoning.

This only changes behavior once thread count exceeds 16 (PACK_IN_MEM_LIMIT / 2) -- below that, packSize() returns the original constant 1000 and
nothing changes.

Validation

24-run sweep, fix/backpressure-scales-with-threads-721 (base) vs this
branch, using the corpus and scripts from #725 (3 datasets x 4 thread counts
x 2 branches, on a 48-vCPU host):

  • Correctness: zero output-digest mismatches across every run.
  • Performance: timings match within noise at -w 1/16 (packSize is
    unchanged there); a small ~5-10% cost at -w 48 on some datasets, the
    expected tradeoff of processing more, smaller packs per unit of gzip
    output -- traded for materially bounded memory growth at high thread
    counts.
dataset -w 1 -w 16 -w 32 -w 48
atac (base / this branch) 31.3s / 31.8s 7.9s / 7.6s 8.7s / 8.5s 9.5s / 10.5s
wgs (base / this branch) 38.4s / 39.0s 26.0s / 26.4s 26.5s / 26.8s 27.0s / 28.9s
synth (base / this branch) 9.1s / 9.0s 3.0s / 3.0s 3.0s / 3.1s 3.2s / 3.3s

Internal unit tests (fastp test) pass, and SE / interleaved-PE / non-interleaved-PE
modes were each spot-checked at 1/8/32/64 threads against the base branch
with matching output digests.

An idea I tried twice and reverted twice

I also attempted distributing packs to whichever worker queue currently
holds the fewest in-flight packs, instead of strict round-robin (the other
follow-up suggested alongside this one). Two attempts, two different
failure modes, both caught before either went anywhere near a PR:

Attempt 1 -- reverted for a correctness bug. The writer's output
ordering -- both the plain round-robin drain in WriterThread::output()
(mWorkingBufferList = (mWorkingBufferList+1) % threads) and the
parallel-pwrite offset-ring path's mNextSeq[t] sequencing -- both assumed
pack round r always lands on worker r % threads. Distributing to the
least-full queue breaks that assumption and silently reorders output
records (caught via a digest diff between two builds).

Attempt 2 -- reworked the writer to fix that, then reverted for a
stall.
I reworked WriterThread to reassemble output by each pack's
explicit global sequence number, then reinstated least-full-queue
distribution on top. It passed a 40-case correctness matrix under Docker on
an oversubscribed ~10-core host, but on real 48-core hardware a workload the
baseline processes in ~4s had not finished after 90s at -w 48 (fine at
-w 8). I initially read the high sys time as mutex contention; a gdb dump
of the stuck process disproved that: all 53 threads were parked in
timed condition-variable waits and none in pthread_mutex_lock. The real
cause is the SPSC list's rule that a queue's newest item isn't consumable
until another item is pushed behind it (the same rule behind #721): once
every queue held one stranded tail pack, all depths tied, ties went to
queue 0, and every other queue's stranded pack blocked ordered output.

A follow-up that replaces per-worker push queues with workers claiming packs
by sequence number (no stranding possible) is on
dougnukem:arch/pull-model-seq-claiming; it's correct but the speedup is
modest, so it isn't proposed here.

🤖 Generated with Claude Code

dougnukem and others added 6 commits September 22, 2026 22:18
Each worker owns one input list, and SingleProducerSingleConsumerList only
lets a consumer take an item once another item has been produced behind it
(or the producer has finished). Reader backpressure used a fixed
PACK_IN_MEM_LIMIT (32) that is smaller than the worker count on machines
with more than 32 cores. With 33+ workers the readers stop after 33 packs,
before any worker list holds two items, so no pack is ever consumed and all
reader and worker threads wait on the backpressure condition variable
forever. The writer-backlog check has the same shape: the writer drains
worker lists round-robin and waits on a list until its worker produces a
second output, so a fixed limit below the worker count can also block the
readers permanently.

Allow at least two in-flight packs per worker in both the reader/processor
and reader/writer backpressure checks (max(PACK_IN_MEM_LIMIT, 2 * threads)),
for PE, interleaved PE and SE readers. The lock-free list semantics are
unchanged, so this does not reintroduce OpenGene#695.
The reader/writer backlog check still gated on the fixed
PACK_IN_MEM_LIMIT * PACK_SIZE period instead of the new
mPackInMemLimit, so it fired far more often than the (now larger)
buffer actually needs at high thread counts. Not a correctness bug —
the comparisons inside already used mPackInMemLimit — just an
inconsistency worth cleaning up alongside it.

Also removed a stray extra blank line left in common.h.
…mory

packInMemLimit(threads) already can't drop below 2*threads without
reintroducing the OpenGene#721 deadlock, so it has no slack left to also cap total
buffered memory. Scale the per-pack read count (packSize) down instead once
thread count exceeds packInMemLimit's own floor, so packInMemLimit(threads)
* packSize(threads) stays roughly constant instead of growing linearly with
thread count, with a floor so packs don't shrink far enough for per-pack
overhead to dominate.

Verified via a 24-run sweep (fix/backpressure-scales-with-threads-721 vs
this branch, 3 corpus datasets x 4 thread counts, on a 48-vCPU host):
zero output-digest mismatches, timings within noise except a ~5-10% cost at
-w 48 on some datasets, the expected tradeoff for materially bounded memory
growth at high thread counts.

A least-full-queue pack-distribution change was also attempted per the
architecture deep-dive, but reverted: the writer's output ordering (both
the plain round-robin drain and the parallel-pwrite offset-ring path) both
assume pack round r always lands on worker r % threads, so distributing
packs to whichever worker queue is least full silently reorders output
records. Fixing that would require reworking the writer to track explicit
per-chunk sequence numbers rather than inferring order from worker index --
a much larger, riskier change than scoped here, caught via a digest diff
before anything shipped.
…d::output

WriterThread::output() polled its per-worker buffer lists with a hardcoded
usleep(100) when nothing was ready, rather than waking on actual input --
the parallel-pwrite path already uses the more responsive published_seq
polling pattern for its own wait, but the plain round-robin drain path had
no wake signal at all. input()/setInputCompleted() now notify a condition
variable that output() waits on (with the same 100us bound as a safety net
in case a notify is ever missed), so the writer wakes immediately when data
arrives instead of on the next poll tick.

No behavioral change to what gets written or in what order -- verified
against the existing digest baseline across -w 1/8/32/64 on synthetic PE
data.
…d-robin

Reworks WriterThread to reassemble output by each pack's global sequence
number instead of inferring order from which worker produced it, which is
what made an earlier attempt at this change (see OpenGene#726's PR description)
silently reorder output. Every ReadPack now carries its round number
(pack->seq), threaded through processPairEnd/processSingleEnd into every
WriterThread::input() call:

- The plain round-robin drain (WriterThread::output()) is replaced with a
  small reorder ring (mOutputRing) keyed by seq, with a single
  next-expected-seq counter -- any worker can publish any seq, and the
  writer just waits for the next one it needs, mirroring the pattern
  the parallel-pwrite path already used for its offset chain.
- The parallel-pwrite path's mNextSeq[tid] (which derived a worker's
  sequence numbers from its own local counter, tid, tid+threads, ...) is
  replaced with the real passed-in seq directly, since it no longer needs
  to be inferred. mMaxPublishedSeq (a single atomic, updated via CAS)
  replaces the per-worker scan in setInputCompletedPwrite() for finding
  where to truncate the output file.
- input()/inputPwrite() now take both tid (still used to index each
  worker's own persistent compressor/buffer in pwrite mode -- worker
  identity hasn't changed, only which pack rounds it handles) and seq
  (determines output order).

With that in place, packs are distributed to whichever worker queue
currently holds the fewest in-flight packs (mQueueDepth), instead of blind
round-robin, in both processors:
- SE and interleaved PE have a single reader thread, so no coordination
  is needed beyond the depth counters themselves.
- Non-interleaved PE reads left/right on two independent threads that must
  agree on the same queue for a given round (a worker's left/right input
  lists are a fixed pair). assignQueueForRound() caches whichever side
  reaches a round first so the other side reuses the same choice, instead
  of each side picking independently and mismatching the pair.

Validated with a 40-case matrix (8 write-path combinations x 5 thread
counts: plain and gzip/pwrite output, SE, interleaved input, --merge,
--failed_out + --unpaired1/2, --overlapped_out) plus --stdout separately,
comparing byte-for-byte (post-decompression) output against
fix/backpressure-scales-with-threads-721 -- all match, no hangs. Internal
`fastp test` suite passes.
@sfchen

sfchen commented Sep 30, 2026

Copy link
Copy Markdown
Member

Hi @dougnukem many thanks for your good catch of the hang issue, and your great fix. I have already tested the PR #723 and merged it.

Is this PR ready to be merged? Memory bound is also important.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants