|
1 | 1 | #include "peprocessor.h" |
| 2 | +#include "hwy/contrib/thread_pool/futex.h" |
2 | 3 | #include "fastqreader.h" |
3 | 4 | #include <iostream> |
4 | 5 | #include <print> |
@@ -42,6 +43,7 @@ PairEndProcessor::PairEndProcessor(Options* opt) |
42 | 43 | mLeftPackReadCounter = 0; |
43 | 44 | mRightPackReadCounter = 0; |
44 | 45 | mPackProcessedCounter = 0; |
| 46 | + mPackProducedCounter = 0; |
45 | 47 |
|
46 | 48 | mLeftReadPool = new ReadPool(mOptions); |
47 | 49 | mRightReadPool = new ReadPool(mOptions); |
@@ -686,7 +688,7 @@ bool PairEndProcessor::processPairEnd(ReadPack* leftPack, ReadPack* rightPack, T |
686 | 688 | delete rightPack; |
687 | 689 |
|
688 | 690 | mPackProcessedCounter.fetch_add(1, std::memory_order_release); |
689 | | - mPackProcessedCounter.notify_all(); |
| 691 | + hwy::WakeAll(mPackProcessedCounter); |
690 | 692 |
|
691 | 693 | return true; |
692 | 694 | } |
@@ -758,6 +760,8 @@ void PairEndProcessor::readerTask(bool isLeft) |
758 | 760 | mRightInputLists[mRightPackReadCounter % mOptions->thread]->produce(pack); |
759 | 761 | mRightPackReadCounter++; |
760 | 762 | } |
| 763 | + mPackProducedCounter.fetch_add(1, std::memory_order_release); |
| 764 | + hwy::WakeAll(mPackProducedCounter); |
761 | 765 | data = NULL; |
762 | 766 | if(read) { |
763 | 767 | delete read; |
@@ -794,32 +798,32 @@ void PairEndProcessor::readerTask(bool isLeft) |
794 | 798 | mRightInputLists[mRightPackReadCounter % mOptions->thread]->produce(pack); |
795 | 799 | mRightPackReadCounter++; |
796 | 800 | } |
| 801 | + mPackProducedCounter.fetch_add(1, std::memory_order_release); |
| 802 | + hwy::WakeAll(mPackProducedCounter); |
797 | 803 |
|
798 | 804 | //re-initialize data for next pack |
799 | 805 | data = new Read*[PACK_SIZE]; |
800 | 806 | memset(data, 0, sizeof(Read*)*PACK_SIZE); |
801 | 807 | // if the processor is far behind this reader, wait to limit memory usage |
802 | 808 | if(isLeft) { |
803 | 809 | while(mLeftPackReadCounter - mPackProcessedCounter.load(std::memory_order_acquire) > PACK_IN_MEM_LIMIT){ |
804 | | - long cur = mPackProcessedCounter.load(std::memory_order_acquire); |
805 | | - mPackProcessedCounter.wait(cur, std::memory_order_acquire); |
| 810 | + uint32_t cur = mPackProcessedCounter.load(std::memory_order_acquire); |
| 811 | + hwy::BlockUntilDifferent(cur, mPackProcessedCounter); |
806 | 812 | slept++; |
807 | 813 | } |
808 | 814 | } else { |
809 | 815 | while(mRightPackReadCounter - mPackProcessedCounter.load(std::memory_order_acquire) > PACK_IN_MEM_LIMIT){ |
810 | | - long cur = mPackProcessedCounter.load(std::memory_order_acquire); |
811 | | - mPackProcessedCounter.wait(cur, std::memory_order_acquire); |
| 816 | + uint32_t cur = mPackProcessedCounter.load(std::memory_order_acquire); |
| 817 | + hwy::BlockUntilDifferent(cur, mPackProcessedCounter); |
812 | 818 | slept++; |
813 | 819 | } |
814 | 820 | } |
815 | 821 | readNum += count; |
816 | 822 | // if the writer threads are far behind this producer, sleep and wait |
817 | 823 | // check this only when necessary |
818 | 824 | if(readNum % (PACK_SIZE * PACK_IN_MEM_LIMIT) == 0 && mLeftWriter) { |
819 | | - while( (mLeftWriter && mLeftWriter->bufferLength() > PACK_IN_MEM_LIMIT) || (mRightWriter && mRightWriter->bufferLength() > PACK_IN_MEM_LIMIT) ){ |
820 | | - std::this_thread::yield(); |
821 | | - slept++; |
822 | | - } |
| 825 | + if(mLeftWriter) mLeftWriter->waitForBufferBelow(PACK_IN_MEM_LIMIT); |
| 826 | + if(mRightWriter) mRightWriter->waitForBufferBelow(PACK_IN_MEM_LIMIT); |
823 | 827 | } |
824 | 828 | // reset count to 0 |
825 | 829 | count = 0; |
@@ -900,6 +904,9 @@ void PairEndProcessor::interleavedReaderTask() |
900 | 904 | mRightInputLists[mRightPackReadCounter % mOptions->thread]->produce(packRight); |
901 | 905 | mRightPackReadCounter++; |
902 | 906 |
|
| 907 | + mPackProducedCounter.fetch_add(1, std::memory_order_release); |
| 908 | + hwy::WakeAll(mPackProducedCounter); |
| 909 | + |
903 | 910 | dataLeft = NULL; |
904 | 911 | dataRight = NULL; |
905 | 912 | break; |
@@ -931,25 +938,26 @@ void PairEndProcessor::interleavedReaderTask() |
931 | 938 | mRightInputLists[mRightPackReadCounter % mOptions->thread]->produce(packRight); |
932 | 939 | mRightPackReadCounter++; |
933 | 940 |
|
| 941 | + mPackProducedCounter.fetch_add(1, std::memory_order_release); |
| 942 | + hwy::WakeAll(mPackProducedCounter); |
| 943 | + |
934 | 944 | //re-initialize data for next pack |
935 | 945 | dataLeft = new Read*[PACK_SIZE]; |
936 | 946 | dataRight = new Read*[PACK_SIZE]; |
937 | 947 | memset(dataLeft, 0, sizeof(Read*)*PACK_SIZE); |
938 | 948 | memset(dataRight, 0, sizeof(Read*)*PACK_SIZE); |
939 | 949 | // if the consumer is far behind this producer, wait to limit memory usage |
940 | 950 | while(mLeftPackReadCounter - mPackProcessedCounter.load(std::memory_order_acquire) > PACK_IN_MEM_LIMIT){ |
941 | | - long cur = mPackProcessedCounter.load(std::memory_order_acquire); |
942 | | - mPackProcessedCounter.wait(cur, std::memory_order_acquire); |
| 951 | + uint32_t cur = mPackProcessedCounter.load(std::memory_order_acquire); |
| 952 | + hwy::BlockUntilDifferent(cur, mPackProcessedCounter); |
943 | 953 | slept++; |
944 | 954 | } |
945 | 955 | readNum += count; |
946 | 956 | // if the writer threads are far behind this producer, sleep and wait |
947 | 957 | // check this only when necessary |
948 | 958 | if(readNum % (PACK_SIZE * PACK_IN_MEM_LIMIT) == 0 && mLeftWriter) { |
949 | | - while( (mLeftWriter && mLeftWriter->bufferLength() > PACK_IN_MEM_LIMIT) || (mRightWriter && mRightWriter->bufferLength() > PACK_IN_MEM_LIMIT) ){ |
950 | | - std::this_thread::yield(); |
951 | | - slept++; |
952 | | - } |
| 959 | + if(mLeftWriter) mLeftWriter->waitForBufferBelow(PACK_IN_MEM_LIMIT); |
| 960 | + if(mRightWriter) mRightWriter->waitForBufferBelow(PACK_IN_MEM_LIMIT); |
953 | 961 | } |
954 | 962 | // reset count to 0 |
955 | 963 | count = 0; |
@@ -1009,7 +1017,8 @@ void PairEndProcessor::processorTask(ThreadConfig* config) |
1009 | 1017 | } else if(inputRight->isProducerFinished() && !inputRight->canBeConsumed()) { |
1010 | 1018 | break; |
1011 | 1019 | } else { |
1012 | | - std::this_thread::yield(); |
| 1020 | + uint32_t cur = mPackProducedCounter.load(std::memory_order_acquire); |
| 1021 | + hwy::BlockUntilDifferent(cur, mPackProducedCounter); |
1013 | 1022 | } |
1014 | 1023 | } |
1015 | 1024 | inputLeft->setConsumerFinished(); |
|
0 commit comments