Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
81a2bf10
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看板
提交
81a2bf10
编写于
6月 28, 2023
作者:
D
dapan1121
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
enh: optimize stable join
上级
8f39b9d2
变更
9
隐藏空白更改
内联
并排
Showing
9 changed file
with
385 addition
and
67 deletion
+385
-67
include/libs/nodes/nodes.h
include/libs/nodes/nodes.h
+1
-0
include/libs/nodes/plannodes.h
include/libs/nodes/plannodes.h
+7
-1
include/libs/nodes/querynodes.h
include/libs/nodes/querynodes.h
+4
-0
source/libs/nodes/src/nodesCloneFuncs.c
source/libs/nodes/src/nodesCloneFuncs.c
+2
-0
source/libs/nodes/src/nodesCodeFuncs.c
source/libs/nodes/src/nodesCodeFuncs.c
+81
-5
source/libs/nodes/src/nodesUtilFuncs.c
source/libs/nodes/src/nodesUtilFuncs.c
+13
-0
source/libs/planner/src/planOptimizer.c
source/libs/planner/src/planOptimizer.c
+208
-33
source/libs/planner/src/planPhysiCreater.c
source/libs/planner/src/planPhysiCreater.c
+68
-27
source/libs/planner/src/planSpliter.c
source/libs/planner/src/planSpliter.c
+1
-1
未找到文件。
include/libs/nodes/nodes.h
浏览文件 @
81a2bf10
...
@@ -233,6 +233,7 @@ typedef enum ENodeType {
...
@@ -233,6 +233,7 @@ typedef enum ENodeType {
QUERY_NODE_LOGIC_PLAN_INDEF_ROWS_FUNC
,
QUERY_NODE_LOGIC_PLAN_INDEF_ROWS_FUNC
,
QUERY_NODE_LOGIC_PLAN_INTERP_FUNC
,
QUERY_NODE_LOGIC_PLAN_INTERP_FUNC
,
QUERY_NODE_LOGIC_PLAN_GROUP_CACHE
,
QUERY_NODE_LOGIC_PLAN_GROUP_CACHE
,
QUERY_NODE_LOGIC_PLAN_DYN_QUERY_CTRL
,
QUERY_NODE_LOGIC_SUBPLAN
,
QUERY_NODE_LOGIC_SUBPLAN
,
QUERY_NODE_LOGIC_PLAN
,
QUERY_NODE_LOGIC_PLAN
,
...
...
include/libs/nodes/plannodes.h
浏览文件 @
81a2bf10
...
@@ -42,6 +42,7 @@ typedef enum EGroupAction {
...
@@ -42,6 +42,7 @@ typedef enum EGroupAction {
typedef
struct
SLogicNode
{
typedef
struct
SLogicNode
{
ENodeType
type
;
ENodeType
type
;
bool
dynamicOp
;
SNodeList
*
pTargets
;
// SColumnNode
SNodeList
*
pTargets
;
// SColumnNode
SNode
*
pConditions
;
SNode
*
pConditions
;
SNodeList
*
pChildren
;
SNodeList
*
pChildren
;
...
@@ -114,6 +115,7 @@ typedef struct SJoinLogicNode {
...
@@ -114,6 +115,7 @@ typedef struct SJoinLogicNode {
SNode
*
pPrimKeyEqCond
;
SNode
*
pPrimKeyEqCond
;
SNode
*
pColEqCond
;
SNode
*
pColEqCond
;
SNode
*
pTagEqCond
;
SNode
*
pTagEqCond
;
SNode
*
pTagOnCond
;
SNode
*
pOtherOnCond
;
SNode
*
pOtherOnCond
;
bool
isSingleTableJoin
;
bool
isSingleTableJoin
;
bool
hasSubQuery
;
bool
hasSubQuery
;
...
@@ -157,9 +159,13 @@ typedef struct SInterpFuncLogicNode {
...
@@ -157,9 +159,13 @@ typedef struct SInterpFuncLogicNode {
typedef
struct
SGroupCacheLogicNode
{
typedef
struct
SGroupCacheLogicNode
{
SLogicNode
node
;
SLogicNode
node
;
SNode
*
pGroupCol
;
SNode
List
*
pGroupCols
;
}
SGroupCacheLogicNode
;
}
SGroupCacheLogicNode
;
typedef
struct
SDynQueryCtrlLogicNode
{
SLogicNode
node
;
EDynQueryType
qType
;
}
SDynQueryCtrlLogicNode
;
typedef
enum
EModifyTableType
{
MODIFY_TABLE_TYPE_INSERT
=
1
,
MODIFY_TABLE_TYPE_DELETE
}
EModifyTableType
;
typedef
enum
EModifyTableType
{
MODIFY_TABLE_TYPE_INSERT
=
1
,
MODIFY_TABLE_TYPE_DELETE
}
EModifyTableType
;
...
...
include/libs/nodes/querynodes.h
浏览文件 @
81a2bf10
...
@@ -180,6 +180,10 @@ typedef enum EJoinAlgorithm {
...
@@ -180,6 +180,10 @@ typedef enum EJoinAlgorithm {
JOIN_ALGO_HASH
,
JOIN_ALGO_HASH
,
}
EJoinAlgorithm
;
}
EJoinAlgorithm
;
typedef
enum
EDynQueryType
{
DYN_QTYPE_STB_HASH
=
1
,
}
EDynQueryType
;
typedef
struct
SJoinTableNode
{
typedef
struct
SJoinTableNode
{
STableNode
table
;
// QUERY_NODE_JOIN_TABLE
STableNode
table
;
// QUERY_NODE_JOIN_TABLE
EJoinType
joinType
;
EJoinType
joinType
;
...
...
source/libs/nodes/src/nodesCloneFuncs.c
浏览文件 @
81a2bf10
...
@@ -403,11 +403,13 @@ static int32_t logicScanCopy(const SScanLogicNode* pSrc, SScanLogicNode* pDst) {
...
@@ -403,11 +403,13 @@ static int32_t logicScanCopy(const SScanLogicNode* pSrc, SScanLogicNode* pDst) {
static
int32_t
logicJoinCopy
(
const
SJoinLogicNode
*
pSrc
,
SJoinLogicNode
*
pDst
)
{
static
int32_t
logicJoinCopy
(
const
SJoinLogicNode
*
pSrc
,
SJoinLogicNode
*
pDst
)
{
COPY_BASE_OBJECT_FIELD
(
node
,
logicNodeCopy
);
COPY_BASE_OBJECT_FIELD
(
node
,
logicNodeCopy
);
COPY_SCALAR_FIELD
(
joinType
);
COPY_SCALAR_FIELD
(
joinType
);
COPY_SCALAR_FIELD
(
joinAlgo
);
CLONE_NODE_FIELD
(
pPrimKeyEqCond
);
CLONE_NODE_FIELD
(
pPrimKeyEqCond
);
CLONE_NODE_FIELD
(
pColEqCond
);
CLONE_NODE_FIELD
(
pColEqCond
);
CLONE_NODE_FIELD
(
pTagEqCond
);
CLONE_NODE_FIELD
(
pTagEqCond
);
CLONE_NODE_FIELD
(
pOtherOnCond
);
CLONE_NODE_FIELD
(
pOtherOnCond
);
COPY_SCALAR_FIELD
(
isSingleTableJoin
);
COPY_SCALAR_FIELD
(
isSingleTableJoin
);
COPY_SCALAR_FIELD
(
hasSubQuery
);
return
TSDB_CODE_SUCCESS
;
return
TSDB_CODE_SUCCESS
;
}
}
...
...
source/libs/nodes/src/nodesCodeFuncs.c
浏览文件 @
81a2bf10
...
@@ -295,6 +295,10 @@ const char* nodesNodeName(ENodeType type) {
...
@@ -295,6 +295,10 @@ const char* nodesNodeName(ENodeType type) {
return
"LogicIndefRowsFunc"
;
return
"LogicIndefRowsFunc"
;
case
QUERY_NODE_LOGIC_PLAN_INTERP_FUNC
:
case
QUERY_NODE_LOGIC_PLAN_INTERP_FUNC
:
return
"LogicInterpFunc"
;
return
"LogicInterpFunc"
;
case
QUERY_NODE_LOGIC_PLAN_GROUP_CACHE
:
return
"LogicGroupCache"
;
case
QUERY_NODE_LOGIC_PLAN_DYN_QUERY_CTRL
:
return
"LogicDynamicQueryCtrl"
;
case
QUERY_NODE_LOGIC_SUBPLAN
:
case
QUERY_NODE_LOGIC_SUBPLAN
:
return
"LogicSubplan"
;
return
"LogicSubplan"
;
case
QUERY_NODE_LOGIC_PLAN
:
case
QUERY_NODE_LOGIC_PLAN
:
...
@@ -1172,6 +1176,55 @@ static int32_t jsonToLogicInterpFuncNode(const SJson* pJson, void* pObj) {
...
@@ -1172,6 +1176,55 @@ static int32_t jsonToLogicInterpFuncNode(const SJson* pJson, void* pObj) {
return
code
;
return
code
;
}
}
static
const
char
*
jkGroupCacheLogicPlanGroupCols
=
"GroupCols"
;
static
int32_t
logicGroupCacheNodeToJson
(
const
void
*
pObj
,
SJson
*
pJson
)
{
const
SGroupCacheLogicNode
*
pNode
=
(
const
SGroupCacheLogicNode
*
)
pObj
;
int32_t
code
=
logicPlanNodeToJson
(
pObj
,
pJson
);
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
nodeListToJson
(
pJson
,
jkGroupCacheLogicPlanGroupCols
,
pNode
->
pGroupCols
);
}
return
code
;
}
static
int32_t
jsonToLogicGroupCacheNode
(
const
SJson
*
pJson
,
void
*
pObj
)
{
SGroupCacheLogicNode
*
pNode
=
(
SGroupCacheLogicNode
*
)
pObj
;
int32_t
code
=
jsonToLogicPlanNode
(
pJson
,
pObj
);
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
jsonToNodeList
(
pJson
,
jkGroupCacheLogicPlanGroupCols
,
&
pNode
->
pGroupCols
);
}
return
code
;
}
static
const
char
*
jkDynQueryCtrlLogicPlanQueryType
=
"QueryType"
;
static
int32_t
logicDynQueryCtrlNodeToJson
(
const
void
*
pObj
,
SJson
*
pJson
)
{
const
SDynQueryCtrlLogicNode
*
pNode
=
(
const
SDynQueryCtrlLogicNode
*
)
pObj
;
int32_t
code
=
logicPlanNodeToJson
(
pObj
,
pJson
);
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAddIntegerToObject
(
pJson
,
jkDynQueryCtrlLogicPlanQueryType
,
pNode
->
qType
);
}
return
code
;
}
static
int32_t
jsonToLogicDynQueryCtrlNode
(
const
SJson
*
pJson
,
void
*
pObj
)
{
SDynQueryCtrlLogicNode
*
pNode
=
(
SDynQueryCtrlLogicNode
*
)
pObj
;
int32_t
code
=
jsonToLogicPlanNode
(
pJson
,
pObj
);
if
(
TSDB_CODE_SUCCESS
==
code
)
{
tjsonGetNumberValue
(
pJson
,
jkDynQueryCtrlLogicPlanQueryType
,
pNode
->
qType
,
code
);
}
return
code
;
}
static
const
char
*
jkSubplanIdQueryId
=
"QueryId"
;
static
const
char
*
jkSubplanIdQueryId
=
"QueryId"
;
static
const
char
*
jkSubplanIdGroupId
=
"GroupId"
;
static
const
char
*
jkSubplanIdGroupId
=
"GroupId"
;
static
const
char
*
jkSubplanIdSubplanId
=
"SubplanId"
;
static
const
char
*
jkSubplanIdSubplanId
=
"SubplanId"
;
...
@@ -1426,9 +1479,11 @@ static int32_t jsonToLogicPlan(const SJson* pJson, void* pObj) {
...
@@ -1426,9 +1479,11 @@ static int32_t jsonToLogicPlan(const SJson* pJson, void* pObj) {
}
}
static
const
char
*
jkJoinLogicPlanJoinType
=
"JoinType"
;
static
const
char
*
jkJoinLogicPlanJoinType
=
"JoinType"
;
static
const
char
*
jkJoinLogicPlanOnConditions
=
"OnConditions"
;
static
const
char
*
jkJoinLogicPlanJoinAlgo
=
"JoinAlgo"
;
static
const
char
*
jkJoinLogicPlanOnConditions
=
"OtherOnCond"
;
static
const
char
*
jkJoinLogicPlanPrimKeyEqCondition
=
"PrimKeyEqCond"
;
static
const
char
*
jkJoinLogicPlanPrimKeyEqCondition
=
"PrimKeyEqCond"
;
static
const
char
*
jkJoinLogicPlanColEqCondition
=
"ColumnEqCond"
;
static
const
char
*
jkJoinLogicPlanColEqCondition
=
"ColumnEqCond"
;
static
const
char
*
jkJoinLogicPlanTagEqCondition
=
"TagEqCond"
;
static
int32_t
logicJoinNodeToJson
(
const
void
*
pObj
,
SJson
*
pJson
)
{
static
int32_t
logicJoinNodeToJson
(
const
void
*
pObj
,
SJson
*
pJson
)
{
const
SJoinLogicNode
*
pNode
=
(
const
SJoinLogicNode
*
)
pObj
;
const
SJoinLogicNode
*
pNode
=
(
const
SJoinLogicNode
*
)
pObj
;
...
@@ -1438,14 +1493,20 @@ static int32_t logicJoinNodeToJson(const void* pObj, SJson* pJson) {
...
@@ -1438,14 +1493,20 @@ static int32_t logicJoinNodeToJson(const void* pObj, SJson* pJson) {
code
=
tjsonAddIntegerToObject
(
pJson
,
jkJoinLogicPlanJoinType
,
pNode
->
joinType
);
code
=
tjsonAddIntegerToObject
(
pJson
,
jkJoinLogicPlanJoinType
,
pNode
->
joinType
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAdd
Object
(
pJson
,
jkJoinLogicPlanPrimKeyEqCondition
,
nodeToJson
,
pNode
->
pPrimKeyEqCond
);
code
=
tjsonAdd
IntegerToObject
(
pJson
,
jkJoinLogicPlanJoinAlgo
,
pNode
->
joinAlgo
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAddObject
(
pJson
,
jkJoinLogicPlan
OnConditions
,
nodeToJson
,
pNode
->
pOtherOn
Cond
);
code
=
tjsonAddObject
(
pJson
,
jkJoinLogicPlan
PrimKeyEqCondition
,
nodeToJson
,
pNode
->
pPrimKeyEq
Cond
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAddObject
(
pJson
,
jkJoinLogicPlanColEqCondition
,
nodeToJson
,
pNode
->
pColEqCond
);
code
=
tjsonAddObject
(
pJson
,
jkJoinLogicPlanColEqCondition
,
nodeToJson
,
pNode
->
pColEqCond
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAddObject
(
pJson
,
jkJoinLogicPlanTagEqCondition
,
nodeToJson
,
pNode
->
pTagEqCond
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAddObject
(
pJson
,
jkJoinLogicPlanOnConditions
,
nodeToJson
,
pNode
->
pOtherOnCond
);
}
return
code
;
return
code
;
}
}
...
@@ -1457,14 +1518,21 @@ static int32_t jsonToLogicJoinNode(const SJson* pJson, void* pObj) {
...
@@ -1457,14 +1518,21 @@ static int32_t jsonToLogicJoinNode(const SJson* pJson, void* pObj) {
tjsonGetNumberValue
(
pJson
,
jkJoinLogicPlanJoinType
,
pNode
->
joinType
,
code
);
tjsonGetNumberValue
(
pJson
,
jkJoinLogicPlanJoinType
,
pNode
->
joinType
,
code
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
jsonToNodeObject
(
pJson
,
jkJoinLogicPlanPrimKeyEqCondition
,
&
pNode
->
pPrimKeyEqCond
);
tjsonGetNumberValue
(
pJson
,
jkJoinLogicPlanJoinAlgo
,
pNode
->
joinAlgo
,
code
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
jsonToNodeObject
(
pJson
,
jkJoinLogicPlan
OnConditions
,
&
pNode
->
pOtherOn
Cond
);
code
=
jsonToNodeObject
(
pJson
,
jkJoinLogicPlan
PrimKeyEqCondition
,
&
pNode
->
pPrimKeyEq
Cond
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
jsonToNodeObject
(
pJson
,
jkJoinLogicPlanColEqCondition
,
&
pNode
->
pColEqCond
);
code
=
jsonToNodeObject
(
pJson
,
jkJoinLogicPlanColEqCondition
,
&
pNode
->
pColEqCond
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
jsonToNodeObject
(
pJson
,
jkJoinLogicPlanTagEqCondition
,
&
pNode
->
pTagEqCond
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
jsonToNodeObject
(
pJson
,
jkJoinLogicPlanOnConditions
,
&
pNode
->
pOtherOnCond
);
}
return
code
;
return
code
;
}
}
...
@@ -6577,6 +6645,10 @@ static int32_t specificNodeToJson(const void* pObj, SJson* pJson) {
...
@@ -6577,6 +6645,10 @@ static int32_t specificNodeToJson(const void* pObj, SJson* pJson) {
return
logicIndefRowsFuncNodeToJson
(
pObj
,
pJson
);
return
logicIndefRowsFuncNodeToJson
(
pObj
,
pJson
);
case
QUERY_NODE_LOGIC_PLAN_INTERP_FUNC
:
case
QUERY_NODE_LOGIC_PLAN_INTERP_FUNC
:
return
logicInterpFuncNodeToJson
(
pObj
,
pJson
);
return
logicInterpFuncNodeToJson
(
pObj
,
pJson
);
case
QUERY_NODE_LOGIC_PLAN_GROUP_CACHE
:
return
logicGroupCacheNodeToJson
(
pObj
,
pJson
);
case
QUERY_NODE_LOGIC_PLAN_DYN_QUERY_CTRL
:
return
logicDynQueryCtrlNodeToJson
(
pObj
,
pJson
);
case
QUERY_NODE_LOGIC_SUBPLAN
:
case
QUERY_NODE_LOGIC_SUBPLAN
:
return
logicSubplanToJson
(
pObj
,
pJson
);
return
logicSubplanToJson
(
pObj
,
pJson
);
case
QUERY_NODE_LOGIC_PLAN
:
case
QUERY_NODE_LOGIC_PLAN
:
...
@@ -6895,6 +6967,10 @@ static int32_t jsonToSpecificNode(const SJson* pJson, void* pObj) {
...
@@ -6895,6 +6967,10 @@ static int32_t jsonToSpecificNode(const SJson* pJson, void* pObj) {
return
jsonToLogicIndefRowsFuncNode
(
pJson
,
pObj
);
return
jsonToLogicIndefRowsFuncNode
(
pJson
,
pObj
);
case
QUERY_NODE_LOGIC_PLAN_INTERP_FUNC
:
case
QUERY_NODE_LOGIC_PLAN_INTERP_FUNC
:
return
jsonToLogicInterpFuncNode
(
pJson
,
pObj
);
return
jsonToLogicInterpFuncNode
(
pJson
,
pObj
);
case
QUERY_NODE_LOGIC_PLAN_GROUP_CACHE
:
return
jsonToLogicGroupCacheNode
(
pJson
,
pObj
);
case
QUERY_NODE_LOGIC_PLAN_DYN_QUERY_CTRL
:
return
jsonToLogicDynQueryCtrlNode
(
pJson
,
pObj
);
case
QUERY_NODE_LOGIC_SUBPLAN
:
case
QUERY_NODE_LOGIC_SUBPLAN
:
return
jsonToLogicSubplan
(
pJson
,
pObj
);
return
jsonToLogicSubplan
(
pJson
,
pObj
);
case
QUERY_NODE_LOGIC_PLAN
:
case
QUERY_NODE_LOGIC_PLAN
:
...
...
source/libs/nodes/src/nodesUtilFuncs.c
浏览文件 @
81a2bf10
...
@@ -492,6 +492,8 @@ SNode* nodesMakeNode(ENodeType type) {
...
@@ -492,6 +492,8 @@ SNode* nodesMakeNode(ENodeType type) {
return
makeNode
(
type
,
sizeof
(
SInterpFuncLogicNode
));
return
makeNode
(
type
,
sizeof
(
SInterpFuncLogicNode
));
case
QUERY_NODE_LOGIC_PLAN_GROUP_CACHE
:
case
QUERY_NODE_LOGIC_PLAN_GROUP_CACHE
:
return
makeNode
(
type
,
sizeof
(
SGroupCacheLogicNode
));
return
makeNode
(
type
,
sizeof
(
SGroupCacheLogicNode
));
case
QUERY_NODE_LOGIC_PLAN_DYN_QUERY_CTRL
:
return
makeNode
(
type
,
sizeof
(
SDynQueryCtrlLogicNode
));
case
QUERY_NODE_LOGIC_SUBPLAN
:
case
QUERY_NODE_LOGIC_SUBPLAN
:
return
makeNode
(
type
,
sizeof
(
SLogicSubplan
));
return
makeNode
(
type
,
sizeof
(
SLogicSubplan
));
case
QUERY_NODE_LOGIC_PLAN
:
case
QUERY_NODE_LOGIC_PLAN
:
...
@@ -1180,6 +1182,17 @@ void nodesDestroyNode(SNode* pNode) {
...
@@ -1180,6 +1182,17 @@ void nodesDestroyNode(SNode* pNode) {
nodesDestroyNode
(
pLogicNode
->
pTimeSeries
);
nodesDestroyNode
(
pLogicNode
->
pTimeSeries
);
break
;
break
;
}
}
case
QUERY_NODE_LOGIC_PLAN_GROUP_CACHE
:
{
SGroupCacheLogicNode
*
pLogicNode
=
(
SGroupCacheLogicNode
*
)
pNode
;
destroyLogicNode
((
SLogicNode
*
)
pLogicNode
);
nodesDestroyList
(
pLogicNode
->
pGroupCols
);
break
;
}
case
QUERY_NODE_LOGIC_PLAN_DYN_QUERY_CTRL
:
{
SDynQueryCtrlLogicNode
*
pLogicNode
=
(
SDynQueryCtrlLogicNode
*
)
pNode
;
destroyLogicNode
((
SLogicNode
*
)
pLogicNode
);
break
;
}
case
QUERY_NODE_LOGIC_SUBPLAN
:
{
case
QUERY_NODE_LOGIC_SUBPLAN
:
{
SLogicSubplan
*
pSubplan
=
(
SLogicSubplan
*
)
pNode
;
SLogicSubplan
*
pSubplan
=
(
SLogicSubplan
*
)
pNode
;
nodesDestroyList
(
pSubplan
->
pChildren
);
nodesDestroyList
(
pSubplan
->
pChildren
);
...
...
source/libs/planner/src/planOptimizer.c
浏览文件 @
81a2bf10
...
@@ -763,14 +763,18 @@ static bool pushDownCondOptIsColEqualOnCond(SJoinLogicNode* pJoin, SNode* pCond,
...
@@ -763,14 +763,18 @@ static bool pushDownCondOptIsColEqualOnCond(SJoinLogicNode* pJoin, SNode* pCond,
return
false
;
return
false
;
}
}
SOperatorNode
*
pOper
=
(
SOperatorNode
*
)
pCond
;
SOperatorNode
*
pOper
=
(
SOperatorNode
*
)
pCond
;
if
(
OP_TYPE_EQUAL
!=
pOper
->
opType
)
{
return
false
;
}
if
(
QUERY_NODE_COLUMN
!=
nodeType
(
pOper
->
pLeft
)
||
QUERY_NODE_COLUMN
!=
nodeType
(
pOper
->
pRight
))
{
if
(
QUERY_NODE_COLUMN
!=
nodeType
(
pOper
->
pLeft
)
||
QUERY_NODE_COLUMN
!=
nodeType
(
pOper
->
pRight
))
{
return
false
;
return
false
;
}
}
SColumnNode
*
pLeft
=
(
SColumnNode
*
)(
pOper
->
pLeft
);
SColumnNode
*
pLeft
=
(
SColumnNode
*
)(
pOper
->
pLeft
);
SColumnNode
*
pRight
=
(
SColumnNode
*
)(
pOper
->
pRight
);
SColumnNode
*
pRight
=
(
SColumnNode
*
)(
pOper
->
pRight
);
*
allTags
=
(
COLUMN_TYPE_TAG
==
pLeft
->
colType
)
&&
(
COLUMN_TYPE_TAG
==
pRight
->
colType
);
if
(
OP_TYPE_EQUAL
!=
pOper
->
opType
)
{
return
false
;
}
//TODO: add cast to operator and remove this restriction of optimization
//TODO: add cast to operator and remove this restriction of optimization
if
(
pLeft
->
node
.
resType
.
type
!=
pRight
->
node
.
resType
.
type
||
pLeft
->
node
.
resType
.
bytes
!=
pRight
->
node
.
resType
.
bytes
)
{
if
(
pLeft
->
node
.
resType
.
type
!=
pRight
->
node
.
resType
.
type
||
pLeft
->
node
.
resType
.
bytes
!=
pRight
->
node
.
resType
.
bytes
)
{
return
false
;
return
false
;
...
@@ -784,7 +788,6 @@ static bool pushDownCondOptIsColEqualOnCond(SJoinLogicNode* pJoin, SNode* pCond,
...
@@ -784,7 +788,6 @@ static bool pushDownCondOptIsColEqualOnCond(SJoinLogicNode* pJoin, SNode* pCond,
isEqual
=
pushDownCondOptIsTableColumn
(
pOper
->
pRight
,
pLeftCols
);
isEqual
=
pushDownCondOptIsTableColumn
(
pOper
->
pRight
,
pLeftCols
);
}
}
if
(
isEqual
)
{
if
(
isEqual
)
{
*
allTags
=
(
COLUMN_TYPE_TAG
==
pLeft
->
colType
)
&&
(
COLUMN_TYPE_TAG
==
pRight
->
colType
);
}
}
return
isEqual
;
return
isEqual
;
}
}
...
@@ -795,30 +798,42 @@ static int32_t pushDownCondOptJoinExtractEqualOnLogicCond(SJoinLogicNode* pJoin)
...
@@ -795,30 +798,42 @@ static int32_t pushDownCondOptJoinExtractEqualOnLogicCond(SJoinLogicNode* pJoin)
int32_t
code
=
TSDB_CODE_SUCCESS
;
int32_t
code
=
TSDB_CODE_SUCCESS
;
SNodeList
*
pColEqOnConds
=
NULL
;
SNodeList
*
pColEqOnConds
=
NULL
;
SNodeList
*
pTagEqOnConds
=
NULL
;
SNodeList
*
pTagEqOnConds
=
NULL
;
SNodeList
*
pTagOnConds
=
NULL
;
SNode
*
pCond
=
NULL
;
SNode
*
pCond
=
NULL
;
bool
allTags
=
false
;
bool
allTags
=
false
;
FOREACH
(
pCond
,
pLogicCond
->
pParameterList
)
{
FOREACH
(
pCond
,
pLogicCond
->
pParameterList
)
{
allTags
=
false
;
if
(
pushDownCondOptIsColEqualOnCond
(
pJoin
,
pCond
,
&
allTags
))
{
if
(
pushDownCondOptIsColEqualOnCond
(
pJoin
,
pCond
,
&
allTags
))
{
if
(
allTags
)
{
if
(
allTags
)
{
code
=
nodesListMakeAppend
(
&
pTagEqOnConds
,
nodesCloneNode
(
pCond
));
code
=
nodesListMakeAppend
(
&
pTagEqOnConds
,
nodesCloneNode
(
pCond
));
}
else
{
}
else
{
code
=
nodesListMakeAppend
(
&
pColEqOnConds
,
nodesCloneNode
(
pCond
));
code
=
nodesListMakeAppend
(
&
pColEqOnConds
,
nodesCloneNode
(
pCond
));
}
}
}
else
if
(
allTags
)
{
code
=
nodesListMakeAppend
(
&
pTagOnConds
,
nodesCloneNode
(
pCond
));
}
if
(
code
)
{
break
;
}
}
}
}
SNode
*
pTempTagEqCond
=
NULL
;
SNode
*
pTempTagEqCond
=
NULL
;
SNode
*
pTempColEqCond
=
NULL
;
SNode
*
pTempColEqCond
=
NULL
;
SNode
*
pTempTagOnCond
=
NULL
;
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
nodesMergeConds
(
&
pTempColEqCond
,
&
pColEqOnConds
);
code
=
nodesMergeConds
(
&
pTempColEqCond
,
&
pColEqOnConds
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
nodesMergeConds
(
&
pTempTagEqCond
,
&
pTagEqOnConds
);
code
=
nodesMergeConds
(
&
pTempTagEqCond
,
&
pTagEqOnConds
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
nodesMergeConds
(
&
pTempTagOnCond
,
&
pTagOnConds
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
pJoin
->
pColEqCond
=
pTempColEqCond
;
pJoin
->
pColEqCond
=
pTempColEqCond
;
pJoin
->
pTagEqCond
=
pTempTagEqCond
;
pJoin
->
pTagEqCond
=
pTempTagEqCond
;
pJoin
->
pTagOnCond
=
pTempTagOnCond
;
return
TSDB_CODE_SUCCESS
;
return
TSDB_CODE_SUCCESS
;
}
else
{
}
else
{
nodesDestroyList
(
pColEqOnConds
);
nodesDestroyList
(
pColEqOnConds
);
...
@@ -846,15 +861,62 @@ static int32_t pushDownCondOptJoinExtractEqualOnCond(SOptimizeContext* pCxt, SJo
...
@@ -846,15 +861,62 @@ static int32_t pushDownCondOptJoinExtractEqualOnCond(SOptimizeContext* pCxt, SJo
}
else
{
}
else
{
pJoin
->
pColEqCond
=
nodesCloneNode
(
pJoin
->
pOtherOnCond
);
pJoin
->
pColEqCond
=
nodesCloneNode
(
pJoin
->
pOtherOnCond
);
}
}
}
else
if
(
allTags
)
{
pJoin
->
pTagOnCond
=
nodesCloneNode
(
pJoin
->
pOtherOnCond
);
}
}
return
TSDB_CODE_SUCCESS
;
return
TSDB_CODE_SUCCESS
;
}
}
static
int32_t
pushDownCondOptAppendFilterCol
(
SOptimizeContext
*
pCxt
,
SJoinLogicNode
*
pJoin
)
{
if
(
NULL
==
pJoin
->
pOtherOnCond
)
{
return
TSDB_CODE_SUCCESS
;
}
int32_t
code
=
TSDB_CODE_SUCCESS
;
SNodeList
*
pCondCols
=
nodesMakeList
();
SNodeList
*
pTargets
=
NULL
;
if
(
NULL
==
pCondCols
)
{
code
=
TSDB_CODE_OUT_OF_MEMORY
;
}
else
{
code
=
nodesCollectColumnsFromNode
(
pJoin
->
pOtherOnCond
,
NULL
,
COLLECT_COL_TYPE_ALL
,
&
pCondCols
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
createColumnByRewriteExprs
(
pCondCols
,
&
pTargets
);
}
nodesDestroyList
(
pCondCols
);
if
(
TSDB_CODE_SUCCESS
==
code
)
{
SNode
*
pNode
=
NULL
;
FOREACH
(
pNode
,
pTargets
)
{
SNode
*
pTmp
=
NULL
;
bool
found
=
false
;
FOREACH
(
pTmp
,
pJoin
->
node
.
pTargets
)
{
if
(
nodesEqualNode
(
pTmp
,
pNode
))
{
found
=
true
;
break
;
}
}
if
(
!
found
)
{
nodesListStrictAppend
(
pJoin
->
node
.
pTargets
,
pNode
);
}
}
}
nodesDestroyList
(
pTargets
);
return
code
;
}
static
int32_t
pushDownCondOptDealJoin
(
SOptimizeContext
*
pCxt
,
SJoinLogicNode
*
pJoin
)
{
static
int32_t
pushDownCondOptDealJoin
(
SOptimizeContext
*
pCxt
,
SJoinLogicNode
*
pJoin
)
{
if
(
OPTIMIZE_FLAG_TEST_MASK
(
pJoin
->
node
.
optimizedFlag
,
OPTIMIZE_FLAG_PUSH_DOWN_CONDE
))
{
if
(
OPTIMIZE_FLAG_TEST_MASK
(
pJoin
->
node
.
optimizedFlag
,
OPTIMIZE_FLAG_PUSH_DOWN_CONDE
))
{
return
TSDB_CODE_SUCCESS
;
return
TSDB_CODE_SUCCESS
;
}
}
if
(
pJoin
->
joinAlgo
!=
JOIN_ALGO_UNKNOWN
)
{
return
TSDB_CODE_SUCCESS
;
}
if
(
NULL
==
pJoin
->
node
.
pConditions
)
{
if
(
NULL
==
pJoin
->
node
.
pConditions
)
{
int32_t
code
=
pushDownCondOptJoinExtractCond
(
pCxt
,
pJoin
);
int32_t
code
=
pushDownCondOptJoinExtractCond
(
pCxt
,
pJoin
);
...
@@ -889,6 +951,10 @@ static int32_t pushDownCondOptDealJoin(SOptimizeContext* pCxt, SJoinLogicNode* p
...
@@ -889,6 +951,10 @@ static int32_t pushDownCondOptDealJoin(SOptimizeContext* pCxt, SJoinLogicNode* p
code
=
pushDownCondOptJoinExtractEqualOnCond
(
pCxt
,
pJoin
);
code
=
pushDownCondOptJoinExtractEqualOnCond
(
pCxt
,
pJoin
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
pushDownCondOptAppendFilterCol
(
pCxt
,
pJoin
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
OPTIMIZE_FLAG_SET_MASK
(
pJoin
->
node
.
optimizedFlag
,
OPTIMIZE_FLAG_PUSH_DOWN_CONDE
);
OPTIMIZE_FLAG_SET_MASK
(
pJoin
->
node
.
optimizedFlag
,
OPTIMIZE_FLAG_PUSH_DOWN_CONDE
);
pCxt
->
optimized
=
true
;
pCxt
->
optimized
=
true
;
...
@@ -2971,13 +3037,23 @@ static bool stbJoinOptShouldBeOptimized(SLogicNode* pNode) {
...
@@ -2971,13 +3037,23 @@ static bool stbJoinOptShouldBeOptimized(SLogicNode* pNode) {
}
}
SJoinLogicNode
*
pJoin
=
(
SJoinLogicNode
*
)
pNode
;
SJoinLogicNode
*
pJoin
=
(
SJoinLogicNode
*
)
pNode
;
if
(
pJoin
->
isSingleTableJoin
||
NULL
==
pJoin
->
pTagEqCond
||
pNode
->
pChildren
->
length
!=
2
||
pJoin
->
hasSubQuery
||
pJoin
->
joinAlgo
!=
UNKNOWN_JOIN_ALGO
)
{
if
(
pJoin
->
isSingleTableJoin
||
NULL
==
pJoin
->
pTagEqCond
||
pNode
->
pChildren
->
length
!=
2
||
pJoin
->
hasSubQuery
||
pJoin
->
joinAlgo
!=
JOIN_ALGO_UNKNOWN
)
{
return
false
;
return
false
;
}
}
return
true
;
return
true
;
}
}
int32_t
stbJoinOptAddFuncToScanNode
(
char
*
funcName
,
SScanLogicNode
*
pScan
)
{
SFunctionNode
*
pUidFunc
=
createFunction
(
funcName
,
NULL
);
snprintf
(
pUidFunc
->
node
.
aliasName
,
sizeof
(
pUidFunc
->
node
.
aliasName
),
"%s.%p"
,
pUidFunc
->
functionName
,
pUidFunc
);
nodesListStrictAppend
(
pScan
->
pScanPseudoCols
,
(
SNode
*
)
pUidFunc
);
return
createColumnByRewriteExpr
((
SNode
*
)
pUidFunc
,
&
pScan
->
node
.
pTargets
);
}
int32_t
stbJoinOptRewriteToTagScan
(
SLogicNode
*
pJoin
,
SNode
*
pNode
)
{
int32_t
stbJoinOptRewriteToTagScan
(
SLogicNode
*
pJoin
,
SNode
*
pNode
)
{
SScanLogicNode
*
pScan
=
(
SScanLogicNode
*
)
pNode
;
SScanLogicNode
*
pScan
=
(
SScanLogicNode
*
)
pNode
;
SJoinLogicNode
*
pJoinNode
=
(
SJoinLogicNode
*
)
pJoin
;
SJoinLogicNode
*
pJoinNode
=
(
SJoinLogicNode
*
)
pJoin
;
...
@@ -2990,6 +3066,9 @@ int32_t stbJoinOptRewriteToTagScan(SLogicNode* pJoin, SNode* pNode) {
...
@@ -2990,6 +3066,9 @@ int32_t stbJoinOptRewriteToTagScan(SLogicNode* pJoin, SNode* pNode) {
SNodeList
*
pTags
=
nodesMakeList
();
SNodeList
*
pTags
=
nodesMakeList
();
int32_t
code
=
nodesCollectColumnsFromNode
(
pJoinNode
->
pTagEqCond
,
NULL
,
COLLECT_COL_TYPE_TAG
,
&
pTags
);
int32_t
code
=
nodesCollectColumnsFromNode
(
pJoinNode
->
pTagEqCond
,
NULL
,
COLLECT_COL_TYPE_TAG
,
&
pTags
);
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
nodesCollectColumnsFromNode
(
pJoinNode
->
pTagOnCond
,
NULL
,
COLLECT_COL_TYPE_TAG
,
&
pTags
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
SNode
*
pTarget
=
NULL
;
SNode
*
pTarget
=
NULL
;
SNode
*
pTag
=
NULL
;
SNode
*
pTag
=
NULL
;
...
@@ -3010,18 +3089,10 @@ int32_t stbJoinOptRewriteToTagScan(SLogicNode* pJoin, SNode* pNode) {
...
@@ -3010,18 +3089,10 @@ int32_t stbJoinOptRewriteToTagScan(SLogicNode* pJoin, SNode* pNode) {
}
}
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
SFunctionNode
*
pUidFunc
=
createFunction
(
"_tbuid"
,
NULL
);
code
=
stbJoinOptAddFuncToScanNode
(
"_tbuid"
,
pScan
);
snprintf
(
pUidFunc
->
node
.
aliasName
,
sizeof
(
pUidFunc
->
node
.
aliasName
),
"%s.%p"
,
pUidFunc
->
functionName
,
pUidFunc
);
nodesListStrictAppend
(
pScan
->
pScanPseudoCols
,
(
SNode
*
)
pUidFunc
);
code
=
createColumnByRewriteExpr
(
pUidFunc
,
&
pScan
->
node
.
pTargets
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
SFunctionNode
*
pVgidFunc
=
createFunction
(
"_vgid"
,
NULL
);
code
=
stbJoinOptAddFuncToScanNode
(
"_vgid"
,
pScan
);
snprintf
(
pVgidFunc
->
node
.
aliasName
,
sizeof
(
pVgidFunc
->
node
.
aliasName
),
"%s.%p"
,
pVgidFunc
->
functionName
,
pVgidFunc
);
nodesListStrictAppend
(
pScan
->
pScanPseudoCols
,
(
SNode
*
)
pVgidFunc
);
code
=
createColumnByRewriteExpr
(
pVgidFunc
,
&
pScan
->
node
.
pTargets
);
}
}
if
(
code
)
{
if
(
code
)
{
...
@@ -3071,6 +3142,7 @@ static int32_t stbJoinOptCreateTagHashJoinNode(SLogicNode* pOrig, SNodeList* pCh
...
@@ -3071,6 +3142,7 @@ static int32_t stbJoinOptCreateTagHashJoinNode(SLogicNode* pOrig, SNodeList* pCh
pJoin
->
node
.
requireDataOrder
=
DATA_ORDER_LEVEL_NONE
;
pJoin
->
node
.
requireDataOrder
=
DATA_ORDER_LEVEL_NONE
;
pJoin
->
node
.
resultDataOrder
=
DATA_ORDER_LEVEL_NONE
;
pJoin
->
node
.
resultDataOrder
=
DATA_ORDER_LEVEL_NONE
;
pJoin
->
pTagEqCond
=
nodesCloneNode
(
pOrigJoin
->
pTagEqCond
);
pJoin
->
pTagEqCond
=
nodesCloneNode
(
pOrigJoin
->
pTagEqCond
);
pJoin
->
pTagOnCond
=
nodesCloneNode
(
pOrigJoin
->
pTagOnCond
);
int32_t
code
=
TSDB_CODE_SUCCESS
;
int32_t
code
=
TSDB_CODE_SUCCESS
;
pJoin
->
node
.
pChildren
=
pChildren
;
pJoin
->
node
.
pChildren
=
pChildren
;
...
@@ -3090,6 +3162,7 @@ static int32_t stbJoinOptCreateTagHashJoinNode(SLogicNode* pOrig, SNodeList* pCh
...
@@ -3090,6 +3162,7 @@ static int32_t stbJoinOptCreateTagHashJoinNode(SLogicNode* pOrig, SNodeList* pCh
if
(
code
)
{
if
(
code
)
{
break
;
break
;
}
}
pScan
->
node
.
pParent
=
(
SLogicNode
*
)
pJoin
;
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
...
@@ -3107,27 +3180,36 @@ static int32_t stbJoinOptCreateTableScanNodes(SLogicNode* pJoin, SNodeList** ppL
...
@@ -3107,27 +3180,36 @@ static int32_t stbJoinOptCreateTableScanNodes(SLogicNode* pJoin, SNodeList** ppL
return
TSDB_CODE_OUT_OF_MEMORY
;
return
TSDB_CODE_OUT_OF_MEMORY
;
}
}
int32_t
code
=
TSDB_CODE_SUCCESS
;
SNode
*
pNode
=
NULL
;
SNode
*
pNode
=
NULL
;
FOREACH
(
pNode
,
pList
)
{
FOREACH
(
pNode
,
pList
)
{
code
=
stbJoinOptAddUidToScan
(
pJoin
,
pNode
);
SScanLogicNode
*
pScan
=
(
SScanLogicNode
*
)
pNode
;
code
=
stbJoinOptAddFuncToScanNode
(
"_tbuid"
,
pScan
);
if
(
code
)
{
if
(
code
)
{
break
;
break
;
}
}
pScan
->
node
.
dynamicOp
=
true
;
}
}
*
ppList
=
pList
;
if
(
TSDB_CODE_SUCCESS
==
code
)
{
*
ppList
=
pList
;
}
else
{
nodesDestroyList
(
pList
);
*
ppList
=
NULL
;
}
return
TSDB_CODE_SUCCESS
;
return
code
;
}
}
static
int32_t
stbJoinOptCreateGroupCacheNode
(
S
LogicNode
*
pOrig
,
S
NodeList
*
pChildren
,
SLogicNode
**
ppLogic
)
{
static
int32_t
stbJoinOptCreateGroupCacheNode
(
SNodeList
*
pChildren
,
SLogicNode
**
ppLogic
)
{
SJoinLogicNode
*
pOrigJoin
=
(
SJoinLogicNode
*
)
pOrig
;
int32_t
code
=
TSDB_CODE_SUCCESS
;
SGroupCacheLogicNode
*
pGrpCache
=
(
S
Join
LogicNode
*
)
nodesMakeNode
(
QUERY_NODE_LOGIC_PLAN_GROUP_CACHE
);
SGroupCacheLogicNode
*
pGrpCache
=
(
S
GroupCache
LogicNode
*
)
nodesMakeNode
(
QUERY_NODE_LOGIC_PLAN_GROUP_CACHE
);
if
(
NULL
==
pGrpCache
)
{
if
(
NULL
==
pGrpCache
)
{
return
TSDB_CODE_OUT_OF_MEMORY
;
return
TSDB_CODE_OUT_OF_MEMORY
;
}
}
int32_t
code
=
TSDB_CODE_SUCCESS
;
pGrpCache
->
node
.
dynamicOp
=
true
;
pGrpCache
->
node
.
pChildren
=
pChildren
;
pGrpCache
->
node
.
pTargets
=
nodesMakeList
();
pGrpCache
->
node
.
pTargets
=
nodesMakeList
();
if
(
NULL
==
pGrpCache
->
node
.
pTargets
)
{
if
(
NULL
==
pGrpCache
->
node
.
pTargets
)
{
code
=
TSDB_CODE_OUT_OF_MEMORY
;
code
=
TSDB_CODE_OUT_OF_MEMORY
;
...
@@ -3136,7 +3218,24 @@ static int32_t stbJoinOptCreateGroupCacheNode(SLogicNode* pOrig, SNodeList* pChi
...
@@ -3136,7 +3218,24 @@ static int32_t stbJoinOptCreateGroupCacheNode(SLogicNode* pOrig, SNodeList* pChi
SScanLogicNode
*
pScan
=
(
SScanLogicNode
*
)
nodesListGetNode
(
pChildren
,
0
);
SScanLogicNode
*
pScan
=
(
SScanLogicNode
*
)
nodesListGetNode
(
pChildren
,
0
);
code
=
nodesListStrictAppendList
(
pGrpCache
->
node
.
pTargets
,
nodesCloneList
(
pScan
->
node
.
pTargets
));
code
=
nodesListStrictAppendList
(
pGrpCache
->
node
.
pTargets
,
nodesCloneList
(
pScan
->
node
.
pTargets
));
}
}
SScanLogicNode
*
pScan
=
(
SScanLogicNode
*
)
nodesListGetNode
(
pChildren
,
0
);
SNode
*
pCol
=
NULL
;
FOREACH
(
pCol
,
pScan
->
pScanPseudoCols
)
{
if
(
QUERY_NODE_FUNCTION
==
nodeType
(
pCol
)
&&
(((
SFunctionNode
*
)
pCol
)
->
funcType
==
FUNCTION_TYPE_TBUID
||
((
SFunctionNode
*
)
pCol
)
->
funcType
==
FUNCTION_TYPE_VGID
))
{
code
=
createColumnByRewriteExpr
(
pCol
,
&
pGrpCache
->
pGroupCols
);
if
(
code
)
{
break
;
}
}
}
SNode
*
pNode
=
NULL
;
FOREACH
(
pNode
,
pChildren
)
{
SScanLogicNode
*
pScan
=
(
SScanLogicNode
*
)
pNode
;
pScan
->
node
.
pParent
=
(
SLogicNode
*
)
pGrpCache
;
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
*
ppLogic
=
(
SLogicNode
*
)
pGrpCache
;
*
ppLogic
=
(
SLogicNode
*
)
pGrpCache
;
}
else
{
}
else
{
...
@@ -3146,20 +3245,95 @@ static int32_t stbJoinOptCreateGroupCacheNode(SLogicNode* pOrig, SNodeList* pChi
...
@@ -3146,20 +3245,95 @@ static int32_t stbJoinOptCreateGroupCacheNode(SLogicNode* pOrig, SNodeList* pChi
return
code
;
return
code
;
}
}
static
void
stbJoinOptRemoveTagEqCond
(
SJoinLogicNode
*
pJoin
)
{
if
(
QUERY_NODE_OPERATOR
==
nodeType
(
pJoin
->
pOtherOnCond
)
&&
nodesEqualNode
(
pJoin
->
pOtherOnCond
,
pJoin
->
pTagEqCond
))
{
NODES_DESTORY_NODE
(
pJoin
->
pOtherOnCond
);
return
;
}
if
(
QUERY_NODE_LOGIC_CONDITION
==
nodeType
(
pJoin
->
pOtherOnCond
))
{
SLogicConditionNode
*
pLogic
=
(
SLogicConditionNode
*
)
pJoin
->
pOtherOnCond
;
SNode
*
pNode
=
NULL
;
FOREACH
(
pNode
,
pLogic
->
pParameterList
)
{
if
(
nodesEqualNode
(
pNode
,
pJoin
->
pTagEqCond
))
{
ERASE_NODE
(
pLogic
->
pParameterList
);
break
;
}
else
if
(
QUERY_NODE_LOGIC_CONDITION
==
nodeType
(
pJoin
->
pTagEqCond
))
{
SLogicConditionNode
*
pTags
=
(
SLogicConditionNode
*
)
pJoin
->
pTagEqCond
;
SNode
*
pTag
=
NULL
;
FOREACH
(
pTag
,
pTags
->
pParameterList
)
{
if
(
nodesEqualNode
(
pTag
,
pNode
))
{
ERASE_NODE
(
pLogic
->
pParameterList
);
break
;
}
}
}
}
if
(
pLogic
->
pParameterList
->
length
<=
0
)
{
NODES_DESTORY_NODE
(
pJoin
->
pOtherOnCond
);
}
}
}
static
int32_t
stbJoinOptCreateMergeJoinNode
(
SLogicNode
*
pOrig
,
SLogicNode
*
pChild
,
SLogicNode
**
ppLogic
)
{
SJoinLogicNode
*
pOrigJoin
=
(
SJoinLogicNode
*
)
pOrig
;
SJoinLogicNode
*
pJoin
=
(
SJoinLogicNode
*
)
nodesCloneNode
((
SNode
*
)
pOrig
);
if
(
NULL
==
pJoin
)
{
return
TSDB_CODE_OUT_OF_MEMORY
;
}
pJoin
->
joinAlgo
=
JOIN_ALGO_MERGE
;
pJoin
->
node
.
dynamicOp
=
true
;
stbJoinOptRemoveTagEqCond
(
pJoin
);
NODES_DESTORY_NODE
(
pJoin
->
pTagEqCond
);
SNode
*
pNode
=
NULL
;
FOREACH
(
pNode
,
pJoin
->
node
.
pChildren
)
{
ERASE_NODE
(
pJoin
->
node
.
pChildren
);
}
nodesListStrictAppend
(
pJoin
->
node
.
pChildren
,
(
SNode
*
)
pChild
);
pChild
->
pParent
=
(
SLogicNode
*
)
pJoin
;
*
ppLogic
=
(
SLogicNode
*
)
pJoin
;
return
TSDB_CODE_SUCCESS
;
}
static
int32_t
stbJoinOptCreateDyn
TaskCtrlNode
(
SLogicNode
*
pJoin
,
SLogicNode
*
pHJoinNode
,
SLogicNode
*
pMJoinNode
,
SLogicNode
**
ppDynNode
)
{
static
int32_t
stbJoinOptCreateDyn
QueryCtrlNode
(
SLogicNode
*
pPrev
,
SLogicNode
*
pPost
,
SLogicNode
**
ppDynNode
)
{
int32_t
code
=
TSDB_CODE_SUCCESS
;
int32_t
code
=
TSDB_CODE_SUCCESS
;
SDynQueryCtrlLogicNode
*
pDynCtrl
=
(
SDynQueryCtrlLogicNode
*
)
nodesMakeNode
(
QUERY_NODE_LOGIC_PLAN_DYN_QUERY_CTRL
);
if
(
NULL
==
pDynCtrl
)
{
return
TSDB_CODE_OUT_OF_MEMORY
;
}
pDynCtrl
->
qType
=
DYN_QTYPE_STB_HASH
;
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
pDyn
Node
->
pChildren
=
nodesMakeList
();
pDyn
Ctrl
->
node
.
pChildren
=
nodesMakeList
();
if
(
NULL
==
pDyn
Node
->
pChildren
)
{
if
(
NULL
==
pDyn
Ctrl
->
node
.
pChildren
)
{
code
=
TSDB_CODE_OUT_OF_MEMORY
;
code
=
TSDB_CODE_OUT_OF_MEMORY
;
}
}
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
nodesListStrictAppend
(
pDynNode
->
pChildren
,
(
SNode
*
)
pHJoinNode
);
nodesListStrictAppend
(
pDynCtrl
->
node
.
pChildren
,
(
SNode
*
)
pPrev
);
nodesListStrictAppend
(
pDynCtrl
->
node
.
pChildren
,
(
SNode
*
)
pPost
);
pDynCtrl
->
node
.
pTargets
=
nodesCloneList
(
pPost
->
pTargets
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
pPrev
->
pParent
=
(
SLogicNode
*
)
pDynCtrl
;
pPost
->
pParent
=
(
SLogicNode
*
)
pDynCtrl
;
*
ppDynNode
=
(
SLogicNode
*
)
pDynCtrl
;
}
else
{
nodesDestroyNode
((
SNode
*
)
pDynCtrl
);
*
ppDynNode
=
NULL
;
}
}
return
code
;
}
}
static
int32_t
stbJoinOptRewriteStableJoin
(
SOptimizeContext
*
pCxt
,
SLogicNode
*
pJoin
,
SLogicSubplan
*
pLogicSubplan
)
{
static
int32_t
stbJoinOptRewriteStableJoin
(
SOptimizeContext
*
pCxt
,
SLogicNode
*
pJoin
,
SLogicSubplan
*
pLogicSubplan
)
{
...
@@ -3174,16 +3348,16 @@ static int32_t stbJoinOptRewriteStableJoin(SOptimizeContext* pCxt, SLogicNode* p
...
@@ -3174,16 +3348,16 @@ static int32_t stbJoinOptRewriteStableJoin(SOptimizeContext* pCxt, SLogicNode* p
code
=
stbJoinOptCreateTagHashJoinNode
(
pJoin
,
pTagScanNodes
,
&
pHJoinNode
);
code
=
stbJoinOptCreateTagHashJoinNode
(
pJoin
,
pTagScanNodes
,
&
pHJoinNode
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
stbJoinOptCreateTableScanNodes
(
pJoin
,
pTbScanNodes
);
code
=
stbJoinOptCreateTableScanNodes
(
pJoin
,
&
pTbScanNodes
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
stbJoinOptCreateGroupCacheNode
(
p
Join
,
p
TbScanNodes
,
&
pGrpCacheNode
);
code
=
stbJoinOptCreateGroupCacheNode
(
pTbScanNodes
,
&
pGrpCacheNode
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
stbJoinOptCreateMergeJoinNode
(
pJoin
,
pGrpCacheNode
,
&
pMJoinNode
);
code
=
stbJoinOptCreateMergeJoinNode
(
pJoin
,
pGrpCacheNode
,
&
pMJoinNode
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
stbJoinOptCreateDyn
TaskCtrlNode
(
pJoin
,
pHJoinNode
,
pMJoinNode
,
&
pDynNode
);
code
=
stbJoinOptCreateDyn
QueryCtrlNode
(
pHJoinNode
,
pMJoinNode
,
&
pDynNode
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
replaceLogicNode
(
pLogicSubplan
,
pJoin
,
(
SLogicNode
*
)
pDynNode
);
code
=
replaceLogicNode
(
pLogicSubplan
,
pJoin
,
(
SLogicNode
*
)
pDynNode
);
...
@@ -3215,14 +3389,14 @@ static const SOptimizeRule optimizeRuleSet[] = {
...
@@ -3215,14 +3389,14 @@ static const SOptimizeRule optimizeRuleSet[] = {
{.
pName
=
"PartitionTags"
,
.
optimizeFunc
=
partTagsOptimize
},
{.
pName
=
"PartitionTags"
,
.
optimizeFunc
=
partTagsOptimize
},
{.
pName
=
"StableJoin"
,
.
optimizeFunc
=
stableJoinOptimize
},
{.
pName
=
"StableJoin"
,
.
optimizeFunc
=
stableJoinOptimize
},
{.
pName
=
"MergeProjects"
,
.
optimizeFunc
=
mergeProjectsOptimize
},
{.
pName
=
"MergeProjects"
,
.
optimizeFunc
=
mergeProjectsOptimize
},
{.
pName
=
"EliminateProject"
,
.
optimizeFunc
=
eliminateProjOptimize
},
{.
pName
=
"EliminateSetOperator"
,
.
optimizeFunc
=
eliminateSetOpOptimize
},
{.
pName
=
"RewriteTail"
,
.
optimizeFunc
=
rewriteTailOptimize
},
{.
pName
=
"RewriteTail"
,
.
optimizeFunc
=
rewriteTailOptimize
},
{.
pName
=
"RewriteUnique"
,
.
optimizeFunc
=
rewriteUniqueOptimize
},
{.
pName
=
"RewriteUnique"
,
.
optimizeFunc
=
rewriteUniqueOptimize
},
{.
pName
=
"LastRowScan"
,
.
optimizeFunc
=
lastRowScanOptimize
},
{.
pName
=
"LastRowScan"
,
.
optimizeFunc
=
lastRowScanOptimize
},
{.
pName
=
"TagScan"
,
.
optimizeFunc
=
tagScanOptimize
},
{.
pName
=
"TagScan"
,
.
optimizeFunc
=
tagScanOptimize
},
{.
pName
=
"PushDownLimit"
,
.
optimizeFunc
=
pushDownLimitOptimize
},
{.
pName
=
"PushDownLimit"
,
.
optimizeFunc
=
pushDownLimitOptimize
},
{.
pName
=
"TableCountScan"
,
.
optimizeFunc
=
tableCountScanOptimize
},
{.
pName
=
"TableCountScan"
,
.
optimizeFunc
=
tableCountScanOptimize
},
{.
pName
=
"EliminateProject"
,
.
optimizeFunc
=
eliminateProjOptimize
},
{.
pName
=
"EliminateSetOperator"
,
.
optimizeFunc
=
eliminateSetOpOptimize
},
};
};
// clang-format on
// clang-format on
...
@@ -3257,6 +3431,7 @@ static int32_t applyOptimizeRule(SPlanContext* pCxt, SLogicSubplan* pLogicSubpla
...
@@ -3257,6 +3431,7 @@ static int32_t applyOptimizeRule(SPlanContext* pCxt, SLogicSubplan* pLogicSubpla
if
(
cxt
.
optimized
)
{
if
(
cxt
.
optimized
)
{
optimized
=
true
;
optimized
=
true
;
dumpLogicSubplan
(
optimizeRuleSet
[
i
].
pName
,
pLogicSubplan
);
dumpLogicSubplan
(
optimizeRuleSet
[
i
].
pName
,
pLogicSubplan
);
break
;
}
}
}
}
}
while
(
optimized
);
}
while
(
optimized
);
...
...
source/libs/planner/src/planPhysiCreater.c
浏览文件 @
81a2bf10
...
@@ -708,7 +708,8 @@ static int32_t mergeEqCond(SNode** ppDst, SNode** ppSrc) {
...
@@ -708,7 +708,8 @@ static int32_t mergeEqCond(SNode** ppDst, SNode** ppSrc) {
return
TSDB_CODE_SUCCESS
;
return
TSDB_CODE_SUCCESS
;
}
}
static
int32_t
createJoinPhysiNode
(
SPhysiPlanContext
*
pCxt
,
SNodeList
*
pChildren
,
SJoinLogicNode
*
pJoinLogicNode
,
static
int32_t
createMergeJoinPhysiNode
(
SPhysiPlanContext
*
pCxt
,
SNodeList
*
pChildren
,
SJoinLogicNode
*
pJoinLogicNode
,
SPhysiNode
**
pPhyNode
)
{
SPhysiNode
**
pPhyNode
)
{
SSortMergeJoinPhysiNode
*
pJoin
=
SSortMergeJoinPhysiNode
*
pJoin
=
(
SSortMergeJoinPhysiNode
*
)
makePhysiNode
(
pCxt
,
(
SLogicNode
*
)
pJoinLogicNode
,
QUERY_NODE_PHYSICAL_PLAN_MERGE_JOIN
);
(
SSortMergeJoinPhysiNode
*
)
makePhysiNode
(
pCxt
,
(
SLogicNode
*
)
pJoinLogicNode
,
QUERY_NODE_PHYSICAL_PLAN_MERGE_JOIN
);
...
@@ -730,32 +731,6 @@ static int32_t createJoinPhysiNode(SPhysiPlanContext* pCxt, SNodeList* pChildren
...
@@ -730,32 +731,6 @@ static int32_t createJoinPhysiNode(SPhysiPlanContext* pCxt, SNodeList* pChildren
&
pJoin
->
pTargets
);
&
pJoin
->
pTargets
);
}
}
if
(
TSDB_CODE_SUCCESS
==
code
&&
NULL
!=
pJoinLogicNode
->
pOtherOnCond
)
{
SNodeList
*
pCondCols
=
nodesMakeList
();
SNodeList
*
pTargets
=
NULL
;
SNodeList
*
pFinTargets
=
NULL
;
if
(
NULL
==
pCondCols
)
{
code
=
TSDB_CODE_OUT_OF_MEMORY
;
}
else
{
code
=
nodesCollectColumnsFromNode
(
pJoinLogicNode
->
pOtherOnCond
,
NULL
,
COLLECT_COL_TYPE_ALL
,
&
pCondCols
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
createColumnByRewriteExprs
(
pCondCols
,
&
pTargets
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
setListSlotId
(
pCxt
,
pLeftDesc
->
dataBlockId
,
pRightDesc
->
dataBlockId
,
pTargets
,
&
pFinTargets
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
nodesListStrictAppendList
(
pJoin
->
pTargets
,
pFinTargets
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
addDataBlockSlots
(
pCxt
,
pCondCols
,
pJoin
->
node
.
pOutputDataBlockDesc
);
}
nodesDestroyList
(
pTargets
);
nodesDestroyList
(
pCondCols
);
}
if
(
TSDB_CODE_SUCCESS
==
code
&&
NULL
!=
pJoinLogicNode
->
pOtherOnCond
)
{
if
(
TSDB_CODE_SUCCESS
==
code
&&
NULL
!=
pJoinLogicNode
->
pOtherOnCond
)
{
code
=
setNodeSlotId
(
pCxt
,
((
SPhysiNode
*
)
pJoin
)
->
pOutputDataBlockDesc
->
dataBlockId
,
-
1
,
code
=
setNodeSlotId
(
pCxt
,
((
SPhysiNode
*
)
pJoin
)
->
pOutputDataBlockDesc
->
dataBlockId
,
-
1
,
pJoinLogicNode
->
pOtherOnCond
,
&
pJoin
->
pOtherOnCond
);
pJoinLogicNode
->
pOtherOnCond
,
&
pJoin
->
pOtherOnCond
);
...
@@ -784,6 +759,72 @@ static int32_t createJoinPhysiNode(SPhysiPlanContext* pCxt, SNodeList* pChildren
...
@@ -784,6 +759,72 @@ static int32_t createJoinPhysiNode(SPhysiPlanContext* pCxt, SNodeList* pChildren
return
code
;
return
code
;
}
}
static
int32_t
createHashJoinPhysiNode
(
SPhysiPlanContext
*
pCxt
,
SNodeList
*
pChildren
,
SJoinLogicNode
*
pJoinLogicNode
,
SPhysiNode
**
pPhyNode
)
{
SHashJoinPhysiNode
*
pJoin
=
(
SHashJoinPhysiNode
*
)
makePhysiNode
(
pCxt
,
(
SLogicNode
*
)
pJoinLogicNode
,
QUERY_NODE_PHYSICAL_PLAN_HASH_JOIN
);
if
(
NULL
==
pJoin
)
{
return
TSDB_CODE_OUT_OF_MEMORY
;
}
SDataBlockDescNode
*
pLeftDesc
=
((
SPhysiNode
*
)
nodesListGetNode
(
pChildren
,
0
))
->
pOutputDataBlockDesc
;
SDataBlockDescNode
*
pRightDesc
=
((
SPhysiNode
*
)
nodesListGetNode
(
pChildren
,
1
))
->
pOutputDataBlockDesc
;
int32_t
code
=
TSDB_CODE_SUCCESS
;
pJoin
->
joinType
=
pJoinLogicNode
->
joinType
;
pJoin
->
node
.
inputTsOrder
=
pJoinLogicNode
->
node
.
inputTsOrder
;
SNode
*
pPrimKeyCond
=
NULL
;
SNode
*
pColEqCond
=
NULL
;
SNode
*
pTagEqCond
=
NULL
;
SNode
*
pTagOnCond
=
NULL
;
code
=
setNodeSlotId
(
pCxt
,
pLeftDesc
->
dataBlockId
,
pRightDesc
->
dataBlockId
,
pJoinLogicNode
->
pPrimKeyEqCond
,
&
pPrimKeyCond
);
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
setNodeSlotId
(
pCxt
,
pLeftDesc
->
dataBlockId
,
pRightDesc
->
dataBlockId
,
pJoinLogicNode
->
pColEqCond
,
&
pColEqCond
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
setNodeSlotId
(
pCxt
,
pLeftDesc
->
dataBlockId
,
pRightDesc
->
dataBlockId
,
pJoinLogicNode
->
pTagEqCond
,
&
pTagEqCond
);
}
if
(
TSDB_CODE_SUCCESS
==
code
&&
NULL
!=
pJoinLogicNode
->
pTagOnCond
)
{
code
=
setNodeSlotId
(
pCxt
,
((
SPhysiNode
*
)
pJoin
)
->
pOutputDataBlockDesc
->
dataBlockId
,
-
1
,
pJoinLogicNode
->
pTagOnCond
,
&
pTagOnCond
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
setListSlotId
(
pCxt
,
pLeftDesc
->
dataBlockId
,
pRightDesc
->
dataBlockId
,
pJoinLogicNode
->
node
.
pTargets
,
&
pJoin
->
pTargets
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
setConditionsSlotId
(
pCxt
,
(
const
SLogicNode
*
)
pJoinLogicNode
,
(
SPhysiNode
*
)
pJoin
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
addDataBlockSlots
(
pCxt
,
pJoin
->
pTargets
,
pJoin
->
node
.
pOutputDataBlockDesc
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
*
pPhyNode
=
(
SPhysiNode
*
)
pJoin
;
}
else
{
nodesDestroyNode
((
SNode
*
)
pJoin
);
}
return
code
;
}
static
int32_t
createJoinPhysiNode
(
SPhysiPlanContext
*
pCxt
,
SNodeList
*
pChildren
,
SJoinLogicNode
*
pJoinLogicNode
,
SPhysiNode
**
pPhyNode
)
{
switch
(
pJoinLogicNode
->
joinAlgo
)
{
case
JOIN_ALGO_MERGE
:
return
createMergeJoinPhysiNode
(
pCxt
,
pChildren
,
pJoinLogicNode
,
pPhyNode
);
case
JOIN_ALGO_HASH
:
return
createHashJoinPhysiNode
(
pCxt
,
pChildren
,
pJoinLogicNode
,
pPhyNode
);
default:
planError
(
"Invalid join algorithm:%d"
,
pJoinLogicNode
->
joinAlgo
);
break
;
}
return
TSDB_CODE_FAILED
;
}
typedef
struct
SRewritePrecalcExprsCxt
{
typedef
struct
SRewritePrecalcExprsCxt
{
int32_t
errCode
;
int32_t
errCode
;
int32_t
planNodeId
;
int32_t
planNodeId
;
...
...
source/libs/planner/src/planSpliter.c
浏览文件 @
81a2bf10
...
@@ -280,7 +280,7 @@ static bool stbSplNeedSplitWindow(bool streamQuery, SLogicNode* pNode) {
...
@@ -280,7 +280,7 @@ static bool stbSplNeedSplitWindow(bool streamQuery, SLogicNode* pNode) {
}
}
static
bool
stbSplNeedSplitJoin
(
bool
streamQuery
,
SJoinLogicNode
*
pJoin
)
{
static
bool
stbSplNeedSplitJoin
(
bool
streamQuery
,
SJoinLogicNode
*
pJoin
)
{
if
(
pJoin
->
isSingleTableJoin
)
{
if
(
pJoin
->
isSingleTableJoin
||
JOIN_ALGO_HASH
==
pJoin
->
joinAlgo
)
{
return
false
;
return
false
;
}
}
SNode
*
pChild
=
NULL
;
SNode
*
pChild
=
NULL
;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录