Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
903d4920
T
TDengine
项目概览
taosdata
/
TDengine
大约 1 年 前同步成功
通知
1184
Star
22015
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看板
体验新版 GitCode,发现更多精彩内容 >>
未验证
提交
903d4920
编写于
3月 04, 2022
作者:
L
Liu Jicong
提交者:
GitHub
3月 04, 2022
浏览文件
操作
浏览文件
下载
差异文件
Merge pull request #10558 from taosdata/feature/tq
remove createRequest in poll
上级
8beb268a
7510a0e7
变更
1
显示空白变更内容
内联
并排
Showing
1 changed file
with
64 addition
and
51 deletion
+64
-51
source/client/src/tmq.c
source/client/src/tmq.c
+64
-51
未找到文件。
source/client/src/tmq.c
浏览文件 @
903d4920
...
...
@@ -202,6 +202,26 @@ int32_t tmq_list_append(tmq_list_t* ptr, const char* src) {
return
0
;
}
void
tmqClearUnhandleMsg
(
tmq_t
*
tmq
)
{
tmq_message_t
*
msg
;
while
(
1
)
{
taosGetQitem
(
tmq
->
qall
,
(
void
**
)
&
msg
);
if
(
msg
)
taosFreeQitem
(
msg
);
else
break
;
}
taosReadAllQitems
(
tmq
->
mqueue
,
tmq
->
qall
);
while
(
1
)
{
taosGetQitem
(
tmq
->
qall
,
(
void
**
)
&
msg
);
if
(
msg
)
taosFreeQitem
(
msg
);
else
break
;
}
}
int32_t
tmqSubscribeCb
(
void
*
param
,
const
SDataBuf
*
pMsg
,
int32_t
code
)
{
SMqSubscribeCbParam
*
pParam
=
(
SMqSubscribeCbParam
*
)
param
;
pParam
->
rspErr
=
code
;
...
...
@@ -845,6 +865,7 @@ tmq_resp_err_t tmq_seek(tmq_t* tmq, const tmq_topic_vgroup_t* offset) {
SMqClientVg
*
pVg
=
taosArrayGet
(
clientTopic
->
vgs
,
j
);
if
(
pVg
->
vgId
==
pOffset
->
vgId
)
{
pVg
->
currentOffset
=
pOffset
->
offset
;
tmqClearUnhandleMsg
(
tmq
);
return
TMQ_RESP_ERR__SUCCESS
;
}
}
...
...
@@ -883,26 +904,6 @@ SMqConsumeReq* tmqBuildConsumeReqImpl(tmq_t* tmq, int64_t blockingTime, SMqClien
return
pReq
;
}
void
tmqClearUnhandleMsg
(
tmq_t
*
tmq
)
{
tmq_message_t
*
msg
;
while
(
1
)
{
taosGetQitem
(
tmq
->
qall
,
(
void
**
)
&
msg
);
if
(
msg
)
taosFreeQitem
(
msg
);
else
break
;
}
taosReadAllQitems
(
tmq
->
mqueue
,
tmq
->
qall
);
while
(
1
)
{
taosGetQitem
(
tmq
->
qall
,
(
void
**
)
&
msg
);
if
(
msg
)
taosFreeQitem
(
msg
);
else
break
;
}
}
tmq_message_t
*
tmqSyncPollImpl
(
tmq_t
*
tmq
,
int64_t
blockingTime
)
{
tmq_message_t
*
msg
=
NULL
;
for
(
int
i
=
0
;
i
<
taosArrayGetSize
(
tmq
->
clientTopics
);
i
++
)
{
...
...
@@ -919,29 +920,35 @@ tmq_message_t* tmqSyncPollImpl(tmq_t* tmq, int64_t blockingTime) {
// TODO: out of mem
return
NULL
;
}
SMqPollCbParam
*
param
=
malloc
(
sizeof
(
SMqPollCbParam
));
if
(
param
==
NULL
)
{
SMqPollCbParam
*
pParam
=
malloc
(
sizeof
(
SMqPollCbParam
));
if
(
pParam
==
NULL
)
{
atomic_store_32
(
&
pVg
->
vgStatus
,
TMQ_VG_STATUS__IDLE
);
// TODO: out of mem
return
NULL
;
}
param
->
tmq
=
tmq
;
param
->
pVg
=
pVg
;
param
->
epoch
=
tmq
->
epoch
;
param
->
sync
=
1
;
param
->
msg
=
&
msg
;
tsem_init
(
&
param
->
rspSem
,
0
,
0
);
SRequestObj
*
pRequest
=
createRequest
(
tmq
->
pTscObj
,
NULL
,
NULL
,
TDMT_VND_CONSUME
);
pRequest
->
body
.
requestMsg
=
(
SDataBuf
){
pParam
->
tmq
=
tmq
;
pParam
->
pVg
=
pVg
;
pParam
->
epoch
=
tmq
->
epoch
;
pParam
->
sync
=
1
;
pParam
->
msg
=
&
msg
;
tsem_init
(
&
pParam
->
rspSem
,
0
,
0
);
SMsgSendInfo
*
sendInfo
=
malloc
(
sizeof
(
SMsgSendInfo
));
if
(
sendInfo
==
NULL
)
{
return
NULL
;
}
sendInfo
->
msgInfo
=
(
SDataBuf
){
.
pData
=
pReq
,
.
len
=
sizeof
(
SMqConsumeReq
),
.
handle
=
NULL
,
};
SMsgSendInfo
*
sendInfo
=
buildMsgInfoImpl
(
pRequest
);
sendInfo
->
requestId
=
generateRequestId
();
sendInfo
->
requestObjRefId
=
0
;
sendInfo
->
param
=
param
;
sendInfo
->
param
=
p
P
aram
;
sendInfo
->
fp
=
tmqPollCb
;
sendInfo
->
msgType
=
TDMT_VND_CONSUME
;
int64_t
transporterId
=
0
;
/*printf("send poll\n");*/
...
...
@@ -950,7 +957,7 @@ tmq_message_t* tmqSyncPollImpl(tmq_t* tmq, int64_t blockingTime) {
pVg
->
pollCnt
++
;
tmq
->
pollCnt
++
;
tsem_wait
(
&
param
->
rspSem
);
tsem_wait
(
&
p
P
aram
->
rspSem
);
tmq_message_t
*
nmsg
=
NULL
;
while
(
1
)
{
taosReadQitem
(
tmq
->
mqueue
,
(
void
**
)
&
nmsg
);
...
...
@@ -978,32 +985,40 @@ int32_t tmqPollImpl(tmq_t* tmq, int64_t blockingTime) {
SMqConsumeReq
*
pReq
=
tmqBuildConsumeReqImpl
(
tmq
,
blockingTime
,
pTopic
,
pVg
);
if
(
pReq
==
NULL
)
{
atomic_store_32
(
&
pVg
->
vgStatus
,
TMQ_VG_STATUS__IDLE
);
// TODO: out of mem
tsem_post
(
&
tmq
->
rspSem
);
return
-
1
;
}
SMqPollCbParam
*
param
=
malloc
(
sizeof
(
SMqPollCbParam
));
if
(
param
==
NULL
)
{
SMqPollCbParam
*
pParam
=
malloc
(
sizeof
(
SMqPollCbParam
));
if
(
pParam
==
NULL
)
{
free
(
pReq
);
atomic_store_32
(
&
pVg
->
vgStatus
,
TMQ_VG_STATUS__IDLE
);
// TODO: out of mem
tsem_post
(
&
tmq
->
rspSem
);
return
-
1
;
}
param
->
tmq
=
tmq
;
param
->
pVg
=
pVg
;
param
->
epoch
=
tmq
->
epoch
;
param
->
sync
=
0
;
SRequestObj
*
pRequest
=
createRequest
(
tmq
->
pTscObj
,
NULL
,
NULL
,
TDMT_VND_CONSUME
);
pRequest
->
body
.
requestMsg
=
(
SDataBuf
){
pParam
->
tmq
=
tmq
;
pParam
->
pVg
=
pVg
;
pParam
->
epoch
=
tmq
->
epoch
;
pParam
->
sync
=
0
;
SMsgSendInfo
*
sendInfo
=
malloc
(
sizeof
(
SMsgSendInfo
));
if
(
sendInfo
==
NULL
)
{
free
(
pReq
);
free
(
pParam
);
atomic_store_32
(
&
pVg
->
vgStatus
,
TMQ_VG_STATUS__IDLE
);
tsem_post
(
&
tmq
->
rspSem
);
return
-
1
;
}
sendInfo
->
msgInfo
=
(
SDataBuf
){
.
pData
=
pReq
,
.
len
=
sizeof
(
SMqConsumeReq
),
.
handle
=
NULL
,
};
SMsgSendInfo
*
sendInfo
=
buildMsgInfoImpl
(
pRequest
);
sendInfo
->
requestId
=
generateRequestId
();
sendInfo
->
requestObjRefId
=
0
;
sendInfo
->
param
=
param
;
sendInfo
->
param
=
p
P
aram
;
sendInfo
->
fp
=
tmqPollCb
;
sendInfo
->
msgType
=
TDMT_VND_CONSUME
;
int64_t
transporterId
=
0
;
/*printf("send poll\n");*/
...
...
@@ -1053,7 +1068,7 @@ tmq_message_t* tmqHandleAllRsp(tmq_t* tmq, int64_t blockingTime, bool pollIfRese
atomic_store_32
(
&
pVg
->
vgStatus
,
TMQ_VG_STATUS__IDLE
);
return
rspMsg
;
}
else
{
printf
(
"epoch mismatch
\n
"
);
/*printf("epoch mismatch\n");*/
taosFreeQitem
(
rspMsg
);
}
}
else
{
...
...
@@ -1107,9 +1122,7 @@ tmq_message_t* tmq_consumer_poll(tmq_t* tmq, int64_t blocking_time) {
while
(
1
)
{
/*printf("cycle\n");*/
if
(
atomic_load_32
(
&
tmq
->
waitingRequest
)
==
0
)
{
tmqPollImpl
(
tmq
,
blocking_time
);
}
tsem_wait
(
&
tmq
->
rspSem
);
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录