diff --git a/src/sort/phase_1_compress.cpp b/src/sort/phase_1_compress.cpp index 8624348..35a116b 100644 --- a/src/sort/phase_1_compress.cpp +++ b/src/sort/phase_1_compress.cpp @@ -27,15 +27,17 @@ uint64_t block_num = 0; uint64_t bam_num = 0; -void CheckBam(uint8_t *addr, int len) { +bool CheckBam(uint8_t *addr, int len) { int bamLen = 0; memcpy(&bamLen, addr, 4); if (nsgv::gIsBigEndian) ed_swap_4p(&bamLen); if (bamLen + 4 != len) { spdlog::error("bam error: {}-{}", bamLen, len); - exit(0); + return false; + // exit(0); } + return true; } static void mtCompressBlock(void* data, long idx, int tid) { diff --git a/src/sort/phase_1_uncompress.cpp b/src/sort/phase_1_uncompress.cpp index 3cd3805..9dda26b 100644 --- a/src/sort/phase_1_uncompress.cpp +++ b/src/sort/phase_1_uncompress.cpp @@ -430,8 +430,12 @@ static void doMemCopy(Phase1PipelineArg& p) { // 开启排序并写入中间文件 PROF_G_BEG(after_full); +#if 1 phase1Sort(&p); phase1MergeCompress(&p); +#else + p.midFileOrder += 1; +#endif PROF_G_END(after_full); p.uncompressData.NextRound(); diff --git a/src/sort/phase_2.h b/src/sort/phase_2.h index 03a1f23..f723531 100644 --- a/src/sort/phase_2.h +++ b/src/sort/phase_2.h @@ -35,6 +35,7 @@ struct CircularBuffer { size_t writeIdx = 0; // 可以写入的开始位置 size_t valueSize = 0; // 有效字节 size_t bufSize = 0; // 缓冲区空间 + size_t skipped = 0; // 跳过的字节数 CircularBuffer() {} CircularBuffer(size_t initSize) { @@ -50,6 +51,7 @@ struct CircularBuffer { writeIdx = 0; valueSize = 0; bufSize = 0; + skipped = 0; } void AllocMem(size_t memSize) { ReAllocMem(memSize); } @@ -95,14 +97,14 @@ struct CircularBuffer { // 返回第一个连续空间的大小 size_t FirstPartWriteSize() { - size_t freeSpace = bufSize - valueSize; + size_t freeSpace = bufSize - valueSize - skipped; size_t firstPart = MIN(freeSpace, bufSize - writeIdx); return firstPart; } // 如果空间不连续,那么返回第二个连续空间的内存大小 size_t SecondPartWriteSize() { - size_t freeSpace = bufSize - valueSize; + size_t freeSpace = bufSize - valueSize - skipped; size_t firstPart = MIN(freeSpace, bufSize - writeIdx); if (firstPart == freeSpace) return 0; @@ -111,7 +113,8 @@ struct CircularBuffer { // 跳过不能完整保存一个bam的空间 void SkipWrite(size_t skipBytes) { - valueSize += skipBytes; + // valueSize += skipBytes; + skipped += skipBytes; writeIdx = (writeIdx + skipBytes) % bufSize; } @@ -120,15 +123,19 @@ struct CircularBuffer { readIdx = 0; // 可以读取的开始位置 writeIdx = 0; // 可以写入的开始位置 valueSize = 0; // 有效字节 + skipped = 0; // 跳过的字节数 } else { if (from == readIdx) { readIdx = (readIdx + skipBytes) % bufSize; valueSize -= skipBytes; } else { - valueSize -= (bufSize - readIdx + skipBytes); + valueSize -= skipBytes; readIdx = skipBytes; } } + //if (valueSize > 1024 * 1024 * 1024) { + // spdlog::info("debug valuesize"); + //} } // 退回一个bam @@ -141,9 +148,6 @@ struct CircularBuffer { size_t ReadBam(uint8_t* out, size_t start, size_t len) { memcpy(out, data + start, len); valueSize -= len; - if (readIdx != start) { - valueSize -= bufSize - readIdx; - } readIdx = start + len; return len; } @@ -163,6 +167,7 @@ struct CircularBuffer { readIdx = 0; // 可以读取的开始位置 writeIdx = 0; // 可以写入的开始位置 valueSize = 0; // 有效字节 + skipped = 0; // 跳过的字节数 } }; @@ -283,7 +288,7 @@ struct Phase2MergeBuffer { // 从一个解压后的block缓冲区拷贝多个bam到循环缓冲区 // 返回实际拷贝的bam数量 - size_t CopyBams(DataBuffer &blockBuf, BamArr &arr, int start) { + size_t CopyBams(DataBuffer &blockBuf, BamArr &arr, int start, size_t copyOrder, int fid) { if (start >= arr.Size()) return 0; int origStart = start; size_t numCopied = 0; @@ -302,6 +307,17 @@ struct Phase2MergeBuffer { size_t secondPartSize = data.SecondPartWriteSize(); size_t needSize = b2->offset - b1->offset + b2->wholeBamLen; + #if 0 + if (copyOrder == 2 && fid == 4) { + spdlog::info("debug"); + } + if (fid == 4) { + spdlog::info("debug"); + } + if (data.valueSize > 1024 * 1024 * 1024) { + spdlog::info("debug"); + } +#endif if (needSize <= firstPartSize) { // 在第一个连续空间里就能放下 // 每个bam的offset需要加上diff,以对应新的buf int64_t diff = (int64_t)data.writeIdx - b1->offset; @@ -310,8 +326,19 @@ struct Phase2MergeBuffer { bams.Push(arr.Get(i)); bams.Back()->offset += diff; bams.Back()->addr = data.data; + // spdlog::info("1"); + #if 0 + auto bam = bams.Back(); + if (!CheckBam(bam->addr + bam->offset, bam->wholeBamLen)) { + spdlog::info("1"); + exit(0); + } + #endif + // auto bam = &arr.Get(i); + // CheckBam(blockBuf.data + bam->offset, bam->wholeBamLen); } } else { + int finalStop = stop; stop = start; size_t firstNeedSize = 0; while (arr.Get(stop).wholeBamLen + firstNeedSize < firstPartSize) { @@ -324,6 +351,14 @@ struct Phase2MergeBuffer { bams.Push(arr.Get(i)); bams.Back()->offset += diff; bams.Back()->addr = data.data; + // spdlog::info("2"); +#if 0 + auto bam = bams.Back(); + if (!CheckBam(bam->addr + bam->offset, bam->wholeBamLen)) { + spdlog::info("2"); + exit(0); + } + #endif } // 跳过first part不能放下完整bam的部分 data.SkipWrite(firstPartSize - firstNeedSize); @@ -332,7 +367,7 @@ struct Phase2MergeBuffer { // 拷贝第二段 size_t secondNeedSize = needSize - firstNeedSize; if (secondNeedSize <= secondPartSize) { - stop = arr.Size(); + stop = finalStop; } else { secondNeedSize = 0; while (arr.Get(stop).wholeBamLen + secondNeedSize < secondPartSize) { @@ -346,6 +381,16 @@ struct Phase2MergeBuffer { bams.Push(arr.Get(i)); bams.Back()->offset += diff; bams.Back()->addr = data.data; +#if 0 + auto bam = bams.Back(); + //auto bam = &arr.Get(i); + // CheckBam(blockBuf.data + bam->offset, bam->wholeBamLen); + if (!CheckBam(bam->addr + bam->offset, bam->wholeBamLen)) { + // if (!CheckBam(blockBuf.data + bam->offset, bam->wholeBamLen)) { + spdlog::info("3-{}-{}", i, copyOrder); + exit(0); + } + #endif } } numCopied = stop - origStart; @@ -410,6 +455,7 @@ struct Phase2File { string fileName; FILE* fp = nullptr; + int fid = 0; // for debug Phase2ReadBuffer readData[READ_BUF_NUM]; UncompressBuffer uncompressBuf[UNCOMPRESSS_BUF_NUM]; DataBuffer halfBlock; // 剩余不完整的压缩的block数据 @@ -548,6 +594,7 @@ struct Phase2PipelineArg { midFiles.resize(midFileNum); for (int i = 0; i < midFileNum; ++i) { midFiles[i].Init(midFilePrefix + std::to_string(i), kFileBufSize); + midFiles[i].fid = i; } readSig = yarn::NEW_LOCK(0); diff --git a/src/sort/phase_2_merge.cpp b/src/sort/phase_2_merge.cpp index d886727..9e4d28d 100644 --- a/src/sort/phase_2_merge.cpp +++ b/src/sort/phase_2_merge.cpp @@ -141,6 +141,13 @@ struct Phase2BamHeap { spdlog::error("null addr"); } minHeap.push({0, minVal.addr, minVal.file, mergeData.Front()}); + // if(minVal.file->fid == 4 && minVal.file->copyOrder == 8) { + // if(minVal.file->fid == 4 && minVal.file->copyOrder == 6) { + // spdlog::info("debug"); + // } + // if (minVal.file->fid == 4 && mergeData.Size() == 1) { + // spdlog::info("debug"); + // } mergeData.Pop(); } else { *emptyFile = minVal.file; diff --git a/src/sort/phase_2_read.cpp b/src/sort/phase_2_read.cpp index 08e0194..c774419 100644 --- a/src/sort/phase_2_read.cpp +++ b/src/sort/phase_2_read.cpp @@ -341,7 +341,7 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) { mergeData.InitSize(uncompressBuf.blockBuf.curLen, uncompressBuf.bamArr.Size()); } - int copiedNum = mergeData.CopyBams(uncompressBuf.blockBuf, uncompressBuf.bamArr, uncompressBuf.startIdx); + int copiedNum = mergeData.CopyBams(uncompressBuf.blockBuf, uncompressBuf.bamArr, uncompressBuf.startIdx, f.copyOrder, i); copied += copiedNum; uncompressBuf.startIdx += copiedNum; if (uncompressBuf.Size() == 0) { @@ -352,7 +352,7 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) { if (f.readyUncompressBufNum > 0) { 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); + int copiedNum = mergeData.CopyBams(uncompressBuf.blockBuf, uncompressBuf.bamArr, uncompressBuf.startIdx, f.copyOrder, i); copied += copiedNum; uncompressBuf.startIdx += copiedNum; if (uncompressBuf.Size() == 0) { diff --git a/src/sort/sort.cpp b/src/sort/sort.cpp index 3e4d7fe..cdd9bbc 100644 --- a/src/sort/sort.cpp +++ b/src/sort/sort.cpp @@ -79,6 +79,7 @@ static void bamSortPipeline() { ///////////////////////////////////////////////////// +#if 1 // 第二阶段参数初始化 Phase2PipelineArg p2(p1.uncompressData, p1.allBams, p1.midFileNamePrefix, p1.midFileOrder, p1.numThread); @@ -94,6 +95,8 @@ static void bamSortPipeline() { // 写结尾 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, p2.outFilePtr); fclose(p2.outFilePtr); +#endif + } // 排序的入口函数,entry function diff --git a/src/sort/sort.h b/src/sort/sort.h index 153d0b7..00bfd07 100644 --- a/src/sort/sort.h +++ b/src/sort/sort.h @@ -194,4 +194,4 @@ struct OneBam { typedef FastVector BamArr; typedef FastVector BamPtrArr; -extern void CheckBam(uint8_t* addr, int len); \ No newline at end of file +extern bool CheckBam(uint8_t* addr, int len); \ No newline at end of file