fix(master): modify logic of updateDecommissionStatus for datanode

Signed-off-by: chihe <chihe@oppo.com>
This commit is contained in:
chihe 2024-07-05 16:01:45 +08:00 committed by AmazingChi
parent 39d2baa37e
commit 8ad898a1e4
5 changed files with 45 additions and 74 deletions

View File

@ -3034,17 +3034,6 @@ func (m *Server) cancelDecommissionDataNode(w http.ResponseWriter, r *http.Reque
dd.cancelDecommission(m.cluster, ns)
}
}
// find all decommission failed data partitions for this node
// leftPartitions := m.cluster.getAllDecommissionDataPartitionByDataNode(offLineAddr)
// for _, dp := range leftPartitions {
// if dp.GetDecommissionStatus() == DecommissionSuccess || dp.IsRollbackFailed() || ns.HasDecommissionToken(dp.PartitionID) {
// continue
// }
// msg := fmt.Sprintf("dp(%v) cancel decommission", dp.decommissionInfo())
// dp.ResetDecommissionStatus()
// dp.setRestoreReplicaStop()
// auditlog.LogMasterOp("CancelDataPartitionDecommission", msg, nil)
// }
rstMsg = fmt.Sprintf("cancel decommission data node [%v] success", offLineAddr)
sendOkReply(w, r, newSuccessHTTPReply(rstMsg))
@ -7117,6 +7106,13 @@ func (m *Server) cancelDecommissionDisk(w http.ResponseWriter, r *http.Request)
return
}
disk := value.(*DecommissionDisk)
status := disk.GetDecommissionStatus()
if status == DecommissionSuccess || status == DecommissionFail {
ret := fmt.Sprintf("action[cancelDecommissionDisk] disk %v status %v do not support cancel",
key, status)
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: ret})
return
}
if err = disk.cancelDecommission(m.cluster, ns); err != nil {
ret := fmt.Sprintf("action[cancelDecommissionDisk] cancel disk %v decommission failed[%v]",
key, err.Error())

View File

@ -4026,7 +4026,7 @@ func (c *Cluster) TryDecommissionDataNode(dataNode *DataNode) {
}
}
}
dataNode.SetDecommissionStatus(DecommissionPrepare)
dataNode.SetDecommissionStatus(DecommissionRunning)
dataNode.ToBeOffline = true
log.LogDebugf("action[TryDecommissionDataNode] dataNode [%s] recover from DecommissionDiskList", dataNode.Addr)
return
@ -4164,9 +4164,9 @@ func (c *Cluster) TryDecommissionDataNode(dataNode *DataNode) {
// c.syncUpdateDataPartition(dp)
// ns.AddToDecommissionDataPartitionList(dp)
// toBeOffLinePartitionIds = append(toBeOffLinePartitionIds, dp.PartitionID)
// }
// disk wait for decommission
dataNode.SetDecommissionStatus(DecommissionPrepare)
//}
//disk wait for decommission
dataNode.SetDecommissionStatus(DecommissionRunning)
// avoid alloc dp on this node
dataNode.ToBeOffline = true
dataNode.DecommissionDiskList = decommissionDiskList

View File

@ -377,15 +377,12 @@ func (dataNode *DataNode) checkDecommissionedDisks(d string) (ok bool) {
func (dataNode *DataNode) updateDecommissionStatus(c *Cluster, debug bool) (uint32, float64) {
var (
partitionIds []uint64
failedPartitionIds []uint64
runningPartitionIds []uint64
preparePartitionIds []uint64
stopPartitionIds []uint64
totalDisk = len(dataNode.DecommissionDiskList)
markDiskNum = 0
successDiskNum = 0
progress float64
totalDisk = len(dataNode.DecommissionDiskList)
markDiskNum = 0
successDiskNum = 0
failedDiskNum = 0
cancelDiskNum = 0
progress float64
)
if dataNode.GetDecommissionStatus() == DecommissionInitial {
return DecommissionInitial, float64(0)
@ -427,6 +424,10 @@ func (dataNode *DataNode) updateDecommissionStatus(c *Cluster, debug bool) (uint
successDiskNum++
} else if status == markDecommission {
markDiskNum++
} else if status == DecommissionFail {
failedDiskNum++
} else if status == DecommissionCancel {
cancelDiskNum++
}
_, diskProgress := dd.updateDecommissionStatus(c, debug)
progress += diskProgress
@ -439,56 +440,29 @@ func (dataNode *DataNode) updateDecommissionStatus(c *Cluster, debug bool) (uint
// only care data node running/prepare/success
// no disk get token
if markDiskNum == totalDisk {
dataNode.SetDecommissionStatus(DecommissionPrepare)
return DecommissionPrepare, float64(0)
} else {
if successDiskNum == totalDisk {
dataNode.SetDecommissionStatus(DecommissionSuccess)
return DecommissionSuccess, float64(1)
}
dataNode.SetDecommissionStatus(markDecommission)
return markDecommission, float64(0)
}
if successDiskNum == totalDisk {
dataNode.SetDecommissionStatus(DecommissionSuccess)
return DecommissionSuccess, float64(1)
}
// update datanode or running status
partitions := dataNode.GetLatestDecommissionDataPartition(c)
// Get all dp on this dataNode
failedNum := 0
runningNum := 0
prepareNum := 0
stopNum := 0
for _, dp := range partitions {
if dp.IsDecommissionFailed() {
failedNum++
failedPartitionIds = append(failedPartitionIds, dp.PartitionID)
}
if dp.GetDecommissionStatus() == DecommissionRunning {
runningNum++
runningPartitionIds = append(runningPartitionIds, dp.PartitionID)
}
if dp.GetDecommissionStatus() == DecommissionPrepare {
prepareNum++
preparePartitionIds = append(preparePartitionIds, dp.PartitionID)
}
// datanode may stop before and will be counted into partitions
if dp.GetDecommissionStatus() == DecommissionPause {
stopNum++
stopPartitionIds = append(stopPartitionIds, dp.PartitionID)
}
partitionIds = append(partitionIds, dp.PartitionID)
}
progress = progress / float64(totalDisk)
if failedNum >= (len(partitions)-stopNum) && failedNum != 0 {
dataNode.markDecommissionFail()
if failedDiskNum == totalDisk {
dataNode.SetDecommissionStatus(DecommissionFail)
return DecommissionFail, progress
}
dataNode.SetDecommissionStatus(DecommissionRunning)
if debug {
log.LogInfof("action[updateDecommissionStatus] dataNode[%v] progress[%v] totalNum[%v] "+
"partitionIds %v FailedNum[%v] failedPartitionIds %v, runningNum[%v] runningDp %v, prepareNum[%v] prepareDp %v "+
"stopNum[%v] stopPartitionIds %v",
dataNode.Addr, progress, len(partitions), partitionIds, failedNum, failedPartitionIds, runningNum, runningPartitionIds,
prepareNum, preparePartitionIds, stopNum, stopPartitionIds)
if cancelDiskNum != 0 {
dataNode.SetDecommissionStatus(DecommissionCancel)
}
return DecommissionRunning, progress
if debug {
log.LogInfof("action[updateDecommissionStatus] dataNode[%v] progress[%v] DecommissionDiskNum[%v] "+
"DecommissionDisks %v markDiskNum[%v] successDiskNum[%v] failedDiskNum[%v] cancelDiskNum[%v]",
dataNode.Addr, progress, len(dataNode.DecommissionDiskList), dataNode.DecommissionDiskList, markDiskNum,
successDiskNum, failedDiskNum, cancelDiskNum)
}
return dataNode.GetDecommissionStatus(), progress
}
func (dataNode *DataNode) GetLatestDecommissionDataPartition(c *Cluster) (partitions []*DataPartition) {

View File

@ -123,11 +123,10 @@ func (c *Cluster) checkDiskRecoveryProgress() {
if partition.DecommissionType == ManualAddReplica {
partition.resetForManualAddReplica()
} else {
partition.SetDecommissionStatus(DecommissionFail)
partition.DecommissionNeedRollback = false
partition.markRollbackFailed(false)
}
partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v cannot finish recover"+
"for host[0] %v is down ", partition.DecommissionDstAddr, masterNode.Addr)
" for host[0] %v is down ", partition.DecommissionDstAddr, masterNode.Addr)
Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] %v",
c.Name, partitionID, partition.DecommissionErrorMessage))
partition.RLock()
@ -376,7 +375,7 @@ func (dd *DecommissionDisk) updateDecommissionStatus(c *Cluster, debug bool) (ui
progress = float64(totalNum-len(partitions)-len(ignorePartitionIds)) / float64(totalNum)
if debug {
log.LogInfof("action[updateDecommissionDiskStatus] disk[%v] progress[%v] totalNum[%v] "+
log.LogInfof("action[updateDecommissionStatus] disk[%v] progress[%v] totalNum[%v] "+
"partitionIds %v FailedNum[%v] failedPartitionIds %v, runningNum[%v] runningDp %v, prepareNum[%v] prepareDp %v "+
"stopNum[%v] stopPartitionIds %v ignorePartitionIds %v term %v",
dd.GenerateKey(), progress, totalNum, partitionIds, failedNum, failedPartitionIds, runningNum, runningPartitionIds,

View File

@ -2109,18 +2109,20 @@ func (l *DecommissionDataPartitionList) traverse(c *Cluster) {
dp.PartitionID)
}
} else if dp.IsDecommissionFailed() {
remove := false
if !dp.tryRollback(c) {
log.LogDebugf("action[DecommissionListTraverse]Remove dp[%v] for fail",
dp.PartitionID)
l.Remove(dp)
// if dp is not removed from decommission list, do not reset RestoreReplica
dp.setRestoreReplicaStop()
remove = true
}
// rollback fail/success need release token
dp.ReleaseDecommissionToken(c)
dp.DecommissionType = InitialDecommission
c.syncUpdateDataPartition(dp)
msg := fmt.Sprintf("dp %v decommission failed", dp.decommissionInfo())
msg := fmt.Sprintf("dp %v decommission failed, remove %v", dp.decommissionInfo(), remove)
auditlog.LogMasterOp("TraverseDataPartition", msg, nil)
} else if dp.IsDecommissionPaused() {
log.LogDebugf("action[DecommissionListTraverse]Remove dp[%v] for paused ",