diff --git a/src/sort/phase_1.h b/src/sort/phase_1.h index 63efe1b..3083701 100644 --- a/src/sort/phase_1.h +++ b/src/sort/phase_1.h @@ -117,6 +117,8 @@ struct MergeCompressData { compressDataArr.resize(blockNum); } + int Size() { return curIdx; } + void Clear() { curIdx = 0; } }; diff --git a/src/sort/phase_1_compress.cpp b/src/sort/phase_1_compress.cpp index faba2fc..8624348 100644 --- a/src/sort/phase_1_compress.cpp +++ b/src/sort/phase_1_compress.cpp @@ -23,6 +23,10 @@ #include "sort.h" #include "util/profiling.h" +// for test +uint64_t block_num = 0; +uint64_t bam_num = 0; + void CheckBam(uint8_t *addr, int len) { int bamLen = 0; memcpy(&bamLen, addr, 4); @@ -42,6 +46,7 @@ static void mtCompressBlock(void* data, long idx, int tid) { auto& compressData = mergeCompressData.compressDataArr[idx]; // spdlog::info("bam size: {}", bams.Size()); + bam_num += bams.Size(); for (int i = 0; i < bams.Size(); ++i) { const OneBam* bp = bams.arr[i]; blockData.MemCopy(p.uncompressData.dataBuf + bp->offset, bp->wholeBamLen); @@ -59,15 +64,16 @@ static void doCompress(Phase1PipelineArg& p) { PROF_G_BEG(compress); DataBuffer& compressBuf = p.compressBuf[p.compressOrder % p.COMPRESS_BUF_NUM]; MergeCompressData& mergeCompressData = p.mergeCompressData[p.compressOrder % p.MERGE_BUF_NUM]; - kt_for(p.numThread, mtCompressBlock, &p, mergeCompressData.blockDataArr.size()); - //kt_for(1, mtCompressBlock, &p, mergeCompressData.blockDataArr.size()); + kt_for(p.numThread, mtCompressBlock, &p, mergeCompressData.Size()); + // kt_for(1, mtCompressBlock, &p, mergeCompressData..Size()); PROF_G_END(compress); compressBuf.Clear(); - for (int i=0; i < mergeCompressData.blockDataArr.size(); ++i) { // 如果太慢,可以考虑并行拷贝,先计算偏移量,然后多线程拷贝 + for (int i = 0; i < mergeCompressData.Size(); ++i) { // 如果太慢,可以考虑并行拷贝,先计算偏移量,然后多线程拷贝 compressBuf.MemCopy(mergeCompressData.compressDataArr[i].data, mergeCompressData.compressDataArr[i].curLen); } - // spdlog::info("compress bytes: {}", compressBuf.curLen); + block_num += mergeCompressData.Size(); + // spdlog::info("block num: {}, {}, bam num: {}", mergeCompressData.Size(), block_num, bam_num); } /* phase1Compress step- 压缩线程 */ @@ -96,6 +102,6 @@ void* phase1Compress(void* data) { yarn::UPDATE_SIG_ORDER(p.compressSig, p.compressOrder); } - spdlog::info("End compress order: {}", p.compressOrder); + spdlog::info("End compress order: {}, blocks: {}, bams: {}", p.compressOrder, block_num, bam_num); return nullptr; } diff --git a/src/sort/phase_1_sort.cpp b/src/sort/phase_1_sort.cpp index 312d94a..1a5d197 100644 --- a/src/sort/phase_1_sort.cpp +++ b/src/sort/phase_1_sort.cpp @@ -23,6 +23,9 @@ #include "sort.h" #include "util/profiling.h" +uint64_t sort_bam_id = 0; +uint64_t block_nums = 0; + /* bam 排序堆 */ struct Phase1BamArrIdIdx { size_t idx = 0; // 下一个待读入数据的idx @@ -206,6 +209,7 @@ bool doMergeSort(Phase1PipelineArg& p, Phase1BamHeap& heap) { if (bamBytes + bam->wholeBamLen > singleBlockBytes) { mergedBlockNum++; if (mergedBlockNum >= p.mergeBlocksThreshold) { + mergeCompressData.curIdx = mergedBlockNum; // 这里和下面finish时候不同,要注意 break; } mergeCompressData.blockDataArr[mergedBlockNum].Clear(); // 清理,为添加bam数据做准备 @@ -214,6 +218,7 @@ bool doMergeSort(Phase1PipelineArg& p, Phase1BamHeap& heap) { #if 0 mergeCompressData.blockDataArr[mergedBlockNum].blockBuf.MemCopy(p.uncompressData.dataBuf + bam->offset, bam->wholeBamLen); #else + sort_bam_id += 1; mergeCompressData.blockDataArr[mergedBlockNum].bamPtrArr.Add(bam); #endif bamBytes += bam->wholeBamLen; // for test @@ -222,8 +227,11 @@ bool doMergeSort(Phase1PipelineArg& p, Phase1BamHeap& heap) { if (bam == nullptr) { // 都处理完 finish = true; + mergeCompressData.curIdx = mergedBlockNum + 1; // 记录有多少个block待压缩,主要是最后一些数据凑不满threshold } - // spdlog::info("mergedBlockNum: {}, bamBytes: {}", mergedBlockNum, bamBytes); + + block_nums += mergedBlockNum; + // spdlog::info("mergedBlockNum: {}, bamBytes: {}, blocks: {} bam id: {}", mergedBlockNum, bamBytes, block_nums, sort_bam_id); return finish; } @@ -255,7 +263,8 @@ void* phase1MergeSort(void* data) { PROF_G_END(merge); if (finish) { - yarn::SIGNAL_FINISH(p.mergeSig, p.mergeFinish); + yarn::SIGNAL_FINISH_ADD_ORDER(p.mergeSig, p.mergeFinish, p.mergeOrder); + // yarn::SIGNAL_FINISH(p.mergeSig, p.mergeFinish); break; } diff --git a/src/sort/phase_1_uncompress.cpp b/src/sort/phase_1_uncompress.cpp index 6de5c5d..3cd3805 100644 --- a/src/sort/phase_1_uncompress.cpp +++ b/src/sort/phase_1_uncompress.cpp @@ -494,8 +494,10 @@ void* phase1MemCopy(void* data) { } // 这里需要再检查一次缓冲区,有数据的话需要处理,只需要线程内排序,再线程间归并排序,不需要压缩写入中间文件了 - if (p.allBams.Size() > 0) + if (p.allBams.Size() > 0) { + spdlog::info("last round bams in buf: {}", p.allBams.Size()); handleLastRoundData(p); + } break; } diff --git a/src/sort/phase_1_write.cpp b/src/sort/phase_1_write.cpp index 110c9ef..afe1a68 100644 --- a/src/sort/phase_1_write.cpp +++ b/src/sort/phase_1_write.cpp @@ -89,7 +89,7 @@ void* phase1Write(void* data) { } // 写结尾,空block,bam文件需要 - //fwrite("\037\213\010\4\0\0\0\0\0\377\6\0\102\103\2\0\033\0\3\0\0\0\0\0\0\0\0\0", 1, 28, p.midFilePtr); + // fwrite("\037\213\010\4\0\0\0\0\0\377\6\0\102\103\2\0\033\0\3\0\0\0\0\0\0\0\0\0", 1, 28, p.midFilePtr); spdlog::info("End write order: {}", p.writeOrder); return nullptr; diff --git a/src/sort/phase_2.cpp b/src/sort/phase_2.cpp index bf6d89f..7087b49 100644 --- a/src/sort/phase_2.cpp +++ b/src/sort/phase_2.cpp @@ -18,6 +18,10 @@ #include "phase_2_write.h" #include "util/profiling.h" +uint64_t block_id = 0; +uint64_t read_bam_id = 0; +uint64_t copy_bam_id = 0; +uint64_t write_bam_id = 0; void phase2Pipeline(Phase2PipelineArg& p) { PROF_G_BEG(phase2); @@ -25,13 +29,13 @@ void phase2Pipeline(Phase2PipelineArg& p) { pthread_t tidArr[6]; // 2-stage pipeline pthread_create(&tidArr[0], NULL, phase2ReadMidFile, &p); pthread_create(&tidArr[1], NULL, phase2Uncompress, &p); - pthread_create(&tidArr[2], NULL, phase2CopyToMergeBuf, &p); - pthread_create(&tidArr[3], NULL, phase2Merge, &p); - pthread_create(&tidArr[4], NULL, phase2Compress, &p); - pthread_create(&tidArr[5], NULL, phase2Write, &p); + //pthread_create(&tidArr[2], NULL, phase2CopyToMergeBuf, &p); + //pthread_create(&tidArr[3], NULL, phase2Merge, &p); + //pthread_create(&tidArr[4], NULL, phase2Compress, &p); + //pthread_create(&tidArr[5], NULL, phase2Write, &p); - for (int i = 0; i < 6; ++i) pthread_join(tidArr[i], NULL); - //for (int i = 0; i < 4; ++i) pthread_join(tidArr[i], NULL); + //for (int i = 0; i < 6; ++i) pthread_join(tidArr[i], NULL); + for (int i = 0; i < 2; ++i) pthread_join(tidArr[i], NULL); spdlog::info("all bams num: {}", p.numBam); diff --git a/src/sort/phase_2.h b/src/sort/phase_2.h index 5ca6e53..7260e89 100644 --- a/src/sort/phase_2.h +++ b/src/sort/phase_2.h @@ -23,6 +23,11 @@ #include "sort.h" using std::string; +extern uint64_t block_id; +extern uint64_t read_bam_id; +extern uint64_t copy_bam_id; +extern uint64_t write_bam_id; + // 循环缓冲区 struct CircularBuffer { uint8_t* data = nullptr; diff --git a/src/sort/phase_2_merge.cpp b/src/sort/phase_2_merge.cpp index 9af2c66..264cd2d 100644 --- a/src/sort/phase_2_merge.cpp +++ b/src/sort/phase_2_merge.cpp @@ -199,6 +199,7 @@ static bool doPhase2Merge(Phase2PipelineArg& p, Phase2BamHeap& h mergeData.blockDataArr[mergeData.curIdx].Clear(); while ((bam = heap.Top()) != nullptr) { + write_bam_id += 1; if (bam->addr == nullptr) { spdlog::info("null addr"); } @@ -242,6 +243,8 @@ static bool doPhase2Merge(Phase2PipelineArg& p, Phase2BamHeap& h *emptyFilePtr = emptyFile; } + spdlog::info("phase2 merge sort order: {}, write bam id: {}", p.mergeOrder, write_bam_id); + return finish; } diff --git a/src/sort/phase_2_read.cpp b/src/sort/phase_2_read.cpp index 505e8f4..bcd03ca 100644 --- a/src/sort/phase_2_read.cpp +++ b/src/sort/phase_2_read.cpp @@ -1,5 +1,5 @@ /* - Description: 中间文件读入和解析等 + Description: 中间文件读入和解析等,包含读入数据,解压,拷贝到merge缓冲区 Copyright : All right reserved by ICT @@ -48,7 +48,7 @@ bool doReadMidFiles(Phase2PipelineArg& p) { finishNum += f.finish; //spdlog::info("fid: {}, readyNum: {}", i, f.readyReadBufNum); } - //spdlog::info("phase2 read order: {}, finishNum: {}, files: {}", p.readOrder, finishNum, p.midFiles.size()); + spdlog::info("phase2 read order: {}, finishNum: {}, files: {}", p.readOrder, finishNum, p.midFiles.size()); finish = finishNum == p.midFiles.size(); PROF_G_END(phase2_read); return finish; @@ -145,7 +145,9 @@ static void mtUncompressBlockBatch(void* data, long idx, int tid) { blockBuf.curLen += dlen; // 解析 - ParseAddAllBams(blockBuf.data, 0, blockBuf.curLen, blockData.bams); + int bam_num = ParseAddAllBams(blockBuf.data, 0, blockBuf.curLen, blockData.bams); + // spdlog::error(" bams: {} - {}", bam_num, blockData.bams.Size()); + read_bam_id += bam_num; } //spdlog::info("tid: {}, blockNum: {}, fid: {}, blocks: {}-{}, block range: {}-{}, {}", tid, blockNum, i, start, stop, startIdx, stopIdx, @@ -255,6 +257,7 @@ bool doPhase2Uncompress(Phase2PipelineArg& p) { } } p.curBlockNum = blockNum; + block_id += blockNum; #if 1 kt_for(p.numThread, mtUncompressBlockBatch, &p, p.numThread); @@ -277,6 +280,9 @@ bool doPhase2Uncompress(Phase2PipelineArg& p) { // f.readyUncompressBufNum = 0; } PROF_G_END(phase2_uncompress); + + spdlog::info("phase2 uncompress order: {}, read bams: {}, blocks: {}", p.uncompressOrder, read_bam_id, block_id); + return hasUncompress; // exit(0); @@ -288,11 +294,11 @@ void* phase2Uncompress(void* data) { while (true) { // previous dependency yarn::DEPENDENCY_NOT_TO_BE(p.readSig, 0); - yarn::DEPENDENCY_NOT_TO_BE(p.uncompressSig, Phase2File::UNCOMPRESSS_BUF_NUM); + //yarn::DEPENDENCY_NOT_TO_BE(p.uncompressSig, Phase2File::UNCOMPRESSS_BUF_NUM); if (p.readFinish) { while (p.uncompressOrder < p.readOrder) { - yarn::DEPENDENCY_NOT_TO_BE(p.uncompressSig, Phase2File::UNCOMPRESSS_BUF_NUM); + //yarn::DEPENDENCY_NOT_TO_BE(p.uncompressSig, Phase2File::UNCOMPRESSS_BUF_NUM); doPhase2Uncompress(p); yarn::UPDATE_SIG_ORDER(p.uncompressSig, p.uncompressOrder); } @@ -317,6 +323,7 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) { PROF_G_BEG(phase2_copyToMerge); int usedUncompress = 0; p.uncompressReadyNum = 0; + size_t copied = 0; #if 1 for (int i = 0; i < p.midFiles.size(); ++i) { int usedUncompressInFile = 0; @@ -331,6 +338,7 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) { } int copiedNum = mergeData.CopyBams(uncompressBuf.blockBuf, uncompressBuf.bamArr, uncompressBuf.startIdx); + copied += copiedNum; uncompressBuf.startIdx += copiedNum; if (uncompressBuf.Size() == 0) { f.readyUncompressBufNum -= 1; @@ -341,6 +349,7 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) { auto& uncompressBuf = f.uncompressBuf[f.copyOrder % f.UNCOMPRESSS_BUF_NUM]; if (mergeData.hasSpace(uncompressBuf.Front()->wholeBamLen)) { // merge有空间,uncompress有数据,那就继续添加到merge int copiedNum = mergeData.CopyBams(uncompressBuf.blockBuf, uncompressBuf.bamArr, uncompressBuf.startIdx); + copied += copiedNum; uncompressBuf.startIdx += copiedNum; if (uncompressBuf.Size() == 0) { usedUncompressInFile += 1; @@ -373,7 +382,8 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) { #else p.midFiles[p.copyMergeOrder % p.midFiles.size()].readyUncompressBufNum -= 1; #endif - spdlog::info("copy to merge order: {}-{}", p.copyMergeOrder, usedUncompress); + copy_bam_id += copied; + spdlog::info("copy to merge order: {}-{}, copied: {}, last id: {}", p.copyMergeOrder, usedUncompress, copied, copy_bam_id); PROF_G_END(phase2_copyToMerge); diff --git a/src/sort/phase_2_write.cpp b/src/sort/phase_2_write.cpp index 714a8d2..fe3e48e 100644 --- a/src/sort/phase_2_write.cpp +++ b/src/sort/phase_2_write.cpp @@ -79,6 +79,7 @@ static void doPhase2Write(Phase2PipelineArg& p) { DataBuffer& compressBuf = p.compressBuf[p.writeOrder % p.COMPRESS_BUF_NUM]; fwrite(compressBuf.data, 1, compressBuf.curLen, p.outFilePtr); PROF_G_END(write_final); + spdlog::info("phase2 write order: {}", p.writeOrder); } void* phase2Write(void* data) { diff --git a/src/util/yarn.h b/src/util/yarn.h index 3265663..a2e996c 100644 --- a/src/util/yarn.h +++ b/src/util/yarn.h @@ -166,6 +166,12 @@ void free_lock_(lock_t *, char const *, long); finish = 1; \ twist_(sig, yarn::BY, 1, __FILE__, __LINE__); +#define SIGNAL_FINISH_ADD_ORDER(sig, finish, order) \ + possess_(sig, __FILE__, __LINE__); \ + finish = 1; \ + order += 1; \ + twist_(sig, yarn::BY, 1, __FILE__, __LINE__); + #define CONSUME_SIGNAL(sig) \ possess_(sig, __FILE__, __LINE__); \ twist_(sig, yarn::BY, -1, __FILE__, __LINE__);