Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
555469c7
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看板
提交
555469c7
编写于
6月 14, 2022
作者:
H
Hongze Cheng
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
more work
上级
c9a3a7a6
变更
4
隐藏空白更改
内联
并排
Showing
4 changed file
with
222 addition
and
352 deletion
+222
-352
source/dnode/vnode/src/inc/tsdb.h
source/dnode/vnode/src/inc/tsdb.h
+5
-11
source/dnode/vnode/src/tsdb/tsdbCommit.c
source/dnode/vnode/src/tsdb/tsdbCommit.c
+121
-171
source/dnode/vnode/src/tsdb/tsdbReaderWriter.c
source/dnode/vnode/src/tsdb/tsdbReaderWriter.c
+94
-59
source/dnode/vnode/src/tsdb/tsdbUtil.c
source/dnode/vnode/src/tsdb/tsdbUtil.c
+2
-111
未找到文件。
source/dnode/vnode/src/inc/tsdb.h
浏览文件 @
555469c7
...
...
@@ -42,7 +42,6 @@ typedef struct SMemTable SMemTable;
typedef
struct
STbDataIter
STbDataIter
;
typedef
struct
SMergeInfo
SMergeInfo
;
typedef
struct
STable
STable
;
typedef
struct
SOffset
SOffset
;
typedef
struct
SMapData
SMapData
;
typedef
struct
SColData
SColData
;
typedef
struct
SColDataBlock
SColDataBlock
;
...
...
@@ -93,10 +92,11 @@ int32_t tsdbFSEnd(STsdbFS *pFS, int8_t rollback);
// SDataFWriter
typedef
struct
SDataFWriter
SDataFWriter
;
int32_t
tsdbDataFWriterOpen
(
SDataFWriter
*
pWriter
,
STsdb
*
pTsdb
,
SDFileSet
*
pSet
);
int32_t
tsdbDataFWriterClose
(
SDataFWriter
*
pWriter
);
int32_t
tsdbDataFWriterOpen
(
SDataFWriter
**
ppWriter
,
STsdb
*
pTsdb
,
SDFileSet
*
pSet
);
int32_t
tsdbDataFWriterClose
(
SDataFWriter
*
pWriter
,
int8_t
sync
);
int32_t
tsdbUpdateDFileSetHeader
(
SDataFWriter
*
pWriter
,
uint8_t
**
ppBuf
);
int32_t
tsdbWriteBlockIdx
(
SDataFWriter
*
pWriter
,
SMapData
*
pMapData
,
uint8_t
**
ppBuf
);
int32_t
tsdbWriteBlock
(
SDataFWriter
*
pWriter
,
SMapData
*
pMapData
,
uint8_t
**
ppBuf
,
int64_t
*
rOffset
,
int64_t
*
rSize
);
int32_t
tsdbWriteBlock
(
SDataFWriter
*
pWriter
,
SMapData
*
pMapData
,
uint8_t
**
ppBuf
,
SBlockIdx
*
pBlockIdx
);
int32_t
tsdbWriteBlockData
(
SDataFWriter
*
pWriter
,
SColDataBlock
*
pBlockData
,
uint8_t
**
ppBuf
,
int64_t
*
rOffset
,
int64_t
*
rSize
);
int32_t
tsdbWriteBlockSMA
(
SDataFWriter
*
pWriter
,
SBlockSMA
*
pBlockSMA
,
int64_t
*
rOffset
,
int64_t
*
rSize
);
...
...
@@ -104,7 +104,7 @@ int32_t tsdbWriteBlockSMA(SDataFWriter *pWriter, SBlockSMA *pBlockSMA, int64_t *
// SDataFReader
typedef
struct
SDataFReader
SDataFReader
;
int32_t
tsdbDataFReaderOpen
(
SDataFReader
*
pReader
,
STsdb
*
pTsdb
,
SDFileSet
*
pSet
);
int32_t
tsdbDataFReaderOpen
(
SDataFReader
*
*
p
pReader
,
STsdb
*
pTsdb
,
SDFileSet
*
pSet
);
int32_t
tsdbDataFReaderClose
(
SDataFReader
*
pReader
);
int32_t
tsdbReadBlockIdx
(
SDataFReader
*
pReader
,
SMapData
*
pMapData
,
uint8_t
**
ppBuf
);
int32_t
tsdbReadBlock
(
SDataFReader
*
pReader
,
SBlockIdx
*
pBlockIdx
,
SMapData
*
pMapData
,
uint8_t
**
ppBuf
);
...
...
@@ -165,12 +165,6 @@ int32_t tPutDelFileHdr(uint8_t *p, SDelFile *pDelFile);
int32_t
tGetDelFileHdr
(
uint8_t
*
p
,
SDelFile
*
pDelFile
);
// structs
struct
SOffset
{
int32_t
nOffset
;
uint8_t
flag
;
uint8_t
*
pOffset
;
};
typedef
struct
{
int
minFid
;
int
midFid
;
...
...
source/dnode/vnode/src/tsdb/tsdbCommit.c
浏览文件 @
555469c7
...
...
@@ -16,34 +16,24 @@
#include "tsdb.h"
typedef
struct
{
STsdb
*
pTsdb
;
uint8_t
*
pBuf1
;
uint8_t
*
pBuf2
;
uint8_t
*
pBuf3
;
uint8_t
*
pBuf4
;
uint8_t
*
pBuf5
;
STsdb
*
pTsdb
;
/* commit data */
int32_t
minutes
;
int8_t
precision
;
int32_t
minRow
;
int32_t
maxRow
;
// --------------
TSKEY
nextKey
;
int32_t
commitFid
;
TSKEY
minKey
;
TSKEY
maxKey
;
// commit file data
TSKEY
nextKey
;
int32_t
commitFid
;
TSKEY
minKey
;
TSKEY
maxKey
;
SDataFReader
*
pReader
;
SMapData
oBlockIdx
;
// SMapData<SBlockIdx>, read from reader
SMapData
oBlock
;
// SMapData<SBlock>, read from reader
SDataFWriter
*
pWriter
;
SMapData
nBlockIdx
;
// SMapData<SBlockIdx>, build by committer
// commit table data
STbDataIter
iter
;
STbDataIter
*
pIter
;
SBlockIdx
*
pBlockIdx
;
SMapData
oBlock
;
SMapData
nBlock
;
SColDataBlock
oColDataBlock
;
SColDataBlock
nColDataBlock
;
SMapData
nBlock
;
// SMapData<SBlock>
/* commit del */
SDelFReader
*
pDelFReader
;
SMapData
oDelIdxMap
;
// SMapData<SDelIdx>, old
...
...
@@ -124,21 +114,6 @@ _err:
return
code
;
}
static
int32_t
tsdbStartCommit
(
STsdb
*
pTsdb
,
SCommitter
*
pCommitter
)
{
int32_t
code
=
0
;
memset
(
pCommitter
,
0
,
sizeof
(
*
pCommitter
));
ASSERT
(
pTsdb
->
mem
&&
pTsdb
->
imem
==
NULL
);
// lock();
pTsdb
->
imem
=
pTsdb
->
mem
;
pTsdb
->
mem
=
NULL
;
// unlock();
pCommitter
->
pTsdb
=
pTsdb
;
return
code
;
}
static
int32_t
tsdbCommitDelStart
(
SCommitter
*
pCommitter
)
{
int32_t
code
=
0
;
STsdb
*
pTsdb
=
pCommitter
->
pTsdb
;
...
...
@@ -371,40 +346,37 @@ _err:
return
code
;
}
static
int32_t
tsdbEndCommit
(
SCommitter
*
pCommitter
,
int32_t
eno
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
static
int32_t
tsdbCommitFileDataStart
(
SCommitter
*
pCommitter
);
static
int32_t
tsdbCommitFileDataImpl
(
SCommitter
*
pCommitter
);
static
int32_t
tsdbCommitFileDataEnd
(
SCommitter
*
pCommitter
);
static
int32_t
tsdbCommitFileData
(
SCommitter
*
pCommitter
)
{
int32_t
code
=
0
;
static
int32_t
tsdbCommitTableData
(
SCommitter
*
pCommitter
,
STbData
*
pTbData
,
SBlockIdx
*
pBlockIdx
)
{
int32_t
code
=
0
;
STbDataIter
iter
;
TSDBROW
row
;
SBlockIdx
blockIdx
;
// commit file data start
code
=
tsdbCommitFileDataStart
(
pCommitter
);
if
(
code
)
{
goto
_err
;
// check: if no memory data and no disk data, exit
if
(
pTbData
)
{
tsdbTbDataIterOpen
(
pTbData
,
&
(
TSDBKEY
){.
ts
=
pCommitter
->
minKey
,
.
version
=
0
},
0
,
&
iter
);
if
((
!
tsdbTbDataIterGet
(
&
iter
,
&
row
)
||
row
.
pTSRow
->
ts
>
pCommitter
->
maxKey
)
&&
pBlockIdx
==
NULL
)
{
goto
_exit
;
}
}
// commit file data impl
code
=
tsdbCommitFileDataImpl
(
pCommitter
);
if
(
code
)
{
goto
_err
;
// start
tMapDataReset
(
&
pCommitter
->
oBlock
);
tMapDataReset
(
&
pCommitter
->
nBlock
);
if
(
pBlockIdx
)
{
code
=
tsdbReadBlock
(
pCommitter
->
pReader
,
pBlockIdx
,
&
pCommitter
->
oBlock
,
NULL
);
if
(
code
)
goto
_err
;
}
// commit file data end
code
=
tsdbCommitFileDataEnd
(
pCommitter
);
if
(
code
)
{
goto
_err
;
}
// impl
// end
_exit:
return
code
;
_err:
tsdbError
(
"vgId:%d commit Table data failed since %s"
,
TD_VID
(
pCommitter
->
pTsdb
->
pVnode
),
tstrerror
(
code
));
return
code
;
}
...
...
@@ -415,20 +387,23 @@ static int32_t tsdbCommitFileDataStart(SCommitter *pCommitter) {
SDFileSet
*
pWSet
=
NULL
;
// TODO
// memory
pCommitter
->
nextKey
=
TSKEY_MAX
;
tMapDataReset
(
&
pCommitter
->
oBlockIdx
);
tMapDataReset
(
&
pCommitter
->
oBlock
);
tMapDataReset
(
&
pCommitter
->
nBlockIdx
);
tMapDataReset
(
&
pCommitter
->
nBlock
);
// load old
if
(
pRSet
)
{
code
=
tsdbD
FileSet
ReaderOpen
(
&
pCommitter
->
pReader
,
pTsdb
,
pRSet
);
code
=
tsdbD
ataF
ReaderOpen
(
&
pCommitter
->
pReader
,
pTsdb
,
pRSet
);
if
(
code
)
goto
_err
;
code
=
tsdbReadBlockIdx
(
pCommitter
->
pReader
,
&
pCommitter
->
oBlockIdx
);
code
=
tsdbReadBlockIdx
(
pCommitter
->
pReader
,
&
pCommitter
->
oBlockIdx
,
NULL
);
if
(
code
)
goto
_err
;
}
// create new
code
=
tsdbD
FileSet
WriterOpen
(
&
pCommitter
->
pWriter
,
pTsdb
,
pWSet
);
code
=
tsdbD
ataF
WriterOpen
(
&
pCommitter
->
pWriter
,
pTsdb
,
pWSet
);
if
(
code
)
goto
_err
;
_exit:
...
...
@@ -439,10 +414,9 @@ _err:
return
code
;
}
static
int32_t
tsdbCommitTableData
(
SCommitter
*
pCommitter
);
static
int32_t
tsdbCommitFileDataImpl
(
SCommitter
*
pCommitter
)
{
int32_t
code
=
0
;
int32_t
c
;
STsdb
*
pTsdb
=
pCommitter
->
pTsdb
;
SMemTable
*
pMemTable
=
pTsdb
->
imem
;
int32_t
iTbData
=
0
;
...
...
@@ -450,77 +424,104 @@ static int32_t tsdbCommitFileDataImpl(SCommitter *pCommitter) {
int32_t
iBlockIdx
=
0
;
int32_t
nBlockIdx
=
pCommitter
->
oBlockIdx
.
nItem
;
STbData
*
pTbData
;
SBlockIdx
*
pBlockIdx
;
SBlockIdx
*
pBlockIdx
=
NULL
;
SBlockIdx
blockIdx
;
int32_t
c
;
while
(
iTbData
<
nTbData
||
iBlockIdx
<
nBlockIdx
)
{
pTbData
=
NULL
;
pBlockIdx
=
NULL
;
if
(
iTbData
<
nTbData
)
{
pTbData
=
(
STbData
*
)
taosArrayGetP
(
pMemTable
->
aTbData
,
iTbData
)
;
}
if
(
iBlockIdx
<
nBlockIdx
)
{
tMapDataGetItemByIdx
(
&
pCommitter
->
oBlockIdx
,
iBlockIdx
,
&
blockIdx
,
NULL
/* TODO */
);
pBlockIdx
=
&
blockIdx
;
}
ASSERT
(
nTbData
>
0
);
pTbData
=
(
STbData
*
)
taosArrayGetP
(
pMemTable
->
aTbData
,
iTbData
)
;
if
(
iBlockIdx
<
nBlockIdx
)
{
code
=
tMapDataGetItemByIdx
(
&
pCommitter
->
oBlockIdx
,
iBlockIdx
,
&
blockIdx
,
tGetBlockIdx
);
if
(
code
)
goto
_err
;
pBlockIdx
=
&
blockIdx
;
}
while
(
true
)
{
if
(
pTbData
==
NULL
&&
pBlockIdx
==
NULL
)
break
;
if
(
pTbData
&&
pBlockIdx
)
{
c
=
tTABLEIDCmprFn
(
pTbData
,
pBlockIdx
);
if
(
c
==
0
)
{
iTbData
++
;
iBlockIdx
++
;
goto
_commit_mem_and_disk_data
;
}
else
if
(
c
<
0
)
{
iTbData
++
;
pBlockIdx
=
NULL
;
goto
_commit_mem_data
;
}
else
{
iBlockIdx
++
;
pTbData
=
NULL
;
goto
_commit_disk_data
;
}
}
else
if
(
pTbData
)
{
goto
_commit_mem_data
;
}
else
{
if
(
pTbData
)
{
iBlockIdx
++
;
}
if
(
pBlockIdx
)
{
iTbData
++
;
}
goto
_commit_disk_data
;
}
if
(
pTbData
&&
!
tsdbTbDataIterOpen
(
pTbData
,
&
(
TSDBKEY
){.
ts
=
pCommitter
->
minKey
,
.
version
=
0
},
0
,
&
pCommitter
->
iter
))
{
_commit_mem_data:
code
=
tsdbCommitTableData
(
pCommitter
,
pTbData
,
NULL
);
if
(
code
)
goto
_err
;
iTbData
++
;
if
(
iTbData
<
nTbData
)
{
pTbData
=
(
STbData
*
)
taosArrayGetP
(
pMemTable
->
aTbData
,
iTbData
);
}
else
{
pTbData
=
NULL
;
}
continue
;
if
(
pTbData
==
NULL
&&
pBlockIdx
==
NULL
)
continue
;
pCommitter
->
pTbData
=
pTbData
;
pCommitter
->
pBlockIdx
=
pBlockIdx
;
_commit_disk_data:
code
=
tsdbCommitTableData
(
pCommitter
,
NULL
,
pBlockIdx
);
if
(
code
)
goto
_err
;
iBlockIdx
++
;
if
(
iBlockIdx
<
nBlockIdx
)
{
code
=
tMapDataGetItemByIdx
(
&
pCommitter
->
oBlockIdx
,
iBlockIdx
,
&
blockIdx
,
tGetBlockIdx
);
if
(
code
)
goto
_err
;
pBlockIdx
=
&
blockIdx
;
}
else
{
pBlockIdx
=
NULL
;
}
continue
;
code
=
tsdbCommitTableData
(
pCommitter
);
_commit_mem_and_disk_data:
code
=
tsdbCommitTableData
(
pCommitter
,
pTbData
,
pBlockIdx
);
if
(
code
)
goto
_err
;
iTbData
++
;
iBlockIdx
++
;
if
(
iTbData
<
nTbData
)
{
pTbData
=
(
STbData
*
)
taosArrayGetP
(
pMemTable
->
aTbData
,
iTbData
);
}
else
{
pTbData
=
NULL
;
}
if
(
iBlockIdx
<
nBlockIdx
)
{
code
=
tMapDataGetItemByIdx
(
&
pCommitter
->
oBlockIdx
,
iBlockIdx
,
&
blockIdx
,
tGetBlockIdx
);
if
(
code
)
goto
_err
;
pBlockIdx
=
&
blockIdx
;
}
else
{
pBlockIdx
=
NULL
;
}
continue
;
}
return
code
;
_err:
tsdbError
(
"vgId:%d commit file data impl failed since %s"
,
TD_VID
(
pTsdb
->
pVnode
),
tstrerror
(
code
));
return
code
;
}
static
int32_t
tsdbCommitFileDataEnd
(
SCommitter
*
pCommitter
)
{
int32_t
code
=
0
;
code
=
tsdbWriteBlockIdx
(
pCommitter
->
pWriter
,
pCommitter
->
nBlockIdx
,
NULL
);
// write blockIdx
code
=
tsdbWriteBlockIdx
(
pCommitter
->
pWriter
,
&
pCommitter
->
nBlockIdx
,
NULL
);
if
(
code
)
goto
_err
;
// update file header
code
=
tsdbUpdateDFileSetHeader
(
pCommitter
->
pWriter
,
NULL
);
if
(
code
)
goto
_err
;
code
=
tsdbDFileSetWriterClose
(
pCommitter
->
pWriter
,
1
);
// close and sync
code
=
tsdbDataFWriterClose
(
pCommitter
->
pWriter
,
1
);
if
(
code
)
goto
_err
;
if
(
pCommitter
->
pReader
)
{
code
=
tsdbD
FileSet
ReaderClose
(
pCommitter
->
pReader
);
code
=
tsdbD
ataF
ReaderClose
(
pCommitter
->
pReader
);
goto
_err
;
}
...
...
@@ -528,107 +529,50 @@ _exit:
return
code
;
_err:
tsdbError
(
"vgId:%d commit file data end failed since %s"
,
TD_VID
(
pCommitter
->
pTsdb
->
pVnode
),
tstrerror
(
code
));
return
code
;
}
static
int32_t
tsdbCommitTableDataStart
(
SCommitter
*
pCommitter
);
static
int32_t
tsdbCommitTableDataImpl
(
SCommitter
*
pCommitter
);
static
int32_t
tsdbCommitTableDataEnd
(
SCommitter
*
pCommitter
);
static
int32_t
tsdbCommitTableData
(
SCommitter
*
pCommitter
)
{
static
int32_t
tsdbCommitFileData
(
SCommitter
*
pCommitter
)
{
int32_t
code
=
0
;
// start
code
=
tsdbCommit
Tab
leDataStart
(
pCommitter
);
//
commit file data
start
code
=
tsdbCommit
Fi
leDataStart
(
pCommitter
);
if
(
code
)
{
goto
_err
;
}
// impl
code
=
tsdbCommit
Tab
leDataImpl
(
pCommitter
);
//
commit file data
impl
code
=
tsdbCommit
Fi
leDataImpl
(
pCommitter
);
if
(
code
)
{
goto
_err
;
}
// end
code
=
tsdbCommit
Tab
leDataEnd
(
pCommitter
);
//
commit file data
end
code
=
tsdbCommit
Fi
leDataEnd
(
pCommitter
);
if
(
code
)
{
goto
_err
;
}
_exit:
return
code
;
_err:
return
code
;
}
static
int32_t
tsdbCommitTableDataStart
(
SCommitter
*
pCommitter
)
{
// ----------------------------------------------------------------------------
static
int32_t
tsdbStartCommit
(
STsdb
*
pTsdb
,
SCommitter
*
pCommitter
)
{
int32_t
code
=
0
;
// old
tMapDataReset
(
&
pCommitter
->
oBlock
);
if
(
pCommitter
->
pBlockIdx
)
{
code
=
tsdbReadBlock
(
pCommitter
->
pReader
,
&
pCommitter
->
oBlock
,
NULL
);
if
(
code
)
goto
_err
;
}
// new
tMapDataReset
(
&
pCommitter
->
nBlock
);
_err:
return
code
;
}
static
int32_t
tsdbCommitTableDataImpl
(
SCommitter
*
pCommitter
)
{
int32_t
code
=
0
;
STsdb
*
pTsdb
=
pCommitter
->
pTsdb
;
STbDataIter
*
pIter
=
NULL
;
int32_t
iBlock
=
0
;
int32_t
nBlock
=
pCommitter
->
nBlock
.
nItem
;
SBlock
*
pBlock
;
SBlock
block
;
TSDBROW
*
pRow
;
TSDBROW
row
;
int32_t
iRow
=
0
;
STSchema
*
pTSchema
=
NULL
;
if
(
pCommitter
->
pTbData
)
{
code
=
tsdbTbDataIterCreate
(
pCommitter
->
pTbData
,
&
(
TSDBKEY
){.
ts
=
pCommitter
->
minKey
,
.
version
=
0
},
0
,
&
pIter
);
if
(
code
)
goto
_err
;
}
if
(
iBlock
<
nBlock
)
{
pBlock
=
&
block
;
}
else
{
pBlock
=
NULL
;
}
tsdbTbDataIterGet
(
pIter
,
pRow
);
// loop to merge memory data and disk data
for
(;
pBlock
==
NULL
||
(
pRow
&&
pRow
->
pTSRow
->
ts
<=
pCommitter
->
maxKey
);)
{
if
(
pRow
==
NULL
||
pRow
->
pTSRow
->
ts
>
pCommitter
->
maxKey
)
{
// only has block data, then move to new index file
}
else
if
(
0
)
{
// only commit memory data
}
else
{
// merge memory and block data
}
}
tsdbTbDataIterDestroy
(
pIter
);
return
code
;
memset
(
pCommitter
,
0
,
sizeof
(
*
pCommitter
));
ASSERT
(
pTsdb
->
mem
&&
pTsdb
->
imem
==
NULL
);
// lock();
pTsdb
->
imem
=
pTsdb
->
mem
;
pTsdb
->
mem
=
NULL
;
// unlock();
_err:
tsdbError
(
"vgId:%d commit table data impl failed since %s"
,
TD_VID
(
pTsdb
->
pVnode
),
tstrerror
(
code
));
tsdbTbDataIterDestroy
(
pIter
);
return
code
;
}
pCommitter
->
pTsdb
=
pTsdb
;
static
int32_t
tsdbCommitTableDataEnd
(
SCommitter
*
pCommitter
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
...
...
@@ -702,3 +646,9 @@ static int32_t tsdbCommitCache(SCommitter *pCommitter) {
// TODO
return
code
;
}
static
int32_t
tsdbEndCommit
(
SCommitter
*
pCommitter
,
int32_t
eno
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
source/dnode/vnode/src/tsdb/tsdbReaderWriter.c
浏览文件 @
555469c7
...
...
@@ -18,64 +18,6 @@
#define TSDB_FHDR_SIZE 512
#define TSDB_FILE_DLMT ((uint32_t)0xF00AFA0F)
// SDFileSetWritter ====================================================
struct
SDFileSetWritter
{
STsdb
*
pTsdb
;
int32_t
szBuf1
;
uint8_t
*
pBuf1
;
int32_t
szBuf2
;
uint8_t
*
pBuf2
;
};
// SDFileSetReader ====================================================
struct
SDFileSetReader
{
STsdb
*
pTsdb
;
int32_t
szBuf1
;
uint8_t
*
pBuf1
;
int32_t
szBuf2
;
uint8_t
*
pBuf2
;
};
int32_t
tsdbDFileSetReaderOpen
(
SDFileSetReader
*
pReader
,
STsdb
*
pTsdb
,
SDFileSet
*
pSet
)
{
int32_t
code
=
0
;
memset
(
pReader
,
0
,
sizeof
(
*
pReader
));
pReader
->
pTsdb
=
pTsdb
;
return
code
;
_err:
tsdbError
(
"vgId:%d failed to open SDFileSetReader since %s"
,
TD_VID
(
pTsdb
->
pVnode
),
tstrerror
(
code
));
return
code
;
}
int32_t
tsdbDFileSetReaderClose
(
SDFileSetReader
*
pReader
)
{
int32_t
code
=
0
;
taosMemoryFreeClear
(
pReader
->
pBuf1
);
taosMemoryFreeClear
(
pReader
->
pBuf2
);
return
code
;
}
int32_t
tsdbLoadSBlockIdx
(
SDFileSetReader
*
pReader
,
SArray
*
pArray
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbLoadSBlockInfo
(
SDFileSetReader
*
pReader
,
SBlockIdx
*
pBlockIdx
,
SBlockInfo
*
pBlockInfo
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbLoadSBlockStatis
(
SDFileSetReader
*
pReader
,
SBlock
*
pBlock
,
SBlockStatis
*
pBlockStatis
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
// SDelFWriter ====================================================
struct
SDelFWriter
{
STsdb
*
pTsdb
;
...
...
@@ -439,4 +381,97 @@ int32_t tsdbReadDelIdx(SDelFReader *pReader, SMapData *pDelIdxMap, uint8_t **ppB
_err:
tsdbError
(
"vgId:%d read del idx failed since %s"
,
TD_VID
(
pReader
->
pTsdb
->
pVnode
),
tstrerror
(
code
));
return
code
;
}
\ No newline at end of file
}
// SDataFReader ====================================================
struct
SDataFReader
{
STsdb
*
pTsdb
;
SDFileSet
*
pSet
;
TdFilePtr
pReadH
;
};
int32_t
tsdbDataFReaderOpen
(
SDataFReader
**
ppReader
,
STsdb
*
pTsdb
,
SDFileSet
*
pSet
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbDataFReaderClose
(
SDataFReader
*
pReader
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbReadBlockIdx
(
SDataFReader
*
pReader
,
SMapData
*
pMapData
,
uint8_t
**
ppBuf
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbReadBlock
(
SDataFReader
*
pReader
,
SBlockIdx
*
pBlockIdx
,
SMapData
*
pMapData
,
uint8_t
**
ppBuf
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbReadBlockData
(
SDataFReader
*
pReader
,
SBlock
*
pBlock
,
SColDataBlock
*
pBlockData
,
uint8_t
**
ppBuf
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbReadBlockSMA
(
SDataFReader
*
pReader
,
SBlockSMA
*
pBlkSMA
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
// SDataFWriter ====================================================
struct
SDataFWriter
{
STsdb
*
pTsdb
;
SDFileSet
*
pSet
;
TdFilePtr
pWriteH
;
};
int32_t
tsdbDataFWriterOpen
(
SDataFWriter
**
ppWriter
,
STsdb
*
pTsdb
,
SDFileSet
*
pSet
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbDataFWriterClose
(
SDataFWriter
*
pWriter
,
int8_t
sync
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbUpdateDFileSetHeader
(
SDataFWriter
*
pWriter
,
uint8_t
**
ppBuf
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbWriteBlockIdx
(
SDataFWriter
*
pWriter
,
SMapData
*
pMapData
,
uint8_t
**
ppBuf
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbWriteBlock
(
SDataFWriter
*
pWriter
,
SMapData
*
pMapData
,
uint8_t
**
ppBuf
,
SBlockIdx
*
pBlockIdx
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbWriteBlockData
(
SDataFWriter
*
pWriter
,
SColDataBlock
*
pBlockData
,
uint8_t
**
ppBuf
,
int64_t
*
rOffset
,
int64_t
*
rSize
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
int32_t
tsdbWriteBlockSMA
(
SDataFWriter
*
pWriter
,
SBlockSMA
*
pBlockSMA
,
int64_t
*
rOffset
,
int64_t
*
rSize
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
source/dnode/vnode/src/tsdb/tsdbUtil.c
浏览文件 @
555469c7
...
...
@@ -15,121 +15,10 @@
#include "tsdb.h"
// SOffset =======================================================================
#define TSDB_OFFSET_I32 ((uint8_t)0)
#define TSDB_OFFSET_I16 ((uint8_t)1)
#define TSDB_OFFSET_I8 ((uint8_t)2)
static
FORCE_INLINE
int32_t
tsdbOffsetSize
(
SOffset
*
pOfst
)
{
switch
(
pOfst
->
flag
)
{
case
TSDB_OFFSET_I32
:
return
sizeof
(
int32_t
);
case
TSDB_OFFSET_I16
:
return
sizeof
(
int16_t
);
case
TSDB_OFFSET_I8
:
return
sizeof
(
int8_t
);
default:
ASSERT
(
0
);
}
}
static
FORCE_INLINE
int32_t
tsdbGetOffset
(
SOffset
*
pOfst
,
int32_t
idx
)
{
int32_t
offset
=
-
1
;
if
(
idx
>=
0
&&
idx
<
pOfst
->
nOffset
)
{
switch
(
pOfst
->
flag
)
{
case
TSDB_OFFSET_I32
:
offset
=
((
int32_t
*
)
pOfst
->
pOffset
)[
idx
];
break
;
case
TSDB_OFFSET_I16
:
offset
=
((
int16_t
*
)
pOfst
->
pOffset
)[
idx
];
break
;
case
TSDB_OFFSET_I8
:
offset
=
((
int8_t
*
)
pOfst
->
pOffset
)[
idx
];
break
;
default:
ASSERT
(
0
);
}
ASSERT
(
offset
>=
0
);
}
return
offset
;
}
static
FORCE_INLINE
int32_t
tsdbAddOffset
(
SOffset
*
pOfst
,
int32_t
offset
)
{
int32_t
code
=
0
;
int32_t
nOffset
=
pOfst
->
nOffset
;
ASSERT
(
pOfst
->
flag
==
TSDB_OFFSET_I32
);
ASSERT
(
offset
>=
0
);
pOfst
->
nOffset
++
;
// alloc
code
=
tsdbRealloc
(
&
pOfst
->
pOffset
,
sizeof
(
int32_t
)
*
pOfst
->
nOffset
);
if
(
code
)
goto
_exit
;
// put
((
int32_t
*
)
pOfst
->
pOffset
)[
nOffset
]
=
offset
;
_exit:
return
code
;
}
static
FORCE_INLINE
int32_t
tPutOffset
(
uint8_t
*
p
,
SOffset
*
pOfst
)
{
int32_t
n
=
0
;
int32_t
maxOffset
;
ASSERT
(
pOfst
->
flag
==
TSDB_OFFSET_I32
);
ASSERT
(
pOfst
->
nOffset
>
0
);
maxOffset
=
tsdbGetOffset
(
pOfst
,
pOfst
->
nOffset
-
1
);
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pOfst
->
nOffset
);
if
(
maxOffset
<=
INT8_MAX
)
{
n
+=
tPutU8
(
p
?
p
+
n
:
p
,
TSDB_OFFSET_I8
);
for
(
int32_t
iOffset
=
0
;
iOffset
<
pOfst
->
nOffset
;
iOffset
++
)
{
n
+=
tPutI8
(
p
?
p
+
n
:
p
,
(
int8_t
)
tsdbGetOffset
(
pOfst
,
iOffset
));
}
}
else
if
(
maxOffset
<=
INT16_MAX
)
{
n
+=
tPutU8
(
p
?
p
+
n
:
p
,
TSDB_OFFSET_I16
);
for
(
int32_t
iOffset
=
0
;
iOffset
<
pOfst
->
nOffset
;
iOffset
++
)
{
n
+=
tPutI16
(
p
?
p
+
n
:
p
,
(
int16_t
)
tsdbGetOffset
(
pOfst
,
iOffset
));
}
}
else
{
n
+=
tPutU8
(
p
?
p
+
n
:
p
,
TSDB_OFFSET_I32
);
for
(
int32_t
iOffset
=
0
;
iOffset
<
pOfst
->
nOffset
;
iOffset
++
)
{
n
+=
tPutI32
(
p
?
p
+
n
:
p
,
(
int32_t
)
tsdbGetOffset
(
pOfst
,
iOffset
));
}
}
return
n
;
}
static
FORCE_INLINE
int32_t
tGetOffset
(
uint8_t
*
p
,
SOffset
*
pOfst
)
{
int32_t
n
=
0
;
n
+=
tGetI32v
(
p
+
n
,
&
pOfst
->
nOffset
);
n
+=
tGetU8
(
p
+
n
,
&
pOfst
->
flag
);
pOfst
->
pOffset
=
p
+
n
;
switch
(
pOfst
->
flag
)
{
case
TSDB_OFFSET_I32
:
n
=
n
+
pOfst
->
nOffset
+
sizeof
(
int32_t
);
break
;
case
TSDB_OFFSET_I16
:
n
=
n
+
pOfst
->
nOffset
+
sizeof
(
int16_t
);
break
;
case
TSDB_OFFSET_I8
:
n
=
n
+
pOfst
->
nOffset
+
sizeof
(
int8_t
);
break
;
default:
ASSERT
(
0
);
}
return
n
;
}
// SMapData =======================================================================
void
tMapDataReset
(
SMapData
*
pMapData
)
{
pMapData
->
flag
=
TSDB_OFFSET_I32
;
...
...
@@ -299,6 +188,7 @@ void tsdbFree(uint8_t *pBuf) {
}
}
// TABLEID =======================================================================
int32_t
tTABLEIDCmprFn
(
const
void
*
p1
,
const
void
*
p2
)
{
TABLEID
*
pId1
=
(
TABLEID
*
)
p1
;
TABLEID
*
pId2
=
(
TABLEID
*
)
p2
;
...
...
@@ -318,6 +208,7 @@ int32_t tTABLEIDCmprFn(const void *p1, const void *p2) {
return
0
;
}
// TSDBKEY =======================================================================
int32_t
tsdbKeyCmprFn
(
const
void
*
p1
,
const
void
*
p2
)
{
TSDBKEY
*
pKey1
=
(
TSDBKEY
*
)
p1
;
TSDBKEY
*
pKey2
=
(
TSDBKEY
*
)
p2
;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录