Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
ca1c961f
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看板
提交
ca1c961f
编写于
6月 30, 2022
作者:
L
Liu Jicong
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix(tmq): set last scan status
上级
8dece648
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
35 addition
and
16 deletion
+35
-16
source/dnode/vnode/src/tq/tqExec.c
source/dnode/vnode/src/tq/tqExec.c
+3
-0
source/libs/executor/src/scanoperator.c
source/libs/executor/src/scanoperator.c
+32
-16
未找到文件。
source/dnode/vnode/src/tq/tqExec.c
浏览文件 @
ca1c961f
...
...
@@ -85,12 +85,15 @@ int32_t tqScanSnapshot(STQ* pTq, const STqExecHandle* pExec, SMqDataRsp* pRsp, S
tqAddBlockDataToRsp
(
pDataBlock
,
pRsp
);
if
(
pRsp
->
withTbName
)
{
pRsp
->
withTbName
=
0
;
#if 1
int64_t
uid
;
int64_t
ts
;
if
(
qGetStreamScanStatus
(
task
,
&
uid
,
&
ts
)
<
0
)
{
ASSERT
(
0
);
}
tqAddTbNameToRsp
(
pTq
,
uid
,
pRsp
,
workerId
);
#endif
}
pRsp
->
blockNum
++
;
...
...
source/libs/executor/src/scanoperator.c
浏览文件 @
ca1c961f
...
...
@@ -392,6 +392,7 @@ static SSDataBlock* doTableScanImpl(SOperatorInfo* pOperator) {
binfo
.
capacity
=
binfo
.
rows
;
blockDataEnsureCapacity
(
pBlock
,
binfo
.
rows
);
pBlock
->
info
=
binfo
;
ASSERT
(
binfo
.
uid
!=
0
);
uint32_t
status
=
0
;
int32_t
code
=
loadDataBlock
(
pOperator
,
pTableScanInfo
,
pBlock
,
&
status
);
...
...
@@ -419,6 +420,7 @@ static SSDataBlock* doTableScanImpl(SOperatorInfo* pOperator) {
pTableScanInfo
->
lastStatus
.
uid
=
pBlock
->
info
.
uid
;
pTableScanInfo
->
lastStatus
.
ts
=
pBlock
->
info
.
window
.
ekey
;
ASSERT
(
pBlock
->
info
.
uid
!=
0
);
return
pBlock
;
}
return
NULL
;
...
...
@@ -438,6 +440,7 @@ static SSDataBlock* doTableScanGroup(SOperatorInfo* pOperator) {
while
(
pTableScanInfo
->
curTWinIdx
<
pTableScanInfo
->
cond
.
numOfTWindows
)
{
SSDataBlock
*
p
=
doTableScanImpl
(
pOperator
);
if
(
p
!=
NULL
)
{
ASSERT
(
p
->
info
.
uid
!=
0
);
return
p
;
}
pTableScanInfo
->
curTWinIdx
+=
1
;
...
...
@@ -517,6 +520,35 @@ static SSDataBlock* doTableScan(SOperatorInfo* pOperator) {
if
(
pInfo
->
scanMode
==
TABLE_SCAN__TABLE_ORDER
)
{
// check status
if
(
pInfo
->
lastStatus
.
uid
==
pInfo
->
expStatus
.
uid
&&
pInfo
->
lastStatus
.
ts
==
pInfo
->
expStatus
.
ts
)
{
while
(
1
)
{
SSDataBlock
*
result
=
doTableScanGroup
(
pOperator
);
if
(
result
)
{
return
result
;
}
// if no data, switch to next table and continue scan
pInfo
->
currentTable
++
;
if
(
pInfo
->
currentTable
>=
taosArrayGetSize
(
pTaskInfo
->
tableqinfoList
.
pTableList
))
{
return
NULL
;
}
STableKeyInfo
*
pTableInfo
=
taosArrayGet
(
pTaskInfo
->
tableqinfoList
.
pTableList
,
pInfo
->
currentTable
);
/*pTableInfo->uid */
tsdbSetTableId
(
pInfo
->
dataReader
,
pTableInfo
->
uid
);
tsdbResetReadHandle
(
pInfo
->
dataReader
,
&
pInfo
->
cond
,
0
);
pInfo
->
scanTimes
=
0
;
pInfo
->
curTWinIdx
=
0
;
}
}
// reset to exp table and window start from ts
tsdbSetTableId
(
pInfo
->
dataReader
,
pInfo
->
expStatus
.
uid
);
SQueryTableDataCond
tmpCond
=
pInfo
->
cond
;
tmpCond
.
twindows
[
0
]
=
(
STimeWindow
){
.
skey
=
pInfo
->
expStatus
.
ts
,
.
ekey
=
INT64_MAX
,
};
tsdbResetReadHandle
(
pInfo
->
dataReader
,
&
tmpCond
,
0
);
pInfo
->
scanTimes
=
0
;
pInfo
->
curTWinIdx
=
0
;
while
(
1
)
{
SSDataBlock
*
result
=
doTableScanGroup
(
pOperator
);
if
(
result
)
{
return
result
;
...
...
@@ -532,23 +564,7 @@ static SSDataBlock* doTableScan(SOperatorInfo* pOperator) {
tsdbResetReadHandle
(
pInfo
->
dataReader
,
&
pInfo
->
cond
,
0
);
pInfo
->
scanTimes
=
0
;
pInfo
->
curTWinIdx
=
0
;
pInfo
->
lastStatus
.
ts
=
pInfo
->
expStatus
.
ts
;
pInfo
->
lastStatus
.
uid
=
pInfo
->
expStatus
.
uid
;
return
doTableScan
(
pOperator
);
}
// reset to exp table and window start from ts
tsdbSetTableId
(
pInfo
->
dataReader
,
pInfo
->
expStatus
.
uid
);
SQueryTableDataCond
tmpCond
=
pInfo
->
cond
;
tmpCond
.
twindows
[
0
]
=
(
STimeWindow
){
.
skey
=
pInfo
->
expStatus
.
ts
,
.
ekey
=
INT64_MAX
,
};
tsdbResetReadHandle
(
pInfo
->
dataReader
,
&
tmpCond
,
0
);
pInfo
->
scanTimes
=
0
;
pInfo
->
curTWinIdx
=
0
;
pInfo
->
lastStatus
.
ts
=
pInfo
->
expStatus
.
ts
;
pInfo
->
lastStatus
.
uid
=
pInfo
->
expStatus
.
uid
;
return
doTableScan
(
pOperator
);
}
if
(
pInfo
->
currentGroupId
==
-
1
)
{
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录