From eb3acb953ad2a5a067079320d49ab6129c8ffda0 Mon Sep 17 00:00:00 2001 From: zzh Date: Mon, 6 Jul 2026 00:06:49 +0800 Subject: [PATCH] =?UTF-8?q?=E6=B7=BB=E5=8A=A0=E4=BA=86=E4=B8=80=E4=BA=9B?= =?UTF-8?q?=E7=BB=9F=E8=AE=A1=E6=97=B6=E9=97=B4=E7=9A=84=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_2_merge.cpp | 27 +++++++++++++++------------ src/sort/phase_2_read.cpp | 4 ++-- src/sort/phase_2_write.cpp | 6 +++--- src/util/profiling.cpp | 6 ++++-- src/util/profiling.h | 1 + 5 files changed, 25 insertions(+), 19 deletions(-) diff --git a/src/sort/phase_2_merge.cpp b/src/sort/phase_2_merge.cpp index 9e4d28d..8b83822 100644 --- a/src/sort/phase_2_merge.cpp +++ b/src/sort/phase_2_merge.cpp @@ -186,6 +186,8 @@ static void mtCopyBams(void* data, long idx, int tid) { Phase2PipelineArg& p = *(Phase2PipelineArg*)data; MergeCompressData& mergeData = p.mergeData[p.mergeOrder % p.MERGE_BUF_NUM]; + idx += mergeData.lastRoundIdx; + auto& bams = mergeData.blockDataArr[idx].bamPtrArr; auto& blockData = mergeData.blockDataArr[idx].blockBuf; @@ -262,20 +264,22 @@ static bool doPhase2Merge(Phase2PipelineArg& p, Phase2BamHeap& h } } // 并行拷贝bam数据 - if (p.mergeOrder == 43) { - spdlog::info("curIdx: {}, bytes: {}", mergeData.curIdx, mergeData.curBytes); - } - spdlog::info("curIdx: {}", mergeData.curIdx); - if (mergeData.lastRoundIdx != p.compressBlocksThreshold) { // 回头再想想 - - } + //if (p.mergeOrder == 43) { + // spdlog::info("curIdx: {}, bytes: {}", mergeData.curIdx, mergeData.curBytes); + //} + spdlog::info("curIdx: {}, lastIdx: {}, diff: {}", mergeData.curIdx, mergeData.lastRoundIdx, mergeData.curIdx - mergeData.lastRoundIdx); + PROF_G_BEG(phase2_copyBam); if (mergeData.curIdx < p.compressBlocksThreshold) { - kt_for(p.numThread, mtCopyBams, &p, mergeData.curIdx + 1); + kt_for(p.numThread, mtCopyBams, &p, mergeData.curIdx + 1 - mergeData.lastRoundIdx); // kt_for(p.numThread, mtCopyBams, &p, mergeData.curIdx); + mergeData.lastRoundIdx = mergeData.curIdx; } else { - kt_for(p.numThread, mtCopyBams, &p, mergeData.curIdx); + kt_for(p.numThread, mtCopyBams, &p, mergeData.curIdx - mergeData.lastRoundIdx); + mergeData.lastRoundIdx = 0; } + PROF_G_END(phase2_copyBam); + // if (mergeFull != nullptr && * mergeFull) { // if (mergeData.curIdx == 2) { // for (int n = 0; n < mergeData.curIdx; ++n) { @@ -310,11 +314,10 @@ static bool doPhase2Merge(Phase2PipelineArg& p, Phase2BamHeap& h for (int i = 0; i <= mergeData.curIdx; ++i) { all_bams += mergeData.blockDataArr[i].bamNum; } - spdlog::info("phase2 merge sort order: {}, write bam id: {}, all bams: {}", p.mergeOrder, write_bam_id, all_bams); + //spdlog::info("phase2 merge sort order: {}, write bam id: {}, all bams: {}", p.mergeOrder, write_bam_id, all_bams); } // spdlog::info("curIdx: {}", mergeData.curIdx); - - return finish; + return finish; } void* phase2Merge(void* data) { diff --git a/src/sort/phase_2_read.cpp b/src/sort/phase_2_read.cpp index c774419..1a3e2f8 100644 --- a/src/sort/phase_2_read.cpp +++ b/src/sort/phase_2_read.cpp @@ -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; @@ -283,7 +283,7 @@ bool doPhase2Uncompress(Phase2PipelineArg& p) { } PROF_G_END(phase2_uncompress); - spdlog::info("phase2 uncompress order: {}, read bams: {}, blocks: {}", p.uncompressOrder, read_bam_id, block_id); + // spdlog::info("phase2 uncompress order: {}, read bams: {}, blocks: {}", p.uncompressOrder, read_bam_id, block_id); return hasUncompress; diff --git a/src/sort/phase_2_write.cpp b/src/sort/phase_2_write.cpp index 9d6ab95..8e02371 100644 --- a/src/sort/phase_2_write.cpp +++ b/src/sort/phase_2_write.cpp @@ -66,7 +66,7 @@ static void mtCompressBlock(void* data, long idx, int tid) { } static void doCompress(Phase2PipelineArg& p) { - PROF_G_BEG(compress); + PROF_G_BEG(phase2_compress); auto& compressBuf = p.compressBuf[p.compressOrder % p.COMPRESS_BUF_NUM]; auto& mergeData = p.mergeData[p.compressOrder % p.MERGE_BUF_NUM]; compressBuf.Clear(); @@ -78,8 +78,8 @@ static void doCompress(Phase2PipelineArg& p) { all_bams += mergeData.blockDataArr[i].bamNum; } mergeData.Clear(); - spdlog::info("compress order: {}, bytes: {}, bams: {}", p.compressOrder, compressBuf.curLen, all_bams); - PROF_G_END(compress); + // spdlog::info("compress order: {}, bytes: {}, bams: {}", p.compressOrder, compressBuf.curLen, all_bams); + PROF_G_END(phase2_compress); // 调试用 #if 0 diff --git a/src/util/profiling.cpp b/src/util/profiling.cpp index d80d0a6..68815cd 100644 --- a/src/util/profiling.cpp +++ b/src/util/profiling.cpp @@ -63,13 +63,15 @@ int displayProfiling(int nthread) { PRINT_GP(merge); PRINT_GP(compress); PRINT_GP(write_mid); - PRINT_GP(phase2_compress); PRINT_GP(read_mid); - PRINT_GP(write_final); PRINT_GP(mid_all); PRINT_GP(phase2_read); PRINT_GP(phase2_uncompress); PRINT_GP(phase2_copyToMerge); + PRINT_GP(phase2_copyBam); + PRINT_GP(phase2_merge); + PRINT_GP(phase2_compress); + PRINT_GP(write_final); PRINT_GP(phase2); PRINT_TP(sort, nthread); diff --git a/src/util/profiling.h b/src/util/profiling.h index 197359a..de82cb5 100644 --- a/src/util/profiling.h +++ b/src/util/profiling.h @@ -79,6 +79,7 @@ enum { GP_phase2_read, GP_phase2_uncompress, GP_phase2_copyToMerge, + GP_phase2_copyBam, GP_phase2_merge, GP_phase2_compress, GP_read_mid,