diff --git a/src/common.h b/src/common.h index 19416a1e..67567885 100644 --- a/src/common.h +++ b/src/common.h @@ -36,6 +36,14 @@ 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..e5a3a21c 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)); } @@ -834,9 +835,9 @@ 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() > 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)); } @@ -970,9 +971,9 @@ 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() > 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..2a0f342b 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)); } @@ -390,9 +391,9 @@ 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() > 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;