Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
a1e18ac8
TDengine
项目概览
taosdata
/
TDengine
大约 2 年 前同步成功
通知
1192
Star
22018
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看板
提交
a1e18ac8
编写于
7月 17, 2023
作者:
D
dapan1121
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix: fix global cache issue
上级
cbc9696f
变更
8
隐藏空白更改
内联
并排
Showing
8 changed file
with
20 addition
and
59 deletion
+20
-59
include/libs/nodes/plannodes.h
include/libs/nodes/plannodes.h
+0
-2
source/libs/executor/inc/groupcache.h
source/libs/executor/inc/groupcache.h
+0
-1
source/libs/executor/src/groupcacheoperator.c
source/libs/executor/src/groupcacheoperator.c
+19
-31
source/libs/nodes/src/nodesCloneFuncs.c
source/libs/nodes/src/nodesCloneFuncs.c
+0
-1
source/libs/nodes/src/nodesCodeFuncs.c
source/libs/nodes/src/nodesCodeFuncs.c
+0
-14
source/libs/nodes/src/nodesMsgFuncs.c
source/libs/nodes/src/nodesMsgFuncs.c
+0
-7
source/libs/planner/src/planOptimizer.c
source/libs/planner/src/planOptimizer.c
+1
-2
source/libs/planner/src/planPhysiCreater.c
source/libs/planner/src/planPhysiCreater.c
+0
-1
未找到文件。
include/libs/nodes/plannodes.h
浏览文件 @
a1e18ac8
...
...
@@ -162,7 +162,6 @@ typedef struct SGroupCacheLogicNode {
bool
grpColsMayBeNull
;
bool
grpByUid
;
bool
globalGrp
;
bool
enableCache
;
SNodeList
*
pGroupCols
;
}
SGroupCacheLogicNode
;
...
...
@@ -446,7 +445,6 @@ typedef struct SGroupCachePhysiNode {
bool
grpColsMayBeNull
;
bool
grpByUid
;
bool
globalGrp
;
bool
enableCache
;
SNodeList
*
pGroupCols
;
}
SGroupCachePhysiNode
;
...
...
source/libs/executor/inc/groupcache.h
浏览文件 @
a1e18ac8
...
...
@@ -131,7 +131,6 @@ typedef struct SGroupCacheOperatorInfo {
SGroupColsInfo
groupColsInfo
;
bool
globalGrp
;
bool
grpByUid
;
bool
enableCache
;
SGcDownstreamCtx
*
pDownstreams
;
SGcBlkCacheInfo
blkCache
;
SHashObj
*
pGrpHash
;
...
...
source/libs/executor/src/groupcacheoperator.c
浏览文件 @
a1e18ac8
...
...
@@ -157,7 +157,7 @@ static int32_t acquireBaseBlockFromList(SGcDownstreamCtx* pCtx, SSDataBlock** pp
taosWUnLockLatch
(
&
pCtx
->
blkLock
);
return
buildGroupCacheBaseBlock
(
ppRes
,
pCtx
->
pBaseBlock
);
}
*
ppRes
=
taosArrayPop
(
pCtx
->
pFreeBlock
);
*
ppRes
=
*
(
SSDataBlock
**
)
taosArrayPop
(
pCtx
->
pFreeBlock
);
taosWUnLockLatch
(
&
pCtx
->
blkLock
);
return
TSDB_CODE_SUCCESS
;
...
...
@@ -270,11 +270,9 @@ static FORCE_INLINE int32_t getBlkFromDownstreamOperator(struct SOperatorInfo* p
SOperatorParam
*
pDownstreamParam
=
NULL
;
SSDataBlock
*
pBlock
=
NULL
;
SGroupCacheOperatorInfo
*
pGCache
=
pOperator
->
info
;
if
(
pGCache
->
enableCache
)
{
code
=
appendNewGroupToDownstream
(
pOperator
,
downstreamIdx
,
&
pDownstreamParam
);
if
(
code
)
{
return
code
;
}
code
=
appendNewGroupToDownstream
(
pOperator
,
downstreamIdx
,
&
pDownstreamParam
);
if
(
code
)
{
return
code
;
}
if
(
pDownstreamParam
)
{
...
...
@@ -283,7 +281,7 @@ static FORCE_INLINE int32_t getBlkFromDownstreamOperator(struct SOperatorInfo* p
pBlock
=
pOperator
->
pDownstream
[
downstreamIdx
]
->
fpSet
.
getNextFn
(
pOperator
->
pDownstream
[
downstreamIdx
]);
}
if
(
pBlock
&&
pGCache
->
enableCache
)
{
if
(
pBlock
)
{
pGCache
->
execInfo
.
pDownstreamBlkNum
[
downstreamIdx
]
++
;
if
(
NULL
==
pGCache
->
pDownstreams
[
downstreamIdx
].
pBaseBlock
)
{
code
=
buildGroupCacheBaseBlock
(
&
pGCache
->
pDownstreams
[
downstreamIdx
].
pBaseBlock
,
pBlock
);
...
...
@@ -389,7 +387,6 @@ static int32_t handleDownstreamFetchDone(struct SOperatorInfo* pOperator, SGcSes
SGroupCacheOperatorInfo
*
pGCache
=
pOperator
->
info
;
SGcDownstreamCtx
*
pCtx
=
&
pGCache
->
pDownstreams
[
pSession
->
downstreamIdx
];
int32_t
uidNum
=
0
;
SHashObj
*
pGrpHash
=
pGCache
->
globalGrp
?
pGCache
->
pGrpHash
:
pCtx
->
pGrpHash
;
SGcVgroupCtx
*
pVgCtx
=
NULL
;
int32_t
iter
=
0
;
while
(
pVgCtx
=
tSimpleHashIterate
(
pCtx
->
pVgTbHash
,
pVgCtx
,
&
iter
))
{
...
...
@@ -420,17 +417,11 @@ static int32_t getCacheBlkFromDownstreamOperator(struct SOperatorInfo* pOperator
if
(
NULL
==
*
ppRes
)
{
code
=
handleDownstreamFetchDone
(
pOperator
,
pSession
);
break
;
}
else
if
(
pGCache
->
enableCache
)
{
code
=
handleGroupCacheRetrievedBlk
(
pOperator
,
*
ppRes
,
pSession
,
&
continueFetch
);
}
else
{
co
ntinueFetch
=
false
;
co
de
=
handleGroupCacheRetrievedBlk
(
pOperator
,
*
ppRes
,
pSession
,
&
continueFetch
)
;
}
}
if
(
!
pGCache
->
enableCache
)
{
return
code
;
}
if
(
!
continueFetch
)
{
SGcSessionCtx
**
ppWaitCtx
=
taosHashIterate
(
pCtx
->
pWaitSessions
,
NULL
);
if
(
ppWaitCtx
)
{
...
...
@@ -514,22 +505,20 @@ static int32_t getBlkFromSessionCacheImpl(struct SOperatorInfo* pOperator, int64
SGcDownstreamCtx
*
pCtx
=
&
pGCache
->
pDownstreams
[
pSession
->
downstreamIdx
];
while
(
true
)
{
if
(
pGCache
->
enableCache
)
{
if
(
pSession
->
lastBlkId
<
0
)
{
int64_t
startBlkId
=
atomic_load_64
(
&
pSession
->
pGroupData
->
startBlkId
);
if
(
startBlkId
>
0
)
{
code
=
retrieveBlkFromBufCache
(
pGCache
,
pSession
->
pGroupData
,
sessionId
,
startBlkId
,
&
pSession
->
nextOffset
,
ppRes
);
pSession
->
lastBlkId
=
startBlkId
;
goto
_return
;
}
}
else
if
(
pSession
->
lastBlkId
<
atomic_load_64
(
&
pSession
->
pGroupData
->
endBlkId
))
{
code
=
retrieveBlkFromBufCache
(
pGCache
,
pSession
->
pGroupData
,
sessionId
,
pSession
->
lastBlkId
+
1
,
&
pSession
->
nextOffset
,
ppRes
);
pSession
->
lastBlkId
++
;
goto
_return
;
}
else
if
(
atomic_load_8
((
int8_t
*
)
&
pSession
->
pGroupData
->
fetchDone
))
{
*
ppRes
=
NULL
;
if
(
pSession
->
lastBlkId
<
0
)
{
int64_t
startBlkId
=
atomic_load_64
(
&
pSession
->
pGroupData
->
startBlkId
);
if
(
startBlkId
>
0
)
{
code
=
retrieveBlkFromBufCache
(
pGCache
,
pSession
->
pGroupData
,
sessionId
,
startBlkId
,
&
pSession
->
nextOffset
,
ppRes
);
pSession
->
lastBlkId
=
startBlkId
;
goto
_return
;
}
}
else
if
(
pSession
->
lastBlkId
<
atomic_load_64
(
&
pSession
->
pGroupData
->
endBlkId
))
{
code
=
retrieveBlkFromBufCache
(
pGCache
,
pSession
->
pGroupData
,
sessionId
,
pSession
->
lastBlkId
+
1
,
&
pSession
->
nextOffset
,
ppRes
);
pSession
->
lastBlkId
++
;
goto
_return
;
}
else
if
(
atomic_load_8
((
int8_t
*
)
&
pSession
->
pGroupData
->
fetchDone
))
{
*
ppRes
=
NULL
;
goto
_return
;
}
if
((
atomic_load_64
(
&
pCtx
->
fetchSessionId
)
==
sessionId
)
...
...
@@ -686,7 +675,7 @@ static int32_t getBlkFromGroupCache(struct SOperatorInfo* pOperator, SSDataBlock
if
(
TSDB_CODE_SUCCESS
!=
code
)
{
return
code
;
}
}
else
if
(
pGCache
->
enableCache
)
{
}
else
{
SSDataBlock
**
ppBlock
=
taosHashGet
(
pGCache
->
blkCache
.
pReadBlk
,
&
pGcParam
->
sessionId
,
sizeof
(
pGcParam
->
sessionId
));
if
(
ppBlock
)
{
releaseBaseBlockToList
(
pCtx
,
*
ppBlock
);
...
...
@@ -788,7 +777,6 @@ SOperatorInfo* createGroupCacheOperatorInfo(SOperatorInfo** pDownstream, int32_t
pInfo
->
maxCacheSize
=
-
1
;
pInfo
->
grpByUid
=
pPhyciNode
->
grpByUid
;
pInfo
->
globalGrp
=
pPhyciNode
->
globalGrp
;
pInfo
->
enableCache
=
pPhyciNode
->
enableCache
;
if
(
!
pInfo
->
grpByUid
)
{
qError
(
"only group cache by uid is supported now"
);
...
...
source/libs/nodes/src/nodesCloneFuncs.c
浏览文件 @
a1e18ac8
...
...
@@ -540,7 +540,6 @@ static int32_t logicGroupCacheCopy(const SGroupCacheLogicNode* pSrc, SGroupCache
COPY_SCALAR_FIELD
(
grpColsMayBeNull
);
COPY_SCALAR_FIELD
(
grpByUid
);
COPY_SCALAR_FIELD
(
globalGrp
);
COPY_SCALAR_FIELD
(
enableCache
);
CLONE_NODE_LIST_FIELD
(
pGroupCols
);
return
TSDB_CODE_SUCCESS
;
}
...
...
source/libs/nodes/src/nodesCodeFuncs.c
浏览文件 @
a1e18ac8
...
...
@@ -1185,7 +1185,6 @@ static int32_t jsonToLogicInterpFuncNode(const SJson* pJson, void* pObj) {
static
const
char
*
jkGroupCacheLogicPlanGrpColsMayBeNull
=
"GroupColsMayBeNull"
;
static
const
char
*
jkGroupCacheLogicPlanGroupByUid
=
"GroupByUid"
;
static
const
char
*
jkGroupCacheLogicPlanGlobalGroup
=
"GlobalGroup"
;
static
const
char
*
jkGroupCacheLogicPlanEnableCache
=
"EnableCache"
;
static
const
char
*
jkGroupCacheLogicPlanGroupCols
=
"GroupCols"
;
static
int32_t
logicGroupCacheNodeToJson
(
const
void
*
pObj
,
SJson
*
pJson
)
{
...
...
@@ -1201,9 +1200,6 @@ static int32_t logicGroupCacheNodeToJson(const void* pObj, SJson* pJson) {
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAddBoolToObject
(
pJson
,
jkGroupCacheLogicPlanGlobalGroup
,
pNode
->
globalGrp
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAddBoolToObject
(
pJson
,
jkGroupCacheLogicPlanEnableCache
,
pNode
->
enableCache
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
nodeListToJson
(
pJson
,
jkGroupCacheLogicPlanGroupCols
,
pNode
->
pGroupCols
);
}
...
...
@@ -1224,9 +1220,6 @@ static int32_t jsonToLogicGroupCacheNode(const SJson* pJson, void* pObj) {
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonGetBoolValue
(
pJson
,
jkGroupCacheLogicPlanGlobalGroup
,
&
pNode
->
globalGrp
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonGetBoolValue
(
pJson
,
jkGroupCacheLogicPlanEnableCache
,
&
pNode
->
enableCache
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
jsonToNodeList
(
pJson
,
jkGroupCacheLogicPlanGroupCols
,
&
pNode
->
pGroupCols
);
}
...
...
@@ -2936,7 +2929,6 @@ static const char* jkGroupCachePhysiPlanGroupCols = "GroupColumns";
static
const
char
*
jkGroupCachePhysiPlanGrpColsMayBeNull
=
"GroupColumnsMayBeNull"
;
static
const
char
*
jkGroupCachePhysiPlanGroupByUid
=
"GroupByUid"
;
static
const
char
*
jkGroupCachePhysiPlanGlobalGroup
=
"GlobalGroup"
;
static
const
char
*
jkGroupCachePhysiPlanEnableCache
=
"EnableCache"
;
static
int32_t
physiGroupCacheNodeToJson
(
const
void
*
pObj
,
SJson
*
pJson
)
{
...
...
@@ -2952,9 +2944,6 @@ static int32_t physiGroupCacheNodeToJson(const void* pObj, SJson* pJson) {
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAddBoolToObject
(
pJson
,
jkGroupCachePhysiPlanGlobalGroup
,
pNode
->
globalGrp
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonAddBoolToObject
(
pJson
,
jkGroupCachePhysiPlanEnableCache
,
pNode
->
enableCache
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
nodeListToJson
(
pJson
,
jkGroupCachePhysiPlanGroupCols
,
pNode
->
pGroupCols
);
}
...
...
@@ -2974,9 +2963,6 @@ static int32_t jsonToPhysiGroupCacheNode(const SJson* pJson, void* pObj) {
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonGetBoolValue
(
pJson
,
jkGroupCachePhysiPlanGlobalGroup
,
&
pNode
->
globalGrp
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tjsonGetBoolValue
(
pJson
,
jkGroupCachePhysiPlanEnableCache
,
&
pNode
->
enableCache
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
jsonToNodeList
(
pJson
,
jkGroupCachePhysiPlanGroupCols
,
&
pNode
->
pGroupCols
);
}
...
...
source/libs/nodes/src/nodesMsgFuncs.c
浏览文件 @
a1e18ac8
...
...
@@ -3530,7 +3530,6 @@ enum {
PHY_GROUP_CACHE_CODE_GROUP_COLS_MAY_BE_NULL
,
PHY_GROUP_CACHE_CODE_GROUP_BY_UID
,
PHY_GROUP_CACHE_CODE_GLOBAL_GROUP
,
PHY_GROUP_CACHE_CODE_ENABLE_CACHE
,
PHY_GROUP_CACHE_CODE_GROUP_COLUMNS
};
...
...
@@ -3550,9 +3549,6 @@ static int32_t physiGroupCacheNodeToMsg(const void* pObj, STlvEncoder* pEncoder)
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tlvEncodeBool
(
pEncoder
,
PHY_GROUP_CACHE_CODE_GLOBAL_GROUP
,
pNode
->
globalGrp
);
}
if
(
TSDB_CODE_SUCCESS
==
code
)
{
code
=
tlvEncodeBool
(
pEncoder
,
PHY_GROUP_CACHE_CODE_ENABLE_CACHE
,
pNode
->
enableCache
);
}
return
code
;
}
...
...
@@ -3579,9 +3575,6 @@ static int32_t msgToPhysiGroupCacheNode(STlvDecoder* pDecoder, void* pObj) {
case
PHY_GROUP_CACHE_CODE_GLOBAL_GROUP
:
code
=
tlvDecodeBool
(
pTlv
,
&
pNode
->
globalGrp
);
break
;
case
PHY_GROUP_CACHE_CODE_ENABLE_CACHE
:
code
=
tlvDecodeBool
(
pTlv
,
&
pNode
->
enableCache
);
break
;
default:
break
;
}
...
...
source/libs/planner/src/planOptimizer.c
浏览文件 @
a1e18ac8
...
...
@@ -3262,8 +3262,7 @@ static int32_t stbJoinOptCreateGroupCacheNode(SNodeList* pChildren, SLogicNode**
}
pScan
->
node
.
pParent
=
(
SLogicNode
*
)
pGrpCache
;
}
pGrpCache
->
globalGrp
=
!
hasCond
;
pGrpCache
->
enableCache
=
pGrpCache
->
globalGrp
;
pGrpCache
->
globalGrp
=
false
;
if
(
TSDB_CODE_SUCCESS
==
code
)
{
*
ppLogic
=
(
SLogicNode
*
)
pGrpCache
;
...
...
source/libs/planner/src/planPhysiCreater.c
浏览文件 @
a1e18ac8
...
...
@@ -988,7 +988,6 @@ static int32_t createGroupCachePhysiNode(SPhysiPlanContext* pCxt, SNodeList* pCh
pGrpCache
->
grpColsMayBeNull
=
pLogicNode
->
grpColsMayBeNull
;
pGrpCache
->
grpByUid
=
pLogicNode
->
grpByUid
;
pGrpCache
->
globalGrp
=
pLogicNode
->
globalGrp
;
pGrpCache
->
enableCache
=
pLogicNode
->
enableCache
;
SDataBlockDescNode
*
pChildDesc
=
((
SPhysiNode
*
)
nodesListGetNode
(
pChildren
,
0
))
->
pOutputDataBlockDesc
;
int32_t
code
=
TSDB_CODE_SUCCESS
;
if
(
TSDB_CODE_SUCCESS
==
code
)
{
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录