Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
d1e5d6e0
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看板
提交
d1e5d6e0
编写于
4月 25, 2023
作者:
wmmhello
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix:pHandle->msg is not null if rebalance
上级
41bec856
变更
3
隐藏空白更改
内联
并排
Showing
3 changed file
with
9 addition
and
9 deletion
+9
-9
source/client/src/clientTmq.c
source/client/src/clientTmq.c
+3
-2
source/dnode/vnode/src/tq/tq.c
source/dnode/vnode/src/tq/tq.c
+1
-1
source/dnode/vnode/src/tq/tqUtil.c
source/dnode/vnode/src/tq/tqUtil.c
+5
-6
未找到文件。
source/client/src/clientTmq.c
浏览文件 @
d1e5d6e0
...
@@ -1363,6 +1363,7 @@ CREATE_MSG_FAIL:
...
@@ -1363,6 +1363,7 @@ CREATE_MSG_FAIL:
typedef
struct
SVgroupSaveInfo
{
typedef
struct
SVgroupSaveInfo
{
STqOffsetVal
offset
;
STqOffsetVal
offset
;
int64_t
numOfRows
;
int64_t
numOfRows
;
int32_t
vgStatus
;
}
SVgroupSaveInfo
;
}
SVgroupSaveInfo
;
static
void
initClientTopicFromRsp
(
SMqClientTopic
*
pTopic
,
SMqSubTopicEp
*
pTopicEp
,
SHashObj
*
pVgOffsetHashMap
,
static
void
initClientTopicFromRsp
(
SMqClientTopic
*
pTopic
,
SMqSubTopicEp
*
pTopicEp
,
SHashObj
*
pVgOffsetHashMap
,
...
@@ -1398,7 +1399,7 @@ static void initClientTopicFromRsp(SMqClientTopic* pTopic, SMqSubTopicEp* pTopic
...
@@ -1398,7 +1399,7 @@ static void initClientTopicFromRsp(SMqClientTopic* pTopic, SMqSubTopicEp* pTopic
.
currentOffset
=
offsetNew
,
.
currentOffset
=
offsetNew
,
.
vgId
=
pVgEp
->
vgId
,
.
vgId
=
pVgEp
->
vgId
,
.
epSet
=
pVgEp
->
epSet
,
.
epSet
=
pVgEp
->
epSet
,
.
vgStatus
=
TMQ_VG_STATUS__IDLE
,
.
vgStatus
=
pInfo
!=
NULL
?
pInfo
->
vgStatus
:
TMQ_VG_STATUS__IDLE
,
.
vgSkipCnt
=
0
,
.
vgSkipCnt
=
0
,
.
emptyBlockReceiveTs
=
0
,
.
emptyBlockReceiveTs
=
0
,
.
numOfRows
=
numOfRows
,
.
numOfRows
=
numOfRows
,
...
@@ -1457,7 +1458,7 @@ static bool doUpdateLocalEp(tmq_t* tmq, int32_t epoch, const SMqAskEpRsp* pRsp)
...
@@ -1457,7 +1458,7 @@ static bool doUpdateLocalEp(tmq_t* tmq, int32_t epoch, const SMqAskEpRsp* pRsp)
tscDebug
(
"consumer:0x%"
PRIx64
", epoch:%d vgId:%d vgKey:%s, offset:%s"
,
tmq
->
consumerId
,
epoch
,
pVgCur
->
vgId
,
tscDebug
(
"consumer:0x%"
PRIx64
", epoch:%d vgId:%d vgKey:%s, offset:%s"
,
tmq
->
consumerId
,
epoch
,
pVgCur
->
vgId
,
vgKey
,
buf
);
vgKey
,
buf
);
SVgroupSaveInfo
info
=
{.
offset
=
pVgCur
->
currentOffset
,
.
numOfRows
=
pVgCur
->
numOfRows
};
SVgroupSaveInfo
info
=
{.
offset
=
pVgCur
->
currentOffset
,
.
numOfRows
=
pVgCur
->
numOfRows
,
.
vgStatus
=
pVgCur
->
vgStatus
};
taosHashPut
(
pVgOffsetHashMap
,
vgKey
,
strlen
(
vgKey
),
&
info
,
sizeof
(
SVgroupSaveInfo
));
taosHashPut
(
pVgOffsetHashMap
,
vgKey
,
strlen
(
vgKey
),
&
info
,
sizeof
(
SVgroupSaveInfo
));
}
}
}
}
...
...
source/dnode/vnode/src/tq/tq.c
浏览文件 @
d1e5d6e0
...
@@ -1079,11 +1079,11 @@ int32_t tqProcessDelReq(STQ* pTq, void* pReq, int32_t len, int64_t ver) {
...
@@ -1079,11 +1079,11 @@ int32_t tqProcessDelReq(STQ* pTq, void* pReq, int32_t len, int64_t ver) {
int32_t
tqProcessSubmitReqForSubscribe
(
STQ
*
pTq
)
{
int32_t
tqProcessSubmitReqForSubscribe
(
STQ
*
pTq
)
{
int32_t
vgId
=
TD_VID
(
pTq
->
pVnode
);
int32_t
vgId
=
TD_VID
(
pTq
->
pVnode
);
tqDebug
(
"vgId:%d start set submit for subscribe"
,
vgId
);
taosWLockLatch
(
&
pTq
->
lock
);
taosWLockLatch
(
&
pTq
->
lock
);
for
(
size_t
i
=
0
;
i
<
taosArrayGetSize
(
pTq
->
pPushArray
);
i
++
){
for
(
size_t
i
=
0
;
i
<
taosArrayGetSize
(
pTq
->
pPushArray
);
i
++
){
STqHandle
*
pHandle
=
(
STqHandle
*
)
taosArrayGetP
(
pTq
->
pPushArray
,
i
);
STqHandle
*
pHandle
=
(
STqHandle
*
)
taosArrayGetP
(
pTq
->
pPushArray
,
i
);
tqDebug
(
"vgId:%d start set submit for pHandle:%p"
,
vgId
,
pHandle
);
if
(
ASSERT
(
pHandle
->
msg
!=
NULL
)){
if
(
ASSERT
(
pHandle
->
msg
!=
NULL
)){
tqError
(
"pHandle->msg should not be null"
);
tqError
(
"pHandle->msg should not be null"
);
break
;
break
;
...
...
source/dnode/vnode/src/tq/tqUtil.c
浏览文件 @
d1e5d6e0
...
@@ -181,17 +181,16 @@ static int32_t extractDataAndRspForNormalSubscribe(STQ* pTq, STqHandle* pHandle,
...
@@ -181,17 +181,16 @@ static int32_t extractDataAndRspForNormalSubscribe(STQ* pTq, STqHandle* pHandle,
// code = tqRegisterPushHandle(pTq, pHandle, pRequest, pMsg, &dataRsp, TMQ_MSG_TYPE__POLL_RSP);
// code = tqRegisterPushHandle(pTq, pHandle, pRequest, pMsg, &dataRsp, TMQ_MSG_TYPE__POLL_RSP);
// lock
// lock
taosWLockLatch
(
&
pTq
->
lock
);
taosWLockLatch
(
&
pTq
->
lock
);
if
(
ASSERT
(
pHandle
->
msg
==
NULL
)){
// tqDebug("data is over, register to handle:%p, msg:%p", pHandle, pHandle->msg);
tqError
(
"pHandle->msg should be null"
);
if
(
pHandle
->
msg
==
NULL
){
taosWUnLockLatch
(
&
pTq
->
lock
);
pHandle
->
msg
=
taosMemoryCalloc
(
1
,
sizeof
(
SRpcMsg
));
goto
end
;
}
}
pHandle
->
msg
=
taosMemoryCalloc
(
1
,
sizeof
(
SRpcMsg
));
memcpy
(
pHandle
->
msg
,
pMsg
,
sizeof
(
SRpcMsg
));
memcpy
(
pHandle
->
msg
,
pMsg
,
sizeof
(
SRpcMsg
));
pHandle
->
msg
->
pCont
=
rpcMallocCont
(
pMsg
->
contLen
);
pHandle
->
msg
->
pCont
=
rpcMallocCont
(
pMsg
->
contLen
);
memcpy
(
pHandle
->
msg
->
pCont
,
pMsg
->
pCont
,
pMsg
->
contLen
);
memcpy
(
pHandle
->
msg
->
pCont
,
pMsg
->
pCont
,
pMsg
->
contLen
);
pHandle
->
msg
->
contLen
=
pMsg
->
contLen
;
pHandle
->
msg
->
contLen
=
pMsg
->
contLen
;
tq
Error
(
"data is over, register to handle:%p, pCont:%p, len:%d"
,
pHandle
,
pHandle
->
msg
->
pCont
,
pHandle
->
msg
->
contLen
);
tq
Debug
(
"data is over, register to handle:%p, pCont:%p, len:%d"
,
pHandle
,
pHandle
->
msg
->
pCont
,
pHandle
->
msg
->
contLen
);
taosArrayPush
(
pTq
->
pPushArray
,
&
pHandle
);
taosArrayPush
(
pTq
->
pPushArray
,
&
pHandle
);
taosWUnLockLatch
(
&
pTq
->
lock
);
taosWUnLockLatch
(
&
pTq
->
lock
);
tDeleteSMqDataRsp
(
&
dataRsp
);
tDeleteSMqDataRsp
(
&
dataRsp
);
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录