解决了阶段一的一个bug,就是最后一次merge之后,如果没达到阈值,merge order不增加,导致中间文件少了一部分数据

This commit is contained in:
zzh 2026-06-27 13:27:13 +08:00
parent eb876d4fe2
commit 36d1aa43dc
11 changed files with 69 additions and 21 deletions

View File

@ -117,6 +117,8 @@ struct MergeCompressData {
compressDataArr.resize(blockNum); compressDataArr.resize(blockNum);
} }
int Size() { return curIdx; }
void Clear() { curIdx = 0; } void Clear() { curIdx = 0; }
}; };

View File

@ -23,6 +23,10 @@
#include "sort.h" #include "sort.h"
#include "util/profiling.h" #include "util/profiling.h"
// for test
uint64_t block_num = 0;
uint64_t bam_num = 0;
void CheckBam(uint8_t *addr, int len) { void CheckBam(uint8_t *addr, int len) {
int bamLen = 0; int bamLen = 0;
memcpy(&bamLen, addr, 4); memcpy(&bamLen, addr, 4);
@ -42,6 +46,7 @@ static void mtCompressBlock(void* data, long idx, int tid) {
auto& compressData = mergeCompressData.compressDataArr[idx]; auto& compressData = mergeCompressData.compressDataArr[idx];
// spdlog::info("bam size: {}", bams.Size()); // spdlog::info("bam size: {}", bams.Size());
bam_num += bams.Size();
for (int i = 0; i < bams.Size(); ++i) { for (int i = 0; i < bams.Size(); ++i) {
const OneBam* bp = bams.arr[i]; const OneBam* bp = bams.arr[i];
blockData.MemCopy(p.uncompressData.dataBuf + bp->offset, bp->wholeBamLen); blockData.MemCopy(p.uncompressData.dataBuf + bp->offset, bp->wholeBamLen);
@ -59,15 +64,16 @@ static void doCompress(Phase1PipelineArg& p) {
PROF_G_BEG(compress); PROF_G_BEG(compress);
DataBuffer& compressBuf = p.compressBuf[p.compressOrder % p.COMPRESS_BUF_NUM]; DataBuffer& compressBuf = p.compressBuf[p.compressOrder % p.COMPRESS_BUF_NUM];
MergeCompressData& mergeCompressData = p.mergeCompressData[p.compressOrder % p.MERGE_BUF_NUM]; MergeCompressData& mergeCompressData = p.mergeCompressData[p.compressOrder % p.MERGE_BUF_NUM];
kt_for(p.numThread, mtCompressBlock, &p, mergeCompressData.blockDataArr.size()); kt_for(p.numThread, mtCompressBlock, &p, mergeCompressData.Size());
//kt_for(1, mtCompressBlock, &p, mergeCompressData.blockDataArr.size()); // kt_for(1, mtCompressBlock, &p, mergeCompressData..Size());
PROF_G_END(compress); PROF_G_END(compress);
compressBuf.Clear(); 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); 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- 压缩线程 */ /* phase1Compress step- 压缩线程 */
@ -96,6 +102,6 @@ void* phase1Compress(void* data) {
yarn::UPDATE_SIG_ORDER(p.compressSig, p.compressOrder); 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; return nullptr;
} }

View File

@ -23,6 +23,9 @@
#include "sort.h" #include "sort.h"
#include "util/profiling.h" #include "util/profiling.h"
uint64_t sort_bam_id = 0;
uint64_t block_nums = 0;
/* bam 排序堆 */ /* bam 排序堆 */
struct Phase1BamArrIdIdx { struct Phase1BamArrIdIdx {
size_t idx = 0; // 下一个待读入数据的idx size_t idx = 0; // 下一个待读入数据的idx
@ -206,6 +209,7 @@ bool doMergeSort(Phase1PipelineArg& p, Phase1BamHeap<BamGreaterThan>& heap) {
if (bamBytes + bam->wholeBamLen > singleBlockBytes) { if (bamBytes + bam->wholeBamLen > singleBlockBytes) {
mergedBlockNum++; mergedBlockNum++;
if (mergedBlockNum >= p.mergeBlocksThreshold) { if (mergedBlockNum >= p.mergeBlocksThreshold) {
mergeCompressData.curIdx = mergedBlockNum; // 这里和下面finish时候不同要注意
break; break;
} }
mergeCompressData.blockDataArr[mergedBlockNum].Clear(); // 清理为添加bam数据做准备 mergeCompressData.blockDataArr[mergedBlockNum].Clear(); // 清理为添加bam数据做准备
@ -214,6 +218,7 @@ bool doMergeSort(Phase1PipelineArg& p, Phase1BamHeap<BamGreaterThan>& heap) {
#if 0 #if 0
mergeCompressData.blockDataArr[mergedBlockNum].blockBuf.MemCopy(p.uncompressData.dataBuf + bam->offset, bam->wholeBamLen); mergeCompressData.blockDataArr[mergedBlockNum].blockBuf.MemCopy(p.uncompressData.dataBuf + bam->offset, bam->wholeBamLen);
#else #else
sort_bam_id += 1;
mergeCompressData.blockDataArr[mergedBlockNum].bamPtrArr.Add(bam); mergeCompressData.blockDataArr[mergedBlockNum].bamPtrArr.Add(bam);
#endif #endif
bamBytes += bam->wholeBamLen; // for test bamBytes += bam->wholeBamLen; // for test
@ -222,8 +227,11 @@ bool doMergeSort(Phase1PipelineArg& p, Phase1BamHeap<BamGreaterThan>& heap) {
if (bam == nullptr) { if (bam == nullptr) {
// 都处理完 // 都处理完
finish = true; 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; return finish;
} }
@ -255,7 +263,8 @@ void* phase1MergeSort(void* data) {
PROF_G_END(merge); PROF_G_END(merge);
if (finish) { 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; break;
} }

View File

@ -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); handleLastRoundData(p);
}
break; break;
} }

View File

@ -18,6 +18,10 @@
#include "phase_2_write.h" #include "phase_2_write.h"
#include "util/profiling.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) { void phase2Pipeline(Phase2PipelineArg& p) {
PROF_G_BEG(phase2); PROF_G_BEG(phase2);
@ -25,13 +29,13 @@ void phase2Pipeline(Phase2PipelineArg& p) {
pthread_t tidArr[6]; // 2-stage pipeline pthread_t tidArr[6]; // 2-stage pipeline
pthread_create(&tidArr[0], NULL, phase2ReadMidFile, &p); pthread_create(&tidArr[0], NULL, phase2ReadMidFile, &p);
pthread_create(&tidArr[1], NULL, phase2Uncompress, &p); pthread_create(&tidArr[1], NULL, phase2Uncompress, &p);
pthread_create(&tidArr[2], NULL, phase2CopyToMergeBuf, &p); //pthread_create(&tidArr[2], NULL, phase2CopyToMergeBuf, &p);
pthread_create(&tidArr[3], NULL, phase2Merge, &p); //pthread_create(&tidArr[3], NULL, phase2Merge, &p);
pthread_create(&tidArr[4], NULL, phase2Compress, &p); //pthread_create(&tidArr[4], NULL, phase2Compress, &p);
pthread_create(&tidArr[5], NULL, phase2Write, &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 < 6; ++i) pthread_join(tidArr[i], NULL);
//for (int i = 0; i < 4; ++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); spdlog::info("all bams num: {}", p.numBam);

View File

@ -23,6 +23,11 @@
#include "sort.h" #include "sort.h"
using std::string; 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 { struct CircularBuffer {
uint8_t* data = nullptr; uint8_t* data = nullptr;

View File

@ -199,6 +199,7 @@ static bool doPhase2Merge(Phase2PipelineArg& p, Phase2BamHeap<BamGreaterThan>& h
mergeData.blockDataArr[mergeData.curIdx].Clear(); mergeData.blockDataArr[mergeData.curIdx].Clear();
while ((bam = heap.Top()) != nullptr) { while ((bam = heap.Top()) != nullptr) {
write_bam_id += 1;
if (bam->addr == nullptr) { if (bam->addr == nullptr) {
spdlog::info("null addr"); spdlog::info("null addr");
} }
@ -242,6 +243,8 @@ static bool doPhase2Merge(Phase2PipelineArg& p, Phase2BamHeap<BamGreaterThan>& h
*emptyFilePtr = emptyFile; *emptyFilePtr = emptyFile;
} }
spdlog::info("phase2 merge sort order: {}, write bam id: {}", p.mergeOrder, write_bam_id);
return finish; return finish;
} }

View File

@ -1,5 +1,5 @@
/* /*
Description: Description: merge
Copyright : All right reserved by ICT Copyright : All right reserved by ICT
@ -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;
@ -145,7 +145,9 @@ static void mtUncompressBlockBatch(void* data, long idx, int tid) {
blockBuf.curLen += dlen; 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, //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; p.curBlockNum = blockNum;
block_id += blockNum;
#if 1 #if 1
kt_for(p.numThread, mtUncompressBlockBatch, &p, p.numThread); kt_for(p.numThread, mtUncompressBlockBatch, &p, p.numThread);
@ -277,6 +280,9 @@ bool doPhase2Uncompress(Phase2PipelineArg& p) {
// f.readyUncompressBufNum = 0; // f.readyUncompressBufNum = 0;
} }
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);
return hasUncompress; return hasUncompress;
// exit(0); // exit(0);
@ -288,11 +294,11 @@ void* phase2Uncompress(void* data) {
while (true) { while (true) {
// previous dependency // previous dependency
yarn::DEPENDENCY_NOT_TO_BE(p.readSig, 0); 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) { if (p.readFinish) {
while (p.uncompressOrder < p.readOrder) { 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); doPhase2Uncompress(p);
yarn::UPDATE_SIG_ORDER(p.uncompressSig, p.uncompressOrder); yarn::UPDATE_SIG_ORDER(p.uncompressSig, p.uncompressOrder);
} }
@ -317,6 +323,7 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) {
PROF_G_BEG(phase2_copyToMerge); PROF_G_BEG(phase2_copyToMerge);
int usedUncompress = 0; int usedUncompress = 0;
p.uncompressReadyNum = 0; p.uncompressReadyNum = 0;
size_t copied = 0;
#if 1 #if 1
for (int i = 0; i < p.midFiles.size(); ++i) { for (int i = 0; i < p.midFiles.size(); ++i) {
int usedUncompressInFile = 0; int usedUncompressInFile = 0;
@ -331,6 +338,7 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) {
} }
int copiedNum = mergeData.CopyBams(uncompressBuf.blockBuf, uncompressBuf.bamArr, uncompressBuf.startIdx); int copiedNum = mergeData.CopyBams(uncompressBuf.blockBuf, uncompressBuf.bamArr, uncompressBuf.startIdx);
copied += copiedNum;
uncompressBuf.startIdx += copiedNum; uncompressBuf.startIdx += copiedNum;
if (uncompressBuf.Size() == 0) { if (uncompressBuf.Size() == 0) {
f.readyUncompressBufNum -= 1; f.readyUncompressBufNum -= 1;
@ -341,6 +349,7 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) {
auto& uncompressBuf = f.uncompressBuf[f.copyOrder % f.UNCOMPRESSS_BUF_NUM]; auto& uncompressBuf = f.uncompressBuf[f.copyOrder % f.UNCOMPRESSS_BUF_NUM];
if (mergeData.hasSpace(uncompressBuf.Front()->wholeBamLen)) { // merge有空间uncompress有数据那就继续添加到merge if (mergeData.hasSpace(uncompressBuf.Front()->wholeBamLen)) { // merge有空间uncompress有数据那就继续添加到merge
int copiedNum = mergeData.CopyBams(uncompressBuf.blockBuf, uncompressBuf.bamArr, uncompressBuf.startIdx); int copiedNum = mergeData.CopyBams(uncompressBuf.blockBuf, uncompressBuf.bamArr, uncompressBuf.startIdx);
copied += copiedNum;
uncompressBuf.startIdx += copiedNum; uncompressBuf.startIdx += copiedNum;
if (uncompressBuf.Size() == 0) { if (uncompressBuf.Size() == 0) {
usedUncompressInFile += 1; usedUncompressInFile += 1;
@ -373,7 +382,8 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) {
#else #else
p.midFiles[p.copyMergeOrder % p.midFiles.size()].readyUncompressBufNum -= 1; p.midFiles[p.copyMergeOrder % p.midFiles.size()].readyUncompressBufNum -= 1;
#endif #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); PROF_G_END(phase2_copyToMerge);

View File

@ -79,6 +79,7 @@ static void doPhase2Write(Phase2PipelineArg& p) {
DataBuffer& compressBuf = p.compressBuf[p.writeOrder % p.COMPRESS_BUF_NUM]; DataBuffer& compressBuf = p.compressBuf[p.writeOrder % p.COMPRESS_BUF_NUM];
fwrite(compressBuf.data, 1, compressBuf.curLen, p.outFilePtr); fwrite(compressBuf.data, 1, compressBuf.curLen, p.outFilePtr);
PROF_G_END(write_final); PROF_G_END(write_final);
spdlog::info("phase2 write order: {}", p.writeOrder);
} }
void* phase2Write(void* data) { void* phase2Write(void* data) {

View File

@ -166,6 +166,12 @@ void free_lock_(lock_t *, char const *, long);
finish = 1; \ finish = 1; \
twist_(sig, yarn::BY, 1, __FILE__, __LINE__); 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) \ #define CONSUME_SIGNAL(sig) \
possess_(sig, __FILE__, __LINE__); \ possess_(sig, __FILE__, __LINE__); \
twist_(sig, yarn::BY, -1, __FILE__, __LINE__); twist_(sig, yarn::BY, -1, __FILE__, __LINE__);