diff --git a/master/api_service.go b/master/api_service.go index a058d33f1..05eb2a7ed 100644 --- a/master/api_service.go +++ b/master/api_service.go @@ -6348,7 +6348,9 @@ func (m *Server) queryDecommissionFirstHostTokenInfo(w http.ResponseWriter, r *h infos := make([]*DataNodeToDecommissionRepairDpInfo, 0) m.cluster.DataNodeToDecommissionRepairDpMap.Range(func(key, value interface{}) bool { info := value.(*DataNodeToDecommissionRepairDpInfo) - infos = append(infos, info) + if atomic.LoadUint64(&info.curParallel) != 0 { + infos = append(infos, info) + } return true }) log.LogDebugf("action[queryDiskToRepairDpInfo] %v", infos) diff --git a/master/cluster.go b/master/cluster.go index 1eb169643..631f4e2ef 100644 --- a/master/cluster.go +++ b/master/cluster.go @@ -5302,7 +5302,9 @@ func (c *Cluster) TryDecommissionDisk(disk *DecommissionDisk) { disk.decommissionInfo(), dp.PartitionID, disk.DecommissionTerm) } } else { - ns.AddToDecommissionDataPartitionList(dp, c) + if dp.GetDecommissionStatus() == markDecommission { + ns.AddToDecommissionDataPartitionList(dp, c) + } } c.syncUpdateDataPartition(dp) badPartitionIds = append(badPartitionIds, dp.PartitionID) @@ -6035,12 +6037,17 @@ func (c *Cluster) markDecommissionDataPartition(dp *DataPartition, src *DataNode return } } + // TODO: handle error err = c.syncUpdateDataPartition(dp) if err != nil { return } - ns.AddToDecommissionDataPartitionList(dp, c) + + if dp.GetDecommissionStatus() == markDecommission { + ns.AddToDecommissionDataPartitionList(dp, c) + } + return } diff --git a/master/data_partition.go b/master/data_partition.go index 856d7419c..35aebc78e 100644 --- a/master/data_partition.go +++ b/master/data_partition.go @@ -1138,38 +1138,45 @@ func (partition *DataPartition) ReleaseDecommissionFirstHostToken(c *Cluster) { if ok { dataNodeToRepairDpInfo := value.(*DataNodeToDecommissionRepairDpInfo) dataNodeToRepairDpInfo.mu.Lock() + defer dataNodeToRepairDpInfo.mu.Unlock() diskToRepairDpInfo, found := dataNodeToRepairDpInfo.diskToDecommissionRepairDpMap[diskPath] if !found { - dataNodeToRepairDpInfo.mu.Unlock() return } if _, isExist := diskToRepairDpInfo.repairingDps[partition.PartitionID]; !isExist { - dataNodeToRepairDpInfo.mu.Unlock() return } delete(diskToRepairDpInfo.repairingDps, partition.PartitionID) - atomic.StoreUint64(&diskToRepairDpInfo.curParallel, uint64(len(diskToRepairDpInfo.repairingDps))) - dataNodeToRepairDpInfo.diskToDecommissionRepairDpMap[diskPath] = diskToRepairDpInfo + if len(diskToRepairDpInfo.repairingDps) == 0 { + delete(dataNodeToRepairDpInfo.diskToDecommissionRepairDpMap, diskPath) + } else { + atomic.StoreUint64(&diskToRepairDpInfo.curParallel, uint64(len(diskToRepairDpInfo.repairingDps))) + dataNodeToRepairDpInfo.diskToDecommissionRepairDpMap[diskPath] = diskToRepairDpInfo + } + if atomic.LoadUint64(&dataNodeToRepairDpInfo.curParallel) > 0 { atomic.AddUint64(&dataNodeToRepairDpInfo.curParallel, -1) } - dataNodeToRepairDpInfo.mu.Unlock() + c.DataNodeToDecommissionRepairDpMap.Store(addr, dataNodeToRepairDpInfo) } } func (partition *DataPartition) AcquireDecommissionFirstHostToken(c *Cluster) bool { + var ok bool + var firstHost string var firstReplica *DataReplica - for _, replica := range partition.Replicas { + for _, host := range partition.Hosts { if partition.DecommissionType == AutoAddReplica || partition.isSpecialReplicaCnt() || - (partition.ReplicaNum == 3 && replica.Addr != partition.DecommissionSrcAddr) { - firstReplica = replica + (partition.ReplicaNum == 3 && host != partition.DecommissionSrcAddr) { + firstReplica, ok = partition.hasReplica(host) + firstHost = host break } } - if firstReplica == nil { - log.LogErrorf("action[AcquireDecommissionFirstHostToken] dp(%v) first replica is nil", partition.PartitionID) + if !ok { + log.LogErrorf("action[AcquireDecommissionFirstHostToken] dp(%v) can not find first host(%v) replica", partition.PartitionID, firstHost) return false } @@ -1190,6 +1197,7 @@ func (partition *DataPartition) AcquireDecommissionFirstHostToken(c *Cluster) bo return false } dataNodeToRepairDpInfo.mu.Lock() + defer dataNodeToRepairDpInfo.mu.Unlock() diskToRepairDpInfo, found := dataNodeToRepairDpInfo.diskToDecommissionRepairDpMap[firstReplica.DiskPath] if !found { diskToRepairDpInfo = &DiskToDecommissionRepairDpInfo{ @@ -1201,7 +1209,6 @@ func (partition *DataPartition) AcquireDecommissionFirstHostToken(c *Cluster) bo if atomic.LoadUint64(&c.DecommissionFirstHostDiskTokenLimit) != 0 && atomic.LoadUint64(&diskToRepairDpInfo.curParallel) >= atomic.LoadUint64(&c.DecommissionFirstHostDiskTokenLimit) { - dataNodeToRepairDpInfo.mu.Unlock() return false } @@ -1209,13 +1216,21 @@ func (partition *DataPartition) AcquireDecommissionFirstHostToken(c *Cluster) bo atomic.StoreUint64(&diskToRepairDpInfo.curParallel, uint64(len(diskToRepairDpInfo.repairingDps))) dataNodeToRepairDpInfo.diskToDecommissionRepairDpMap[firstReplica.DiskPath] = diskToRepairDpInfo atomic.AddUint64(&dataNodeToRepairDpInfo.curParallel, 1) - dataNodeToRepairDpInfo.mu.Unlock() c.DataNodeToDecommissionRepairDpMap.Store(firstReplica.Addr, dataNodeToRepairDpInfo) key := fmt.Sprintf("%v_%v", firstReplica.Addr, firstReplica.DiskPath) partition.DecommissionFirstHostDiskTokenKey = key return true } +func isReplicasContainsHost(replicas []*DataReplica, host string) bool { + for _, replica := range replicas { + if replica.Addr == host { + return true + } + } + return false +} + func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk string, raftForce bool, term uint64, migrateType uint32, c *Cluster, ns *nodeSet, ) (err error) { @@ -1276,6 +1291,55 @@ func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk " cannot handle in auto decommission mode", partition.decommissionInfo()) return proto.ErrAllReplicaUnavailable } + + if partition.ReplicaNum == 3 && len(partition.Hosts) == 3 { + diskErrReplicas := partition.getAllDiskErrorReplica() + if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) && isReplicasContainsHost(diskErrReplicas, partition.Hosts[1]) { + //raftForce delete host0 and host1 + toDeleteHosts := partition.Hosts[:2] + for _, toDeleteHost := range toDeleteHosts { + if err = c.removeDataReplica(partition, toDeleteHost, false, true); err != nil { + log.LogWarnf("action[MarkDecommissionStatus] dp[%v] replicaNum[%v] remove first data replica[%v] failed, err: %v", + partition.PartitionID, partition.ReplicaNum, toDeleteHosts, err) + msg := fmt.Sprintf("dp(%v) replicaNum(%v) mark decommission found host(%v) unavailable, raftForce delete it", + partition.decommissionInfo(), partition.ReplicaNum, toDeleteHost) + auditlog.LogMasterOp("DataPartitionDecommission", msg, err) + return + } + } + //decommission success, reset status + partition.ResetDecommissionStatus() + partition.setRestoreReplicaStop() + msg := fmt.Sprintf("dp(%v) replicaNum(%v) mark decommission found host0(%v) and host1(%v) unavailable, raftForce delete them", + partition.decommissionInfo(), partition.ReplicaNum, toDeleteHosts[0], toDeleteHosts[1]) + auditlog.LogMasterOp("DataPartitionDecommission", msg, nil) + return + } + } + + if partition.ReplicaNum == 2 && len(partition.Hosts) == 2 { + diskErrReplicas := partition.getAllDiskErrorReplica() + if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) { + //raftForce delete host0 + toDeleteHost := partition.Hosts[0] + if err = c.removeDataReplica(partition, toDeleteHost, false, true); err != nil { + log.LogWarnf("action[MarkDecommissionStatus] dp[%v] replicaNum[%v] remove first data replica[%v] failed, err: %v", + partition.PartitionID, partition.ReplicaNum, toDeleteHost, err) + msg := fmt.Sprintf("dp(%v) replicaNum(%v) mark decommission found host0(%v) unavailable, raftForce delete it", + partition.decommissionInfo(), partition.ReplicaNum, toDeleteHost) + auditlog.LogMasterOp("DataPartitionDecommission", msg, err) + return + } + //decommission success, reset status + partition.ResetDecommissionStatus() + partition.setRestoreReplicaStop() + msg := fmt.Sprintf("dp(%v) replicaNum(%v) mark decommission found host0(%v) unavailable, raftForce delete it", + partition.decommissionInfo(), partition.ReplicaNum, toDeleteHost) + auditlog.LogMasterOp("DataPartitionDecommission", msg, nil) + return + } + } + raftForce = true diskErrReplica := partition.getDiskErrorReplica() if diskErrReplica != nil { @@ -1315,6 +1379,15 @@ func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk return proto.ErrAllReplicaUnavailable } } + if migrateType == ManualDecommission && partition.ReplicaNum == 2 && len(partition.Hosts) >= 1 { + diskErrReplicas := partition.getAllDiskErrorReplica() + if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) { + // mark decommission failed + log.LogWarnf("action[MarkDecommissionStatus] dp[%v] replicaNum[%v] host0[%v] is unavaliable, cannot handle in manual decommission mode", + partition.PartitionID, partition.ReplicaNum, partition.Replicas[0].Addr) + return proto.ErrFirstHostUnavailable + } + } } directly: waitTimes := 0 @@ -2222,6 +2295,18 @@ func (partition *DataPartition) getDiskErrorReplica() *DataReplica { return nil } +func (partition *DataPartition) getAllDiskErrorReplica() []*DataReplica { + partition.RLock() + defer partition.RUnlock() + diskErrReplicas := make([]*DataReplica, 0) + for _, replica := range partition.Replicas { + if replica.TriggerDiskError { + diskErrReplicas = append(diskErrReplicas, replica) + } + } + return diskErrReplicas +} + func (partition *DataPartition) checkReplicaMeta(c *Cluster) (err error) { var auditMsg string diff --git a/master/disk_manager.go b/master/disk_manager.go index c58de6fea..6cdcffafe 100644 --- a/master/disk_manager.go +++ b/master/disk_manager.go @@ -120,7 +120,25 @@ func (c *Cluster) checkDiskRecoveryProgress() { if !partition.isSpecialReplicaCnt() || (partition.isSpecialReplicaCnt() && partition.DecommissionRaftForce) { masterNode, _ := partition.getReplica(partition.Hosts[0]) duration := time.Unix(masterNode.ReportTime, 0).Sub(time.Unix(newReplica.ReportTime, 0)) - if math.Abs(duration.Minutes()) > 10 { + diskErrReplicas := partition.getAllDiskErrorReplica() + if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) { + if partition.DecommissionType == ManualAddReplica { + partition.resetForManualAddReplica() + } else { + partition.markRollbackFailed(false) + } + partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v cannot finish recover"+ + " for host[0] %v is unavailable", partition.DecommissionDstAddr, partition.Hosts[0]) + Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] %v", + c.Name, partitionID, partition.DecommissionErrorMessage)) + partition.RLock() + err = c.syncUpdateDataPartition(partition) + if err != nil { + log.LogErrorf("[checkDiskRecoveryProgress] update dp(%v) fail, err(%v)", partitionID, err) + } + partition.RUnlock() + continue + } else if math.Abs(duration.Minutes()) > 10 { if partition.DecommissionType == ManualAddReplica { partition.resetForManualAddReplica() } else { diff --git a/master/monitor_metrics.go b/master/monitor_metrics.go index 3bff56f3f..af67b1c7b 100644 --- a/master/monitor_metrics.go +++ b/master/monitor_metrics.go @@ -18,6 +18,7 @@ import ( "fmt" "math" "strconv" + "strings" "sync" "time" @@ -52,6 +53,7 @@ const ( MetricDiskError = "disk_error" MetricFlashNodesDiskError = "flashNodes_disk_error" MetricDiskLost = "disk_lost" + MetricDpUnableDecommission = "dp_unable_decommission" MetricDataNodesInactive = "dataNodes_inactive" MetricInactiveDataNodeInfo = "inactive_dataNodes_info" MetricMetaNodesInactive = "metaNodes_inactive" @@ -124,6 +126,7 @@ type monitorMetrics struct { diskError *exporter.GaugeVec flashNodesDiskError *exporter.GaugeVec diskLost *exporter.GaugeVec + dpUnableDecommission *exporter.GaugeVec dataNodesNotWritable *exporter.Gauge dataNodesAllocable *exporter.Gauge metaNodesNotWritable *exporter.Gauge @@ -510,6 +513,7 @@ func (mm *monitorMetrics) start() { mm.diskError = exporter.NewGaugeVec(MetricDiskError, "", []string{"addr", "path"}) mm.flashNodesDiskError = exporter.NewGaugeVec(MetricFlashNodesDiskError, "", []string{"addr", "path"}) mm.diskLost = exporter.NewGaugeVec(MetricDiskLost, "", []string{"addr", "path"}) + mm.dpUnableDecommission = exporter.NewGaugeVec(MetricDpUnableDecommission, "", []string{"dpId"}) mm.nodeStat = exporter.NewGaugeVec(MetricNodeStat, "", []string{"type", "addr", "stat"}) mm.dataNodesInactive = exporter.NewGauge(MetricDataNodesInactive) mm.InactiveDataNodeInfo = exporter.NewGaugeVec(MetricInactiveDataNodeInfo, "", []string{"clusterName", "addr"}) @@ -612,6 +616,7 @@ func (mm *monitorMetrics) doStat() { mm.setDiskErrorMetric() mm.setDiskLostMetric() mm.setFlashNodesDiskErrorMetric() + mm.setDpUnableDecommissionMetric() mm.setNotWritableDataNodesCount() mm.setNotWritableMetaNodesCount() mm.setMpInconsistentErrorMetric() @@ -918,6 +923,21 @@ func (mm *monitorMetrics) setDiskLostMetric() { }) } +func (mm *monitorMetrics) setDpUnableDecommissionMetric() { + mm.dpUnableDecommission.Reset() + + vols := mm.cluster.allVols() + for _, vol := range vols { + partitions := vol.dataPartitions.clonePartitions() + for _, dp := range partitions { + if dp.GetDecommissionStatus() == DecommissionFail && strings.Contains(dp.DecommissionErrorMessage, proto.ErrAllReplicaUnavailable.Error()) { + idStr := strconv.FormatUint(dp.PartitionID, 10) + mm.dpUnableDecommission.SetWithLabelValues(1, idStr) + } + } + } +} + func (mm *monitorMetrics) setDiskDecommissionedMetric() { mm.diskDecommissioned.Reset() @@ -1300,6 +1320,7 @@ func (mm *monitorMetrics) resetAllLeaderMetrics() { mm.metaNodesIncreased.Set(0) // mm.diskError.Set(0) mm.diskLost.Reset() + mm.dpUnableDecommission.Reset() mm.diskDecommissioned.Reset() mm.dataNodesInactive.Set(0) mm.metaNodesInactive.Set(0) diff --git a/proto/errors.go b/proto/errors.go index f6dd688cf..1c27ce88d 100644 --- a/proto/errors.go +++ b/proto/errors.go @@ -100,6 +100,7 @@ var ( ErrDecompressFailed = errors.New("decompress data failed") ErrDecommissionDiskErrDPFirst = errors.New("decommission disk error data partition first") ErrAllReplicaUnavailable = errors.New("all replica unavailable") + ErrFirstHostUnavailable = errors.New("first host unavailable") ErrDiskNotExists = errors.New("disk not exists") ErrPerformingRestoreReplica = errors.New("is performing restore replica") ErrPerformingDecommission = errors.New("one replica is performing decommission")