55 lines
1.3 KiB
C++
55 lines
1.3 KiB
C++
/*
|
|
Description: 第一阶段写入中间文件
|
|
|
|
Copyright : All right reserved by ICT
|
|
|
|
Author : Zhang Zhonghai
|
|
Date : 2026/06/02
|
|
*/
|
|
|
|
#include "phase_1_write.h"
|
|
|
|
#include <klib/kthread.h>
|
|
#include <spdlog/spdlog.h>
|
|
|
|
#include <algorithm>
|
|
#include <string>
|
|
|
|
#include "common_data.h"
|
|
#include "const_val.h"
|
|
#include "phase_1.h"
|
|
#include "phase_1_compress.h"
|
|
#include "phase_1_write.h"
|
|
#include "sort.h"
|
|
#include "util/profiling.h"
|
|
|
|
static void doWrite(Phase1PipelineArg& p) {
|
|
PROF_G_BEG(write_mid);
|
|
DataBuffer& compressBuf = p.compressBuf[p.compressOrder % p.COMPRESS_BUF_NUM];
|
|
fwrite(compressBuf.data, 1, compressBuf.curLen, p.midFilePtr);
|
|
PROF_G_END(write_mid);
|
|
}
|
|
|
|
void* phase1Write(void* data) {
|
|
Phase1PipelineArg& p = *(Phase1PipelineArg*)data;
|
|
/* do the work */
|
|
while (true) {
|
|
// previous dependency
|
|
yarn::DEPENDENCY_NOT_TO_BE(p.compressSig, 0);
|
|
|
|
if (p.compressFinish) {
|
|
while (p.writeOrder < p.compressOrder) {
|
|
doWrite(p);
|
|
p.writeOrder += 1;
|
|
}
|
|
break;
|
|
}
|
|
doWrite(p);
|
|
// update status
|
|
yarn::CONSUME_SIGNAL(p.compressSig);
|
|
p.writeOrder += 1;
|
|
}
|
|
|
|
spdlog::info("End write order: {}", p.writeOrder);
|
|
return nullptr;
|
|
} |