添加了一些统计时间的代码
Co-authored-by: Copilot <copilot@github.com>
This commit is contained in:
parent
1827d89663
commit
eb3acb953a
|
|
@ -186,6 +186,8 @@ static void mtCopyBams(void* data, long idx, int tid) {
|
||||||
Phase2PipelineArg& p = *(Phase2PipelineArg*)data;
|
Phase2PipelineArg& p = *(Phase2PipelineArg*)data;
|
||||||
MergeCompressData& mergeData = p.mergeData[p.mergeOrder % p.MERGE_BUF_NUM];
|
MergeCompressData& mergeData = p.mergeData[p.mergeOrder % p.MERGE_BUF_NUM];
|
||||||
|
|
||||||
|
idx += mergeData.lastRoundIdx;
|
||||||
|
|
||||||
auto& bams = mergeData.blockDataArr[idx].bamPtrArr;
|
auto& bams = mergeData.blockDataArr[idx].bamPtrArr;
|
||||||
auto& blockData = mergeData.blockDataArr[idx].blockBuf;
|
auto& blockData = mergeData.blockDataArr[idx].blockBuf;
|
||||||
|
|
||||||
|
|
@ -262,20 +264,22 @@ static bool doPhase2Merge(Phase2PipelineArg& p, Phase2BamHeap<BamGreaterThan>& h
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// 并行拷贝bam数据
|
// 并行拷贝bam数据
|
||||||
if (p.mergeOrder == 43) {
|
//if (p.mergeOrder == 43) {
|
||||||
spdlog::info("curIdx: {}, bytes: {}", mergeData.curIdx, mergeData.curBytes);
|
// spdlog::info("curIdx: {}, bytes: {}", mergeData.curIdx, mergeData.curBytes);
|
||||||
}
|
//}
|
||||||
spdlog::info("curIdx: {}", mergeData.curIdx);
|
spdlog::info("curIdx: {}, lastIdx: {}, diff: {}", mergeData.curIdx, mergeData.lastRoundIdx, mergeData.curIdx - mergeData.lastRoundIdx);
|
||||||
if (mergeData.lastRoundIdx != p.compressBlocksThreshold) { // 回头再想想
|
|
||||||
|
|
||||||
}
|
|
||||||
|
|
||||||
|
PROF_G_BEG(phase2_copyBam);
|
||||||
if (mergeData.curIdx < p.compressBlocksThreshold) {
|
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);
|
// kt_for(p.numThread, mtCopyBams, &p, mergeData.curIdx);
|
||||||
|
mergeData.lastRoundIdx = mergeData.curIdx;
|
||||||
} else {
|
} 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 (mergeFull != nullptr && * mergeFull) {
|
||||||
// if (mergeData.curIdx == 2) {
|
// if (mergeData.curIdx == 2) {
|
||||||
// for (int n = 0; n < mergeData.curIdx; ++n) {
|
// for (int n = 0; n < mergeData.curIdx; ++n) {
|
||||||
|
|
@ -310,10 +314,9 @@ static bool doPhase2Merge(Phase2PipelineArg& p, Phase2BamHeap<BamGreaterThan>& h
|
||||||
for (int i = 0; i <= mergeData.curIdx; ++i) {
|
for (int i = 0; i <= mergeData.curIdx; ++i) {
|
||||||
all_bams += mergeData.blockDataArr[i].bamNum;
|
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);
|
// spdlog::info("curIdx: {}", mergeData.curIdx);
|
||||||
|
|
||||||
return finish;
|
return finish;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -48,7 +48,7 @@ bool doReadMidFiles(Phase2PipelineArg& p) {
|
||||||
finishNum += f.finish;
|
finishNum += f.finish;
|
||||||
//spdlog::info("fid: {}, readyNum: {}", i, f.readyReadBufNum);
|
//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();
|
finish = finishNum == p.midFiles.size();
|
||||||
PROF_G_END(phase2_read);
|
PROF_G_END(phase2_read);
|
||||||
return finish;
|
return finish;
|
||||||
|
|
@ -283,7 +283,7 @@ bool doPhase2Uncompress(Phase2PipelineArg& p) {
|
||||||
}
|
}
|
||||||
PROF_G_END(phase2_uncompress);
|
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;
|
return hasUncompress;
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -66,7 +66,7 @@ static void mtCompressBlock(void* data, long idx, int tid) {
|
||||||
}
|
}
|
||||||
|
|
||||||
static void doCompress(Phase2PipelineArg& p) {
|
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& compressBuf = p.compressBuf[p.compressOrder % p.COMPRESS_BUF_NUM];
|
||||||
auto& mergeData = p.mergeData[p.compressOrder % p.MERGE_BUF_NUM];
|
auto& mergeData = p.mergeData[p.compressOrder % p.MERGE_BUF_NUM];
|
||||||
compressBuf.Clear();
|
compressBuf.Clear();
|
||||||
|
|
@ -78,8 +78,8 @@ static void doCompress(Phase2PipelineArg& p) {
|
||||||
all_bams += mergeData.blockDataArr[i].bamNum;
|
all_bams += mergeData.blockDataArr[i].bamNum;
|
||||||
}
|
}
|
||||||
mergeData.Clear();
|
mergeData.Clear();
|
||||||
spdlog::info("compress order: {}, bytes: {}, bams: {}", p.compressOrder, compressBuf.curLen, all_bams);
|
// spdlog::info("compress order: {}, bytes: {}, bams: {}", p.compressOrder, compressBuf.curLen, all_bams);
|
||||||
PROF_G_END(compress);
|
PROF_G_END(phase2_compress);
|
||||||
|
|
||||||
// 调试用
|
// 调试用
|
||||||
#if 0
|
#if 0
|
||||||
|
|
|
||||||
|
|
@ -63,13 +63,15 @@ int displayProfiling(int nthread) {
|
||||||
PRINT_GP(merge);
|
PRINT_GP(merge);
|
||||||
PRINT_GP(compress);
|
PRINT_GP(compress);
|
||||||
PRINT_GP(write_mid);
|
PRINT_GP(write_mid);
|
||||||
PRINT_GP(phase2_compress);
|
|
||||||
PRINT_GP(read_mid);
|
PRINT_GP(read_mid);
|
||||||
PRINT_GP(write_final);
|
|
||||||
PRINT_GP(mid_all);
|
PRINT_GP(mid_all);
|
||||||
PRINT_GP(phase2_read);
|
PRINT_GP(phase2_read);
|
||||||
PRINT_GP(phase2_uncompress);
|
PRINT_GP(phase2_uncompress);
|
||||||
PRINT_GP(phase2_copyToMerge);
|
PRINT_GP(phase2_copyToMerge);
|
||||||
|
PRINT_GP(phase2_copyBam);
|
||||||
|
PRINT_GP(phase2_merge);
|
||||||
|
PRINT_GP(phase2_compress);
|
||||||
|
PRINT_GP(write_final);
|
||||||
PRINT_GP(phase2);
|
PRINT_GP(phase2);
|
||||||
|
|
||||||
PRINT_TP(sort, nthread);
|
PRINT_TP(sort, nthread);
|
||||||
|
|
|
||||||
|
|
@ -79,6 +79,7 @@ enum {
|
||||||
GP_phase2_read,
|
GP_phase2_read,
|
||||||
GP_phase2_uncompress,
|
GP_phase2_uncompress,
|
||||||
GP_phase2_copyToMerge,
|
GP_phase2_copyToMerge,
|
||||||
|
GP_phase2_copyBam,
|
||||||
GP_phase2_merge,
|
GP_phase2_merge,
|
||||||
GP_phase2_compress,
|
GP_phase2_compress,
|
||||||
GP_read_mid,
|
GP_read_mid,
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue