Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
e587cc50
T
TDengine
项目概览
taosdata
/
TDengine
大约 2 年 前同步成功
通知
1192
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看板
提交
e587cc50
编写于
8月 10, 2023
作者:
W
wangjiaming0909
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
feat: optimize select agg_func partition by tag slimit
上级
68e0fb92
变更
9
隐藏空白更改
内联
并排
Showing
9 changed file
with
125 addition
and
31 deletion
+125
-31
include/libs/nodes/plannodes.h
include/libs/nodes/plannodes.h
+2
-0
source/libs/executor/src/executor.c
source/libs/executor/src/executor.c
+8
-1
source/libs/nodes/src/nodesCodeFuncs.c
source/libs/nodes/src/nodesCodeFuncs.c
+14
-0
source/libs/nodes/src/nodesMsgFuncs.c
source/libs/nodes/src/nodesMsgFuncs.c
+12
-0
source/libs/planner/inc/planInt.h
source/libs/planner/inc/planInt.h
+8
-2
source/libs/planner/src/planOptimizer.c
source/libs/planner/src/planOptimizer.c
+13
-21
source/libs/planner/src/planPhysiCreater.c
source/libs/planner/src/planPhysiCreater.c
+8
-2
source/libs/planner/src/planSpliter.c
source/libs/planner/src/planSpliter.c
+9
-1
source/libs/planner/src/planUtil.c
source/libs/planner/src/planUtil.c
+51
-4
未找到文件。
include/libs/nodes/plannodes.h
浏览文件 @
e587cc50
...
@@ -603,6 +603,8 @@ typedef struct SSubplan {
...
@@ -603,6 +603,8 @@ typedef struct SSubplan {
SNode
*
pTagCond
;
SNode
*
pTagCond
;
SNode
*
pTagIndexCond
;
SNode
*
pTagIndexCond
;
bool
showRewrite
;
bool
showRewrite
;
int32_t
rowsThreshold
;
bool
dynamicRowThreshold
;
}
SSubplan
;
}
SSubplan
;
typedef
enum
EExplainMode
{
EXPLAIN_MODE_DISABLE
=
1
,
EXPLAIN_MODE_STATIC
,
EXPLAIN_MODE_ANALYZE
}
EExplainMode
;
typedef
enum
EExplainMode
{
EXPLAIN_MODE_DISABLE
=
1
,
EXPLAIN_MODE_STATIC
,
EXPLAIN_MODE_ANALYZE
}
EExplainMode
;
...
...
source/libs/executor/src/executor.c
浏览文件 @
e587cc50
...
@@ -589,6 +589,10 @@ int32_t qExecTaskOpt(qTaskInfo_t tinfo, SArray* pResList, uint64_t* useconds, bo
...
@@ -589,6 +589,10 @@ int32_t qExecTaskOpt(qTaskInfo_t tinfo, SArray* pResList, uint64_t* useconds, bo
int64_t
st
=
taosGetTimestampUs
();
int64_t
st
=
taosGetTimestampUs
();
int32_t
blockIndex
=
0
;
int32_t
blockIndex
=
0
;
int32_t
rowsThreshold
=
pTaskInfo
->
pSubplan
->
rowsThreshold
;
if
(
!
pTaskInfo
->
pSubplan
->
dynamicRowThreshold
||
4096
<=
pTaskInfo
->
pSubplan
->
rowsThreshold
)
{
rowsThreshold
=
4096
;
}
while
((
pRes
=
pTaskInfo
->
pRoot
->
fpSet
.
getNextFn
(
pTaskInfo
->
pRoot
))
!=
NULL
)
{
while
((
pRes
=
pTaskInfo
->
pRoot
->
fpSet
.
getNextFn
(
pTaskInfo
->
pRoot
))
!=
NULL
)
{
SSDataBlock
*
p
=
NULL
;
SSDataBlock
*
p
=
NULL
;
if
(
blockIndex
>=
taosArrayGetSize
(
pTaskInfo
->
pResultBlockList
))
{
if
(
blockIndex
>=
taosArrayGetSize
(
pTaskInfo
->
pResultBlockList
))
{
...
@@ -606,10 +610,13 @@ int32_t qExecTaskOpt(qTaskInfo_t tinfo, SArray* pResList, uint64_t* useconds, bo
...
@@ -606,10 +610,13 @@ int32_t qExecTaskOpt(qTaskInfo_t tinfo, SArray* pResList, uint64_t* useconds, bo
ASSERT
(
p
->
info
.
rows
>
0
);
ASSERT
(
p
->
info
.
rows
>
0
);
taosArrayPush
(
pResList
,
&
p
);
taosArrayPush
(
pResList
,
&
p
);
if
(
current
>=
4096
)
{
if
(
current
>=
rowsThreshold
)
{
break
;
break
;
}
}
}
}
if
(
pTaskInfo
->
pSubplan
->
dynamicRowThreshold
)
{
pTaskInfo
->
pSubplan
->
rowsThreshold
-=
current
;
}
*
hasMore
=
(
pRes
!=
NULL
);
*
hasMore
=
(
pRes
!=
NULL
);
uint64_t
el
=
(
taosGetTimestampUs
()
-
st
);
uint64_t
el
=
(
taosGetTimestampUs
()
-
st
);
...
...
source/libs/nodes/src/nodesCodeFuncs.c
浏览文件 @
e587cc50
...
@@ -2814,6 +2814,8 @@ static const char* jkSubplanDataSink = "DataSink";
...
@@ -2814,6 +2814,8 @@ static const char* jkSubplanDataSink = "DataSink";
static
const
char
*
jkSubplanTagCond
=
"TagCond"
;
static
const
char
*
jkSubplanTagCond
=
"TagCond"
;
static
const
char
*
jkSubplanTagIndexCond
=
"TagIndexCond"
;
static
const
char
*
jkSubplanTagIndexCond
=
"TagIndexCond"
;
static
const
char
*
jkSubplanShowRewrite
=
"ShowRewrite"
;
static
const
char
*
jkSubplanShowRewrite
=
"ShowRewrite"
;
static
const
char
*
jkSubplanRowsThreshold
=
"RowThreshold"
;
static
const
char
*
jkSubplanDynamicRowsThreshold
=
"DyRowThreshold"
;
static
int32_t
subplanToJson
(
const
void
*
pObj
,
SJson
*
pJson
)
{
static
int32_t
subplanToJson
(
const
void
*
pObj
,
SJson
*
pJson
)
{
const
SSubplan
*
pNode
=
(
const
SSubplan
*
)
pObj
;
const
SSubplan
*
pNode
=
(
const
SSubplan
*
)
pObj
;
...
@@ -2852,6 +2854,12 @@ static int32_t subplanToJson(const void* pObj, SJson* pJson) {
...
@@ -2852,6 +2854,12 @@ static int32_t subplanToJson(const void* pObj, SJson* pJson) {
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAddBoolToObject
(
pJson
,
jkSubplanShowRewrite
,
pNode
->
showRewrite
);
code
=
tjsonAddBoolToObject
(
pJson
,
jkSubplanShowRewrite
,
pNode
->
showRewrite
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAddIntegerToObject
(
pJson
,
jkSubplanRowsThreshold
,
pNode
->
rowsThreshold
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAddBoolToObject
(
pJson
,
jkSubplanDynamicRowsThreshold
,
pNode
->
dynamicRowThreshold
);
}
return
code
;
return
code
;
}
}
...
@@ -2893,6 +2901,12 @@ static int32_t jsonToSubplan(const SJson* pJson, void* pObj) {
...
@@ -2893,6 +2901,12 @@ static int32_t jsonToSubplan(const SJson* pJson, void* pObj) {
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonGetBoolValue
(
pJson
,
jkSubplanShowRewrite
,
&
pNode
->
showRewrite
);
code
=
tjsonGetBoolValue
(
pJson
,
jkSubplanShowRewrite
,
&
pNode
->
showRewrite
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonGetIntValue
(
pJson
,
jkSubplanRowsThreshold
,
&
pNode
->
rowsThreshold
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonGetBoolValue
(
pJson
,
jkSubplanDynamicRowsThreshold
,
&
pNode
->
dynamicRowThreshold
);
}
return
code
;
return
code
;
}
}
...
...
source/libs/nodes/src/nodesMsgFuncs.c
浏览文件 @
e587cc50
...
@@ -3538,6 +3538,12 @@ static int32_t subplanInlineToMsg(const void* pObj, STlvEncoder* pEncoder) {
...
@@ -3538,6 +3538,12 @@ static int32_t subplanInlineToMsg(const void* pObj, STlvEncoder* pEncoder) {
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tlvEncodeValueBool
(
pEncoder
,
pNode
->
showRewrite
);
code
=
tlvEncodeValueBool
(
pEncoder
,
pNode
->
showRewrite
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tlvEncodeValueI32
(
pEncoder
,
pNode
->
rowsThreshold
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tlvEncodeValueBool
(
pEncoder
,
pNode
->
dynamicRowThreshold
);
}
return
code
;
return
code
;
}
}
...
@@ -3587,6 +3593,12 @@ static int32_t msgToSubplanInline(STlvDecoder* pDecoder, void* pObj) {
...
@@ -3587,6 +3593,12 @@ static int32_t msgToSubplanInline(STlvDecoder* pDecoder, void* pObj) {
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tlvDecodeValueBool
(
pDecoder
,
&
pNode
->
showRewrite
);
code
=
tlvDecodeValueBool
(
pDecoder
,
&
pNode
->
showRewrite
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tlvDecodeValueI32
(
pDecoder
,
&
pNode
->
rowsThreshold
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tlvDecodeValueBool
(
pDecoder
,
&
pNode
->
dynamicRowThreshold
);
}
return
code
;
return
code
;
}
}
...
...
source/libs/planner/inc/planInt.h
浏览文件 @
e587cc50
...
@@ -43,8 +43,14 @@ int32_t splitLogicPlan(SPlanContext* pCxt, SLogicSubplan* pLogicSubplan);
...
@@ -43,8 +43,14 @@ int32_t splitLogicPlan(SPlanContext* pCxt, SLogicSubplan* pLogicSubplan);
int32_t
scaleOutLogicPlan
(
SPlanContext
*
pCxt
,
SLogicSubplan
*
pLogicSubplan
,
SQueryLogicPlan
**
pLogicPlan
);
int32_t
scaleOutLogicPlan
(
SPlanContext
*
pCxt
,
SLogicSubplan
*
pLogicSubplan
,
SQueryLogicPlan
**
pLogicPlan
);
int32_t
createPhysiPlan
(
SPlanContext
*
pCxt
,
SQueryLogicPlan
*
pLogicPlan
,
SQueryPlan
**
pPlan
,
SArray
*
pExecNodeList
);
int32_t
createPhysiPlan
(
SPlanContext
*
pCxt
,
SQueryLogicPlan
*
pLogicPlan
,
SQueryPlan
**
pPlan
,
SArray
*
pExecNodeList
);
bool
isPartTableAgg
(
SAggLogicNode
*
pAgg
);
bool
isPartTableAgg
(
SAggLogicNode
*
pAgg
);
bool
isPartTableWinodw
(
SWindowLogicNode
*
pWindow
);
bool
isPartTagAgg
(
SAggLogicNode
*
pAgg
);
bool
isPartTableWinodw
(
SWindowLogicNode
*
pWindow
);
#define CLONE_LIMIT 1
#define CLONE_SLIMIT 1 << 1
#define CLONE_LIMIT_SLIMIT (CLONE_LIMIT | CLONE_SLIMIT)
bool
cloneLimit
(
SLogicNode
*
pParent
,
SLogicNode
*
pChild
,
uint8_t
cloneWhat
);
#ifdef __cplusplus
#ifdef __cplusplus
}
}
...
...
source/libs/planner/src/planOptimizer.c
浏览文件 @
e587cc50
...
@@ -368,8 +368,8 @@ static void scanPathOptSetGroupOrderScan(SScanLogicNode* pScan) {
...
@@ -368,8 +368,8 @@ static void scanPathOptSetGroupOrderScan(SScanLogicNode* pScan) {
if
(
pScan
->
node
.
pParent
&&
nodeType
(
pScan
->
node
.
pParent
)
==
QUERY_NODE_LOGIC_PLAN_AGG
)
{
if
(
pScan
->
node
.
pParent
&&
nodeType
(
pScan
->
node
.
pParent
)
==
QUERY_NODE_LOGIC_PLAN_AGG
)
{
SAggLogicNode
*
pAgg
=
(
SAggLogicNode
*
)
pScan
->
node
.
pParent
;
SAggLogicNode
*
pAgg
=
(
SAggLogicNode
*
)
pScan
->
node
.
pParent
;
bool
withSlimit
=
pAgg
->
node
.
pSlimit
!=
NULL
||
(
pAgg
->
node
.
pParent
&&
pAgg
->
node
.
pParent
->
pSlimit
)
;
bool
withSlimit
=
pAgg
->
node
.
pSlimit
!=
NULL
;
if
(
withSlimit
&&
isPartTableAgg
(
pAgg
))
{
if
(
withSlimit
&&
(
isPartTableAgg
(
pAgg
)
||
isPartTagAgg
(
pAgg
)
))
{
pScan
->
groupOrderScan
=
pAgg
->
node
.
forceCreateNonBlockingOptr
=
true
;
pScan
->
groupOrderScan
=
pAgg
->
node
.
forceCreateNonBlockingOptr
=
true
;
}
}
}
}
...
@@ -2698,39 +2698,31 @@ static void swapLimit(SLogicNode* pParent, SLogicNode* pChild) {
...
@@ -2698,39 +2698,31 @@ static void swapLimit(SLogicNode* pParent, SLogicNode* pChild) {
pParent
->
pLimit
=
NULL
;
pParent
->
pLimit
=
NULL
;
}
}
static
void
cloneLimit
(
SLogicNode
*
pParent
,
SLogicNode
*
pChild
)
{
SLimitNode
*
pLimit
=
NULL
;
if
(
pParent
->
pLimit
)
{
pChild
->
pLimit
=
nodesCloneNode
(
pParent
->
pLimit
);
pLimit
=
(
SLimitNode
*
)
pChild
->
pLimit
;
pLimit
->
limit
+=
pLimit
->
offset
;
pLimit
->
offset
=
0
;
}
if
(
pParent
->
pSlimit
)
{
pChild
->
pSlimit
=
nodesCloneNode
(
pParent
->
pSlimit
);
pLimit
=
(
SLimitNode
*
)
pChild
->
pSlimit
;
pLimit
->
limit
+=
pLimit
->
offset
;
pLimit
->
offset
=
0
;
}
}
static
bool
pushDownLimitHow
(
SLogicNode
*
pNodeWithLimit
,
SLogicNode
*
pNodeLimitPushTo
);
static
bool
pushDownLimitHow
(
SLogicNode
*
pNodeWithLimit
,
SLogicNode
*
pNodeLimitPushTo
);
static
bool
pushDownLimitTo
(
SLogicNode
*
pNodeWithLimit
,
SLogicNode
*
pNodeLimitPushTo
)
{
static
bool
pushDownLimitTo
(
SLogicNode
*
pNodeWithLimit
,
SLogicNode
*
pNodeLimitPushTo
)
{
switch
(
nodeType
(
pNodeLimitPushTo
))
{
switch
(
nodeType
(
pNodeLimitPushTo
))
{
case
QUERY_NODE_LOGIC_PLAN_WINDOW
:
{
case
QUERY_NODE_LOGIC_PLAN_WINDOW
:
{
SWindowLogicNode
*
pWindow
=
(
SWindowLogicNode
*
)
pNodeLimitPushTo
;
SWindowLogicNode
*
pWindow
=
(
SWindowLogicNode
*
)
pNodeLimitPushTo
;
if
(
pWindow
->
winType
!=
WINDOW_TYPE_INTERVAL
)
break
;
if
(
pWindow
->
winType
!=
WINDOW_TYPE_INTERVAL
)
break
;
cloneLimit
(
pNodeWithLimit
,
pNodeLimitPushTo
);
cloneLimit
(
pNodeWithLimit
,
pNodeLimitPushTo
,
CLONE_LIMIT_SLIMIT
);
return
true
;
return
true
;
}
}
case
QUERY_NODE_LOGIC_PLAN_FILL
:
case
QUERY_NODE_LOGIC_PLAN_FILL
:
case
QUERY_NODE_LOGIC_PLAN_SORT
:
{
case
QUERY_NODE_LOGIC_PLAN_SORT
:
{
cloneLimit
(
pNodeWithLimit
,
pNodeLimitPushTo
);
cloneLimit
(
pNodeWithLimit
,
pNodeLimitPushTo
,
CLONE_LIMIT_SLIMIT
);
SNode
*
pChild
=
NULL
;
SNode
*
pChild
=
NULL
;
FOREACH
(
pChild
,
pNodeLimitPushTo
->
pChildren
)
{
pushDownLimitHow
(
pNodeLimitPushTo
,
(
SLogicNode
*
)
pChild
);
}
FOREACH
(
pChild
,
pNodeLimitPushTo
->
pChildren
)
{
pushDownLimitHow
(
pNodeLimitPushTo
,
(
SLogicNode
*
)
pChild
);
}
return
true
;
return
true
;
}
}
case
QUERY_NODE_LOGIC_PLAN_AGG
:
{
if
(
nodeType
(
pNodeWithLimit
)
==
QUERY_NODE_LOGIC_PLAN_PROJECT
&&
(
isPartTagAgg
((
SAggLogicNode
*
)
pNodeLimitPushTo
)
||
isPartTableAgg
((
SAggLogicNode
*
)
pNodeLimitPushTo
)))
{
// when part by tag, slimit will be cloned to agg, and it will be pipelined.
// The scan below will do scanning with group order
return
cloneLimit
(
pNodeWithLimit
,
pNodeLimitPushTo
,
CLONE_SLIMIT
);
}
break
;
}
case
QUERY_NODE_LOGIC_PLAN_SCAN
:
case
QUERY_NODE_LOGIC_PLAN_SCAN
:
if
(
nodeType
(
pNodeWithLimit
)
==
QUERY_NODE_LOGIC_PLAN_PROJECT
&&
pNodeWithLimit
->
pLimit
)
{
if
(
nodeType
(
pNodeWithLimit
)
==
QUERY_NODE_LOGIC_PLAN_PROJECT
&&
pNodeWithLimit
->
pLimit
)
{
swapLimit
(
pNodeWithLimit
,
pNodeLimitPushTo
);
swapLimit
(
pNodeWithLimit
,
pNodeLimitPushTo
);
...
...
source/libs/planner/src/planPhysiCreater.c
浏览文件 @
e587cc50
...
@@ -872,12 +872,16 @@ static int32_t rewritePrecalcExpr(SPhysiPlanContext* pCxt, SNode* pNode, SNodeLi
...
@@ -872,12 +872,16 @@ static int32_t rewritePrecalcExpr(SPhysiPlanContext* pCxt, SNode* pNode, SNodeLi
}
}
static
int32_t
createAggPhysiNode
(
SPhysiPlanContext
*
pCxt
,
SNodeList
*
pChildren
,
SAggLogicNode
*
pAggLogicNode
,
static
int32_t
createAggPhysiNode
(
SPhysiPlanContext
*
pCxt
,
SNodeList
*
pChildren
,
SAggLogicNode
*
pAggLogicNode
,
SPhysiNode
**
pPhyNode
)
{
SPhysiNode
**
pPhyNode
,
SSubplan
*
pSubPlan
)
{
SAggPhysiNode
*
pAgg
=
SAggPhysiNode
*
pAgg
=
(
SAggPhysiNode
*
)
makePhysiNode
(
pCxt
,
(
SLogicNode
*
)
pAggLogicNode
,
QUERY_NODE_PHYSICAL_PLAN_HASH_AGG
);
(
SAggPhysiNode
*
)
makePhysiNode
(
pCxt
,
(
SLogicNode
*
)
pAggLogicNode
,
QUERY_NODE_PHYSICAL_PLAN_HASH_AGG
);
if
(
NULL
==
pAgg
)
{
if
(
NULL
==
pAgg
)
{
return
TSDB_CODE_OUT_OF_MEMORY
;
return
TSDB_CODE_OUT_OF_MEMORY
;
}
}
if
(
pAgg
->
node
.
pSlimit
)
{
pSubPlan
->
dynamicRowThreshold
=
true
;
pSubPlan
->
rowsThreshold
=
((
SLimitNode
*
)
pAgg
->
node
.
pSlimit
)
->
limit
;
}
pAgg
->
mergeDataBlock
=
(
GROUP_ACTION_KEEP
==
pAggLogicNode
->
node
.
groupAction
?
false
:
true
);
pAgg
->
mergeDataBlock
=
(
GROUP_ACTION_KEEP
==
pAggLogicNode
->
node
.
groupAction
?
false
:
true
);
pAgg
->
groupKeyOptimized
=
pAggLogicNode
->
hasGroupKeyOptimized
;
pAgg
->
groupKeyOptimized
=
pAggLogicNode
->
hasGroupKeyOptimized
;
...
@@ -1617,7 +1621,7 @@ static int32_t doCreatePhysiNode(SPhysiPlanContext* pCxt, SLogicNode* pLogicNode
...
@@ -1617,7 +1621,7 @@ static int32_t doCreatePhysiNode(SPhysiPlanContext* pCxt, SLogicNode* pLogicNode
case
QUERY_NODE_LOGIC_PLAN_JOIN
:
case
QUERY_NODE_LOGIC_PLAN_JOIN
:
return
createJoinPhysiNode
(
pCxt
,
pChildren
,
(
SJoinLogicNode
*
)
pLogicNode
,
pPhyNode
);
return
createJoinPhysiNode
(
pCxt
,
pChildren
,
(
SJoinLogicNode
*
)
pLogicNode
,
pPhyNode
);
case
QUERY_NODE_LOGIC_PLAN_AGG
:
case
QUERY_NODE_LOGIC_PLAN_AGG
:
return
createAggPhysiNode
(
pCxt
,
pChildren
,
(
SAggLogicNode
*
)
pLogicNode
,
pPhyNode
);
return
createAggPhysiNode
(
pCxt
,
pChildren
,
(
SAggLogicNode
*
)
pLogicNode
,
pPhyNode
,
pSubplan
);
case
QUERY_NODE_LOGIC_PLAN_PROJECT
:
case
QUERY_NODE_LOGIC_PLAN_PROJECT
:
return
createProjectPhysiNode
(
pCxt
,
pChildren
,
(
SProjectLogicNode
*
)
pLogicNode
,
pPhyNode
);
return
createProjectPhysiNode
(
pCxt
,
pChildren
,
(
SProjectLogicNode
*
)
pLogicNode
,
pPhyNode
);
case
QUERY_NODE_LOGIC_PLAN_EXCHANGE
:
case
QUERY_NODE_LOGIC_PLAN_EXCHANGE
:
...
@@ -1721,6 +1725,8 @@ static SSubplan* makeSubplan(SPhysiPlanContext* pCxt, SLogicSubplan* pLogicSubpl
...
@@ -1721,6 +1725,8 @@ static SSubplan* makeSubplan(SPhysiPlanContext* pCxt, SLogicSubplan* pLogicSubpl
pSubplan
->
id
=
pLogicSubplan
->
id
;
pSubplan
->
id
=
pLogicSubplan
->
id
;
pSubplan
->
subplanType
=
pLogicSubplan
->
subplanType
;
pSubplan
->
subplanType
=
pLogicSubplan
->
subplanType
;
pSubplan
->
level
=
pLogicSubplan
->
level
;
pSubplan
->
level
=
pLogicSubplan
->
level
;
pSubplan
->
rowsThreshold
=
4096
;
pSubplan
->
dynamicRowThreshold
=
false
;
if
(
NULL
!=
pCxt
->
pPlanCxt
->
pUser
)
{
if
(
NULL
!=
pCxt
->
pPlanCxt
->
pUser
)
{
snprintf
(
pSubplan
->
user
,
sizeof
(
pSubplan
->
user
),
"%s"
,
pCxt
->
pPlanCxt
->
pUser
);
snprintf
(
pSubplan
->
user
,
sizeof
(
pSubplan
->
user
),
"%s"
,
pCxt
->
pPlanCxt
->
pUser
);
}
}
...
...
source/libs/planner/src/planSpliter.c
浏览文件 @
e587cc50
...
@@ -867,8 +867,16 @@ static int32_t stbSplSplitAggNodeForPartTable(SSplitContext* pCxt, SStableSplitI
...
@@ -867,8 +867,16 @@ static int32_t stbSplSplitAggNodeForPartTable(SSplitContext* pCxt, SStableSplitI
static
int32_t
stbSplSplitAggNodeForCrossTable
(
SSplitContext
*
pCxt
,
SStableSplitInfo
*
pInfo
)
{
static
int32_t
stbSplSplitAggNodeForCrossTable
(
SSplitContext
*
pCxt
,
SStableSplitInfo
*
pInfo
)
{
SLogicNode
*
pPartAgg
=
NULL
;
SLogicNode
*
pPartAgg
=
NULL
;
int32_t
code
=
stbSplCreatePartAggNode
((
SAggLogicNode
*
)
pInfo
->
pSplitNode
,
&
pPartAgg
);
int32_t
code
=
stbSplCreatePartAggNode
((
SAggLogicNode
*
)
pInfo
->
pSplitNode
,
&
pPartAgg
);
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
stbSplCreateExchangeNode
(
pCxt
,
pInfo
->
pSplitNode
,
pPartAgg
);
// if slimit was pushed down to agg, agg will be pipelined mode, add sort merge before parent agg
if
((
SAggLogicNode
*
)
pInfo
->
pSplitNode
->
pSlimit
)
code
=
stbSplCreateMergeNode
(
pCxt
,
NULL
,
pInfo
->
pSplitNode
,
NULL
,
pPartAgg
,
true
);
else
code
=
stbSplCreateExchangeNode
(
pCxt
,
pInfo
->
pSplitNode
,
pPartAgg
);
}
else
{
nodesDestroyNode
((
SNode
*
)
pPartAgg
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
nodesListMakeStrictAppend
(
&
pInfo
->
pSubplan
->
pChildren
,
code
=
nodesListMakeStrictAppend
(
&
pInfo
->
pSubplan
->
pChildren
,
...
...
source/libs/planner/src/planUtil.c
浏览文件 @
e587cc50
...
@@ -349,7 +349,7 @@ static bool stbHasPartTbname(SNodeList* pPartKeys) {
...
@@ -349,7 +349,7 @@ static bool stbHasPartTbname(SNodeList* pPartKeys) {
return
false
;
return
false
;
}
}
static
SNodeList
*
stb
Spl
GetPartKeys
(
SLogicNode
*
pNode
)
{
static
SNodeList
*
stbGetPartKeys
(
SLogicNode
*
pNode
)
{
if
(
QUERY_NODE_LOGIC_PLAN_SCAN
==
nodeType
(
pNode
))
{
if
(
QUERY_NODE_LOGIC_PLAN_SCAN
==
nodeType
(
pNode
))
{
return
((
SScanLogicNode
*
)
pNode
)
->
pGroupTags
;
return
((
SScanLogicNode
*
)
pNode
)
->
pGroupTags
;
}
else
if
(
QUERY_NODE_LOGIC_PLAN_PARTITION
==
nodeType
(
pNode
))
{
}
else
if
(
QUERY_NODE_LOGIC_PLAN_PARTITION
==
nodeType
(
pNode
))
{
...
@@ -367,11 +367,58 @@ bool isPartTableAgg(SAggLogicNode* pAgg) {
...
@@ -367,11 +367,58 @@ bool isPartTableAgg(SAggLogicNode* pAgg) {
return
stbHasPartTbname
(
pAgg
->
pGroupKeys
)
&&
return
stbHasPartTbname
(
pAgg
->
pGroupKeys
)
&&
stbNotSystemScan
((
SLogicNode
*
)
nodesListGetNode
(
pAgg
->
node
.
pChildren
,
0
));
stbNotSystemScan
((
SLogicNode
*
)
nodesListGetNode
(
pAgg
->
node
.
pChildren
,
0
));
}
}
return
stbHasPartTbname
(
stb
Spl
GetPartKeys
((
SLogicNode
*
)
nodesListGetNode
(
pAgg
->
node
.
pChildren
,
0
)));
return
stbHasPartTbname
(
stbGetPartKeys
((
SLogicNode
*
)
nodesListGetNode
(
pAgg
->
node
.
pChildren
,
0
)));
}
}
bool
isPartTableWinodw
(
SWindowLogicNode
*
pWindow
)
{
static
bool
stbHasPartTag
(
SNodeList
*
pPartKeys
)
{
return
stbHasPartTbname
(
stbSplGetPartKeys
((
SLogicNode
*
)
nodesListGetNode
(
pWindow
->
node
.
pChildren
,
0
)));
if
(
NULL
==
pPartKeys
)
{
return
false
;
}
SNode
*
pPartKey
=
NULL
;
FOREACH
(
pPartKey
,
pPartKeys
)
{
if
(
QUERY_NODE_GROUPING_SET
==
nodeType
(
pPartKey
))
{
pPartKey
=
nodesListGetNode
(((
SGroupingSetNode
*
)
pPartKey
)
->
pParameterList
,
0
);
}
if
((
QUERY_NODE_FUNCTION
==
nodeType
(
pPartKey
)
&&
FUNCTION_TYPE_TAGS
==
((
SFunctionNode
*
)
pPartKey
)
->
funcType
)
||
(
QUERY_NODE_COLUMN
==
nodeType
(
pPartKey
)
&&
COLUMN_TYPE_TAG
==
((
SColumnNode
*
)
pPartKey
)
->
colType
))
{
return
true
;
}
}
return
false
;
}
}
bool
isPartTagAgg
(
SAggLogicNode
*
pAgg
)
{
if
(
1
!=
LIST_LENGTH
(
pAgg
->
node
.
pChildren
))
{
return
false
;
}
if
(
pAgg
->
pGroupKeys
)
{
return
stbHasPartTag
(
pAgg
->
pGroupKeys
)
&&
stbNotSystemScan
((
SLogicNode
*
)
nodesListGetNode
(
pAgg
->
node
.
pChildren
,
0
));
}
return
stbHasPartTag
(
stbGetPartKeys
((
SLogicNode
*
)
nodesListGetNode
(
pAgg
->
node
.
pChildren
,
0
)));
}
bool
isPartTableWinodw
(
SWindowLogicNode
*
pWindow
)
{
return
stbHasPartTbname
(
stbGetPartKeys
((
SLogicNode
*
)
nodesListGetNode
(
pWindow
->
node
.
pChildren
,
0
)));
}
bool
cloneLimit
(
SLogicNode
*
pParent
,
SLogicNode
*
pChild
,
uint8_t
cloneWhat
)
{
SLimitNode
*
pLimit
;
bool
cloned
=
false
;
if
(
pParent
->
pLimit
&&
(
cloneWhat
&
CLONE_LIMIT
))
{
pChild
->
pLimit
=
nodesCloneNode
(
pParent
->
pLimit
);
pLimit
=
(
SLimitNode
*
)
pChild
->
pLimit
;
pLimit
->
limit
+=
pLimit
->
offset
;
pLimit
->
offset
=
0
;
cloned
=
true
;
}
if
(
pParent
->
pSlimit
&&
(
cloneWhat
&
CLONE_SLIMIT
))
{
pChild
->
pSlimit
=
nodesCloneNode
(
pParent
->
pSlimit
);
pLimit
=
(
SLimitNode
*
)
pChild
->
pSlimit
;
pLimit
->
limit
+=
pLimit
->
offset
;
pLimit
->
offset
=
0
;
cloned
=
true
;
}
return
cloned
;
}
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录