fix:remove report offset flag

This commit is contained in:
wangmm0220 2023-06-25 15:24:52 +08:00
parent 610aa0c8ba
commit 5b88dfb90b
1 changed files with 6 additions and 6 deletions

View File

@ -99,7 +99,7 @@ struct tmq_t {
// poll info // poll info
int64_t pollCnt; int64_t pollCnt;
int64_t totalRows; int64_t totalRows;
bool needReportOffsetRows; // bool needReportOffsetRows;
// timer // timer
tmr_h hbLiveTimer; tmr_h hbLiveTimer;
@ -797,7 +797,7 @@ void tmqSendHbReq(void* param, void* tmrId) {
SMqHbReq req = {0}; SMqHbReq req = {0};
req.consumerId = tmq->consumerId; req.consumerId = tmq->consumerId;
req.epoch = tmq->epoch; req.epoch = tmq->epoch;
if(tmq->needReportOffsetRows){ // if(tmq->needReportOffsetRows){
req.topics = taosArrayInit(taosArrayGetSize(tmq->clientTopics), sizeof(TopicOffsetRows)); req.topics = taosArrayInit(taosArrayGetSize(tmq->clientTopics), sizeof(TopicOffsetRows));
for(int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++){ for(int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++){
SMqClientTopic* pTopic = taosArrayGet(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); 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); int32_t tlen = tSerializeSMqHbReq(NULL, 0, &req);
if (tlen < 0) { 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->status = TMQ_CONSUMER_STATUS__INIT;
pTmq->pollCnt = 0; pTmq->pollCnt = 0;
pTmq->epoch = 0; pTmq->epoch = 0;
pTmq->needReportOffsetRows = true; // pTmq->needReportOffsetRows = true;
// set conf // set conf
strcpy(pTmq->clientId, conf->clientId); strcpy(pTmq->clientId, conf->clientId);
@ -2452,7 +2452,7 @@ int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet) {
// if no more waiting rsp // if no more waiting rsp
pParamSet->callbackFn(tmq, pParamSet->code, pParamSet->userParam); pParamSet->callbackFn(tmq, pParamSet->code, pParamSet->userParam);
taosMemoryFree(pParamSet); taosMemoryFree(pParamSet);
tmq->needReportOffsetRows = true; // tmq->needReportOffsetRows = true;
taosReleaseRef(tmqMgmt.rsetId, refId); taosReleaseRef(tmqMgmt.rsetId, refId);
return 0; return 0;