diff --git a/master/admin_task_manager.go b/master/admin_task_manager.go index c9afa39b4..73fd87a60 100644 --- a/master/admin_task_manager.go +++ b/master/admin_task_manager.go @@ -17,6 +17,7 @@ package master import ( "encoding/json" "fmt" + "github.com/google/uuid" "net" "sync" "time" @@ -108,7 +109,9 @@ func (sender *AdminTaskManager) getToBeDeletedTasks() (delTasks []*proto.AdminTa } func (sender *AdminTaskManager) doSendTasks() { - tasks := sender.getToDoTasks() + id := uuid.New() + log.LogDebugf("doSendTasks %v", id.String()) + tasks := sender.getToDoTasks(id.String()) if len(tasks) == 0 { return } @@ -252,7 +255,7 @@ func (sender *AdminTaskManager) AddTask(t *proto.AdminTask) { } } -func (sender *AdminTaskManager) getToDoTasks() (tasks []*proto.AdminTask) { +func (sender *AdminTaskManager) getToDoTasks(id string) (tasks []*proto.AdminTask) { sender.RLock() defer sender.RUnlock() tasks = make([]*proto.AdminTask, 0) @@ -262,6 +265,7 @@ func (sender *AdminTaskManager) getToDoTasks() (tasks []*proto.AdminTask) { if t.IsHeartbeatTask() && t.CheckTaskNeedSend() { tasks = append(tasks, t) t.SendTime = time.Now().Unix() + log.LogDebugf("getToDoTasks get heartbeatTask %v %v", t.RequestID, id) } } // send urgent task immediately diff --git a/master/api_service.go b/master/api_service.go index 3dda6f262..534fd3cf8 100644 --- a/master/api_service.go +++ b/master/api_service.go @@ -1635,7 +1635,6 @@ func (m *Server) addDataReplica(w http.ResponseWriter, r *http.Request) { if !dp.setRestoreReplicaForbidden() { retry++ if retry > defaultDecommissionRetryLimit { - } else { err = errors.NewErrorf("set RestoreReplicaMetaForbidden failed") sendErrReply(w, r, newErrHTTPReply(err)) return @@ -5680,6 +5679,7 @@ func (m *Server) queryDecommissionToken(w http.ResponseWriter, r *http.Request) CurTokenNum: s.CurTokenNum, MaxTokenNum: s.MaxTokenNum, RunningDp: s.RunningDp, + TotalDP: s.TotalDP, }) } } diff --git a/master/cluster.go b/master/cluster.go index cd22af67c..a11c4a006 100644 --- a/master/cluster.go +++ b/master/cluster.go @@ -766,10 +766,15 @@ func (c *Cluster) checkLeaderAddr() { func (c *Cluster) checkDataNodeHeartbeat() { tasks := make([]*proto.AdminTask, 0) + id := uuid.New() + log.LogDebugf("checkDataNodeHeartbeat start %v", id.String()) c.dataNodes.Range(func(addr, dataNode interface{}) bool { node := dataNode.(*DataNode) node.checkLiveness() + log.LogDebugf("checkDataNodeHeartbeat checkLiveness for data node %v %v", node.Addr, id.String()) task := node.createHeartbeatTask(c.masterAddr(), c.diskQosEnable) + log.LogDebugf("checkDataNodeHeartbeat createHeartbeatTask for data node %v task %v %v", node.Addr, + task.RequestID, id.String()) hbReq := task.Request.(*proto.HeartBeatRequest) c.volMutex.RLock() defer c.volMutex.RUnlock() @@ -784,7 +789,9 @@ func (c *Cluster) checkDataNodeHeartbeat() { tasks = append(tasks, task) return true }) + log.LogDebugf("checkDataNodeHeartbeat add task %v", id.String()) c.addDataNodeTasks(tasks) + log.LogDebugf("checkDataNodeHeartbeat end %v", id.String()) } func (c *Cluster) checkMetaNodeHeartbeat() { diff --git a/master/data_node.go b/master/data_node.go index cd9b24a55..abd22f0f7 100644 --- a/master/data_node.go +++ b/master/data_node.go @@ -417,6 +417,10 @@ func (dataNode *DataNode) updateDecommissionStatus(c *Cluster, debug bool) (uint successDiskNum = 0 failedDiskNum = 0 cancelDiskNum = 0 + markDisks = make([]string, 0) + successDisks = make([]string, 0) + failedDisks = make([]string, 0) + cancelDisks = make([]string, 0) progress float64 ) if dataNode.GetDecommissionStatus() == DecommissionInitial { @@ -457,12 +461,16 @@ func (dataNode *DataNode) updateDecommissionStatus(c *Cluster, debug bool) (uint status := dd.GetDecommissionStatus() if status == DecommissionSuccess { successDiskNum++ + successDisks = append(successDisks, dd.DiskPath) } else if status == markDecommission { markDiskNum++ + markDisks = append(markDisks, dd.DiskPath) } else if status == DecommissionFail { failedDiskNum++ + failedDisks = append(failedDisks, dd.DiskPath) } else if status == DecommissionCancel { cancelDiskNum++ + cancelDisks = append(cancelDisks, dd.DiskPath) } _, diskProgress := dd.updateDecommissionStatus(c, debug) progress += diskProgress @@ -492,9 +500,9 @@ func (dataNode *DataNode) updateDecommissionStatus(c *Cluster, debug bool) (uint if debug { log.LogInfof("action[updateDecommissionStatus] dataNode[%v] progress[%v] DecommissionDiskNum[%v] "+ - "DecommissionDisks %v markDiskNum[%v] successDiskNum[%v] failedDiskNum[%v] cancelDiskNum[%v]", + "DecommissionDisks %v markDiskNum[%v] %v successDiskNum[%v] %v failedDiskNum[%v] %v cancelDiskNum[%v] %v", dataNode.Addr, progress/float64(totalDisk), len(dataNode.DecommissionDiskList), dataNode.DecommissionDiskList, markDiskNum, - successDiskNum, failedDiskNum, cancelDiskNum) + markDisks, successDiskNum, successDisks, failedDiskNum, failedDisks, cancelDiskNum, cancelDisks) } return dataNode.GetDecommissionStatus(), progress / float64(totalDisk) } diff --git a/master/data_partition.go b/master/data_partition.go index 82acedba4..8bbe80c6f 100644 --- a/master/data_partition.go +++ b/master/data_partition.go @@ -1875,10 +1875,16 @@ func (partition *DataPartition) needRollback(c *Cluster) bool { log.LogWarnf("action[rollback]dp[%v] rollback to del from bad dataPartitionIDs failed:%v", partition.PartitionID, err) } partition.DecommissionNeedRollback = false - err = partition.removeReplicaByForce(c, partition.DecommissionDstAddr) + removeAddr := partition.DecommissionDstAddr + // when special replica partition enter SpecialDecommissionWaitAddResFin, new replica is recoverd, so only + // need to delete DecommissionSrcAddr + if partition.isSpecialReplicaCnt() && partition.GetSpecialReplicaDecommissionStep() >= SpecialDecommissionWaitAddResFin { + removeAddr = partition.DecommissionSrcAddr + } + err = partition.removeReplicaByForce(c, removeAddr) if err != nil { log.LogWarnf("action[needRollback]dp[%v] remove decommission dst replica %v failed: %v", - partition.PartitionID, partition.DecommissionDstAddr, err) + partition.PartitionID, removeAddr, err) } c.syncUpdateDataPartition(partition) auditlog.LogMasterOp("DataPartitionDecommissionRollback", diff --git a/master/disk_manager.go b/master/disk_manager.go index 2fde28497..1c3d1ae09 100644 --- a/master/disk_manager.go +++ b/master/disk_manager.go @@ -527,6 +527,7 @@ func (dd *DecommissionDisk) cancelDecommission(cluster *Cluster, ns *nodeSet) (e msg := fmt.Sprintf("dp(%v) cancel decommission", dp.decommissionInfo()) dp.ResetDecommissionStatus() dp.setRestoreReplicaStop() + cluster.syncUpdateDataPartition(dp) auditlog.LogMasterOp("CancelDataPartitionDecommission", msg, nil) } dd.SetDecommissionStatus(DecommissionCancel) diff --git a/master/topology.go b/master/topology.go index e25542889..761d4a4cc 100644 --- a/master/topology.go +++ b/master/topology.go @@ -960,6 +960,7 @@ type nodeSetDecommissionParallelStatus struct { CurTokenNum int32 MaxTokenNum int32 RunningDp []uint64 + TotalDP int } func newNodeSet(c *Cluster, id uint64, cap int, zoneName string) *nodeSet { @@ -1085,7 +1086,7 @@ func (ns *nodeSet) UpdateMaxParallel(maxParallel int32) { atomic.StoreInt32(&ns.decommissionParallelLimit, maxParallel) } -func (ns *nodeSet) getDecommissionParallelStatus() (int32, int32, []uint64) { +func (ns *nodeSet) getDecommissionParallelStatus() (int32, int32, []uint64, int) { return ns.decommissionDataPartitionList.getDecommissionParallelStatus() } @@ -1867,12 +1868,13 @@ func (zone *Zone) queryDecommissionParallelStatus() (err error, stats []nodeSetD } for _, ns := range nodeSets { - curToken, maxToken, dps := ns.getDecommissionParallelStatus() + curToken, maxToken, dps, total := ns.getDecommissionParallelStatus() stat := nodeSetDecommissionParallelStatus{ ID: ns.ID, CurTokenNum: curToken, MaxTokenNum: maxToken, RunningDp: dps, + TotalDP: total, } stats = append(stats, stat) } @@ -2005,15 +2007,15 @@ func (l *DecommissionDataPartitionList) Remove(value *DataPartition) { } } -func (l *DecommissionDataPartitionList) getDecommissionParallelStatus() (int32, int32, []uint64) { +func (l *DecommissionDataPartitionList) getDecommissionParallelStatus() (int32, int32, []uint64, int) { l.mu.Lock() defer l.mu.Unlock() dps := make([]uint64, 0) for id := range l.runningMap { dps = append(dps, id) } - - return atomic.LoadInt32(&l.curParallel), atomic.LoadInt32(&l.parallelLimit), dps + total := l.decommissionList.Len() + return atomic.LoadInt32(&l.curParallel), atomic.LoadInt32(&l.parallelLimit), dps, total } func (l *DecommissionDataPartitionList) updateMaxParallel(maxParallel int32) { diff --git a/proto/model.go b/proto/model.go index 83c2ffb67..c4d59cbde 100644 --- a/proto/model.go +++ b/proto/model.go @@ -413,6 +413,7 @@ type DecommissionTokenStatus struct { CurTokenNum int32 MaxTokenNum int32 RunningDp []uint64 + TotalDP int } type VolVersionInfo struct {