From 4ee79c3081322bc463aa70ad6c6d599581372446 Mon Sep 17 00:00:00 2001 From: Douglas Daniels Date: Tue, 22 Sep 2026 22:18:39 -0500 Subject: [PATCH 1/2] fix: scale reader backpressure with worker thread count (#721) 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 #695. --- src/common.h | 9 +++++++++ src/peprocessor.cpp | 11 ++++++----- src/peprocessor.h | 1 + src/seprocessor.cpp | 5 +++-- src/seprocessor.h | 1 + 5 files changed, 20 insertions(+), 7 deletions(-) diff --git a/src/common.h b/src/common.h index 19416a1e..df259edc 100644 --- a/src/common.h +++ b/src/common.h @@ -36,6 +36,15 @@ static const int PACK_SIZE = 1000; // similar peak memory: 32 × 1000 reads ≈ 32K reads in flight ≈ old 128 × 256. static const int PACK_IN_MEM_LIMIT = 32; +// Each worker thread owns one input list, and a list's newest item only becomes +// consumable once another item is produced behind it (or the producer finishes). +// Reader backpressure must therefore allow at least one in-flight pack per worker; +// otherwise, with more workers than PACK_IN_MEM_LIMIT, readers stop before any +// worker has a consumable pack and all threads wait forever (#721). +inline long packInMemLimit(int threads) { + return threads * 2L > PACK_IN_MEM_LIMIT ? threads * 2L : PACK_IN_MEM_LIMIT; +} + // different filtering results, bigger number means worse // if r1 and r2 are both failed, then the bigger one of the two results will be recorded diff --git a/src/peprocessor.cpp b/src/peprocessor.cpp index c3355525..18d0ae6e 100644 --- a/src/peprocessor.cpp +++ b/src/peprocessor.cpp @@ -15,6 +15,7 @@ PairEndProcessor::PairEndProcessor(Options* opt){ mOptions = opt; + mPackInMemLimit = packInMemLimit(mOptions->thread); mLeftReaderFinished = false; mRightReaderFinished = false; mFinishedThreads = 0; @@ -820,12 +821,12 @@ void PairEndProcessor::readerTask(bool isLeft) { std::unique_lock lk(mBackpressureMtx); if(isLeft) { - while(mLeftPackReadCounter - mPackProcessedCounter.load(std::memory_order_acquire) > PACK_IN_MEM_LIMIT){ + while(mLeftPackReadCounter - mPackProcessedCounter.load(std::memory_order_acquire) > mPackInMemLimit){ slept++; mBackpressureCV.wait_for(lk, std::chrono::milliseconds(1)); } } else { - while(mRightPackReadCounter - mPackProcessedCounter.load(std::memory_order_acquire) > PACK_IN_MEM_LIMIT){ + while(mRightPackReadCounter - mPackProcessedCounter.load(std::memory_order_acquire) > mPackInMemLimit){ slept++; mBackpressureCV.wait_for(lk, std::chrono::milliseconds(1)); } @@ -836,7 +837,7 @@ void PairEndProcessor::readerTask(bool isLeft) // check this only when necessary if(readNum % (PACK_SIZE * PACK_IN_MEM_LIMIT) == 0 && mLeftWriter) { std::unique_lock lk(mBackpressureMtx); - while( (mLeftWriter && mLeftWriter->bufferLength() > PACK_IN_MEM_LIMIT) || (mRightWriter && mRightWriter->bufferLength() > PACK_IN_MEM_LIMIT) ){ + while( (mLeftWriter && mLeftWriter->bufferLength() > mPackInMemLimit) || (mRightWriter && mRightWriter->bufferLength() > mPackInMemLimit) ){ slept++; mBackpressureCV.wait_for(lk, std::chrono::milliseconds(1)); } @@ -962,7 +963,7 @@ void PairEndProcessor::interleavedReaderTask() // if the consumer is far behind this producer, sleep and wait to limit memory usage { std::unique_lock lk(mBackpressureMtx); - while(mLeftPackReadCounter - mPackProcessedCounter.load(std::memory_order_acquire) > PACK_IN_MEM_LIMIT){ + while(mLeftPackReadCounter - mPackProcessedCounter.load(std::memory_order_acquire) > mPackInMemLimit){ slept++; mBackpressureCV.wait_for(lk, std::chrono::milliseconds(1)); } @@ -972,7 +973,7 @@ void PairEndProcessor::interleavedReaderTask() // check this only when necessary if(readNum % (PACK_SIZE * PACK_IN_MEM_LIMIT) == 0 && mLeftWriter) { std::unique_lock lk(mBackpressureMtx); - while( (mLeftWriter && mLeftWriter->bufferLength() > PACK_IN_MEM_LIMIT) || (mRightWriter && mRightWriter->bufferLength() > PACK_IN_MEM_LIMIT) ){ + while( (mLeftWriter && mLeftWriter->bufferLength() > mPackInMemLimit) || (mRightWriter && mRightWriter->bufferLength() > mPackInMemLimit) ){ slept++; mBackpressureCV.wait_for(lk, std::chrono::milliseconds(1)); } diff --git a/src/peprocessor.h b/src/peprocessor.h index 707e41cb..887b5984 100644 --- a/src/peprocessor.h +++ b/src/peprocessor.h @@ -64,6 +64,7 @@ class PairEndProcessor{ size_t mLeftPackReadCounter; size_t mRightPackReadCounter; alignas(128) atomic_long mPackProcessedCounter; + long mPackInMemLimit; ReadPool* mLeftReadPool; ReadPool* mRightReadPool; atomic_bool shouldStopReading; diff --git a/src/seprocessor.cpp b/src/seprocessor.cpp index 76663b02..eed0411a 100644 --- a/src/seprocessor.cpp +++ b/src/seprocessor.cpp @@ -14,6 +14,7 @@ SingleEndProcessor::SingleEndProcessor(Options* opt){ mOptions = opt; + mPackInMemLimit = packInMemLimit(mOptions->thread); mReaderFinished = false; mFinishedThreads = 0; mFilter = new Filter(opt); @@ -382,7 +383,7 @@ void SingleEndProcessor::readerTask() // if the processor is far behind this reader, sleep and wait to limit memory usage { std::unique_lock lk(mBackpressureMtx); - while( mPackReadCounter - mPackProcessedCounter.load(std::memory_order_acquire) > PACK_IN_MEM_LIMIT){ + while( mPackReadCounter - mPackProcessedCounter.load(std::memory_order_acquire) > mPackInMemLimit){ slept++; mBackpressureCV.wait_for(lk, std::chrono::milliseconds(1)); } @@ -392,7 +393,7 @@ void SingleEndProcessor::readerTask() // check this only when necessary if(readNum % (PACK_SIZE * PACK_IN_MEM_LIMIT) == 0 && mLeftWriter) { std::unique_lock lk(mBackpressureMtx); - while(mLeftWriter->bufferLength() > PACK_IN_MEM_LIMIT) { + while(mLeftWriter->bufferLength() > mPackInMemLimit) { slept++; mBackpressureCV.wait_for(lk, std::chrono::milliseconds(1)); } diff --git a/src/seprocessor.h b/src/seprocessor.h index b3c71521..8d32a0eb 100644 --- a/src/seprocessor.h +++ b/src/seprocessor.h @@ -50,6 +50,7 @@ class SingleEndProcessor{ SingleProducerSingleConsumerList** mInputLists; size_t mPackReadCounter; alignas(128) atomic_long mPackProcessedCounter; + long mPackInMemLimit; ReadPool* mReadPool; std::mutex mBackpressureMtx; std::condition_variable mBackpressureCV; From 2a14542ffedbf7fd637b2142b52264b8e2a11d11 Mon Sep 17 00:00:00 2001 From: Douglas Daniels Date: Sun, 27 Sep 2026 15:39:11 -0500 Subject: [PATCH 2/2] fix: use the scaled backpressure limit for the writer-check period too MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- src/common.h | 1 - src/peprocessor.cpp | 4 ++-- src/seprocessor.cpp | 2 +- 3 files changed, 3 insertions(+), 4 deletions(-) diff --git a/src/common.h b/src/common.h index df259edc..67567885 100644 --- a/src/common.h +++ b/src/common.h @@ -45,7 +45,6 @@ inline long packInMemLimit(int threads) { return threads * 2L > PACK_IN_MEM_LIMIT ? threads * 2L : PACK_IN_MEM_LIMIT; } - // different filtering results, bigger number means worse // if r1 and r2 are both failed, then the bigger one of the two results will be recorded // we reserve some gaps for future types to be added diff --git a/src/peprocessor.cpp b/src/peprocessor.cpp index 18d0ae6e..e5a3a21c 100644 --- a/src/peprocessor.cpp +++ b/src/peprocessor.cpp @@ -835,7 +835,7 @@ void PairEndProcessor::readerTask(bool isLeft) readNum += count; // if the writer threads are far behind this producer, sleep and wait // check this only when necessary - if(readNum % (PACK_SIZE * PACK_IN_MEM_LIMIT) == 0 && mLeftWriter) { + if(readNum % (PACK_SIZE * mPackInMemLimit) == 0 && mLeftWriter) { std::unique_lock lk(mBackpressureMtx); while( (mLeftWriter && mLeftWriter->bufferLength() > mPackInMemLimit) || (mRightWriter && mRightWriter->bufferLength() > mPackInMemLimit) ){ slept++; @@ -971,7 +971,7 @@ void PairEndProcessor::interleavedReaderTask() readNum += count; // if the writer threads are far behind this producer, sleep and wait // check this only when necessary - if(readNum % (PACK_SIZE * PACK_IN_MEM_LIMIT) == 0 && mLeftWriter) { + if(readNum % (PACK_SIZE * mPackInMemLimit) == 0 && mLeftWriter) { std::unique_lock lk(mBackpressureMtx); while( (mLeftWriter && mLeftWriter->bufferLength() > mPackInMemLimit) || (mRightWriter && mRightWriter->bufferLength() > mPackInMemLimit) ){ slept++; diff --git a/src/seprocessor.cpp b/src/seprocessor.cpp index eed0411a..2a0f342b 100644 --- a/src/seprocessor.cpp +++ b/src/seprocessor.cpp @@ -391,7 +391,7 @@ void SingleEndProcessor::readerTask() readNum += count; // if the writer threads are far behind this reader, sleep and wait // check this only when necessary - if(readNum % (PACK_SIZE * PACK_IN_MEM_LIMIT) == 0 && mLeftWriter) { + if(readNum % (PACK_SIZE * mPackInMemLimit) == 0 && mLeftWriter) { std::unique_lock lk(mBackpressureMtx); while(mLeftWriter->bufferLength() > mPackInMemLimit) { slept++;