refactor: update logs.
This commit is contained in:
parent
8825065364
commit
1e388cace7
|
@ -1486,16 +1486,18 @@ int32_t streamProcessDispatchRsp(SStreamTask* pTask, SStreamDispatchRsp* pRsp, i
|
||||||
int32_t numOfFailed = 0;
|
int32_t numOfFailed = 0;
|
||||||
bool triggerDispatchRsp = false;
|
bool triggerDispatchRsp = false;
|
||||||
SActiveCheckpointInfo* pInfo = pTask->chkInfo.pActiveInfo;
|
SActiveCheckpointInfo* pInfo = pTask->chkInfo.pActiveInfo;
|
||||||
|
|
||||||
// we only set the dispatch msg info for current checkpoint trans
|
|
||||||
int64_t tmpCheckpointId = -1;
|
int64_t tmpCheckpointId = -1;
|
||||||
int32_t tmpTranId = -1;
|
int32_t tmpTranId = -1;
|
||||||
|
const char* pStatus = NULL;
|
||||||
|
|
||||||
|
// we only set the dispatch msg info for current checkpoint trans
|
||||||
streamMutexLock(&pTask->lock);
|
streamMutexLock(&pTask->lock);
|
||||||
triggerDispatchRsp = (streamTaskGetStatus(pTask).state == TASK_STATUS__CK) &&
|
SStreamTaskState s = streamTaskGetStatus(pTask);
|
||||||
(pInfo->activeId == pMsgInfo->checkpointId) && (pInfo->transId != pMsgInfo->transId);
|
triggerDispatchRsp = (s.state == TASK_STATUS__CK) && (pInfo->activeId == pMsgInfo->checkpointId) &&
|
||||||
|
(pInfo->transId != pMsgInfo->transId);
|
||||||
tmpCheckpointId = pInfo->activeId;
|
tmpCheckpointId = pInfo->activeId;
|
||||||
tmpTranId = pInfo->transId;
|
tmpTranId = pInfo->transId;
|
||||||
|
pStatus = s.name;
|
||||||
streamMutexUnlock(&pTask->lock);
|
streamMutexUnlock(&pTask->lock);
|
||||||
|
|
||||||
streamMutexLock(&pMsgInfo->lock);
|
streamMutexLock(&pMsgInfo->lock);
|
||||||
|
@ -1561,8 +1563,9 @@ int32_t streamProcessDispatchRsp(SStreamTask* pTask, SStreamDispatchRsp* pRsp, i
|
||||||
streamTaskSetTriggerDispatchConfirmed(pTask, pRsp->downstreamNodeId);
|
streamTaskSetTriggerDispatchConfirmed(pTask, pRsp->downstreamNodeId);
|
||||||
} else {
|
} else {
|
||||||
stWarn("s-task:%s checkpoint-trigger msg rsp for checkpointId:%" PRId64
|
stWarn("s-task:%s checkpoint-trigger msg rsp for checkpointId:%" PRId64
|
||||||
" transId:%d discard, current active checkpointId:%" PRId64 " active transId:%d, since expired",
|
" transId:%d discard, current status:%s, active checkpointId:%" PRId64
|
||||||
pTask->id.idStr, pMsgInfo->checkpointId, pMsgInfo->transId, tmpCheckpointId, tmpTranId);
|
" active transId:%d, since expired",
|
||||||
|
pTask->id.idStr, pMsgInfo->checkpointId, pMsgInfo->transId, pStatus, tmpCheckpointId, tmpTranId);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
Loading…
Reference in New Issue