Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
53226eae
TDengine
项目概览
taosdata
/
TDengine
1 年多 前同步成功
通知
1185
Star
22016
Fork
4786
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
1
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
TDengine
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
1
Issue
1
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
提交
53226eae
编写于
12月 26, 2022
作者:
H
Hongze Cheng
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
more code
上级
07e89142
变更
1
隐藏空白更改
内联
并排
Showing
1 changed file
with
132 addition
and
46 deletion
+132
-46
source/dnode/vnode/src/tsdb/tsdbCompact.c
source/dnode/vnode/src/tsdb/tsdbCompact.c
+132
-46
未找到文件。
source/dnode/vnode/src/tsdb/tsdbCompact.c
浏览文件 @
53226eae
...
...
@@ -15,6 +15,10 @@
#include "tsdb.h"
#define TSDB_ITER_TYPE_MEM 0x0
#define TSDB_ITER_TYPE_DAT 0x1
#define TSDB_ITER_TYPE_STT 0x2
typedef
struct
{
}
SMemDIter
;
...
...
@@ -37,7 +41,9 @@ typedef struct {
int32_t
iRow
;
}
SSttDIter
;
typedef
struct
{
typedef
struct
STsdbDataIter
{
struct
STsdbDataIter
*
next
;
int32_t
flag
;
SRowInfo
rowInfo
;
SRBTreeNode
n
;
...
...
@@ -45,19 +51,22 @@ typedef struct {
}
STsdbDataIter
;
typedef
struct
{
STsdb
*
pTsdb
;
STsdbFS
fs
;
int64_t
cid
;
int32_t
fid
;
SDataFReader
*
pReader
;
SDFileSet
*
pDFileSet
;
SRBTree
rtree
;
SBlockData
bData
;
STsdb
*
pTsdb
;
STsdbFS
fs
;
int64_t
cid
;
int32_t
fid
;
SDFileSet
*
pDFileSet
;
SDataFReader
*
pReader
;
STsdbDataIter
*
iterList
;
// list of iterators
SRBTree
rtree
;
SBlockData
bData
;
}
STsdbCompactor
;
#define TSDB_FLG_DEEP_COMPACT 0x1
// ITER =========================
static
int32_t
tsdbDataIterNext
(
STsdbDataIter
*
pIter
);
static
int32_t
tsdbDataIterCmprFn
(
const
SRBTreeNode
*
n1
,
const
SRBTreeNode
*
n2
)
{
const
STsdbDataIter
*
pIter1
=
(
STsdbDataIter
*
)((
char
*
)
n1
-
offsetof
(
STsdbDataIter
,
n
));
const
STsdbDataIter
*
pIter2
=
(
STsdbDataIter
*
)((
char
*
)
n2
-
offsetof
(
STsdbDataIter
,
n
));
...
...
@@ -95,6 +104,7 @@ static int32_t tsdbDataDIterOpen(SDataFReader *pReader, STsdbDataIter **ppIter)
code
=
TSDB_CODE_OUT_OF_MEMORY
;
goto
_exit
;
}
pIter
->
flag
=
TSDB_ITER_TYPE_DAT
;
SDataDIter
*
pDataDIter
=
(
SDataDIter
*
)
pIter
->
handle
;
pDataDIter
->
pReader
=
pReader
;
...
...
@@ -109,21 +119,16 @@ static int32_t tsdbDataDIterOpen(SDataFReader *pReader, STsdbDataIter **ppIter)
if
(
taosArrayGetSize
(
pDataDIter
->
aBlockIdx
)
==
0
)
goto
_clear_exit
;
// TODO
code
=
tBlockDataCreate
(
&
pDataDIter
->
bData
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
// read first data block
pDataDIter
->
iBlockIdx
=
0
;
code
=
tsdbReadDataBlk
(
pReader
,
taosArrayGet
(
pDataDIter
->
aBlockIdx
,
pDataDIter
->
iBlockIdx
),
&
pDataDIter
->
mDataBlk
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
pDataDIter
->
iBlockIdx
=
-
1
;
pDataDIter
->
iDataBlk
=
0
;
// code = tsdbReadDataBlock(pReader, tMapDat);
// TSDB_CHECK_CODE(code, lino, _exit);
pDataDIter
->
iRow
=
0
;
// TODO
code
=
tsdbDataIterNext
(
pIter
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
_exit:
if
(
code
)
{
...
...
@@ -150,6 +155,7 @@ static int32_t tsdbSttDIterOpen(SDataFReader *pReader, int32_t iStt, STsdbDataIt
code
=
TSDB_CODE_OUT_OF_MEMORY
;
goto
_exit
;
}
pIter
->
flag
=
TSDB_ITER_TYPE_STT
;
SSttDIter
*
pSttDIter
=
(
SSttDIter
*
)
pIter
->
handle
;
pSttDIter
->
pReader
=
pReader
;
...
...
@@ -168,13 +174,11 @@ static int32_t tsdbSttDIterOpen(SDataFReader *pReader, int32_t iStt, STsdbDataIt
code
=
tBlockDataCreate
(
&
pSttDIter
->
bData
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
pSttDIter
->
iSttBlk
=
0
;
// code = tsdbReadSttBlock(pReader, taosArrayGet(pSttDIter->aSttBlk, pSttDIter->iSttBlk), &pSttDIter->bData);
// TSDB_CHECK_CODE(code, lino, _exit);
pSttDIter
->
iSttBlk
=
-
1
;
pSttDIter
->
iRow
=
-
1
;
pSttDIter
->
iRow
=
0
;
// TODO
code
=
tsdbDataIterNext
(
pIter
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
_exit:
if
(
code
)
{
...
...
@@ -193,12 +197,42 @@ _exit:
static
void
tsdbDataIterClose
(
STsdbDataIter
*
pIter
)
{
// TODO
ASSERT
(
0
);
}
static
int32_t
tsdbDataIterNext
(
STsdbDataIter
*
pIter
)
{
int32_t
code
=
0
;
int32_t
lino
=
0
;
// TODO
if
(
pIter
->
flag
&
TSDB_ITER_TYPE_MEM
)
{
// TODO
ASSERT
(
0
);
}
else
if
(
pIter
->
flag
&
TSDB_ITER_TYPE_DAT
)
{
// TODO
ASSERT
(
0
);
}
else
if
(
pIter
->
flag
&
TSDB_ITER_TYPE_STT
)
{
SSttDIter
*
pSttDIter
=
(
SSttDIter
*
)
pIter
->
handle
;
pSttDIter
->
iRow
++
;
if
(
pSttDIter
->
iRow
<
pSttDIter
->
bData
.
nRow
)
{
ASSERT
(
0
);
}
else
{
pSttDIter
->
iSttBlk
++
;
if
(
pSttDIter
->
iSttBlk
<
taosArrayGetSize
(
pSttDIter
->
aSttBlk
))
{
code
=
tsdbReadSttBlock
(
pSttDIter
->
pReader
,
pSttDIter
->
iStt
,
taosArrayGet
(
pSttDIter
->
aSttBlk
,
pSttDIter
->
iSttBlk
),
&
pSttDIter
->
bData
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
pSttDIter
->
iRow
=
0
;
}
else
{
// code = TSDB_CODE_TDB_NO_DATA;
// goto _exit;
}
}
}
else
{
ASSERT
(
0
);
}
_exit:
return
code
;
}
...
...
@@ -289,52 +323,102 @@ _exit:
return
code
;
}
int32_t
tsdbCompact
(
STsdb
*
pTsdb
,
int32_t
flag
)
{
static
int32_t
tsdbOpenCompactor
(
STsdbCompactor
*
pCompactor
)
{
int32_t
code
=
0
;
int32_t
lino
=
0
;
// Check if can do compact (TODO)
STsdb
*
pTsdb
=
pCompactor
->
pTsdb
;
// Do compact
STsdbCompactor
compactor
=
{
0
};
// next compact file
pCompactor
->
pDFileSet
=
(
SDFileSet
*
)
taosArraySearch
(
pCompactor
->
fs
.
aDFileSet
,
&
(
SDFileSet
){.
fid
=
pCompactor
->
fid
},
tDFileSetCmprFn
,
TD_GT
);
if
(
pCompactor
->
pDFileSet
==
NULL
)
goto
_exit
;
code
=
tsdbBeginCompact
(
pTsdb
,
&
compactor
);
pCompactor
->
fid
=
pCompactor
->
pDFileSet
->
fid
;
code
=
tsdbDataFReaderOpen
(
&
pCompactor
->
pReader
,
pTsdb
,
pCompactor
->
pDFileSet
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
while
(
true
)
{
compactor
.
pDFileSet
=
(
SDFileSet
*
)
taosArraySearch
(
compactor
.
fs
.
aDFileSet
,
&
compactor
.
fid
,
tDFileSetCmprFn
,
TD_GT
);
if
(
compactor
.
pDFileSet
==
NULL
)
break
;
// open iters
STsdbDataIter
*
pIter
;
compactor
.
fid
=
compactor
.
pDFileSet
->
fid
;
pCompactor
->
iterList
=
NULL
;
tRBTreeCreate
(
&
pCompactor
->
rtree
,
tsdbDataIterCmprFn
);
code
=
tsdbDataFReaderOpen
(
&
compactor
.
pReader
,
pTsdb
,
compactor
.
pDFileSet
);
code
=
tsdbDataDIterOpen
(
pCompactor
->
pReader
,
&
pIter
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
if
(
pIter
)
{
pIter
->
next
=
pCompactor
->
iterList
;
pCompactor
->
iterList
=
pIter
;
tRBTreePut
(
&
pCompactor
->
rtree
,
&
pIter
->
n
);
}
for
(
int32_t
iStt
=
0
;
iStt
<
pCompactor
->
pReader
->
pSet
->
nSttF
;
iStt
++
)
{
code
=
tsdbSttDIterOpen
(
pCompactor
->
pReader
,
iStt
,
&
pIter
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
// open those iterators
tRBTreeCreate
(
&
compactor
.
rtree
,
tsdbDataIterCmprFn
);
if
(
pIter
)
{
pIter
->
next
=
pCompactor
->
iterList
;
pCompactor
->
iterList
=
pIter
;
tRBTreePut
(
&
pCompactor
->
rtree
,
&
pIter
->
n
);
}
}
_exit:
if
(
code
)
{
tsdbError
(
"vgId:%d %s failed at line %d since %s"
,
TD_VID
(
pTsdb
->
pVnode
),
__func__
,
lino
,
tstrerror
(
code
));
}
else
{
tsdbDebug
(
"vgId:%d %s done"
,
TD_VID
(
pTsdb
->
pVnode
),
__func__
);
}
return
code
;
}
STsdbDataIter
*
pIter
;
static
void
tsdbCloseCompactor
(
STsdbCompactor
*
pCompactor
)
{
STsdb
*
pTsdb
=
pCompactor
->
pTsdb
;
code
=
tsdbDataDIterOpen
(
compactor
.
pReader
,
&
pIter
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
for
(
STsdbDataIter
*
pIter
=
pCompactor
->
iterList
;
pIter
;)
{
STsdbDataIter
*
pIterNext
=
pIter
->
next
;
tsdbDataIterClose
(
pIter
);
pIter
=
pIterNext
;
}
if
(
pIter
)
tRBTreePut
(
&
compactor
.
rtree
,
&
pIter
->
n
);
// TODO
ASSERT
(
0
);
for
(
int32_t
iStt
=
0
;
iStt
<
compactor
.
pReader
->
pSet
->
nSttF
;
iStt
++
)
{
code
=
tsdbSttDIterOpen
(
compactor
.
pReader
,
iStt
,
&
pIter
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
_exit:
tsdbDebug
(
"vgId:%d %s done"
,
TD_VID
(
pTsdb
->
pVnode
),
__func__
);
}
if
(
pIter
)
tRBTreePut
(
&
compactor
.
rtree
,
&
pIter
->
n
);
}
int32_t
tsdbCompact
(
STsdb
*
pTsdb
,
int32_t
flag
)
{
int32_t
code
=
0
;
int32_t
lino
=
0
;
// Check if can do compact (TODO)
// Do compact
STsdbCompactor
compactor
=
{
0
};
code
=
tsdbBeginCompact
(
pTsdb
,
&
compactor
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
while
(
true
)
{
code
=
tsdbOpenCompactor
(
&
compactor
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
if
(
compactor
.
pDFileSet
==
NULL
)
break
;
// loop to merge row by row
TSDBROW
*
pRow
=
NULL
;
int64_t
nRow
=
0
;
for
(;;)
{
code
=
tsdbCompactNextRow
(
&
compactor
,
&
pRow
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
if
(
pRow
==
NULL
)
break
;
nRow
++
;
// code = tBlockDataAppendRow(&compactor.bData, pRow, pRow, NULL, 0);
// TSDB_CHECK_CODE(code, lino, _exit);
...
...
@@ -343,6 +427,8 @@ int32_t tsdbCompact(STsdb *pTsdb, int32_t flag) {
// TSDB_CHECK_CODE(code, lino, _exit);
// }
}
tsdbCloseCompactor
(
&
compactor
);
}
_exit:
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录