小数据结果差不多了,大数据还有bug,再调试一下

Co-authored-by: Copilot <copilot@github.com>
This commit is contained in:
zzh 2026-07-02 03:12:05 +08:00
parent 36d1aa43dc
commit e477905410
7 changed files with 200 additions and 44 deletions

View File

@ -99,10 +99,12 @@ struct ThreadUncompressWrap {
struct BlockBams {
DataBuffer blockBuf; // 解压的block放在这里
BamPtrArr bamPtrArr; // 解析后的bam数据放在这里
int bamNum = 0; // 解析后的bam数量调试用
void Clear() {
blockBuf.Clear();
bamPtrArr.Clear();
bamNum = 0;
}
};
@ -111,6 +113,8 @@ struct MergeCompressData {
vector<BlockBams> blockDataArr; // 待压缩的数据
vector<DataBuffer> compressDataArr; // 压缩后的数据
int curIdx = 0;
int curBytes = 0;
int lastRoundIdx = 0; // 上一轮已经处理的block数量
void Resize(int blockNum) {
blockDataArr.resize(blockNum);
@ -119,7 +123,11 @@ struct MergeCompressData {
int Size() { return curIdx; }
void Clear() { curIdx = 0; }
void Clear() {
curIdx = 0;
curBytes = 0;
lastRoundIdx = 0;
}
};
/* 第一阶段的多线程流水线参数 */

View File

@ -29,13 +29,13 @@ void phase2Pipeline(Phase2PipelineArg& p) {
pthread_t tidArr[6]; // 2-stage pipeline
pthread_create(&tidArr[0], NULL, phase2ReadMidFile, &p);
pthread_create(&tidArr[1], NULL, phase2Uncompress, &p);
//pthread_create(&tidArr[2], NULL, phase2CopyToMergeBuf, &p);
//pthread_create(&tidArr[3], NULL, phase2Merge, &p);
//pthread_create(&tidArr[4], NULL, phase2Compress, &p);
//pthread_create(&tidArr[5], NULL, phase2Write, &p);
pthread_create(&tidArr[2], NULL, phase2CopyToMergeBuf, &p);
pthread_create(&tidArr[3], NULL, phase2Merge, &p);
pthread_create(&tidArr[4], NULL, phase2Compress, &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 < 2; ++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);
spdlog::info("all bams num: {}", p.numBam);

View File

@ -196,7 +196,7 @@ struct CircularArray {
// 推入元素(自动扩容)
void Push(const T& value) {
if (Full()) {
ReAllocArr(bufSize * 1.5);
ReAllocArr(bufSize * 1.5); // 这是bug的来源不能扩展或者把原来的数据搬运到后面的空间否则writeIdx对应的位置和ReadIdx冲突
}
arr[writeIdx] = value;
writeIdx = (writeIdx + 1) % bufSize;
@ -222,7 +222,7 @@ struct CircularArray {
void Revert() {
++valueSize;
readIdx = (readIdx - 1) % bufSize;
readIdx = (readIdx + bufSize - 1) % bufSize;
}
// 查看头部/尾部(不弹出)
@ -235,7 +235,7 @@ struct CircularArray {
T* Back() {
if (Empty())
return nullptr;
return &arr[(writeIdx - 1) % bufSize];
return &arr[(writeIdx + bufSize - 1) % bufSize];
}
// 是否为空
@ -287,13 +287,21 @@ struct Phase2MergeBuffer {
if (start >= arr.Size()) return 0;
int origStart = start;
size_t numCopied = 0;
int bamArrFree = bams.Free();
if (bamArrFree == 0)
return 0;
int stop = arr.Size();
if (stop - start > bamArrFree) {
stop = start + bamArrFree;
}
OneBam* b1 = &arr.Get(start);
OneBam* b2 = arr.Back();
OneBam* b2 = &arr.Get(stop - 1);
size_t firstPartSize = data.FirstPartWriteSize();
size_t secondPartSize = data.SecondPartWriteSize();
size_t needSize = b2->offset - b1->offset + b2->wholeBamLen;
int stop = arr.Size();
if (needSize <= firstPartSize) { // 在第一个连续空间里就能放下
// 每个bam的offset需要加上diff以对应新的buf
int64_t diff = (int64_t)data.writeIdx - b1->offset;

View File

@ -16,6 +16,8 @@
#include "phase_2.h"
#include "util/profiling.h"
static size_t all_bams = 0;
/* bam 排序堆 */
struct Phase2BamArrIdIdx {
int idx = 0; // 如果是第一阶段转移过来的buffer用这个此时file为nullptr
@ -161,6 +163,13 @@ struct Phase2BamHeap {
const OneBam* ret = nullptr;
if (!minHeap.empty()) {
ret = minHeap.top().bam;
auto &top = minHeap.top();
int idx = 0; // 如果是第一阶段转移过来的buffer用这个此时file为nullptr
uint8_t* addr = nullptr; // bam的原始地址就是减去offset
Phase2File* file = nullptr;
const OneBam* bam = nullptr;
// spdlog::info("idx: {}, addr: {}, file: {}, offset: {}, qname: {}", top.idx, (void*)top.addr, top.file != nullptr ? top.file->fileName : "null", top.bam->offset, std::string((char*)(top.addr + top.bam->offset + OneBam::QnameOffset), top.bam->qnameLen));
}
return ret;
}
@ -178,8 +187,12 @@ static void mtCopyBams(void* data, long idx, int tid) {
const OneBam* bp = bams.arr[i];
blockData.MemCopy(bp->addr + bp->offset, bp->wholeBamLen);
// check bam
// CheckBam(p.uncompressData.dataBuf + bp->offset, bp->wholeBamLen);
// CheckBam(bp->addr + bp->offset, bp->wholeBamLen);
}
// 要清理bams数组
mergeData.blockDataArr[idx].bamNum += bams.Size();
bams.Clear();
}
template <class BamGreaterThan>
@ -189,61 +202,110 @@ static bool doPhase2Merge(Phase2PipelineArg& p, Phase2BamHeap<BamGreaterThan>& h
auto& mergeData = p.mergeData[p.mergeOrder % p.MERGE_BUF_NUM];
const OneBam* bam = nullptr;
size_t bamBytes = 0;
int singleBlockBytes = 0xff00;
Phase2File* emptyFile = nullptr;
if (emptyFilePtr != nullptr && *emptyFilePtr != nullptr) {
heap.Init(&p);
heap.ReInitFile(*emptyFilePtr);
}
// spdlog::error("curIdx: {}, bytes: {}", mergeData.curIdx, mergeData.curBytes);
if (mergeData.curBytes == 0) { // 如果上一轮用完了某个文件则此时curBytes可能大于0而且没满此时不能清空
mergeData.blockDataArr[mergeData.curIdx].Clear();
}
// mergeData.curBytes = 0;
while ((bam = heap.Top()) != nullptr) {
write_bam_id += 1;
if (bam->addr == nullptr) {
spdlog::info("null addr");
}
if (bamBytes + bam->wholeBamLen > singleBlockBytes) {
// spdlog::info("bam: {}, pos: {}, tid: {}, qname: {}", bam->offset, bam->pos, bam->tid, std::string((char*)(bam->addr + bam->offset + OneBam::QnameOffset), bam->qnameLen));
// CheckBam(bam->addr + bam->offset, bam->wholeBamLen);
if (mergeData.curBytes + bam->wholeBamLen > singleBlockBytes) {
mergeData.curIdx++;
if (mergeData.curIdx >= p.compressBlocksThreshold) {
if (mergeFull != nullptr)
if (mergeFull != nullptr) {
*mergeFull = true;
}
// mergeData.curIdx = 0; // 调试用的
break;
}
// spdlog::error("bamBytes: {}", mergeData.curBytes);
mergeData.blockDataArr[mergeData.curIdx].Clear(); // 清理为添加bam数据做准备
bamBytes = 0;
mergeData.curBytes = 0;
}
#if 0
mergeData.blockDataArr[mergeData.curIdx].blockBuf.MemCopy(p.uncompressData.dataBuf + bam->offset, bam->wholeBamLen);
#else
mergeData.blockDataArr[mergeData.curIdx].bamPtrArr.Add(bam);
// spdlog::error("bamBytes-1: {} - {}", mergeData.curIdx, mergeData.curBytes);
write_bam_id += 1;
#endif
bamBytes += bam->wholeBamLen; // for test
mergeData.curBytes += bam->wholeBamLen; // for test
if (copyFinish) {
heap.Pop();
} else {
heap.Pop(&emptyFile);
if (emptyFile != nullptr)
if (emptyFile != nullptr) {
// spdlog::error("emptyfile, curIdx: {}, bytes: {}", mergeData.curIdx, mergeData.blockDataArr[mergeData.curIdx].blockBuf.curLen);
// for (int n = 0; n < mergeData.curIdx; ++n) {
// spdlog::info("emptyfile, block len: {}, bam len: {}", mergeData.blockDataArr[n].blockBuf.curLen, bam->wholeBamLen);
// }
break;
}
}
}
// 并行拷贝bam数据
kt_for(p.numThread, mtCopyBams, &p, mergeData.curIdx);
if (p.mergeOrder == 43) {
spdlog::info("curIdx: {}, bytes: {}", mergeData.curIdx, mergeData.curBytes);
}
spdlog::info("curIdx: {}", mergeData.curIdx);
if (mergeData.lastRoundIdx != p.compressBlocksThreshold) { // 回头再想想
}
if (mergeData.curIdx < p.compressBlocksThreshold) {
kt_for(p.numThread, mtCopyBams, &p, mergeData.curIdx + 1);
// kt_for(p.numThread, mtCopyBams, &p, mergeData.curIdx);
} else {
kt_for(p.numThread, mtCopyBams, &p, mergeData.curIdx);
}
// if (mergeFull != nullptr && * mergeFull) {
// if (mergeData.curIdx == 2) {
// for (int n = 0; n < mergeData.curIdx; ++n) {
// spdlog::info("n: {}, block len: {}, bam len: {}", n, mergeData.blockDataArr[n].blockBuf.curLen, bam->wholeBamLen);
// }
// }
if (bam == nullptr) {
// 都处理完
finish = true;
if (mergeData.curIdx < p.compressBlocksThreshold) {
mergeData.curIdx++; // 最后一次需要包含curIdx
}
}
if (emptyFile != nullptr) {
heap.Revert();
// 这里应该不需要revert因为弹出empty file的最后一个后没有在调用top
// heap.Revert();
emptyFile->readyMergeBufNum--;
}
if (emptyFilePtr != nullptr) {
*emptyFilePtr = emptyFile;
}
spdlog::info("phase2 merge sort order: {}, write bam id: {}", p.mergeOrder, write_bam_id);
// 调试用
#if 0
yarn::CONSUME_SIGNAL(p.mergeSig);
#endif
// spdlog::info("phase2 merge sort order: {}, write bam id: {}", p.mergeOrder, write_bam_id);
if (mergeFull != nullptr && *mergeFull) {
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("curIdx: {}", mergeData.curIdx);
return finish;
}
@ -266,7 +328,7 @@ void* phase2Merge(void* data) {
// bool finish = doMergeSort(p, heap);
bool finish = false;
bool mergeFull = false;
emptyFile = nullptr;
// emptyFile = nullptr;
if (firstInit) {
if (nsgv::gSortArg.SORT_COORIDINATE) {
@ -278,6 +340,7 @@ void* phase2Merge(void* data) {
}
if (p.copyMergeFinish) {
while (!finish) {
mergeFull = false;
PROF_G_BEG(phase2_merge);
yarn::DEPENDENCY_NOT_TO_BE(p.mergeSig, p.MERGE_BUF_NUM);
if (nsgv::gSortArg.SORT_COORIDINATE) {
@ -285,10 +348,13 @@ void* phase2Merge(void* data) {
} else {
finish = doPhase2Merge(p, nameHeap, &mergeFull, &emptyFile, true);
}
if (mergeFull)
if (mergeFull) {
yarn::UPDATE_SIG_ORDER(p.mergeSig, p.mergeOrder);
}
PROF_G_END(phase2_merge);
}
yarn::UPDATE_SIG_ORDER(p.mergeSig, p.mergeOrder);
spdlog::info("last phase2 merge sort order: {}, write bam id: {}", p.mergeOrder, write_bam_id);
yarn::SIGNAL_FINISH(p.mergeSig, p.mergeFinish);
break;
}
@ -301,11 +367,13 @@ void* phase2Merge(void* data) {
PROF_G_END(phase2_merge);
// update self status
if (emptyFile != nullptr) // 只有某个文件的buf消耗完了才读入
if (emptyFile != nullptr) {// 只有某个文件的buf消耗完了才读入
yarn::CONSUME_SIGNAL(p.copyMerge);
if (mergeFull)
}
if (mergeFull) {
yarn::UPDATE_SIG_ORDER(p.mergeSig, p.mergeOrder);
}
}
#endif
spdlog::info("End phase2 merge sort order: {}", p.mergeOrder);
return nullptr;

View File

@ -270,8 +270,10 @@ bool doPhase2Uncompress(Phase2PipelineArg& p) {
if (f.needUncompress) { // 还有空余缓冲区
f.readyReadBufNum -= 1;
f.readyUncompressBufNum += 1;
f.uncompressBuf[f.uncompressOrder % Phase2File::UNCOMPRESSS_BUF_NUM].startIdx = 0; // 刚读入的全部bam都还没拷贝到merge
f.uncompressOrder += 1;
f.uncompressBuf->startIdx = 0; // 刚读入的全部bam都还没拷贝到merge
// f.uncompressBuf->startIdx = 0; // 刚读入的全部bam都还没拷贝到merge, 这个是错误的buf是个数组这个一直把第一个buf的startidx设置为0bug来源
hasUncompress = true;
}
@ -294,7 +296,7 @@ void* phase2Uncompress(void* data) {
while (true) {
// previous dependency
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) {
while (p.uncompressOrder < p.readOrder) {
@ -313,7 +315,7 @@ void* phase2Uncompress(void* data) {
// }
}
spdlog::info("phase2 uncompress end order: {}", p.uncompressOrder);
// spdlog::info("phase2 uncompress end order: {}", p.uncompressOrder);
return nullptr;
}
@ -323,10 +325,11 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) {
PROF_G_BEG(phase2_copyToMerge);
int usedUncompress = 0;
p.uncompressReadyNum = 0;
int maxUncompressReady = 0;
size_t copied = 0;
#if 1
for (int i = 0; i < p.midFiles.size(); ++i) {
int usedUncompressInFile = 0;
int usedUncompressInFile = 0; // 使用了多少uncompress缓冲区
auto& f = p.midFiles[i];
if (f.readyUncompressBufNum > 0 && f.readyMergeBufNum < f.COPY_BUF_NUM) { // 有数据
// 拷贝到merge数据里的循环数组和循环缓冲区
@ -334,6 +337,7 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) {
auto& mergeData = f.mergeData;
if (!mergeData.initialized) {
// 初始化的时候确定merge buf的大小已经对应的bam个数以后每轮按照buf和bam的最小值拷贝
mergeData.InitSize(uncompressBuf.blockBuf.curLen, uncompressBuf.bamArr.Size());
}
@ -361,13 +365,21 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) {
}
f.readyMergeBufNum += 1; // 当前文件的merge data数据准备好了
//spdlog::info("copy num: {}", copiedNum);
// check bam
// for (int j = 0; j < mergeData.Size(); ++j) {
// auto& bam = mergeData.bams[j];
// CheckBam(bam.addr + bam.offset, bam.wholeBamLen);
// }
#if 1
//spdlog::info("copy num: {}", copiedNum);
maxUncompressReady = MAX(maxUncompressReady, f.readyUncompressBufNum);
// 调试用
#if 0
// for test消耗mergedata
//f.mergeData.bams.Clear();
//f.mergeData.data.Clear();
int clearSize = mergeData.Size() / 3;
//int clearSize = 110000;
//int clearSize = mergeData.Size();
for (int j = 0; j < clearSize; ++j) {
mergeData.Pop();
@ -383,10 +395,21 @@ int doPhase2CopyToMerge(Phase2PipelineArg& p) {
p.midFiles[p.copyMergeOrder % p.midFiles.size()].readyUncompressBufNum -= 1;
#endif
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);
// for test模拟消耗merge缓冲区的数据
// 调试用
#if 0
yarn::CONSUME_SIGNAL(p.copyMerge);
#endif
usedUncompress = MAX(usedUncompress, Phase2File::UNCOMPRESSS_BUF_NUM - maxUncompressReady);
//spdlog::info("copy to merge order: {}-{}, copied: {}, ready num: {}, max ready: {}, last id: {}", p.copyMergeOrder, usedUncompress, copied,
// p.uncompressReadyNum, maxUncompressReady, copy_bam_id);
return usedUncompress;
}
@ -399,15 +422,18 @@ void* phase2CopyToMergeBuf(void* data) {
yarn::DEPENDENCY_NOT_TO_BE(p.copyMerge, Phase2File::COPY_BUF_NUM);
if (p.uncompressFinish) {
while (p.uncompressReadyNum > 0) { // 一直有没处理完的解压数据
while (p.uncompressReadyNum > 0 || yarn::PEEK_LOCK(p.uncompressSig) > 0) { // 一直有没处理完的解压数据
yarn::DEPENDENCY_NOT_TO_BE(p.copyMerge, Phase2File::COPY_BUF_NUM);
doPhase2CopyToMerge(p);
int usedUncompress = doPhase2CopyToMerge(p);
usedUncompress = MIN(usedUncompress, yarn::PEEK_LOCK(p.uncompressSig));
yarn::CONSUME_SIGNAL_BY(p.uncompressSig, usedUncompress);
yarn::UPDATE_SIG_ORDER(p.copyMerge, p.copyMergeOrder);
}
yarn::SIGNAL_FINISH(p.copyMerge, p.copyMergeFinish);
break;
}
int usedUncompress = doPhase2CopyToMerge(p);
usedUncompress = MIN(usedUncompress, yarn::PEEK_LOCK(p.uncompressSig));
// update status
yarn::CONSUME_SIGNAL_BY(p.uncompressSig, usedUncompress);

View File

@ -11,10 +11,46 @@
#include <klib/kthread.h>
#include <spdlog/spdlog.h>
#include <zlib.h>
#include "phase_2.h"
#include "util/profiling.h"
static size_t all_bams = 0;
static void checkCompress(uint8_t* addr, uint64_t len) {
size_t curReadPos = 0;
int blockLen = 0;
int maxBlockLen = 0;
DataBuffer buf;
buf.AllocMem(SINGLE_BLOCK_SIZE);
while (curReadPos + BLOCK_HEADER_LENGTH <= len) { /* 确保能解析block长度 */
blockLen = unpackInt16(&addr[curReadPos + 16]) + 1;
if (blockLen > maxBlockLen) {
maxBlockLen = blockLen;
}
if (curReadPos + blockLen <= len) { /* 完整的block数据在buf里 */
size_t dlen = SINGLE_BLOCK_SIZE; // 65535
uint32_t crc = le_to_u32(addr + curReadPos + blockLen - 8);
int ret = bgzfUncompress(buf.data, &dlen, (Bytef*)(addr + curReadPos) + BLOCK_HEADER_LENGTH, blockLen - BLOCK_HEADER_LENGTH, crc);
if (ret != 0) {
spdlog::error("block len: {}, uncompressed len: {}", blockLen, dlen);
exit(0);
}
curReadPos += blockLen;
} else {
spdlog::error("not valid compressed block: {}, {}", curReadPos + blockLen, len);
break; /* 当前block数据不完整一部分在还没读入的file数据里 */
}
}
if (curReadPos != len) {
spdlog::error("addr: {}, len: {}", curReadPos, len);
exit(0);
}
}
//////////////////////////////压缩
static void mtCompressBlock(void* data, long idx, int tid) {
Phase2PipelineArg& p = *(Phase2PipelineArg*)data;
@ -25,6 +61,8 @@ static void mtCompressBlock(void* data, long idx, int tid) {
compressedData.ReAllocMem(SINGLE_BLOCK_SIZE); // 压缩后的block数据不会超过单个block的大小
compressedData.curLen = SINGLE_BLOCK_SIZE;
bgzfCompress(compressedData.data, &compressedData.curLen, blockData.blockBuf.data, blockData.blockBuf.curLen, p.compressLevel);
// spdlog::info("bytes: {}", blockData.blockBuf.curLen);
// checkCompress(compressedData.data, compressedData.curLen);
}
static void doCompress(Phase2PipelineArg& p) {
@ -33,14 +71,21 @@ static void doCompress(Phase2PipelineArg& p) {
auto& mergeData = p.mergeData[p.compressOrder % p.MERGE_BUF_NUM];
compressBuf.Clear();
kt_for(p.numThread, mtCompressBlock, &p, mergeData.blockDataArr.size());
kt_for(p.numThread, mtCompressBlock, &p, mergeData.Size());
for (int i = 0; i < mergeData.blockDataArr.size(); ++i) {
for (int i = 0; i < mergeData.Size(); ++i) {
compressBuf.MemCopy(mergeData.compressDataArr[i].data, mergeData.compressDataArr[i].curLen);
all_bams += mergeData.blockDataArr[i].bamNum;
}
mergeData.Clear();
// spdlog::info("compress bytes: {}", compressBuf.curLen);
spdlog::info("compress order: {}, bytes: {}, bams: {}", p.compressOrder, compressBuf.curLen, all_bams);
PROF_G_END(compress);
// 调试用
#if 0
yarn::CONSUME_SIGNAL(p.compressSig);
#endif
}
/* phase1Compress step- 压缩线程 */
@ -77,9 +122,10 @@ void* phase2Compress(void* data) {
static void doPhase2Write(Phase2PipelineArg& p) {
PROF_G_BEG(write_final);
DataBuffer& compressBuf = p.compressBuf[p.writeOrder % p.COMPRESS_BUF_NUM];
// checkCompress(compressBuf.data, compressBuf.curLen);
fwrite(compressBuf.data, 1, compressBuf.curLen, p.outFilePtr);
PROF_G_END(write_final);
spdlog::info("phase2 write order: {}", p.writeOrder);
// spdlog::info("phase2 write order: {}, bytes: {}", p.writeOrder, compressBuf.curLen);
}
void* phase2Write(void* data) {

View File

@ -178,7 +178,7 @@ void free_lock_(lock_t *, char const *, long);
#define CONSUME_SIGNAL_BY(sig, num) \
possess_(sig, __FILE__, __LINE__); \
twist_(sig, yarn::BY, num, __FILE__, __LINE__);
twist_(sig, yarn::BY, (0 - num), __FILE__, __LINE__);
#define INIT_SIG(sig) \
possess_(sig, __FILE__, __LINE__); \