Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
71df0a00
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看板
提交
71df0a00
编写于
6月 12, 2023
作者:
H
Hongze Cheng
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
more code
上级
438d0cae
变更
2
显示空白变更内容
内联
并排
Showing
2 changed file
with
73 addition
and
88 deletion
+73
-88
source/dnode/vnode/src/tsdb/dev/tsdbCommit.c
source/dnode/vnode/src/tsdb/dev/tsdbCommit.c
+71
-88
source/dnode/vnode/src/tsdb/dev/tsdbMerge.c
source/dnode/vnode/src/tsdb/dev/tsdbMerge.c
+2
-0
未找到文件。
source/dnode/vnode/src/tsdb/dev/tsdbCommit.c
浏览文件 @
71df0a00
...
@@ -39,6 +39,7 @@ typedef struct {
...
@@ -39,6 +39,7 @@ typedef struct {
TSKEY
nextKey
;
TSKEY
nextKey
;
int32_t
fid
;
int32_t
fid
;
int32_t
expLevel
;
int32_t
expLevel
;
SDiskID
did
;
TSKEY
minKey
;
TSKEY
minKey
;
TSKEY
maxKey
;
TSKEY
maxKey
;
STFileSet
*
fset
;
STFileSet
*
fset
;
...
@@ -60,12 +61,6 @@ static int32_t tsdbCommitOpenNewSttWriter(SCommitter2 *committer) {
...
@@ -60,12 +61,6 @@ static int32_t tsdbCommitOpenNewSttWriter(SCommitter2 *committer) {
int32_t
code
=
0
;
int32_t
code
=
0
;
int32_t
lino
=
0
;
int32_t
lino
=
0
;
SDiskID
did
[
1
];
if
(
tfsAllocDisk
(
committer
->
tsdb
->
pVnode
->
pTfs
,
committer
->
ctx
->
expLevel
,
did
)
<
0
)
{
code
=
TSDB_CODE_FS_NO_VALID_DISK
;
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
}
SSttFileWriterConfig
config
[
1
]
=
{{
SSttFileWriterConfig
config
[
1
]
=
{{
.
tsdb
=
committer
->
tsdb
,
.
tsdb
=
committer
->
tsdb
,
.
maxRow
=
committer
->
maxRow
,
.
maxRow
=
committer
->
maxRow
,
...
@@ -75,7 +70,7 @@ static int32_t tsdbCommitOpenNewSttWriter(SCommitter2 *committer) {
...
@@ -75,7 +70,7 @@ static int32_t tsdbCommitOpenNewSttWriter(SCommitter2 *committer) {
.
file
=
.
file
=
{
{
.
type
=
TSDB_FTYPE_STT
,
.
type
=
TSDB_FTYPE_STT
,
.
did
=
did
[
0
]
,
.
did
=
committer
->
ctx
->
did
,
.
fid
=
committer
->
ctx
->
fid
,
.
fid
=
committer
->
ctx
->
fid
,
.
cid
=
committer
->
ctx
->
cid
,
.
cid
=
committer
->
ctx
->
cid
,
},
},
...
@@ -122,31 +117,62 @@ static int32_t tsdbCommitOpenWriter(SCommitter2 *committer) {
...
@@ -122,31 +117,62 @@ static int32_t tsdbCommitOpenWriter(SCommitter2 *committer) {
int32_t
code
=
0
;
int32_t
code
=
0
;
int32_t
lino
=
0
;
int32_t
lino
=
0
;
// if (committer->sttTrigger == 1) {
// SDataFileWriterConfig config = {
// // TODO
// };
// code = tsdbDataFileWriterOpen(&config, &committer->dataWriter);
// TSDB_CHECK_CODE(code, lino, _exit);
// // TODO
// }
// stt writer
// stt writer
if
(
!
committer
->
ctx
->
fset
)
{
if
(
committer
->
ctx
->
fset
==
NULL
)
{
return
tsdbCommitOpenNewSttWriter
(
committer
);
code
=
tsdbCommitOpenNewSttWriter
(
committer
);
}
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
}
else
{
const
SSttLvl
*
lvl0
=
tsdbTFileSetGetSttLvl
(
committer
->
ctx
->
fset
,
0
);
const
SSttLvl
*
lvl0
=
tsdbTFileSetGetSttLvl
(
committer
->
ctx
->
fset
,
0
);
if
(
lvl0
==
NULL
||
TARRAY2_SIZE
(
lvl0
->
fobjArr
)
==
0
)
{
if
(
lvl0
==
NULL
||
TARRAY2_SIZE
(
lvl0
->
fobjArr
)
==
0
)
{
return
tsdbCommitOpenNewSttWriter
(
committer
);
code
=
tsdbCommitOpenNewSttWriter
(
committer
);
}
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
}
else
{
STFileObj
*
fobj
=
TARRAY2_LAST
(
lvl0
->
fobjArr
);
STFileObj
*
fobj
=
TARRAY2_LAST
(
lvl0
->
fobjArr
);
if
(
fobj
->
f
->
stt
->
nseg
>=
committer
->
sttTrigger
)
{
if
(
fobj
->
f
->
stt
->
nseg
>=
committer
->
sttTrigger
)
{
return
tsdbCommitOpenNewSttWriter
(
committer
);
code
=
tsdbCommitOpenNewSttWriter
(
committer
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
if
(
committer
->
sttTrigger
==
1
)
{
SSttFileReaderConfig
sttFileReaderConfig
=
{
.
tsdb
=
committer
->
tsdb
,
.
szPage
=
committer
->
szPage
,
.
file
=
fobj
->
f
[
0
],
};
code
=
tsdbSttFileReaderOpen
(
NULL
,
&
sttFileReaderConfig
,
&
committer
->
sttReader
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
}
}
else
{
}
else
{
return
tsdbCommitOpenExistSttWriter
(
committer
,
fobj
->
f
);
code
=
tsdbCommitOpenExistSttWriter
(
committer
,
fobj
->
f
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
}
}
}
// data writer
if
(
committer
->
sttTrigger
==
1
)
{
// data writer
SDataFileWriterConfig
config
=
{
.
tsdb
=
committer
->
tsdb
,
.
cmprAlg
=
committer
->
cmprAlg
,
.
maxRow
=
committer
->
maxRow
,
.
szPage
=
committer
->
szPage
,
.
fid
=
committer
->
ctx
->
fid
,
.
cid
=
committer
->
ctx
->
cid
,
.
did
=
committer
->
ctx
->
did
,
.
compactVersion
=
committer
->
compactVersion
,
};
if
(
committer
->
ctx
->
fset
)
{
for
(
int32_t
ftype
=
TSDB_FTYPE_MIN
;
ftype
<
TSDB_FTYPE_MAX
;
ftype
++
)
{
if
(
committer
->
ctx
->
fset
->
farr
[
ftype
]
!=
NULL
)
{
config
.
files
[
ftype
].
exist
=
true
;
config
.
files
[
ftype
].
file
=
committer
->
ctx
->
fset
->
farr
[
ftype
]
->
f
[
0
];
}
}
}
code
=
tsdbDataFileWriterOpen
(
&
config
,
&
committer
->
dataWriter
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
}
}
_exit:
_exit:
...
@@ -156,13 +182,6 @@ _exit:
...
@@ -156,13 +182,6 @@ _exit:
return
code
;
return
code
;
}
}
static
int32_t
tsdbCommitWriteDelData
(
SCommitter2
*
committer
,
int64_t
suid
,
int64_t
uid
,
int64_t
version
,
int64_t
sKey
,
int64_t
eKey
)
{
int32_t
code
=
0
;
// TODO
return
code
;
}
static
int32_t
tsdbCommitTSDataOpenIterMerger
(
SCommitter2
*
committer
)
{
static
int32_t
tsdbCommitTSDataOpenIterMerger
(
SCommitter2
*
committer
)
{
int32_t
code
=
0
;
int32_t
code
=
0
;
int32_t
lino
=
0
;
int32_t
lino
=
0
;
...
@@ -455,54 +474,6 @@ _exit:
...
@@ -455,54 +474,6 @@ _exit:
return
code
;
return
code
;
}
}
static
int32_t
tsdbCommitTombDataToStt
(
SCommitter2
*
committer
)
{
int32_t
code
=
0
;
int32_t
lino
=
0
;
for
(
STombRecord
*
record
;
(
record
=
tsdbIterMergerGetTombRecord
(
committer
->
iterMerger
));)
{
code
=
tsdbSttFileWriteTombRecord
(
committer
->
sttWriter
,
record
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
code
=
tsdbIterMergerNext
(
committer
->
iterMerger
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
}
_exit:
if
(
code
)
{
TSDB_ERROR_LOG
(
TD_VID
(
committer
->
tsdb
->
pVnode
),
lino
,
code
);
}
return
code
;
}
static
int32_t
tsdbCommitTombDataToData
(
SCommitter2
*
committer
)
{
int32_t
code
=
0
;
int32_t
lino
=
0
;
if
(
committer
->
dataWriter
==
NULL
||
tsdbSttFileWriterIsOpened
(
committer
->
sttWriter
))
{
for
(
STombRecord
*
record
;
(
record
=
tsdbIterMergerGetTombRecord
(
committer
->
iterMerger
));)
{
code
=
tsdbSttFileWriteTombRecord
(
committer
->
sttWriter
,
record
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
code
=
tsdbIterMergerNext
(
committer
->
iterMerger
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
}
}
else
{
for
(
STombRecord
*
record
;
(
record
=
tsdbIterMergerGetTombRecord
(
committer
->
iterMerger
));)
{
code
=
tsdbDataFileWriteTombRecord
(
committer
->
dataWriter
,
record
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
code
=
tsdbIterMergerNext
(
committer
->
iterMerger
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
}
}
_exit:
if
(
code
)
{
TSDB_ERROR_LOG
(
TD_VID
(
committer
->
tsdb
->
pVnode
),
lino
,
code
);
}
return
code
;
}
static
int32_t
tsdbCommitTombDataOpenIter
(
SCommitter2
*
committer
)
{
static
int32_t
tsdbCommitTombDataOpenIter
(
SCommitter2
*
committer
)
{
int32_t
code
=
0
;
int32_t
code
=
0
;
int32_t
lino
=
0
;
int32_t
lino
=
0
;
...
@@ -561,13 +532,23 @@ static int32_t tsdbCommitTombData(SCommitter2 *committer) {
...
@@ -561,13 +532,23 @@ static int32_t tsdbCommitTombData(SCommitter2 *committer) {
code
=
tsdbCommitTombDataOpenIter
(
committer
);
code
=
tsdbCommitTombDataOpenIter
(
committer
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
if
(
committer
->
sttTrigger
>
1
)
{
if
(
committer
->
dataWriter
==
NULL
||
tsdbSttFileWriterIsOpened
(
committer
->
sttWriter
))
{
code
=
tsdbCommitTombDataToStt
(
committer
);
for
(
STombRecord
*
record
;
(
record
=
tsdbIterMergerGetTombRecord
(
committer
->
iterMerger
));)
{
code
=
tsdbSttFileWriteTombRecord
(
committer
->
sttWriter
,
record
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
code
=
tsdbIterMergerNext
(
committer
->
iterMerger
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
}
}
else
{
}
else
{
code
=
tsdbCommitTombDataToData
(
committer
);
for
(
STombRecord
*
record
;
(
record
=
tsdbIterMergerGetTombRecord
(
committer
->
iterMerger
));)
{
code
=
tsdbDataFileWriteTombRecord
(
committer
->
dataWriter
,
record
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
code
=
tsdbIterMergerNext
(
committer
->
iterMerger
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
}
}
}
code
=
tsdbCommitTombDataCloseIter
(
committer
);
code
=
tsdbCommitTombDataCloseIter
(
committer
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
...
@@ -588,6 +569,8 @@ static int32_t tsdbCommitFileSetBegin(SCommitter2 *committer) {
...
@@ -588,6 +569,8 @@ static int32_t tsdbCommitFileSetBegin(SCommitter2 *committer) {
committer
->
ctx
->
expLevel
=
tsdbFidLevel
(
committer
->
ctx
->
fid
,
&
tsdb
->
keepCfg
,
committer
->
ctx
->
now
);
committer
->
ctx
->
expLevel
=
tsdbFidLevel
(
committer
->
ctx
->
fid
,
&
tsdb
->
keepCfg
,
committer
->
ctx
->
now
);
tsdbFidKeyRange
(
committer
->
ctx
->
fid
,
committer
->
minutes
,
committer
->
precision
,
&
committer
->
ctx
->
minKey
,
tsdbFidKeyRange
(
committer
->
ctx
->
fid
,
committer
->
minutes
,
committer
->
precision
,
&
committer
->
ctx
->
minKey
,
&
committer
->
ctx
->
maxKey
);
&
committer
->
ctx
->
maxKey
);
code
=
tfsAllocDisk
(
committer
->
tsdb
->
pVnode
->
pTfs
,
committer
->
ctx
->
expLevel
,
&
committer
->
ctx
->
did
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
STFileSet
fset
=
{.
fid
=
committer
->
ctx
->
fid
};
STFileSet
fset
=
{.
fid
=
committer
->
ctx
->
fid
};
committer
->
ctx
->
fset
=
&
fset
;
committer
->
ctx
->
fset
=
&
fset
;
committer
->
ctx
->
fset
=
TARRAY2_SEARCH_EX
(
committer
->
fsetArr
,
&
committer
->
ctx
->
fset
,
tsdbTFileSetCmprFn
,
TD_EQ
);
committer
->
ctx
->
fset
=
TARRAY2_SEARCH_EX
(
committer
->
fsetArr
,
&
committer
->
ctx
->
fset
,
tsdbTFileSetCmprFn
,
TD_EQ
);
...
...
source/dnode/vnode/src/tsdb/dev/tsdbMerge.c
浏览文件 @
71df0a00
...
@@ -661,6 +661,8 @@ int32_t tsdbMerge(void *arg) {
...
@@ -661,6 +661,8 @@ int32_t tsdbMerge(void *arg) {
.
sttTrigger
=
tsdb
->
pVnode
->
config
.
sttTrigger
,
.
sttTrigger
=
tsdb
->
pVnode
->
config
.
sttTrigger
,
}};
}};
ASSERT
(
merger
->
sttTrigger
>
1
);
code
=
tsdbFSCreateCopySnapshot
(
tsdb
->
pFS
,
&
merger
->
fsetArr
);
code
=
tsdbFSCreateCopySnapshot
(
tsdb
->
pFS
,
&
merger
->
fsetArr
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
TSDB_CHECK_CODE
(
code
,
lino
,
_exit
);
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录