Conversation
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.
… of round-robin" This reverts commit 3ab24aa.
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. |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Stacked on #723: this branch is built on top of
fix/backpressure-scales-with-threads-721, so the diff below includes #723'scommits until that one merges -- the only new commit here is the last one
("scale pack size down at high thread counts..."). Opening against
masterdirectly (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 below2*threadswithout 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, sopackInMemLimit(threads) * packSize(threads)(roughly the total bufferedread 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()insrc/common.hfor 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 andnothing changes.
Validation
24-run sweep,
fix/backpressure-scales-with-threads-721(base) vs thisbranch, using the corpus and scripts from #725 (3 datasets x 4 thread counts
x 2 branches, on a 48-vCPU host):
-w 1/16(packSize isunchanged there); a small ~5-10% cost at
-w 48on some datasets, theexpected tradeoff of processing more, smaller packs per unit of gzip
output -- traded for materially bounded memory growth at high thread
counts.
Internal unit tests (
fastp test) pass, and SE / interleaved-PE / non-interleaved-PEmodes 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 theparallel-pwrite offset-ring path's
mNextSeq[t]sequencing -- both assumedpack round
ralways lands on workerr % threads. Distributing to theleast-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
WriterThreadto reassemble output by each pack'sexplicit 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 highsystime as mutex contention; a gdb dumpof the stuck process disproved that: all 53 threads were parked in
timed condition-variable waits and none in
pthread_mutex_lock. The realcause 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 ismodest, so it isn't proposed here.
🤖 Generated with Claude Code