找到了一个bug,是循环缓冲区的skipwrite函数导致的,skip和valuesize混淆会导致pop的时候valuesize和pop出去的bam size不一致,导致bug,新增了skipped字段来解决这一问题

Co-authored-by: Copilot <copilot@github.com>
This commit is contained in:
zzh 2026-07-05 23:35:33 +08:00
parent e477905410
commit 1827d89663
7 changed files with 77 additions and 14 deletions

View File

@ -27,15 +27,17 @@
uint64_t block_num = 0; uint64_t block_num = 0;
uint64_t bam_num = 0; uint64_t bam_num = 0;
void CheckBam(uint8_t *addr, int len) { bool CheckBam(uint8_t *addr, int len) {
int bamLen = 0; int bamLen = 0;
memcpy(&bamLen, addr, 4); memcpy(&bamLen, addr, 4);
if (nsgv::gIsBigEndian) if (nsgv::gIsBigEndian)
ed_swap_4p(&bamLen); ed_swap_4p(&bamLen);
if (bamLen + 4 != len) { if (bamLen + 4 != len) {
spdlog::error("bam error: {}-{}", bamLen, len); spdlog::error("bam error: {}-{}", bamLen, len);
exit(0); return false;
// exit(0);
} }
return true;
} }
static void mtCompressBlock(void* data, long idx, int tid) { static void mtCompressBlock(void* data, long idx, int tid) {

View File

@ -430,8 +430,12 @@ static void doMemCopy(Phase1PipelineArg& p) {
// 开启排序并写入中间文件 // 开启排序并写入中间文件
PROF_G_BEG(after_full); PROF_G_BEG(after_full);
#if 1
phase1Sort(&p); phase1Sort(&p);
phase1MergeCompress(&p); phase1MergeCompress(&p);
#else
p.midFileOrder += 1;
#endif
PROF_G_END(after_full); PROF_G_END(after_full);
p.uncompressData.NextRound(); p.uncompressData.NextRound();

View File

@ -35,6 +35,7 @@ struct CircularBuffer {
size_t writeIdx = 0; // 可以写入的开始位置 size_t writeIdx = 0; // 可以写入的开始位置
size_t valueSize = 0; // 有效字节 size_t valueSize = 0; // 有效字节
size_t bufSize = 0; // 缓冲区空间 size_t bufSize = 0; // 缓冲区空间
size_t skipped = 0; // 跳过的字节数
CircularBuffer() {} CircularBuffer() {}
CircularBuffer(size_t initSize) { CircularBuffer(size_t initSize) {
@ -50,6 +51,7 @@ struct CircularBuffer {
writeIdx = 0; writeIdx = 0;
valueSize = 0; valueSize = 0;
bufSize = 0; bufSize = 0;
skipped = 0;
} }
void AllocMem(size_t memSize) { ReAllocMem(memSize); } void AllocMem(size_t memSize) { ReAllocMem(memSize); }
@ -95,14 +97,14 @@ struct CircularBuffer {
// 返回第一个连续空间的大小 // 返回第一个连续空间的大小
size_t FirstPartWriteSize() { size_t FirstPartWriteSize() {
size_t freeSpace = bufSize - valueSize; size_t freeSpace = bufSize - valueSize - skipped;
size_t firstPart = MIN(freeSpace, bufSize - writeIdx); size_t firstPart = MIN(freeSpace, bufSize - writeIdx);
return firstPart; return firstPart;
} }
// 如果空间不连续,那么返回第二个连续空间的内存大小 // 如果空间不连续,那么返回第二个连续空间的内存大小
size_t SecondPartWriteSize() { size_t SecondPartWriteSize() {
size_t freeSpace = bufSize - valueSize; size_t freeSpace = bufSize - valueSize - skipped;
size_t firstPart = MIN(freeSpace, bufSize - writeIdx); size_t firstPart = MIN(freeSpace, bufSize - writeIdx);
if (firstPart == freeSpace) if (firstPart == freeSpace)
return 0; return 0;
@ -111,7 +113,8 @@ struct CircularBuffer {
// 跳过不能完整保存一个bam的空间 // 跳过不能完整保存一个bam的空间
void SkipWrite(size_t skipBytes) { void SkipWrite(size_t skipBytes) {
valueSize += skipBytes; // valueSize += skipBytes;
skipped += skipBytes;
writeIdx = (writeIdx + skipBytes) % bufSize; writeIdx = (writeIdx + skipBytes) % bufSize;
} }
@ -120,15 +123,19 @@ struct CircularBuffer {
readIdx = 0; // 可以读取的开始位置 readIdx = 0; // 可以读取的开始位置
writeIdx = 0; // 可以写入的开始位置 writeIdx = 0; // 可以写入的开始位置
valueSize = 0; // 有效字节 valueSize = 0; // 有效字节
skipped = 0; // 跳过的字节数
} else { } else {
if (from == readIdx) { if (from == readIdx) {
readIdx = (readIdx + skipBytes) % bufSize; readIdx = (readIdx + skipBytes) % bufSize;
valueSize -= skipBytes; valueSize -= skipBytes;
} else { } else {
valueSize -= (bufSize - readIdx + skipBytes); valueSize -= skipBytes;
readIdx = skipBytes; readIdx = skipBytes;
} }
} }
//if (valueSize > 1024 * 1024 * 1024) {
// spdlog::info("debug valuesize");
//}
} }
// 退回一个bam // 退回一个bam
@ -141,9 +148,6 @@ struct CircularBuffer {
size_t ReadBam(uint8_t* out, size_t start, size_t len) { size_t ReadBam(uint8_t* out, size_t start, size_t len) {
memcpy(out, data + start, len); memcpy(out, data + start, len);
valueSize -= len; valueSize -= len;
if (readIdx != start) {
valueSize -= bufSize - readIdx;
}
readIdx = start + len; readIdx = start + len;
return len; return len;
} }
@ -163,6 +167,7 @@ struct CircularBuffer {
readIdx = 0; // 可以读取的开始位置 readIdx = 0; // 可以读取的开始位置
writeIdx = 0; // 可以写入的开始位置 writeIdx = 0; // 可以写入的开始位置
valueSize = 0; // 有效字节 valueSize = 0; // 有效字节
skipped = 0; // 跳过的字节数
} }
}; };
@ -283,7 +288,7 @@ struct Phase2MergeBuffer {
// 从一个解压后的block缓冲区拷贝多个bam到循环缓冲区 // 从一个解压后的block缓冲区拷贝多个bam到循环缓冲区
// 返回实际拷贝的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; if (start >= arr.Size()) return 0;
int origStart = start; int origStart = start;
size_t numCopied = 0; size_t numCopied = 0;
@ -302,6 +307,17 @@ struct Phase2MergeBuffer {
size_t secondPartSize = data.SecondPartWriteSize(); size_t secondPartSize = data.SecondPartWriteSize();
size_t needSize = b2->offset - b1->offset + b2->wholeBamLen; 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) { // 在第一个连续空间里就能放下 if (needSize <= firstPartSize) { // 在第一个连续空间里就能放下
// 每个bam的offset需要加上diff以对应新的buf // 每个bam的offset需要加上diff以对应新的buf
int64_t diff = (int64_t)data.writeIdx - b1->offset; int64_t diff = (int64_t)data.writeIdx - b1->offset;
@ -310,8 +326,19 @@ struct Phase2MergeBuffer {
bams.Push(arr.Get(i)); bams.Push(arr.Get(i));
bams.Back()->offset += diff; bams.Back()->offset += diff;
bams.Back()->addr = data.data; 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 { } else {
int finalStop = stop;
stop = start; stop = start;
size_t firstNeedSize = 0; size_t firstNeedSize = 0;
while (arr.Get(stop).wholeBamLen + firstNeedSize < firstPartSize) { while (arr.Get(stop).wholeBamLen + firstNeedSize < firstPartSize) {
@ -324,6 +351,14 @@ struct Phase2MergeBuffer {
bams.Push(arr.Get(i)); bams.Push(arr.Get(i));
bams.Back()->offset += diff; bams.Back()->offset += diff;
bams.Back()->addr = data.data; 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的部分 // 跳过first part不能放下完整bam的部分
data.SkipWrite(firstPartSize - firstNeedSize); data.SkipWrite(firstPartSize - firstNeedSize);
@ -332,7 +367,7 @@ struct Phase2MergeBuffer {
// 拷贝第二段 // 拷贝第二段
size_t secondNeedSize = needSize - firstNeedSize; size_t secondNeedSize = needSize - firstNeedSize;
if (secondNeedSize <= secondPartSize) { if (secondNeedSize <= secondPartSize) {
stop = arr.Size(); stop = finalStop;
} else { } else {
secondNeedSize = 0; secondNeedSize = 0;
while (arr.Get(stop).wholeBamLen + secondNeedSize < secondPartSize) { while (arr.Get(stop).wholeBamLen + secondNeedSize < secondPartSize) {
@ -346,6 +381,16 @@ struct Phase2MergeBuffer {
bams.Push(arr.Get(i)); bams.Push(arr.Get(i));
bams.Back()->offset += diff; bams.Back()->offset += diff;
bams.Back()->addr = data.data; 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; numCopied = stop - origStart;
@ -410,6 +455,7 @@ struct Phase2File {
string fileName; string fileName;
FILE* fp = nullptr; FILE* fp = nullptr;
int fid = 0; // for debug
Phase2ReadBuffer readData[READ_BUF_NUM]; Phase2ReadBuffer readData[READ_BUF_NUM];
UncompressBuffer uncompressBuf[UNCOMPRESSS_BUF_NUM]; UncompressBuffer uncompressBuf[UNCOMPRESSS_BUF_NUM];
DataBuffer halfBlock; // 剩余不完整的压缩的block数据 DataBuffer halfBlock; // 剩余不完整的压缩的block数据
@ -548,6 +594,7 @@ struct Phase2PipelineArg {
midFiles.resize(midFileNum); midFiles.resize(midFileNum);
for (int i = 0; i < midFileNum; ++i) { for (int i = 0; i < midFileNum; ++i) {
midFiles[i].Init(midFilePrefix + std::to_string(i), kFileBufSize); midFiles[i].Init(midFilePrefix + std::to_string(i), kFileBufSize);
midFiles[i].fid = i;
} }
readSig = yarn::NEW_LOCK(0); readSig = yarn::NEW_LOCK(0);

View File

@ -141,6 +141,13 @@ struct Phase2BamHeap {
spdlog::error("null addr"); spdlog::error("null addr");
} }
minHeap.push({0, minVal.addr, minVal.file, mergeData.Front()}); 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(); mergeData.Pop();
} else { } else {
*emptyFile = minVal.file; *emptyFile = minVal.file;

View File

@ -341,7 +341,7 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) {
mergeData.InitSize(uncompressBuf.blockBuf.curLen, uncompressBuf.bamArr.Size()); 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; copied += copiedNum;
uncompressBuf.startIdx += copiedNum; uncompressBuf.startIdx += copiedNum;
if (uncompressBuf.Size() == 0) { if (uncompressBuf.Size() == 0) {
@ -352,7 +352,7 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) {
if (f.readyUncompressBufNum > 0) { if (f.readyUncompressBufNum > 0) {
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, f.copyOrder, i);
copied += copiedNum; copied += copiedNum;
uncompressBuf.startIdx += copiedNum; uncompressBuf.startIdx += copiedNum;
if (uncompressBuf.Size() == 0) { if (uncompressBuf.Size() == 0) {

View File

@ -79,6 +79,7 @@ static void bamSortPipeline() {
///////////////////////////////////////////////////// /////////////////////////////////////////////////////
#if 1
// 第二阶段参数初始化 // 第二阶段参数初始化
Phase2PipelineArg p2(p1.uncompressData, p1.allBams, p1.midFileNamePrefix, p1.midFileOrder, p1.numThread); 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); 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); fclose(p2.outFilePtr);
#endif
} }
// 排序的入口函数entry function // 排序的入口函数entry function

View File

@ -194,4 +194,4 @@ struct OneBam {
typedef FastVector<OneBam> BamArr; typedef FastVector<OneBam> BamArr;
typedef FastVector<const OneBam*> BamPtrArr; typedef FastVector<const OneBam*> BamPtrArr;
extern void CheckBam(uint8_t* addr, int len); extern bool CheckBam(uint8_t* addr, int len);