Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
5b88dfb9
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,发现更多精彩内容 >>
提交
5b88dfb9
编写于
6月 25, 2023
作者:
wmmhello
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix:remove report offset flag
上级
610aa0c8
变更
1
隐藏空白更改
内联
并排
Showing
1 changed file
with
6 addition
and
6 deletion
+6
-6
source/client/src/clientTmq.c
source/client/src/clientTmq.c
+6
-6
未找到文件。
source/client/src/clientTmq.c
浏览文件 @
5b88dfb9
...
...
@@ -99,7 +99,7 @@ struct tmq_t {
// poll info
int64_t
pollCnt
;
int64_t
totalRows
;
bool
needReportOffsetRows
;
//
bool needReportOffsetRows;
// timer
tmr_h
hbLiveTimer
;
...
...
@@ -797,7 +797,7 @@ void tmqSendHbReq(void* param, void* tmrId) {
SMqHbReq
req
=
{
0
};
req
.
consumerId
=
tmq
->
consumerId
;
req
.
epoch
=
tmq
->
epoch
;
if
(
tmq
->
needReportOffsetRows
){
//
if(tmq->needReportOffsetRows){
req
.
topics
=
taosArrayInit
(
taosArrayGetSize
(
tmq
->
clientTopics
),
sizeof
(
TopicOffsetRows
));
for
(
int
i
=
0
;
i
<
taosArrayGetSize
(
tmq
->
clientTopics
);
i
++
){
SMqClientTopic
*
pTopic
=
taosArrayGet
(
tmq
->
clientTopics
,
i
);
...
...
@@ -816,8 +816,8 @@ void tmqSendHbReq(void* param, void* tmrId) {
tscInfo
(
"report offset: vgId:%d, offset:%s, rows:%"
PRId64
,
offRows
->
vgId
,
buf
,
offRows
->
rows
);
}
}
tmq
->
needReportOffsetRows
=
false
;
}
//
tmq->needReportOffsetRows = false;
//
}
int32_t
tlen
=
tSerializeSMqHbReq
(
NULL
,
0
,
&
req
);
if
(
tlen
<
0
)
{
...
...
@@ -1094,7 +1094,7 @@ tmq_t* tmq_consumer_new(tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
pTmq
->
status
=
TMQ_CONSUMER_STATUS__INIT
;
pTmq
->
pollCnt
=
0
;
pTmq
->
epoch
=
0
;
pTmq
->
needReportOffsetRows
=
true
;
//
pTmq->needReportOffsetRows = true;
// set conf
strcpy
(
pTmq
->
clientId
,
conf
->
clientId
);
...
...
@@ -2452,7 +2452,7 @@ int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet) {
// if no more waiting rsp
pParamSet
->
callbackFn
(
tmq
,
pParamSet
->
code
,
pParamSet
->
userParam
);
taosMemoryFree
(
pParamSet
);
tmq
->
needReportOffsetRows
=
true
;
//
tmq->needReportOffsetRows = true;
taosReleaseRef
(
tmqMgmt
.
rsetId
,
refId
);
return
0
;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录