Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
fe4e4564
TDengine
项目概览
taosdata
/
TDengine
1 年多 前同步成功
通知
1185
Star
22016
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看板
提交
fe4e4564
编写于
8月 22, 2023
作者:
D
dapan1121
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix: global data sink manager issue
上级
b5bd8f7c
变更
5
隐藏空白更改
内联
并排
Showing
5 changed file
with
22 addition
and
10 deletion
+22
-10
source/libs/executor/src/dataDeleter.c
source/libs/executor/src/dataDeleter.c
+2
-0
source/libs/executor/src/dataDispatcher.c
source/libs/executor/src/dataDispatcher.c
+1
-0
source/libs/executor/src/dataInserter.c
source/libs/executor/src/dataInserter.c
+2
-0
source/libs/executor/src/dataSinkMgt.c
source/libs/executor/src/dataSinkMgt.c
+14
-8
source/libs/executor/src/executor.c
source/libs/executor/src/executor.c
+3
-2
未找到文件。
source/libs/executor/src/dataDeleter.c
浏览文件 @
fe4e4564
...
...
@@ -224,6 +224,8 @@ static int32_t destroyDataSinker(SDataSinkHandle* pHandle) {
}
taosCloseQueue
(
pDeleter
->
pDataBlocks
);
taosThreadMutexDestroy
(
&
pDeleter
->
mutex
);
taosMemoryFree
(
pDeleter
->
pManager
);
return
TSDB_CODE_SUCCESS
;
}
...
...
source/libs/executor/src/dataDispatcher.c
浏览文件 @
fe4e4564
...
...
@@ -226,6 +226,7 @@ static int32_t destroyDataSinker(SDataSinkHandle* pHandle) {
}
taosCloseQueue
(
pDispatcher
->
pDataBlocks
);
taosThreadMutexDestroy
(
&
pDispatcher
->
mutex
);
taosMemoryFree
(
pDispatcher
->
pManager
);
return
TSDB_CODE_SUCCESS
;
}
...
...
source/libs/executor/src/dataInserter.c
浏览文件 @
fe4e4564
...
...
@@ -395,6 +395,8 @@ static int32_t destroyDataSinker(SDataSinkHandle* pHandle) {
taosMemoryFree
(
pInserter
->
pParam
);
taosHashCleanup
(
pInserter
->
pCols
);
taosThreadMutexDestroy
(
&
pInserter
->
mutex
);
taosMemoryFree
(
pInserter
->
pManager
);
return
TSDB_CODE_SUCCESS
;
}
...
...
source/libs/executor/src/dataSinkMgt.c
浏览文件 @
fe4e4564
...
...
@@ -18,12 +18,17 @@
#include "planner.h"
#include "tarray.h"
static
SDataSinkManager
gDataSinkManager
=
{
0
};
SDataSinkStat
gDataSinkStat
=
{
0
};
int32_t
dsDataSinkMgtInit
(
SDataSinkMgtCfg
*
cfg
,
SStorageAPI
*
pAPI
)
{
gDataSinkManager
.
cfg
=
*
cfg
;
gDataSinkManager
.
pAPI
=
pAPI
;
int32_t
dsDataSinkMgtInit
(
SDataSinkMgtCfg
*
cfg
,
SStorageAPI
*
pAPI
,
void
**
ppSinkManager
)
{
SDataSinkManager
*
pSinkManager
=
taosMemoryMalloc
(
sizeof
(
SDataSinkManager
));
if
(
NULL
==
pSinkManager
)
{
return
TSDB_CODE_OUT_OF_MEMORY
;
}
pSinkManager
->
cfg
=
*
cfg
;
pSinkManager
->
pAPI
=
pAPI
;
*
ppSinkManager
=
pSinkManager
;
return
0
;
// to avoid compiler eror
}
...
...
@@ -33,15 +38,16 @@ int32_t dsDataSinkGetCacheSize(SDataSinkStat* pStat) {
return
0
;
}
int32_t
dsCreateDataSinker
(
const
SDataSinkNode
*
pDataSink
,
DataSinkHandle
*
pHandle
,
void
*
pParam
,
const
char
*
id
)
{
int32_t
dsCreateDataSinker
(
void
*
pSinkManager
,
const
SDataSinkNode
*
pDataSink
,
DataSinkHandle
*
pHandle
,
void
*
pParam
,
const
char
*
id
)
{
SDataSinkManager
*
pManager
=
pSinkManager
;
switch
((
int
)
nodeType
(
pDataSink
))
{
case
QUERY_NODE_PHYSICAL_PLAN_DISPATCH
:
return
createDataDispatcher
(
&
gDataSink
Manager
,
pDataSink
,
pHandle
);
return
createDataDispatcher
(
p
Manager
,
pDataSink
,
pHandle
);
case
QUERY_NODE_PHYSICAL_PLAN_DELETE
:
{
return
createDataDeleter
(
&
gDataSink
Manager
,
pDataSink
,
pHandle
,
pParam
);
return
createDataDeleter
(
p
Manager
,
pDataSink
,
pHandle
,
pParam
);
}
case
QUERY_NODE_PHYSICAL_PLAN_QUERY_INSERT
:
{
return
createDataInserter
(
&
gDataSink
Manager
,
pDataSink
,
pHandle
,
pParam
);
return
createDataInserter
(
p
Manager
,
pDataSink
,
pHandle
,
pParam
);
}
}
...
...
source/libs/executor/src/executor.c
浏览文件 @
fe4e4564
...
...
@@ -512,7 +512,8 @@ int32_t qCreateExecTask(SReadHandle* readHandle, int32_t vgId, uint64_t taskId,
}
SDataSinkMgtCfg
cfg
=
{.
maxDataBlockNum
=
500
,
.
maxDataBlockNumPerQuery
=
50
};
code
=
dsDataSinkMgtInit
(
&
cfg
,
&
(
*
pTask
)
->
storageAPI
);
void
*
pSinkManager
=
NULL
;
code
=
dsDataSinkMgtInit
(
&
cfg
,
&
(
*
pTask
)
->
storageAPI
,
&
pSinkManager
);
if
(
code
!=
TSDB_CODE_SUCCESS
)
{
qError
(
"failed to dsDataSinkMgtInit, code:%s, %s"
,
tstrerror
(
code
),
(
*
pTask
)
->
id
.
str
);
goto
_error
;
...
...
@@ -527,7 +528,7 @@ int32_t qCreateExecTask(SReadHandle* readHandle, int32_t vgId, uint64_t taskId,
}
// pSinkParam has been freed during create sinker.
code
=
dsCreateDataSinker
(
pSubplan
->
pDataSink
,
handle
,
pSinkParam
,
(
*
pTask
)
->
id
.
str
);
code
=
dsCreateDataSinker
(
pS
inkManager
,
pS
ubplan
->
pDataSink
,
handle
,
pSinkParam
,
(
*
pTask
)
->
id
.
str
);
}
qDebug
(
"subplan task create completed, TID:0x%"
PRIx64
" QID:0x%"
PRIx64
,
taskId
,
pSubplan
->
id
.
queryId
);
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录