Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
506c192b
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看板
提交
506c192b
编写于
3月 20, 2023
作者:
wmmhello
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix:error in TD-23218 & remove useless logic
上级
f8007803
变更
4
隐藏空白更改
内联
并排
Showing
4 changed file
with
22 addition
and
32 deletion
+22
-32
include/libs/executor/executor.h
include/libs/executor/executor.h
+0
-2
source/dnode/vnode/src/tq/tqExec.c
source/dnode/vnode/src/tq/tqExec.c
+16
-19
source/libs/executor/src/executor.c
source/libs/executor/src/executor.c
+0
-5
source/libs/executor/src/scanoperator.c
source/libs/executor/src/scanoperator.c
+6
-6
未找到文件。
include/libs/executor/executor.h
浏览文件 @
506c192b
...
...
@@ -198,8 +198,6 @@ int32_t qStreamSetScanMemData(qTaskInfo_t tinfo, SPackedData submit);
void
qStreamExtractOffset
(
qTaskInfo_t
tinfo
,
STqOffsetVal
*
pOffset
);
int64_t
qStreamExtractOffsetUid
(
qTaskInfo_t
tinfo
);
SMqMetaRsp
*
qStreamExtractMetaMsg
(
qTaskInfo_t
tinfo
);
const
SSchemaWrapper
*
qExtractSchemaFromTask
(
qTaskInfo_t
tinfo
);
...
...
source/dnode/vnode/src/tq/tqExec.c
浏览文件 @
506c192b
...
...
@@ -156,36 +156,33 @@ int32_t tqScanTaosx(STQ* pTq, const STqHandle* pHandle, STaosxRsp* pRsp, SMqMeta
}
}
if
(
pDataBlock
==
NULL
&&
pOffset
->
type
==
TMQ_OFFSET__SNAPSHOT_DATA
)
{
if
(
qStreamExtractOffsetUid
(
task
)
!=
0
)
{
// get meta
SMqMetaRsp
*
tmp
=
qStreamExtractMetaMsg
(
task
);
if
(
tmp
->
metaRspLen
>
0
)
{
qStreamExtractOffset
(
task
,
&
tmp
->
rspOffset
);
*
pMetaRsp
=
*
tmp
;
tqDebug
(
"tmqsnap task get data"
);
break
;
}
if
(
pDataBlock
==
NULL
)
{
qStreamExtractOffset
(
task
,
pOffset
);
if
(
pOffset
->
type
==
TMQ_OFFSET__SNAPSHOT_DATA
)
{
continue
;
}
tqDebug
(
"tmqsnap vgId: %d, tsdb consume over, switch to wal, ver %"
PRId64
,
TD_VID
(
pTq
->
pVnode
),
pHandle
->
snapshotVer
+
1
);
tqDebug
(
"tmqsnap vgId: %d, tsdb consume over, switch to wal, ver %"
PRId64
,
TD_VID
(
pTq
->
pVnode
),
pHandle
->
snapshotVer
+
1
);
qStreamExtractOffset
(
task
,
&
pRsp
->
rspOffset
);
break
;
}
if
(
pRsp
->
blockNum
>
0
)
{
tqDebug
(
"tmqsnap task exec exited, get data"
);
qStreamExtractOffset
(
task
,
&
pRsp
->
rspOffset
);
break
;
}
SMqMetaRsp
*
tmp
=
qStreamExtractMetaMsg
(
task
);
if
(
tmp
->
rspOffset
.
type
==
TMQ_OFFSET__SNAPSHOT_DATA
)
{
*
pOffset
=
tmp
->
rspOffset
;
qStreamPrepareScan
(
task
,
pOffset
,
pHandle
->
execHandle
.
subType
);
tmp
->
rspOffset
.
type
=
TMQ_OFFSET__SNAPSHOT_META
;
tqDebug
(
"tmqsnap task exec change to get data"
);
continue
;
}
*
pMetaRsp
=
*
tmp
;
tqDebug
(
"tmqsnap task exec exited, get meta"
);
break
;
}
qStreamExtractOffset
(
task
,
&
pRsp
->
rspOffset
);
return
0
;
}
...
...
source/libs/executor/src/executor.c
浏览文件 @
506c192b
...
...
@@ -993,11 +993,6 @@ SMqMetaRsp* qStreamExtractMetaMsg(qTaskInfo_t tinfo) {
return
&
pTaskInfo
->
streamInfo
.
metaRsp
;
}
int64_t
qStreamExtractOffsetUid
(
qTaskInfo_t
tinfo
)
{
SExecTaskInfo
*
pTaskInfo
=
(
SExecTaskInfo
*
)
tinfo
;
return
pTaskInfo
->
streamInfo
.
currentOffset
.
uid
;
}
void
qStreamExtractOffset
(
qTaskInfo_t
tinfo
,
STqOffsetVal
*
pOffset
)
{
SExecTaskInfo
*
pTaskInfo
=
(
SExecTaskInfo
*
)
tinfo
;
memcpy
(
pOffset
,
&
pTaskInfo
->
streamInfo
.
currentOffset
,
sizeof
(
STqOffsetVal
));
...
...
source/libs/executor/src/scanoperator.c
浏览文件 @
506c192b
...
...
@@ -2087,17 +2087,16 @@ static SSDataBlock* doRawScan(SOperatorInfo* pOperator) {
}
SMetaTableInfo
mtInfo
=
getUidfromSnapShot
(
pInfo
->
sContext
);
STqOffsetVal
offset
=
{
0
};
if
(
mtInfo
.
uid
==
0
)
{
// read snapshot done, change to get data from wal
qDebug
(
"tmqsnap read snapshot done, change to get data from wal"
);
tqOffsetResetToLog
(
&
pTaskInfo
->
streamInfo
.
currentO
ffset
,
pInfo
->
sContext
->
snapVersion
);
tqOffsetResetToLog
(
&
o
ffset
,
pInfo
->
sContext
->
snapVersion
);
}
else
{
STqOffsetVal
offset
=
{
0
};
tqOffsetResetToData
(
&
offset
,
mtInfo
.
uid
,
INT64_MIN
);
qStreamPrepareScan
(
pTaskInfo
,
&
offset
,
pInfo
->
sContext
->
subType
);
qDebug
(
"tmqsnap change get data uid:%"
PRId64
""
,
mtInfo
.
uid
);
}
qStreamPrepareScan
(
pTaskInfo
,
&
offset
,
pInfo
->
sContext
->
subType
);
tDeleteSSchemaWrapper
(
mtInfo
.
schema
);
qDebug
(
"tmqsnap stream scan tsdb return null"
);
return
NULL
;
}
else
if
(
pTaskInfo
->
streamInfo
.
currentOffset
.
type
==
TMQ_OFFSET__SNAPSHOT_META
)
{
SSnapContext
*
sContext
=
pInfo
->
sContext
;
...
...
@@ -2112,10 +2111,11 @@ static SSDataBlock* doRawScan(SOperatorInfo* pOperator) {
}
if
(
!
sContext
->
queryMeta
)
{
// change to get data next poll request
tqOffsetResetToData
(
&
pTaskInfo
->
streamInfo
.
metaRsp
.
rspOffset
,
0
,
INT64_MIN
);
STqOffsetVal
offset
=
{
0
};
tqOffsetResetToData
(
&
offset
,
0
,
INT64_MIN
);
qStreamPrepareScan
(
pTaskInfo
,
&
offset
,
pInfo
->
sContext
->
subType
);
}
else
{
tqOffsetResetToMeta
(
&
pTaskInfo
->
streamInfo
.
currentOffset
,
uid
);
pTaskInfo
->
streamInfo
.
metaRsp
.
rspOffset
=
pTaskInfo
->
streamInfo
.
currentOffset
;
pTaskInfo
->
streamInfo
.
metaRsp
.
resMsgType
=
type
;
pTaskInfo
->
streamInfo
.
metaRsp
.
metaRspLen
=
dataLen
;
pTaskInfo
->
streamInfo
.
metaRsp
.
metaRsp
=
data
;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录