Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
e3e42cdf
T
TDengine
项目概览
taosdata
/
TDengine
1 年多 前同步成功
通知
1187
Star
22018
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看板
提交
e3e42cdf
编写于
12月 13, 2022
作者:
H
Haojun Liao
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix(query): set update ts flag for stream.
上级
d7edcfd2
变更
5
隐藏空白更改
内联
并排
Showing
5 changed file
with
13 addition
and
0 deletion
+13
-0
source/common/src/tdatablock.c
source/common/src/tdatablock.c
+4
-0
source/dnode/vnode/src/tq/tqRead.c
source/dnode/vnode/src/tq/tqRead.c
+1
-0
source/libs/executor/src/exchangeoperator.c
source/libs/executor/src/exchangeoperator.c
+1
-0
source/libs/executor/src/executorimpl.c
source/libs/executor/src/executorimpl.c
+2
-0
source/libs/executor/src/scanoperator.c
source/libs/executor/src/scanoperator.c
+5
-0
未找到文件。
source/common/src/tdatablock.c
浏览文件 @
e3e42cdf
...
...
@@ -358,6 +358,10 @@ size_t blockDataGetNumOfCols(const SSDataBlock* pBlock) { return taosArrayGetSiz
size_t
blockDataGetNumOfRows
(
const
SSDataBlock
*
pBlock
)
{
return
pBlock
->
info
.
rows
;
}
int32_t
blockDataUpdateTsWindow
(
SSDataBlock
*
pDataBlock
,
int32_t
tsColumnIndex
)
{
if
(
pDataBlock
->
info
.
rows
>
0
)
{
ASSERT
(
pDataBlock
->
info
.
dataLoad
==
1
);
}
if
(
pDataBlock
==
NULL
||
pDataBlock
->
info
.
rows
<=
0
||
pDataBlock
->
info
.
dataLoad
==
0
)
{
return
0
;
}
...
...
source/dnode/vnode/src/tq/tqRead.c
浏览文件 @
e3e42cdf
...
...
@@ -533,6 +533,7 @@ int32_t tqRetrieveDataBlock(SSDataBlock* pBlock, STqReader* pReader) {
pBlock
->
info
.
id
.
uid
=
pReader
->
msgIter
.
uid
;
pBlock
->
info
.
rows
=
pReader
->
msgIter
.
numOfRows
;
pBlock
->
info
.
version
=
pReader
->
pMsg
->
version
;
pBlock
->
info
.
dataLoad
=
1
;
while
((
row
=
tGetSubmitBlkNext
(
&
pReader
->
blkIter
))
!=
NULL
)
{
tdSTSRowIterReset
(
&
iter
,
row
);
...
...
source/libs/executor/src/exchangeoperator.c
浏览文件 @
e3e42cdf
...
...
@@ -510,6 +510,7 @@ int32_t extractDataBlockFromFetchRsp(SSDataBlock* pRes, char* pData, SArray* pCo
blockDataEnsureCapacity
(
pRes
,
pBlock
->
info
.
rows
);
// data from mnode
pRes
->
info
.
dataLoad
=
1
;
pRes
->
info
.
rows
=
pBlock
->
info
.
rows
;
relocateColumnData
(
pRes
,
pColList
,
pBlock
->
pDataBlock
,
false
);
blockDataDestroy
(
pBlock
);
...
...
source/libs/executor/src/executorimpl.c
浏览文件 @
e3e42cdf
...
...
@@ -2546,6 +2546,7 @@ int32_t buildDataBlockFromGroupRes(SOperatorInfo* pOperator, SStreamState* pStat
pBlock
->
info
.
rows
+=
pRow
->
numOfRows
;
releaseOutputBuf
(
pState
,
&
key
,
pRow
);
}
pBlock
->
info
.
dataLoad
=
1
;
blockDataUpdateTsWindow
(
pBlock
,
0
);
return
TSDB_CODE_SUCCESS
;
}
...
...
@@ -2635,6 +2636,7 @@ int32_t buildSessionResultDataBlock(SOperatorInfo* pOperator, SStreamState* pSta
}
}
pBlock
->
info
.
dataLoad
=
1
;
pBlock
->
info
.
rows
+=
pRow
->
numOfRows
;
// saveSessionDiscBuf(pState, pKey, pVal, size);
releaseOutputBuf
(
pState
,
NULL
,
pRow
);
...
...
source/libs/executor/src/scanoperator.c
浏览文件 @
e3e42cdf
...
...
@@ -1313,6 +1313,7 @@ static int32_t generateScanRange(SStreamScanInfo* pInfo, SSDataBlock* pSrcBlock,
}
pDestBlock
->
info
.
type
=
STREAM_CLEAR
;
pDestBlock
->
info
.
version
=
pSrcBlock
->
info
.
version
;
pDestBlock
->
info
.
dataLoad
=
1
;
blockDataUpdateTsWindow
(
pDestBlock
,
0
);
return
code
;
}
...
...
@@ -1421,6 +1422,7 @@ static void checkUpdateData(SStreamScanInfo* pInfo, bool invertible, SSDataBlock
}
if
(
out
&&
pInfo
->
pUpdateDataRes
->
info
.
rows
>
0
)
{
pInfo
->
pUpdateDataRes
->
info
.
version
=
pBlock
->
info
.
version
;
pInfo
->
pUpdateDataRes
->
info
.
dataLoad
=
1
;
blockDataUpdateTsWindow
(
pInfo
->
pUpdateDataRes
,
0
);
pInfo
->
pUpdateDataRes
->
info
.
type
=
pInfo
->
partitionSup
.
needCalc
?
STREAM_DELETE_DATA
:
STREAM_CLEAR
;
}
...
...
@@ -1483,6 +1485,7 @@ static int32_t setBlockIntoRes(SStreamScanInfo* pInfo, const SSDataBlock* pBlock
doFilter
(
pInfo
->
pRes
,
pOperator
->
exprSupp
.
pFilterInfo
,
NULL
);
}
pInfo
->
pRes
->
info
.
dataLoad
=
1
;
blockDataUpdateTsWindow
(
pInfo
->
pRes
,
pInfo
->
primaryTsIndex
);
blockDataFreeRes
((
SSDataBlock
*
)
pBlock
);
...
...
@@ -1771,6 +1774,7 @@ FETCH_NEXT_BLOCK:
// TODO move into scan
pBlock
->
info
.
calWin
.
skey
=
INT64_MIN
;
pBlock
->
info
.
calWin
.
ekey
=
INT64_MAX
;
pBlock
->
info
.
dataLoad
=
1
;
blockDataUpdateTsWindow
(
pBlock
,
0
);
switch
(
pBlock
->
info
.
type
)
{
case
STREAM_NORMAL
:
...
...
@@ -1948,6 +1952,7 @@ FETCH_NEXT_BLOCK:
}
doFilter
(
pInfo
->
pRes
,
pOperator
->
exprSupp
.
pFilterInfo
,
NULL
);
pInfo
->
pRes
->
info
.
dataLoad
=
1
;
blockDataUpdateTsWindow
(
pInfo
->
pRes
,
pInfo
->
primaryTsIndex
);
if
(
pBlockInfo
->
rows
>
0
||
pInfo
->
pUpdateDataRes
->
info
.
rows
>
0
)
{
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录