改进了merge,开始第二阶段过程的工作,需要clean一下代码

Co-authored-by: Copilot <copilot@github.com>
This commit is contained in:
zzh 2026-06-03 14:09:51 +08:00
parent 438375572d
commit f5776a795a
10 changed files with 57 additions and 8 deletions

View File

@ -68,5 +68,8 @@ void phase1Pipeline() {
PROF_G_END(mid_all); PROF_G_END(mid_all);
#endif #endif
} }

View File

@ -99,9 +99,19 @@ struct ThreadUncompressWrap {
} }
}; };
struct BlockBams {
DataBuffer blockBuf; // 解压的block放在这里
BamPtrArr bamPtrArr; // 解析后的bam数据放在这里
void Clear() {
blockBuf.Clear();
bamPtrArr.Clear();
}
};
// 用于合并压缩的数据结构 // 用于合并压缩的数据结构
struct MergeCompressData { struct MergeCompressData {
vector<DataBuffer> blockDataArr; // 待压缩的数据 vector<BlockBams> blockDataArr; // 待压缩的数据
vector<DataBuffer> compressDataArr; // 压缩后的数据 vector<DataBuffer> compressDataArr; // 压缩后的数据
void Resize(int blockNum) { void Resize(int blockNum) {

View File

@ -26,9 +26,15 @@
static void mtCompressBlock(void* data, long idx, int tid) { static void mtCompressBlock(void* data, long idx, int tid) {
Phase1PipelineArg& p = *(Phase1PipelineArg*)data; Phase1PipelineArg& p = *(Phase1PipelineArg*)data;
MergeCompressData& mergeCompressData = p.mergeCompressData[p.compressOrder % p.COMPRESS_BUF_NUM]; MergeCompressData& mergeCompressData = p.mergeCompressData[p.compressOrder % p.COMPRESS_BUF_NUM];
auto& blockData = mergeCompressData.blockDataArr[idx]; auto& bams = mergeCompressData.blockDataArr[idx].bamPtrArr;
auto& blockData = mergeCompressData.blockDataArr[idx].blockBuf;
auto& compressData = mergeCompressData.compressDataArr[idx]; auto& compressData = mergeCompressData.compressDataArr[idx];
for (int i = 0; i < bams.Size(); ++i) {
const OneBam* bp = bams.arr[i];
blockData.MemCopy(p.uncompressData.dataBuf + bp->offset, bp->wholeBamLen);
}
// bgzfCompress(void* _dst, size_t* dlen, const void* src, size_t slen, int level) // bgzfCompress(void* _dst, size_t* dlen, const void* src, size_t slen, int level)
compressData.ReAllocMem(SINGLE_BLOCK_SIZE); // 压缩后的block数据不会超过单个block的大小 compressData.ReAllocMem(SINGLE_BLOCK_SIZE); // 压缩后的block数据不会超过单个block的大小
compressData.curLen = SINGLE_BLOCK_SIZE; compressData.curLen = SINGLE_BLOCK_SIZE;

View File

@ -153,15 +153,16 @@ bool doMergeSort(Phase1PipelineArg& p, Phase1BamHeap<BamGreaterThan>& heap) {
mergeCompressData.blockDataArr[mergedBlockNum].Clear(); // 清理为添加bam数据做准备 mergeCompressData.blockDataArr[mergedBlockNum].Clear(); // 清理为添加bam数据做准备
bamBytes = 0; bamBytes = 0;
} }
mergeCompressData.blockDataArr[mergedBlockNum].MemCopy(p.uncompressData.dataBuf + bam->offset, bam->wholeBamLen); // mergeCompressData.blockDataArr[mergedBlockNum].MemCopy(p.uncompressData.dataBuf + bam->offset, bam->wholeBamLen);
bamBytes += bam->wholeBamLen; // for test mergeCompressData.blockDataArr[mergedBlockNum].bamPtrArr.Add(bam);
bamBytes += bam->wholeBamLen; // for test
heap.Pop(); heap.Pop();
} }
if (bam == nullptr) { if (bam == nullptr) {
// 都处理完 // 都处理完
finish = true; finish = true;
} }
spdlog::info("mergedBlockNum: {}, bamBytes: {}", mergedBlockNum, bamBytes); // spdlog::info("mergedBlockNum: {}, bamBytes: {}", mergedBlockNum, bamBytes);
return finish; return finish;
} }

View File

@ -381,8 +381,10 @@ static void doMemCopy(Phase1PipelineArg& p) {
#endif #endif
// 开启排序并写入中间文件 // 开启排序并写入中间文件
PROF_G_BEG(after_full);
phase1Sort(&p); phase1Sort(&p);
phase1MergeCompress(&p); phase1MergeCompress(&p);
PROF_G_END(after_full);
p.uncompressData.NextRound(); p.uncompressData.NextRound();
p.allBams.Clear(); p.allBams.Clear();
@ -390,7 +392,7 @@ static void doMemCopy(Phase1PipelineArg& p) {
PROF_G_BEG(mem_copy); PROF_G_BEG(mem_copy);
p.allBams.Add(uncompressWrap.GetTotalBamNum()); p.allBams.AddSize(uncompressWrap.GetTotalBamNum());
p.bamNum += uncompressWrap.GetTotalBamNum(); p.bamNum += uncompressWrap.GetTotalBamNum();
p.blockNum += uncompressWrap.GetTotalBlockNum(); p.blockNum += uncompressWrap.GetTotalBlockNum();

View File

@ -24,8 +24,10 @@
#include "util/profiling.h" #include "util/profiling.h"
static void doWrite(Phase1PipelineArg& p) { static void doWrite(Phase1PipelineArg& p) {
PROF_G_BEG(write_mid);
DataBuffer& compressBuf = p.compressBuf[p.compressOrder % p.COMPRESS_BUF_NUM]; DataBuffer& compressBuf = p.compressBuf[p.compressOrder % p.COMPRESS_BUF_NUM];
fwrite(compressBuf.data, 1, compressBuf.curLen, p.midFilePtr); fwrite(compressBuf.data, 1, compressBuf.curLen, p.midFilePtr);
PROF_G_END(write_mid);
} }
void* phase1Write(void* data) { void* phase1Write(void* data) {

View File

@ -12,4 +12,12 @@
Date : 2026/02/08 Date : 2026/02/08
*/ */
#pragma once #pragma once
struct Phase2File {
FILE* fp;
// 双buffer
// 当前读入的buffer指针df
// ReadBuffer
};

View File

@ -68,6 +68,12 @@ struct UncompressBlockBuffer {
lastEndPos = 0; lastEndPos = 0;
} }
uint8_t *HandOverBuf() {
uint8_t *handOverBuf = dataBuf;
dataBuf = nullptr; // 交出buf的所有权外部需要负责释放内存
return handOverBuf;
}
void NextRound() { void NextRound() {
usedBufSize -= lastEndPos; usedBufSize -= lastEndPos;
memcpy(dataBuf, dataBuf + lastEndPos, usedBufSize); memcpy(dataBuf, dataBuf + lastEndPos, usedBufSize);
@ -95,7 +101,13 @@ struct FastVector {
#endif #endif
} }
} }
void Add(size_t num) {
void Add(const T& item) {
T &newItem = Add();
newItem = item;
}
void AddSize(size_t num) {
if (curIdx + num > arr.size()) { if (curIdx + num > arr.size()) {
arr.resize(curIdx + num); arr.resize(curIdx + num);
} }
@ -383,6 +395,7 @@ struct MergeSortData {
}; };
typedef FastVector<OneBam> BamArr; typedef FastVector<OneBam> BamArr;
typedef FastVector<const OneBam*> BamPtrArr;
typedef FastVector<OneBlock> BlockArr; typedef FastVector<OneBlock> BlockArr;
@ -413,9 +426,11 @@ struct FirstPipeArg {
ReadBuffer readData[READ_BUF_NUM]; // 用来读如数据,双缓冲 ReadBuffer readData[READ_BUF_NUM]; // 用来读如数据,双缓冲
UncompressData uncompressData[UNCOMPRESS_BUF_NUM]; // 每个线程内保留一些自己的数据,比如解压缩的数据 UncompressData uncompressData[UNCOMPRESS_BUF_NUM]; // 每个线程内保留一些自己的数据,比如解压缩的数据
UncompressBlockBuffer unCompblockDataBuf; // 所有线程共用一个串行往这里添加解压后的block数据 UncompressBlockBuffer unCompblockDataBuf; // 所有线程共用一个串行往这里添加解压后的block数据
vector<const OneBam*> sortedBamArr; vector<const OneBam*> sortedBamArr;
vector<ArrayInterval> taskArr; vector<ArrayInterval> taskArr;
vector<DataBuffer> threadBuf; vector<DataBuffer> threadBuf;
MergeSortData mergeSortData; MergeSortData mergeSortData;
FirstPipeArg() FirstPipeArg()

View File

@ -58,6 +58,7 @@ int displayProfiling(int nthread) {
PRINT_GP(mem_copy); PRINT_GP(mem_copy);
PRINT_GP(parse_block); PRINT_GP(parse_block);
PRINT_GP(uncompress); PRINT_GP(uncompress);
PRINT_GP(after_full);
PRINT_GP(sort); PRINT_GP(sort);
PRINT_GP(merge); PRINT_GP(merge);
PRINT_GP(compress); PRINT_GP(compress);

View File

@ -71,6 +71,7 @@ enum {
GP_mem_copy, GP_mem_copy,
GP_parse_block, GP_parse_block,
GP_uncompress, GP_uncompress,
GP_after_full,
GP_sort, GP_sort,
GP_merge, GP_merge,
GP_compress, GP_compress,