Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
92ca817b
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看板
提交
92ca817b
编写于
8月 28, 2022
作者:
H
Hongze Cheng
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
more code
上级
d31459dc
变更
1
隐藏空白更改
内联
并排
Showing
1 changed file
with
126 addition
and
44 deletion
+126
-44
source/dnode/vnode/src/tsdb/tsdbCommit.c
source/dnode/vnode/src/tsdb/tsdbCommit.c
+126
-44
未找到文件。
source/dnode/vnode/src/tsdb/tsdbCommit.c
浏览文件 @
92ca817b
...
...
@@ -559,18 +559,14 @@ _err:
return
code
;
}
static
int32_t
tsdbCommitDataBlock
(
SCommitter
*
pCommitter
,
SBlock
*
pBlock
)
{
static
int32_t
tsdbCommitDataBlock
(
SCommitter
*
pCommitter
)
{
int32_t
code
=
0
;
SBlockData
*
pBlockData
=
&
pCommitter
->
dWriter
.
bData
;
SBlock
block
;
ASSERT
(
pBlockData
->
nRow
>
0
);
if
(
pBlock
)
{
block
=
*
pBlock
;
// as a subblock
}
else
{
tBlockReset
(
&
block
);
// as a new block
}
tBlockReset
(
&
block
);
// info
block
.
nRow
+=
pBlockData
->
nRow
;
...
...
@@ -1547,14 +1543,14 @@ static int32_t tsdbCommitAheadBlock(SCommitter *pCommitter, SBlock *pBlock) {
}
}
if
(
pBlockData
->
nRow
>=
pCommitter
->
maxRow
*
4
/
5
)
{
code
=
tsdbCommitDataBlock
(
pCommitter
,
NULL
);
if
(
pBlockData
->
nRow
>=
pCommitter
->
maxRow
)
{
code
=
tsdbCommitDataBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
}
}
if
(
pBlockData
->
nRow
)
{
code
=
tsdbCommitDataBlock
(
pCommitter
,
NULL
);
code
=
tsdbCommitDataBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
}
...
...
@@ -1616,8 +1612,8 @@ static int32_t tsdbCommitMergeBlock(SCommitter *pCommitter, SBlock *pBlock) {
ASSERT
(
0
);
}
if
(
pBDataW
->
nRow
>=
pCommitter
->
maxRow
*
4
/
5
)
{
code
=
tsdbCommitDataBlock
(
pCommitter
,
NULL
);
if
(
pBDataW
->
nRow
>=
pCommitter
->
maxRow
)
{
code
=
tsdbCommitDataBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
}
}
...
...
@@ -1633,14 +1629,14 @@ static int32_t tsdbCommitMergeBlock(SCommitter *pCommitter, SBlock *pBlock) {
pRow
=
NULL
;
}
if
(
pBDataW
->
nRow
>=
pCommitter
->
maxRow
*
4
/
5
)
{
code
=
tsdbCommitDataBlock
(
pCommitter
,
NULL
);
if
(
pBDataW
->
nRow
>=
pCommitter
->
maxRow
)
{
code
=
tsdbCommitDataBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
}
}
if
(
pBDataW
->
nRow
)
{
code
=
tsdbCommitDataBlock
(
pCommitter
,
NULL
);
code
=
tsdbCommitDataBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
}
...
...
@@ -1724,59 +1720,146 @@ _err:
return
code
;
}
static
int32_t
tsdbAppendLastBlock
(
SCommitter
*
pCommitter
)
{
int32_t
code
=
0
;
SBlockData
*
pBData
=
&
pCommitter
->
dWriter
.
bData
;
SBlockData
*
pBDatal
=
&
pCommitter
->
dWriter
.
bDatal
;
if
(
pBDatal
->
suid
||
pBDatal
->
uid
)
{
if
(
pBDatal
->
suid
!=
pBData
->
suid
||
pBDatal
->
suid
==
0
)
{
if
(
pBDatal
->
nRow
)
{
code
=
tsdbCommitLastBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
}
tBlockDataReset
(
pBDatal
);
}
}
if
(
!
pBDatal
->
suid
&&
!
pBDatal
->
uid
)
{
ASSERT
(
pCommitter
->
skmTable
.
suid
==
pBData
->
suid
);
ASSERT
(
pCommitter
->
skmTable
.
uid
==
pBData
->
uid
);
code
=
tBlockDataInit
(
pBDatal
,
pBData
->
suid
,
0
,
pCommitter
->
skmTable
.
pTSchema
);
if
(
code
)
goto
_err
;
}
for
(
int32_t
iRow
=
0
;
iRow
<
pBData
->
nRow
;
iRow
++
)
{
TSDBROW
row
=
tsdbRowFromBlockData
(
pBData
,
iRow
);
code
=
tBlockDataAppendRow
(
pBDatal
,
&
row
,
NULL
,
pBData
->
uid
);
if
(
code
)
goto
_err
;
if
(
pBDatal
->
nRow
>=
pCommitter
->
maxRow
)
{
code
=
tsdbCommitLastBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
}
}
return
code
;
_err:
return
code
;
}
static
int32_t
tsdbCommitTableData
(
SCommitter
*
pCommitter
,
TABLEID
id
)
{
int32_t
code
=
0
;
SRowInfo
*
pRowInfo
=
tsdbGetCommitRow
(
pCommitter
);
int32_t
code
=
0
;
SRowInfo
*
pRowInfo
=
tsdbGetCommitRow
(
pCommitter
);
if
(
pRowInfo
&&
(
pRowInfo
->
suid
!=
id
.
suid
||
pRowInfo
->
uid
!=
id
.
uid
))
{
pRowInfo
=
NULL
;
}
if
(
pRowInfo
==
NULL
)
goto
_e
rr
;
if
(
pRowInfo
==
NULL
)
goto
_e
xit
;
#if 0
if (pBlockData->suid || pBlockData->uid) {
if (pBlockData->suid != pTbData->suid || pBlockData->suid == 0) {
if (pBlockData->nRow > 0) {
code = tsdbCommitLastBlock(pCommitter);
if (code) goto _err;
if
(
pCommitter
->
toLastOnly
)
{
SBlockData
*
pBDatal
=
&
pCommitter
->
dWriter
.
bDatal
;
if
(
pBDatal
->
suid
||
pBDatal
->
uid
)
{
if
(
pBDatal
->
suid
!=
id
.
suid
||
pBDatal
->
suid
==
0
)
{
if
(
pBDatal
->
nRow
)
{
code
=
tsdbCommitLastBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
}
tBlockDataReset
(
pBDatal
);
}
}
tBlockDataReset(pBlockData);
if
(
!
pBDatal
->
suid
&&
!
pBDatal
->
uid
)
{
ASSERT
(
pCommitter
->
skmTable
.
suid
==
id
.
suid
);
ASSERT
(
pCommitter
->
skmTable
.
uid
==
id
.
uid
);
code
=
tBlockDataInit
(
pBDatal
,
id
.
suid
,
0
,
pCommitter
->
skmTable
.
pTSchema
);
if
(
code
)
goto
_err
;
}
}
if (!pBlockData->suid && !pBlockData->uid) {
code = tBlockDataInit(pBlockData, pTbData->suid, 0, pCommitter->skmTable.pTSchema);
if (code) goto _err;
}
#endif
while
(
pRowInfo
)
{
STSchema
*
pTSchema
=
NULL
;
if
(
pRowInfo
->
row
.
type
==
0
)
{
code
=
tsdbCommitterUpdateRowSchema
(
pCommitter
,
id
.
suid
,
id
.
uid
,
TSDBROW_SVERSION
(
&
pRowInfo
->
row
));
if
(
code
)
goto
_err
;
pTSchema
=
pCommitter
->
skmRow
.
pTSchema
;
}
SBlockData
*
pBlockData
=
NULL
;
// TODO
while
(
pRowInfo
)
{
STSchema
*
pTSchema
=
NULL
;
if
(
pRowInfo
->
row
.
type
==
0
)
{
code
=
tsdbCommitterUpdateRowSchema
(
pCommitter
,
id
.
suid
,
id
.
uid
,
TSDBROW_SVERSION
(
&
pRowInfo
->
row
));
code
=
tBlockDataAppendRow
(
pBDatal
,
&
pRowInfo
->
row
,
pTSchema
,
id
.
uid
);
if
(
code
)
goto
_err
;
code
=
tsdbNextCommitRow
(
pCommitter
);
if
(
code
)
goto
_err
;
pTSchema
=
pCommitter
->
skmRow
.
pTSchema
;
pRowInfo
=
tsdbGetCommitRow
(
pCommitter
);
if
(
pRowInfo
&&
(
pRowInfo
->
suid
!=
id
.
suid
||
pRowInfo
->
uid
!=
id
.
uid
))
{
pRowInfo
=
NULL
;
}
if
(
pBDatal
->
nRow
>=
pCommitter
->
maxRow
)
{
code
=
tsdbCommitLastBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
}
}
}
else
{
SBlockData
*
pBData
=
&
pCommitter
->
dWriter
.
bData
;
code
=
tBlockDataAppendRow
(
pBlockData
,
&
pRowInfo
->
row
,
pTSchema
,
id
.
uid
);
if
(
code
)
goto
_err
;
ASSERT
(
pBData
->
nRow
==
0
);
code
=
tsdbNextCommitRow
(
pCommitter
);
if
(
code
)
goto
_err
;
while
(
pRowInfo
)
{
STSchema
*
pTSchema
=
NULL
;
if
(
pRowInfo
->
row
.
type
==
0
)
{
code
=
tsdbCommitterUpdateRowSchema
(
pCommitter
,
id
.
suid
,
id
.
uid
,
TSDBROW_SVERSION
(
&
pRowInfo
->
row
));
if
(
code
)
goto
_err
;
pTSchema
=
pCommitter
->
skmRow
.
pTSchema
;
}
pRowInfo
=
tsdbGetCommitRow
(
pCommitter
);
if
(
pRowInfo
&&
(
pRowInfo
->
suid
!=
id
.
suid
||
pRowInfo
->
uid
!=
id
.
uid
))
{
pRowInfo
=
NULL
;
code
=
tBlockDataAppendRow
(
pBData
,
&
pRowInfo
->
row
,
pTSchema
,
id
.
uid
);
if
(
code
)
goto
_err
;
code
=
tsdbNextCommitRow
(
pCommitter
);
if
(
code
)
goto
_err
;
pRowInfo
=
tsdbGetCommitRow
(
pCommitter
);
if
(
pRowInfo
&&
(
pRowInfo
->
suid
!=
id
.
suid
||
pRowInfo
->
uid
!=
id
.
uid
))
{
pRowInfo
=
NULL
;
}
if
(
pBData
->
nRow
>=
pCommitter
->
maxRow
)
{
code
=
tsdbCommitDataBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
}
}
if
(
pBData
->
nRow
)
{
if
(
pBData
->
nRow
>
pCommitter
->
minRow
)
{
code
=
tsdbCommitDataBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
}
else
{
code
=
tsdbAppendLastBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
}
}
}
_exit:
return
code
;
_err:
tsdbError
(
"vgId:%d tsdb commit table data failed since %s"
,
TD_VID
(
pCommitter
->
pTsdb
->
pVnode
),
tstrerror
(
code
));
return
code
;
}
...
...
@@ -1830,7 +1913,6 @@ static int32_t tsdbCommitFileDataImpl(SCommitter *pCommitter) {
code
=
tsdbMoveCommitData
(
pCommitter
,
id
);
if
(
code
)
goto
_err
;
// TODO: here may have problem
if
(
pCommitter
->
dWriter
.
bDatal
.
nRow
>
0
)
{
code
=
tsdbCommitLastBlock
(
pCommitter
);
if
(
code
)
goto
_err
;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录