Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
7032405e
T
TDengine
项目概览
taosdata
/
TDengine
1 年多 前同步成功
通知
1185
Star
22016
Fork
4786
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
1
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
T
TDengine
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
1
Issue
1
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
提交
7032405e
编写于
3月 27, 2020
作者:
H
hzcheng
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
TD-34
上级
497300f1
变更
3
显示空白变更内容
内联
并排
Showing
3 changed file
with
163 addition
and
99 deletion
+163
-99
src/vnode/tsdb/inc/tsdbFile.h
src/vnode/tsdb/inc/tsdbFile.h
+14
-2
src/vnode/tsdb/src/tsdbFile.c
src/vnode/tsdb/src/tsdbFile.c
+31
-13
src/vnode/tsdb/src/tsdbMain.c
src/vnode/tsdb/src/tsdbMain.c
+118
-84
未找到文件。
src/vnode/tsdb/inc/tsdbFile.h
浏览文件 @
7032405e
...
@@ -25,6 +25,8 @@
...
@@ -25,6 +25,8 @@
extern
"C"
{
extern
"C"
{
#endif
#endif
#define TSDB_FILE_HEAD_SIZE 512
#define tsdbGetKeyFileId(key, daysPerFile, precision) ((key) / tsMsPerDay[(precision)] / (daysPerFile))
#define tsdbGetKeyFileId(key, daysPerFile, precision) ((key) / tsMsPerDay[(precision)] / (daysPerFile))
#define tsdbGetMaxNumOfFiles(keep, daysPerFile) ((keep) / (daysPerFile) + 3)
#define tsdbGetMaxNumOfFiles(keep, daysPerFile) ((keep) / (daysPerFile) + 3)
...
@@ -69,6 +71,7 @@ typedef struct {
...
@@ -69,6 +71,7 @@ typedef struct {
STsdbFileH
*
tsdbInitFileH
(
char
*
dataDir
,
int
maxFiles
);
STsdbFileH
*
tsdbInitFileH
(
char
*
dataDir
,
int
maxFiles
);
void
tsdbCloseFileH
(
STsdbFileH
*
pFileH
);
void
tsdbCloseFileH
(
STsdbFileH
*
pFileH
);
int
tsdbCreateFile
(
char
*
dataDir
,
int
fileId
,
char
*
suffix
,
int
maxTables
,
SFile
*
pFile
,
int
writeHeader
,
int
toClose
);
int
tsdbCreateFGroup
(
STsdbFileH
*
pFileH
,
char
*
dataDir
,
int
fid
,
int
maxTables
);
int
tsdbCreateFGroup
(
STsdbFileH
*
pFileH
,
char
*
dataDir
,
int
fid
,
int
maxTables
);
int
tsdbOpenFile
(
SFile
*
pFile
,
int
oflag
);
int
tsdbOpenFile
(
SFile
*
pFile
,
int
oflag
);
SFileGroup
*
tsdbOpenFilesForCommit
(
STsdbFileH
*
pFileH
,
int
fid
);
SFileGroup
*
tsdbOpenFilesForCommit
(
STsdbFileH
*
pFileH
,
int
fid
);
...
@@ -104,6 +107,9 @@ typedef struct {
...
@@ -104,6 +107,9 @@ typedef struct {
TSKEY
keyLast
;
TSKEY
keyLast
;
}
SCompBlock
;
}
SCompBlock
;
#define IS_SUPER_BLOCK(pBlock) ((pBlock)->numOfSubBlocks >= 1)
#define IS_SUB_BLOCK(pBlock) ((pBlock)->numOfSubBlocks == 0)
typedef
struct
{
typedef
struct
{
int32_t
delimiter
;
// For recovery usage
int32_t
delimiter
;
// For recovery usage
int32_t
checksum
;
// TODO: decide if checksum logic in this file or make it one API
int32_t
checksum
;
// TODO: decide if checksum logic in this file or make it one API
...
@@ -111,8 +117,7 @@ typedef struct {
...
@@ -111,8 +117,7 @@ typedef struct {
SCompBlock
blocks
[];
SCompBlock
blocks
[];
}
SCompInfo
;
}
SCompInfo
;
int
tsdbLoadCompIdx
(
SFileGroup
*
pGroup
,
void
*
buf
,
int
maxTables
);
#define TSDB_COMPBLOCK_AT(pCompInfo, idx) ((pCompInfo)->blocks + (idx))
int
tsdbLoadCompBlocks
(
SFileGroup
*
pGroup
,
SCompIdx
*
pIdx
,
void
*
buf
);
// TODO: take pre-calculation into account
// TODO: take pre-calculation into account
typedef
struct
{
typedef
struct
{
...
@@ -129,6 +134,13 @@ typedef struct {
...
@@ -129,6 +134,13 @@ typedef struct {
int64_t
uid
;
// For recovery usage
int64_t
uid
;
// For recovery usage
SCompCol
cols
[];
SCompCol
cols
[];
}
SCompData
;
}
SCompData
;
int
tsdbCopyCompBlockToFile
(
SFile
*
outFile
,
SFile
*
inFile
,
SCompInfo
*
pCompInfo
,
int
index
,
int
isOutLast
);
int
tsdbLoadCompIdx
(
SFileGroup
*
pGroup
,
void
*
buf
,
int
maxTables
);
int
tsdbLoadCompBlocks
(
SFileGroup
*
pGroup
,
SCompIdx
*
pIdx
,
void
*
buf
);
int
tsdbLoadCompCols
(
SFile
*
pFile
,
SCompBlock
*
pBlock
,
void
*
buf
);
int
tsdbLoadColData
(
SFile
*
pFile
,
SCompCol
*
pCol
,
int64_t
blockBaseOffset
,
void
*
buf
);
// TODO: need an API to merge all sub-block data into one
int
tsdbWriteBlockToFile
(
SFileGroup
*
pGroup
,
SCompInfo
*
pCompInfo
,
SCompIdx
*
pIdx
,
int
isMerge
,
SCompBlock
*
pBlock
,
SDataCols
*
pCols
);
int
tsdbWriteBlockToFile
(
SFileGroup
*
pGroup
,
SCompInfo
*
pCompInfo
,
SCompIdx
*
pIdx
,
int
isMerge
,
SCompBlock
*
pBlock
,
SDataCols
*
pCols
);
...
...
src/vnode/tsdb/src/tsdbFile.c
浏览文件 @
7032405e
...
@@ -24,7 +24,6 @@
...
@@ -24,7 +24,6 @@
#include "tsdbFile.h"
#include "tsdbFile.h"
#define TSDB_FILE_HEAD_SIZE 512
#define TSDB_FILE_DELIMITER 0xF00AFA0F
#define TSDB_FILE_DELIMITER 0xF00AFA0F
const
char
*
tsdbFileSuffix
[]
=
{
const
char
*
tsdbFileSuffix
[]
=
{
...
@@ -35,8 +34,7 @@ const char *tsdbFileSuffix[] = {
...
@@ -35,8 +34,7 @@ const char *tsdbFileSuffix[] = {
static
int
compFGroupKey
(
const
void
*
key
,
const
void
*
fgroup
);
static
int
compFGroupKey
(
const
void
*
key
,
const
void
*
fgroup
);
static
int
compFGroup
(
const
void
*
arg1
,
const
void
*
arg2
);
static
int
compFGroup
(
const
void
*
arg1
,
const
void
*
arg2
);
static
int
tsdbGetFileName
(
char
*
dataDir
,
int
fileId
,
int8_t
type
,
char
*
fname
);
static
int
tsdbGetFileName
(
char
*
dataDir
,
int
fileId
,
char
*
suffix
,
char
*
fname
);
static
int
tsdbCreateFile
(
char
*
dataDir
,
int
fileId
,
int8_t
type
,
int
maxTables
,
SFile
*
pFile
);
static
int
tsdbWriteFileHead
(
SFile
*
pFile
);
static
int
tsdbWriteFileHead
(
SFile
*
pFile
);
static
int
tsdbWriteHeadFileIdx
(
SFile
*
pFile
,
int
maxTables
);
static
int
tsdbWriteHeadFileIdx
(
SFile
*
pFile
,
int
maxTables
);
static
SFileGroup
*
tsdbSearchFGroup
(
STsdbFileH
*
pFileH
,
int
fid
);
static
SFileGroup
*
tsdbSearchFGroup
(
STsdbFileH
*
pFileH
,
int
fid
);
...
@@ -71,10 +69,10 @@ int tsdbCreateFGroup(STsdbFileH *pFileH, char *dataDir, int fid, int maxTables)
...
@@ -71,10 +69,10 @@ int tsdbCreateFGroup(STsdbFileH *pFileH, char *dataDir, int fid, int maxTables)
SFileGroup
fGroup
;
SFileGroup
fGroup
;
SFileGroup
*
pFGroup
=
&
fGroup
;
SFileGroup
*
pFGroup
=
&
fGroup
;
if
(
tsdbSearchFGroup
(
pFileH
,
fid
)
==
NULL
)
{
if
(
tsdbSearchFGroup
(
pFileH
,
fid
)
==
NULL
)
{
// if not exists, create one
pFGroup
->
fileId
=
fid
;
pFGroup
->
fileId
=
fid
;
for
(
int
type
=
TSDB_FILE_TYPE_HEAD
;
type
<
TSDB_FILE_TYPE_MAX
;
type
++
)
{
for
(
int
type
=
TSDB_FILE_TYPE_HEAD
;
type
<
TSDB_FILE_TYPE_MAX
;
type
++
)
{
if
(
tsdbCreateFile
(
dataDir
,
fid
,
t
ype
,
maxTables
,
&
(
pFGroup
->
files
[
type
])
)
<
0
)
{
if
(
tsdbCreateFile
(
dataDir
,
fid
,
t
sdbFileSuffix
[
type
],
maxTables
,
&
(
pFGroup
->
files
[
type
]),
type
==
TSDB_FILE_TYPE_HEAD
?
1
:
0
,
1
)
<
0
)
{
// TODO: deal with the ERROR here, remove those creaed file
// TODO: deal with the ERROR here, remove those creaed file
return
-
1
;
return
-
1
;
}
}
...
@@ -105,6 +103,10 @@ int tsdbRemoveFileGroup(STsdbFileH *pFileH, int fid) {
...
@@ -105,6 +103,10 @@ int tsdbRemoveFileGroup(STsdbFileH *pFileH, int fid) {
return
0
;
return
0
;
}
}
int
tsdbCopyCompBlockToFile
(
SFile
*
outFile
,
SFile
*
inFile
,
SCompInfo
*
pCompInfo
,
int
index
,
int
isOutLast
)
{
// TODO
return
0
;
}
int
tsdbLoadCompIdx
(
SFileGroup
*
pGroup
,
void
*
buf
,
int
maxTables
)
{
int
tsdbLoadCompIdx
(
SFileGroup
*
pGroup
,
void
*
buf
,
int
maxTables
)
{
SFile
*
pFile
=
&
(
pGroup
->
files
[
TSDB_FILE_TYPE_HEAD
]);
SFile
*
pFile
=
&
(
pGroup
->
files
[
TSDB_FILE_TYPE_HEAD
]);
...
@@ -127,6 +129,22 @@ int tsdbLoadCompBlocks(SFileGroup *pGroup, SCompIdx *pIdx, void *buf) {
...
@@ -127,6 +129,22 @@ int tsdbLoadCompBlocks(SFileGroup *pGroup, SCompIdx *pIdx, void *buf) {
return
0
;
return
0
;
}
}
int
tsdbLoadCompCols
(
SFile
*
pFile
,
SCompBlock
*
pBlock
,
void
*
buf
)
{
// assert(pBlock->numOfSubBlocks == 0 || pBlock->numOfSubBlocks == 1);
if
(
lseek
(
pFile
->
fd
,
pBlock
->
offset
,
SEEK_SET
)
<
0
)
return
-
1
;
size_t
size
=
sizeof
(
SCompData
)
+
sizeof
(
SCompCol
)
*
pBlock
->
numOfCols
;
if
(
read
(
pFile
->
fd
,
buf
,
size
)
<
0
)
return
-
1
;
return
0
;
}
int
tsdbLoadColData
(
SFile
*
pFile
,
SCompCol
*
pCol
,
int64_t
blockBaseOffset
,
void
*
buf
)
{
if
(
lseek
(
pFile
->
fd
,
blockBaseOffset
+
pCol
->
offset
,
SEEK_SET
)
<
0
)
return
-
1
;
if
(
read
(
pFile
->
fd
,
buf
,
pCol
->
len
)
<
0
)
return
-
1
;
return
0
;
}
static
int
tsdbWriteBlockToFileImpl
(
SFile
*
pFile
,
// File to write
static
int
tsdbWriteBlockToFileImpl
(
SFile
*
pFile
,
// File to write
SDataCols
*
pCols
,
// Data column buffer
SDataCols
*
pCols
,
// Data column buffer
int
numOfPointsToWrie
,
// Number of points to write to the file
int
numOfPointsToWrie
,
// Number of points to write to the file
...
@@ -229,10 +247,10 @@ static int tsdbWriteHeadFileIdx(SFile *pFile, int maxTables) {
...
@@ -229,10 +247,10 @@ static int tsdbWriteHeadFileIdx(SFile *pFile, int maxTables) {
return
0
;
return
0
;
}
}
static
int
tsdbGetFileName
(
char
*
dataDir
,
int
fileId
,
int8_t
type
,
char
*
fname
)
{
static
int
tsdbGetFileName
(
char
*
dataDir
,
int
fileId
,
char
*
suffix
,
char
*
fname
)
{
if
(
dataDir
==
NULL
||
fname
==
NULL
||
!
IS_VALID_TSDB_FILE_TYPE
(
type
)
)
return
-
1
;
if
(
dataDir
==
NULL
||
fname
==
NULL
)
return
-
1
;
sprintf
(
fname
,
"%s/f%d%s"
,
dataDir
,
fileId
,
tsdbFileSuffix
[
type
]
);
sprintf
(
fname
,
"%s/f%d%s"
,
dataDir
,
fileId
,
suffix
);
return
0
;
return
0
;
}
}
...
@@ -264,12 +282,12 @@ static int tsdbCloseFile(SFile *pFile) {
...
@@ -264,12 +282,12 @@ static int tsdbCloseFile(SFile *pFile) {
return
ret
;
return
ret
;
}
}
static
int
tsdbCreateFile
(
char
*
dataDir
,
int
fileId
,
int8_t
type
,
int
maxTables
,
SFile
*
pFil
e
)
{
int
tsdbCreateFile
(
char
*
dataDir
,
int
fileId
,
char
*
suffix
,
int
maxTables
,
SFile
*
pFile
,
int
writeHeader
,
int
toClos
e
)
{
memset
((
void
*
)
pFile
,
0
,
sizeof
(
SFile
));
memset
((
void
*
)
pFile
,
0
,
sizeof
(
SFile
));
pFile
->
type
=
type
;
pFile
->
fd
=
-
1
;
pFile
->
fd
=
-
1
;
tsdbGetFileName
(
dataDir
,
fileId
,
type
,
pFile
->
fname
);
tsdbGetFileName
(
dataDir
,
fileId
,
suffix
,
pFile
->
fname
);
if
(
access
(
pFile
->
fname
,
F_OK
)
==
0
)
{
if
(
access
(
pFile
->
fname
,
F_OK
)
==
0
)
{
// File already exists
// File already exists
return
-
1
;
return
-
1
;
...
@@ -280,7 +298,7 @@ static int tsdbCreateFile(char *dataDir, int fileId, int8_t type, int maxTables,
...
@@ -280,7 +298,7 @@ static int tsdbCreateFile(char *dataDir, int fileId, int8_t type, int maxTables,
return
-
1
;
return
-
1
;
}
}
if
(
type
==
TSDB_FILE_TYPE_HEAD
)
{
if
(
writeHeader
)
{
if
(
tsdbWriteHeadFileIdx
(
pFile
,
maxTables
)
<
0
)
{
if
(
tsdbWriteHeadFileIdx
(
pFile
,
maxTables
)
<
0
)
{
tsdbCloseFile
(
pFile
);
tsdbCloseFile
(
pFile
);
return
-
1
;
return
-
1
;
...
@@ -292,7 +310,7 @@ static int tsdbCreateFile(char *dataDir, int fileId, int8_t type, int maxTables,
...
@@ -292,7 +310,7 @@ static int tsdbCreateFile(char *dataDir, int fileId, int8_t type, int maxTables,
return
-
1
;
return
-
1
;
}
}
tsdbCloseFile
(
pFile
);
if
(
toClose
)
tsdbCloseFile
(
pFile
);
return
0
;
return
0
;
}
}
...
...
src/vnode/tsdb/src/tsdbMain.c
浏览文件 @
7032405e
...
@@ -8,6 +8,7 @@
...
@@ -8,6 +8,7 @@
#include <string.h>
#include <string.h>
#include <sys/stat.h>
#include <sys/stat.h>
#include <sys/types.h>
#include <sys/types.h>
#include <sys/sendfile.h>
#include <unistd.h>
#include <unistd.h>
// #include "taosdef.h"
// #include "taosdef.h"
...
@@ -45,6 +46,7 @@
...
@@ -45,6 +46,7 @@
#define TSDB_CFG_FILE_NAME "CONFIG"
#define TSDB_CFG_FILE_NAME "CONFIG"
#define TSDB_DATA_DIR_NAME "data"
#define TSDB_DATA_DIR_NAME "data"
#define TSDB_DEFAULT_FILE_BLOCK_ROW_OPTION 0.7
#define TSDB_DEFAULT_FILE_BLOCK_ROW_OPTION 0.7
#define TSDB_MAX_LAST_FILE_SIZE (1024 * 1024 * 10) // 10M
enum
{
TSDB_REPO_STATE_ACTIVE
,
TSDB_REPO_STATE_CLOSED
,
TSDB_REPO_STATE_CONFIGURING
};
enum
{
TSDB_REPO_STATE_ACTIVE
,
TSDB_REPO_STATE_CLOSED
,
TSDB_REPO_STATE_CONFIGURING
};
...
@@ -775,7 +777,7 @@ static int tsdbReadRowsFromCache(SSkipListIterator *pIter, TSKEY maxKey, int max
...
@@ -775,7 +777,7 @@ static int tsdbReadRowsFromCache(SSkipListIterator *pIter, TSKEY maxKey, int max
tdAppendDataRowToDataCol
(
row
,
pCols
);
tdAppendDataRowToDataCol
(
row
,
pCols
);
numOfRows
++
;
numOfRows
++
;
if
(
numOfRows
>
maxRowsToRead
)
break
;
if
(
numOfRows
>
=
maxRowsToRead
)
break
;
}
while
(
tSkipListIterNext
(
pIter
));
}
while
(
tSkipListIterNext
(
pIter
));
return
numOfRows
;
return
numOfRows
;
...
@@ -842,7 +844,10 @@ static void *tsdbCommitData(void *arg) {
...
@@ -842,7 +844,10 @@ static void *tsdbCommitData(void *arg) {
int
efid
=
tsdbGetKeyFileId
(
pCache
->
imem
->
keyLast
,
pCfg
->
daysPerFile
,
pCfg
->
precision
);
int
efid
=
tsdbGetKeyFileId
(
pCache
->
imem
->
keyLast
,
pCfg
->
daysPerFile
,
pCfg
->
precision
);
for
(
int
fid
=
sfid
;
fid
<=
efid
;
fid
++
)
{
for
(
int
fid
=
sfid
;
fid
<=
efid
;
fid
++
)
{
tsdbCommitToFile
(
pRepo
,
fid
,
iters
,
pCols
);
if
(
tsdbCommitToFile
(
pRepo
,
fid
,
iters
,
pCols
)
<
0
)
{
// TODO: deal with the error here
// assert(0);
}
}
}
tdFreeDataCols
(
pCols
);
tdFreeDataCols
(
pCols
);
...
@@ -867,7 +872,8 @@ static void *tsdbCommitData(void *arg) {
...
@@ -867,7 +872,8 @@ static void *tsdbCommitData(void *arg) {
}
}
static
int
tsdbCommitToFile
(
STsdbRepo
*
pRepo
,
int
fid
,
SSkipListIterator
**
iters
,
SDataCols
*
pCols
)
{
static
int
tsdbCommitToFile
(
STsdbRepo
*
pRepo
,
int
fid
,
SSkipListIterator
**
iters
,
SDataCols
*
pCols
)
{
int
flag
=
0
;
int
hasDataToCommit
=
0
;
int
isNewLastFile
=
0
;
STsdbMeta
*
pMeta
=
pRepo
->
tsdbMeta
;
STsdbMeta
*
pMeta
=
pRepo
->
tsdbMeta
;
STsdbFileH
*
pFileH
=
pRepo
->
tsdbFileH
;
STsdbFileH
*
pFileH
=
pRepo
->
tsdbFileH
;
...
@@ -887,97 +893,125 @@ static int tsdbCommitToFile(STsdbRepo *pRepo, int fid, SSkipListIterator **iters
...
@@ -887,97 +893,125 @@ static int tsdbCommitToFile(STsdbRepo *pRepo, int fid, SSkipListIterator **iters
STable
*
pTable
=
pMeta
->
tables
[
tid
];
STable
*
pTable
=
pMeta
->
tables
[
tid
];
SSkipListIterator
*
pIter
=
iters
[
tid
];
SSkipListIterator
*
pIter
=
iters
[
tid
];
int
isLoadCompBlocks
=
0
;
int
isLoadCompBlocks
=
0
;
char
dataDir
[
128
]
=
"
\0
"
;
if
(
pIter
==
NULL
)
continue
;
if
(
pIter
==
NULL
)
continue
;
tdInitDataCols
(
pCols
,
pTable
->
schema
);
tdInitDataCols
(
pCols
,
pTable
->
schema
);
int
numOfWrites
=
0
;
int
numOfWrites
=
0
;
// while (tsdbReadRowsFromCache(pIter, maxKey, pCfg->maxRowsPerFileBlock, pCols)) {
int
maxRowsToRead
=
pCfg
->
maxRowsPerFileBlock
*
4
/
5
;
// We keep 20% of space for merge purpose
// break;
// Loop to read columns from cache
// if (!flag) {
while
(
tsdbReadRowsFromCache
(
pIter
,
maxKey
,
maxRowsToRead
,
pCols
))
{
// // There are data to commit to this file, we need to create/open it for read/write.
if
(
!
hasDataToCommit
)
{
// // At the meantime, we set the flag to prevent further create/open operations
// There are data to commit to this fileId, we need to create/open it for read/write.
// if (tsdbCreateFGroup(pFileH, pRepo->rootDir, fid, pCfg->maxTables) < 0) {
// At the meantime, we set the flag to prevent further create/open operations
// // TODO: deal with the ERROR here
tsdbGetDataDirName
(
pRepo
,
dataDir
);
// }
// // Open files for commit
// pGroup = tsdbOpenFilesForCommit(pFileH, fid);
// if (pGroup == NULL) {
// // TODO: deal with the ERROR here
// }
// // TODO: open .h file and if neccessary, open .l file
// {}
// pIndices = (SCompIdx *)malloc(sizeof(SCompIdx) * pCfg->maxTables);
// if (pIndices == NULL) {
// // TODO: deal with the ERROR
// }
// // load the SCompIdx part
// if (tsdbLoadCompIdx(pGroup, (void *)pIndices, pCfg->maxTables) < 0) {
// // TODO: deal with the ERROR here
// }
// // TODO: sendfile those not need changed table content
if
(
tsdbCreateFGroup
(
pFileH
,
dataDir
,
fid
,
pCfg
->
maxTables
)
<
0
)
{
// for (int ttid = 0; ttid < tid; ttid++) {
// TODO: deal with the ERROR here
// // SCompIdx *pIdx = &pIndices[ttid];
}
// // if (pIdx->len > 0) {
// Open files for commit
// // lseek(pGroup->files[TSDB_FILE_TYPE_HEAD].fd, pIdx->offset, 0, SEEK_CUR);
pGroup
=
tsdbOpenFilesForCommit
(
pFileH
,
fid
);
// // sendfile(fd, pGroup->files[TSDB_FILE_TYPE_HEAD].fd, NULL, pIdx->len);
if
(
pGroup
==
NULL
)
{
// // }
// TODO: deal with the ERROR here
// }
}
// flag = 1;
// TODO: open .h file and if neccessary, open .l file
// }
tsdbCreateFile
(
dataDir
,
fid
,
".h"
,
pCfg
->
maxTables
,
&
tFile
,
1
,
0
);
if
(
1
/*pGroup->files[TSDB_FILE_TYPE_LAST].size > TSDB_MAX_LAST_FILE_SIZE*/
)
{
// TODO: make it not to write the last file every time
tsdbCreateFile
(
dataDir
,
fid
,
".l"
,
pCfg
->
maxTables
,
&
lFile
,
0
,
0
);
isNewLastFile
=
1
;
}
// load the SCompIdx part
pIndices
=
(
SCompIdx
*
)
malloc
(
sizeof
(
SCompIdx
)
*
pCfg
->
maxTables
);
if
(
pIndices
==
NULL
)
{
// TODO: deal with the ERROR
}
if
(
tsdbLoadCompIdx
(
pGroup
,
(
void
*
)
pIndices
,
pCfg
->
maxTables
)
<
0
)
{
// TODO: deal with the ERROR here
}
// sendfile those not need to changed table content
lseek
(
pGroup
->
files
[
TSDB_FILE_TYPE_HEAD
].
fd
,
TSDB_FILE_HEAD_SIZE
+
sizeof
(
SCompIdx
)
*
pCfg
->
maxTables
,
SEEK_SET
);
lseek
(
tFile
.
fd
,
TSDB_FILE_HEAD_SIZE
+
sizeof
(
SCompIdx
)
*
pCfg
->
maxTables
,
SEEK_SET
);
for
(
int
ttid
=
0
;
ttid
<
tid
;
ttid
++
)
{
SCompIdx
*
tIdx
=
&
pIndices
[
ttid
];
if
(
tIdx
->
len
<=
0
)
continue
;
if
(
isNewLastFile
&&
tIdx
->
hasLast
)
{
// TODO: Need to load the SCompBlock part and copy to new last file
pCompInfo
=
(
SCompInfo
*
)
realloc
((
void
*
)
pCompInfo
,
tIdx
->
len
);
if
(
pCompInfo
==
NULL
)
{
/* TODO */
}
if
(
tsdbLoadCompBlocks
(
pGroup
,
tIdx
,
(
void
*
)
pCompInfo
)
<
0
)
{
/* TODO */
}
SCompBlock
*
pLastBlock
=
TSDB_COMPBLOCK_AT
(
pCompInfo
,
tIdx
->
numOfSuperBlocks
-
1
);
int
numOfSubBlocks
=
pLastBlock
->
numOfSubBlocks
;
assert
(
pLastBlock
->
last
);
if
(
tsdbCopyCompBlockToFile
(
&
pGroup
->
files
[
TSDB_FILE_TYPE_LAST
],
&
lFile
,
pCompInfo
,
tIdx
->
numOfSuperBlocks
,
1
/* isOutLast*/
)
<
0
)
{
/* TODO */
}
{
if
(
numOfSubBlocks
>
1
)
tIdx
->
len
-=
(
sizeof
(
SCompBlock
)
*
numOfSubBlocks
);
tIdx
->
checksum
=
0
;
}
write
(
tFile
.
fd
,
(
void
*
)
pCompInfo
,
tIdx
->
len
);
tFile
.
size
+=
tIdx
->
len
;
}
else
{
sendfile
(
pGroup
->
files
[
TSDB_FILE_TYPE_HEAD
].
fd
,
tFile
.
fd
,
NULL
,
tIdx
->
len
);
tFile
.
size
+=
tIdx
->
len
;
}
}
// SCompIdx *pIdx = &pIndices[tid];
hasDataToCommit
=
1
;
}
// /* The first time to write to the table, need to decide
// * if it is neccessary to load the SComplock part. If it
// * is needed, just load it, or, just use sendfile and
// * append it.
// */
// if (numOfWrites == 0 && pIdx->offset > 0) {
// if (dataColsKeyFirst(pCols) <= pIdx->maxKey || pIdx->hasLast) {
// pCompInfo = (SCompInfo *)realloc((void *)pCompInfo, pIdx->len);
// if (tsdbLoadCompBlocks(pGroup, pIdx, (void *)pCompInfo) < 0) {
// // TODO: deal with the ERROR here
// }
// if (pCompInfo->uid == pTable->tableId.uid) isLoadCompBlocks = 1;
// } else {
// // TODO: sendfile the prefix part
// }
// }
// // if (tsdbWriteBlockToFile(pGroup, pCompInfo, pIdx, isLoadCompBlocks, pBlock, pCols) < 0) {
SCompIdx
*
pIdx
=
&
pIndices
[
tid
];
// // // TODO: deal with the ERROR here
// // }
// // pCompInfo = tsdbMergeBlock(pCompInfo, pBlock);
/* The first time to write to the table, need to decide
* if it is neccessary to load the SComplock part. If it
* is needed, just load it, or, just use sendfile and
* append it.
*/
if
(
numOfWrites
==
0
&&
pIdx
->
offset
>
0
)
{
if
(
dataColsKeyFirst
(
pCols
)
<=
pIdx
->
maxKey
||
pIdx
->
hasLast
)
{
pCompInfo
=
(
SCompInfo
*
)
realloc
((
void
*
)
pCompInfo
,
pIdx
->
len
);
if
(
tsdbLoadCompBlocks
(
pGroup
,
pIdx
,
(
void
*
)
pCompInfo
)
<
0
)
{
// TODO: deal with the ERROR here
}
if
(
pCompInfo
->
uid
==
pTable
->
tableId
.
uid
)
isLoadCompBlocks
=
1
;
}
else
{
// TODO: sendfile the prefix part
}
}
// if (tsdbWriteBlockToFile(pGroup, pCompInfo, pIdx, isLoadCompBlocks, pBlock, pCols) < 0) {
// // TODO: deal with the ERROR here
// }
// // if (1 /* the SCompBlock part is not loaded*/) {
// pCompInfo = tsdbMergeBlock(pCompInfo, pBlock);
// // // Append to .data file generate a SCompBlock and record it
// // } else {
// // }
// // // TODO: need to reset the pCols
// numOfWrites++;
// if (1 /* the SCompBlock part is not loaded*/) {
// // Append to .data file generate a SCompBlock and record it
// } else {
// }
// }
// if (pCols->numOfPoints > 0) {
// // TODO: need to reset the pCols
// // TODO: still has data to commit, commit it
// }
// if (1/* SCompBlock part is loaded, write it to .head file*/) {
numOfWrites
++
;
// // TODO
}
// } else {
// // TODO: use sendfile send the old part and append the newly added part
if
(
pCols
->
numOfPoints
>
0
)
{
// }
// TODO: still has data to commit, commit it
}
if
(
1
/* SCompBlock part is loaded, write it to .head file*/
)
{
// TODO
}
else
{
// TODO: use sendfile send the old part and append the newly added part
}
}
}
// Write the SCompIdx part
// Write the SCompIdx part
// Close all files and return
// Close all files and return
if
(
flag
)
{
if
(
hasDataToCommit
)
{
// TODO
// TODO
}
}
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录