From f5776a795a8cd99c4479730e3a092aa3c34f04fe Mon Sep 17 00:00:00 2001 From: zzh Date: Wed, 3 Jun 2026 14:09:51 +0800 Subject: [PATCH] =?UTF-8?q?=E6=94=B9=E8=BF=9B=E4=BA=86merge=EF=BC=8C?= =?UTF-8?q?=E5=BC=80=E5=A7=8B=E7=AC=AC=E4=BA=8C=E9=98=B6=E6=AE=B5=E8=BF=87?= =?UTF-8?q?=E7=A8=8B=E7=9A=84=E5=B7=A5=E4=BD=9C=EF=BC=8C=E9=9C=80=E8=A6=81?= =?UTF-8?q?clean=E4=B8=80=E4=B8=8B=E4=BB=A3=E7=A0=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Copilot --- src/sort/phase_1.cpp | 3 +++ src/sort/phase_1.h | 12 +++++++++++- src/sort/phase_1_compress.cpp | 8 +++++++- src/sort/phase_1_sort.cpp | 7 ++++--- src/sort/phase_1_uncompress.cpp | 4 +++- src/sort/phase_1_write.cpp | 2 ++ src/sort/phase_2.h | 10 +++++++++- src/sort/sort.h | 17 ++++++++++++++++- src/util/profiling.cpp | 1 + src/util/profiling.h | 1 + 10 files changed, 57 insertions(+), 8 deletions(-) diff --git a/src/sort/phase_1.cpp b/src/sort/phase_1.cpp index 860a1c1..e081359 100644 --- a/src/sort/phase_1.cpp +++ b/src/sort/phase_1.cpp @@ -68,5 +68,8 @@ void phase1Pipeline() { PROF_G_END(mid_all); + + + #endif } \ No newline at end of file diff --git a/src/sort/phase_1.h b/src/sort/phase_1.h index c290420..26dfec3 100644 --- a/src/sort/phase_1.h +++ b/src/sort/phase_1.h @@ -99,9 +99,19 @@ struct ThreadUncompressWrap { } }; +struct BlockBams { + DataBuffer blockBuf; // 解压的block放在这里 + BamPtrArr bamPtrArr; // 解析后的bam数据放在这里 + + void Clear() { + blockBuf.Clear(); + bamPtrArr.Clear(); + } +}; + // 用于合并压缩的数据结构 struct MergeCompressData { - vector blockDataArr; // 待压缩的数据 + vector blockDataArr; // 待压缩的数据 vector compressDataArr; // 压缩后的数据 void Resize(int blockNum) { diff --git a/src/sort/phase_1_compress.cpp b/src/sort/phase_1_compress.cpp index aaed600..cfacf16 100644 --- a/src/sort/phase_1_compress.cpp +++ b/src/sort/phase_1_compress.cpp @@ -26,9 +26,15 @@ static void mtCompressBlock(void* data, long idx, int tid) { Phase1PipelineArg& p = *(Phase1PipelineArg*)data; MergeCompressData& mergeCompressData = p.mergeCompressData[p.compressOrder % p.COMPRESS_BUF_NUM]; - auto& blockData = mergeCompressData.blockDataArr[idx]; + auto& bams = mergeCompressData.blockDataArr[idx].bamPtrArr; + auto& blockData = mergeCompressData.blockDataArr[idx].blockBuf; auto& compressData = mergeCompressData.compressDataArr[idx]; + for (int i = 0; i < bams.Size(); ++i) { + const OneBam* bp = bams.arr[i]; + blockData.MemCopy(p.uncompressData.dataBuf + bp->offset, bp->wholeBamLen); + } + // bgzfCompress(void* _dst, size_t* dlen, const void* src, size_t slen, int level) compressData.ReAllocMem(SINGLE_BLOCK_SIZE); // 压缩后的block数据不会超过单个block的大小 compressData.curLen = SINGLE_BLOCK_SIZE; diff --git a/src/sort/phase_1_sort.cpp b/src/sort/phase_1_sort.cpp index c3e14fb..5911ac3 100644 --- a/src/sort/phase_1_sort.cpp +++ b/src/sort/phase_1_sort.cpp @@ -153,15 +153,16 @@ bool doMergeSort(Phase1PipelineArg& p, Phase1BamHeap& heap) { mergeCompressData.blockDataArr[mergedBlockNum].Clear(); // 清理,为添加bam数据做准备 bamBytes = 0; } - mergeCompressData.blockDataArr[mergedBlockNum].MemCopy(p.uncompressData.dataBuf + bam->offset, bam->wholeBamLen); - bamBytes += bam->wholeBamLen; // for test + // mergeCompressData.blockDataArr[mergedBlockNum].MemCopy(p.uncompressData.dataBuf + bam->offset, bam->wholeBamLen); + mergeCompressData.blockDataArr[mergedBlockNum].bamPtrArr.Add(bam); + bamBytes += bam->wholeBamLen; // for test heap.Pop(); } if (bam == nullptr) { // 都处理完 finish = true; } - spdlog::info("mergedBlockNum: {}, bamBytes: {}", mergedBlockNum, bamBytes); + // spdlog::info("mergedBlockNum: {}, bamBytes: {}", mergedBlockNum, bamBytes); return finish; } diff --git a/src/sort/phase_1_uncompress.cpp b/src/sort/phase_1_uncompress.cpp index 209f94f..287ff25 100644 --- a/src/sort/phase_1_uncompress.cpp +++ b/src/sort/phase_1_uncompress.cpp @@ -381,8 +381,10 @@ static void doMemCopy(Phase1PipelineArg& p) { #endif // 开启排序并写入中间文件 + PROF_G_BEG(after_full); phase1Sort(&p); phase1MergeCompress(&p); + PROF_G_END(after_full); p.uncompressData.NextRound(); p.allBams.Clear(); @@ -390,7 +392,7 @@ static void doMemCopy(Phase1PipelineArg& p) { PROF_G_BEG(mem_copy); - p.allBams.Add(uncompressWrap.GetTotalBamNum()); + p.allBams.AddSize(uncompressWrap.GetTotalBamNum()); p.bamNum += uncompressWrap.GetTotalBamNum(); p.blockNum += uncompressWrap.GetTotalBlockNum(); diff --git a/src/sort/phase_1_write.cpp b/src/sort/phase_1_write.cpp index eed7cf4..4bc6dd7 100644 --- a/src/sort/phase_1_write.cpp +++ b/src/sort/phase_1_write.cpp @@ -24,8 +24,10 @@ #include "util/profiling.h" static void doWrite(Phase1PipelineArg& p) { + PROF_G_BEG(write_mid); DataBuffer& compressBuf = p.compressBuf[p.compressOrder % p.COMPRESS_BUF_NUM]; fwrite(compressBuf.data, 1, compressBuf.curLen, p.midFilePtr); + PROF_G_END(write_mid); } void* phase1Write(void* data) { diff --git a/src/sort/phase_2.h b/src/sort/phase_2.h index f562f55..46c7f16 100644 --- a/src/sort/phase_2.h +++ b/src/sort/phase_2.h @@ -12,4 +12,12 @@ Date : 2026/02/08 */ -#pragma once \ No newline at end of file +#pragma once + +struct Phase2File { + FILE* fp; + // 双buffer + // 当前读入的buffer指针df + // ReadBuffer + +}; \ No newline at end of file diff --git a/src/sort/sort.h b/src/sort/sort.h index 52d3da4..01f8a1c 100644 --- a/src/sort/sort.h +++ b/src/sort/sort.h @@ -68,6 +68,12 @@ struct UncompressBlockBuffer { lastEndPos = 0; } + uint8_t *HandOverBuf() { + uint8_t *handOverBuf = dataBuf; + dataBuf = nullptr; // 交出buf的所有权,外部需要负责释放内存 + return handOverBuf; + } + void NextRound() { usedBufSize -= lastEndPos; memcpy(dataBuf, dataBuf + lastEndPos, usedBufSize); @@ -95,7 +101,13 @@ struct FastVector { #endif } } - void Add(size_t num) { + + void Add(const T& item) { + T &newItem = Add(); + newItem = item; + } + + void AddSize(size_t num) { if (curIdx + num > arr.size()) { arr.resize(curIdx + num); } @@ -383,6 +395,7 @@ struct MergeSortData { }; typedef FastVector BamArr; +typedef FastVector BamPtrArr; typedef FastVector BlockArr; @@ -413,9 +426,11 @@ struct FirstPipeArg { ReadBuffer readData[READ_BUF_NUM]; // 用来读如数据,双缓冲 UncompressData uncompressData[UNCOMPRESS_BUF_NUM]; // 每个线程内保留一些自己的数据,比如解压缩的数据 UncompressBlockBuffer unCompblockDataBuf; // 所有线程共用一个,串行往这里添加解压后的block数据 + vector sortedBamArr; vector taskArr; vector threadBuf; + MergeSortData mergeSortData; FirstPipeArg() diff --git a/src/util/profiling.cpp b/src/util/profiling.cpp index 08f18b7..43b9f92 100644 --- a/src/util/profiling.cpp +++ b/src/util/profiling.cpp @@ -58,6 +58,7 @@ int displayProfiling(int nthread) { PRINT_GP(mem_copy); PRINT_GP(parse_block); PRINT_GP(uncompress); + PRINT_GP(after_full); PRINT_GP(sort); PRINT_GP(merge); PRINT_GP(compress); diff --git a/src/util/profiling.h b/src/util/profiling.h index 68f3511..8cd87bc 100644 --- a/src/util/profiling.h +++ b/src/util/profiling.h @@ -71,6 +71,7 @@ enum { GP_mem_copy, GP_parse_block, GP_uncompress, + GP_after_full, GP_sort, GP_merge, GP_compress,