Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
1b263602
T
TDengine
项目概览
taosdata
/
TDengine
1 年多 前同步成功
通知
1185
Star
22016
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看板
提交
1b263602
编写于
7月 21, 2023
作者:
H
Haojun Liao
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix(stream): fix memory leak.
上级
16d7707b
变更
5
隐藏空白更改
内联
并排
Showing
5 changed file
with
10 addition
and
1 deletion
+10
-1
include/libs/stream/tstream.h
include/libs/stream/tstream.h
+1
-0
source/dnode/snode/src/snode.c
source/dnode/snode/src/snode.c
+1
-0
source/dnode/vnode/src/tq/tq.c
source/dnode/vnode/src/tq/tq.c
+1
-0
source/libs/stream/src/streamDispatch.c
source/libs/stream/src/streamDispatch.c
+2
-1
source/libs/stream/src/streamTask.c
source/libs/stream/src/streamTask.c
+5
-0
未找到文件。
include/libs/stream/tstream.h
浏览文件 @
1b263602
...
...
@@ -337,6 +337,7 @@ struct SStreamTask {
SMsgCb
*
pMsgCb
;
// msg handle
SStreamState
*
pState
;
// state backend
SArray
*
pRspMsgList
;
TdThreadMutex
lock
;
// the followings attributes don't be serialized
int32_t
notReadyTasks
;
...
...
source/dnode/snode/src/snode.c
浏览文件 @
1b263602
...
...
@@ -91,6 +91,7 @@ int32_t sndExpandTask(SSnode *pSnode, SStreamTask *pTask, int64_t ver) {
pTask
->
exec
.
pExecutor
=
qCreateStreamExecTaskInfo
(
pTask
->
exec
.
qmsg
,
&
handle
,
0
);
ASSERT
(
pTask
->
exec
.
pExecutor
);
taosThreadMutexInit
(
&
pTask
->
lock
,
NULL
);
streamSetupScheduleTrigger
(
pTask
);
qDebug
(
"snode:%d expand stream task on snode, s-task:%s, checkpoint ver:%"
PRId64
" child id:%d, level:%d"
,
SNODE_HANDLE
,
...
...
source/dnode/vnode/src/tq/tq.c
浏览文件 @
1b263602
...
...
@@ -921,6 +921,7 @@ int32_t tqExpandTask(STQ* pTq, SStreamTask* pTask, int64_t ver) {
pTask
->
status
.
taskStatus
=
TASK_STATUS__NORMAL
;
}
taosThreadMutexInit
(
&
pTask
->
lock
,
NULL
);
streamSetupScheduleTrigger
(
pTask
);
tqInfo
(
"vgId:%d expand stream task, s-task:%s, checkpoint ver:%"
PRId64
...
...
source/libs/stream/src/streamDispatch.c
浏览文件 @
1b263602
...
...
@@ -691,10 +691,11 @@ int32_t streamAddEndScanHistoryMsg(SStreamTask* pTask, SRpcHandleInfo* pRpcInfo,
initRpcMsg
(
&
info
.
msg
,
0
,
pBuf
,
sizeof
(
SMsgHead
)
+
len
);
info
.
msg
.
info
=
*
pRpcInfo
;
// todo: fix race condition here
taosThreadMutexLock
(
&
pTask
->
lock
);
if
(
pTask
->
pRspMsgList
==
NULL
)
{
pTask
->
pRspMsgList
=
taosArrayInit
(
4
,
sizeof
(
SStreamContinueExecInfo
));
}
taosThreadMutexUnlock
(
&
pTask
->
lock
);
taosArrayPush
(
pTask
->
pRspMsgList
,
&
info
);
...
...
source/libs/stream/src/streamTask.c
浏览文件 @
1b263602
...
...
@@ -251,5 +251,10 @@ void tFreeStreamTask(SStreamTask* pTask) {
tSimpleHashCleanup
(
pTask
->
pNameMap
);
}
if
(
pTask
->
pRspMsgList
!=
NULL
)
{
pTask
->
pRspMsgList
=
taosArrayDestroy
(
pTask
->
pRspMsgList
);
}
taosThreadMutexDestroy
(
&
pTask
->
lock
);
taosMemoryFree
(
pTask
);
}
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录