diff --git a/datanode/wrap_operator.go b/datanode/wrap_operator.go index 02a57f2eb..e9362989f 100644 --- a/datanode/wrap_operator.go +++ b/datanode/wrap_operator.go @@ -475,36 +475,6 @@ end: } } -func (s *DataNode) checkVolumeForbidden(volNames []string) { - s.space.RangePartitions(func(partition *DataPartition) bool { - for _, volName := range volNames { - if volName == partition.volumeID { - partition.SetForbidden(true) - return true - } - } - partition.SetForbidden(false) - return true - }, "") -} - -func (s *DataNode) checkVolumeDpRepairBlockSize(dpRepairBlockSize map[string]uint64) { - s.space.RangePartitions(func(partition *DataPartition) bool { - size := uint64(proto.DefaultDpRepairBlockSize) - if len(dpRepairBlockSize) != 0 { - var ok bool - if size, ok = dpRepairBlockSize[partition.volumeID]; !ok { - size = proto.DefaultDpRepairBlockSize - } - } - log.LogDebugf("[checkVolumeDpRepairBlockSize] volume(%v) dp(%v) repair block size(%v) current size(%v)", partition.volumeID, partition.partitionID, size, partition.GetRepairBlockSize()) - if partition.GetRepairBlockSize() != size { - partition.SetRepairBlockSize(size) - } - return true - }, "") -} - func (s *DataNode) checkDecommissionDisks(decommissionDisks []string) { decommissionDiskSet := util.NewSet() for _, disk := range decommissionDisks { diff --git a/master/api_service.go b/master/api_service.go index 56b7d000a..dd870da4c 100644 --- a/master/api_service.go +++ b/master/api_service.go @@ -4286,13 +4286,14 @@ func (m *Server) queryAllDecommissionDisk(w http.ResponseWriter, r *http.Request var ( err error decommissoinType int + showAll bool ) metric := exporter.NewTPCnt("req_queryAllDecommissionDisk") defer func() { metric.Set(err) }() - if decommissoinType, err = parseReqToQueryDecoDisk(r); err != nil { + if decommissoinType, showAll, err = parseReqToQueryDecoDisk(r); err != nil { sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()}) return } @@ -4302,6 +4303,9 @@ func (m *Server) queryAllDecommissionDisk(w http.ResponseWriter, r *http.Request disk := value.(*DecommissionDisk) if decommissoinType == int(QueryDecommission) || (decommissoinType != int(QueryDecommission) && disk.Type == uint32(decommissoinType)) { status, progress := disk.updateDecommissionStatus(m.cluster, true) + if !showAll && (status != DecommissionFail && status != DecommissionRunning) { + return true + } progress, _ = FormatFloatFloor(progress, 4) decommissionProgress := proto.DecommissionProgress{ Status: status, @@ -4876,7 +4880,7 @@ func parseNodeAddrAndDisk(r *http.Request) (nodeAddr, diskPath string, err error return } -func parseReqToQueryDecoDisk(r *http.Request) (decommissionType int, err error) { +func parseReqToQueryDecoDisk(r *http.Request) (decommissionType int, showAll bool, err error) { if err = r.ParseForm(); err != nil { return } @@ -4884,6 +4888,10 @@ func parseReqToQueryDecoDisk(r *http.Request) (decommissionType int, err error) if err != nil { return } + showAll, err = pareseBoolWithDefault(r, ShowAll, false) + if err != nil { + return + } return } diff --git a/master/cluster.go b/master/cluster.go index f99d7ba30..f9bb19413 100644 --- a/master/cluster.go +++ b/master/cluster.go @@ -4214,6 +4214,8 @@ func (c *Cluster) migrateDisk(dataNode *DataNode, diskPath, dstPath string, raft c.DecommissionDisks.Store(disk.GenerateKey(), disk) } disk.Type = migrateType + disk.DiskDisable = diskDisable + disk.ResidualDecommissionDps = make([]proto.IgnoreDecommissionDP, 0) // disk should be decommission all the dp disk.markDecommission(dstPath, raftForce, limit) if err = c.syncAddDecommissionDisk(disk); err != nil { diff --git a/master/const.go b/master/const.go index c4089f047..e4038e824 100644 --- a/master/const.go +++ b/master/const.go @@ -144,6 +144,7 @@ const ( autoDpMetaRepairKey = "autoDpMetaRepair" autoDpMetaRepairParallelCntKey = "autoDpMetaRepairParallelCnt" dpTimeoutKey = "dpTimeout" + ShowAll = "showAll" ) const ( diff --git a/master/disk_manager.go b/master/disk_manager.go index 21ea57e6d..679ac3e5b 100644 --- a/master/disk_manager.go +++ b/master/disk_manager.go @@ -299,14 +299,15 @@ func (dd *DecommissionDisk) GenerateKey() string { func (dd *DecommissionDisk) updateDecommissionStatus(c *Cluster, debug bool) (uint32, float64) { var ( - progress float64 - totalNum = dd.DecommissionDpTotal - partitionIds []uint64 - failedPartitionIds []uint64 - runningPartitionIds []uint64 - preparePartitionIds []uint64 - stopPartitionIds []uint64 - ignorePartitionIds []uint64 + progress float64 + totalNum = dd.DecommissionDpTotal + partitionIds []uint64 + failedPartitionIds []uint64 + runningPartitionIds []uint64 + preparePartitionIds []uint64 + stopPartitionIds []uint64 + ignorePartitionIds []uint64 + residualPartitionIds []uint64 ) if dd.GetDecommissionStatus() == DecommissionInitial { @@ -321,10 +322,6 @@ func (dd *DecommissionDisk) updateDecommissionStatus(c *Cluster, debug bool) (ui return DecommissionFail, float64(0) } - if dd.GetDecommissionStatus() == DecommissionSuccess { - return DecommissionSuccess, float64(1) - } - if dd.GetDecommissionStatus() == DecommissionPause { return DecommissionPause, float64(0) } @@ -349,7 +346,12 @@ func (dd *DecommissionDisk) updateDecommissionStatus(c *Cluster, debug bool) (ui failedNum++ } - if len(partitions)+len(ignorePartitionIds) == 0 { + for _, info := range dd.ResidualDecommissionDps { + residualPartitionIds = append(residualPartitionIds, info.PartitionID) + failedNum++ + } + + if len(partitions)+len(ignorePartitionIds)+len(residualPartitionIds) == 0 { log.LogDebugf("action[updateDecommissionDiskStatus]no partitions left:%v", dd.GenerateKey()) dd.markDecommissionSuccess() return DecommissionSuccess, float64(1) @@ -377,7 +379,7 @@ func (dd *DecommissionDisk) updateDecommissionStatus(c *Cluster, debug bool) (ui partitionIds = append(partitionIds, dp.PartitionID) } - progress = float64(totalNum-len(partitions)-len(ignorePartitionIds)) / float64(totalNum) + progress = float64(totalNum-len(partitions)-len(ignorePartitionIds)-len(residualPartitionIds)) / float64(totalNum) if debug { log.LogInfof("action[updateDecommissionStatus] disk[%v] progress[%v] totalNum[%v] "+ "partitionIds %v left %v FailedNum[%v] failedPartitionIds %v, runningNum[%v] runningDp %v, prepareNum[%v] prepareDp %v "+ @@ -389,7 +391,7 @@ func (dd *DecommissionDisk) updateDecommissionStatus(c *Cluster, debug bool) (ui if dd.GetDecommissionStatus() == DecommissionCancel { return DecommissionCancel, progress } - if failedNum >= (len(partitions)+len(ignorePartitionIds)-stopNum) && failedNum != 0 { + if failedNum >= (len(partitions)+len(ignorePartitionIds)+len(residualPartitionIds)-stopNum) && failedNum != 0 { dd.markDecommissionFailed() return DecommissionFail, progress } @@ -514,10 +516,10 @@ func (dd *DecommissionDisk) CanBePaused() bool { } func (dd *DecommissionDisk) decommissionInfo() string { - return fmt.Sprintf("disk(%v_%v)_dst(%v)_total(%v)_term(%v)_type(%v)_force(%v)_retry(%v)_status(%v)", + return fmt.Sprintf("disk(%v_%v)_dst(%v)_total(%v)_term(%v)_type(%v)_force(%v)_retry(%v)_status(%v)_disable(%v)", dd.SrcAddr, dd.DiskPath, dd.DstAddr, dd.DecommissionDpTotal, dd.DecommissionTerm, GetDecommissionTypeMessage(dd.Type), dd.DecommissionRaftForce, dd.DecommissionTimes, - GetDecommissionStatusMessage(dd.DecommissionStatus)) + GetDecommissionStatusMessage(dd.DecommissionStatus), dd.DiskDisable) } func (dd *DecommissionDisk) cancelDecommission(cluster *Cluster, ns *nodeSet) (err error) { diff --git a/master/metadata_fsm_op.go b/master/metadata_fsm_op.go index f553c2c33..81da21e79 100644 --- a/master/metadata_fsm_op.go +++ b/master/metadata_fsm_op.go @@ -1726,6 +1726,7 @@ type decommissionDiskValue struct { DecommissionLimit int IgnoreDecommissionDps []bsProto.IgnoreDecommissionDP ResidualDecommissionDps []bsProto.IgnoreDecommissionDP + DiskDisable bool } func newDecommissionDiskValue(disk *DecommissionDisk) *decommissionDiskValue { @@ -1743,6 +1744,7 @@ func newDecommissionDiskValue(disk *DecommissionDisk) *decommissionDiskValue { DecommissionLimit: disk.DecommissionDpCount, IgnoreDecommissionDps: disk.IgnoreDecommissionDps, ResidualDecommissionDps: disk.ResidualDecommissionDps, + DiskDisable: disk.DiskDisable, } } @@ -1761,6 +1763,7 @@ func (ddv *decommissionDiskValue) Restore() *DecommissionDisk { DecommissionDpCount: ddv.DecommissionLimit, IgnoreDecommissionDps: ddv.IgnoreDecommissionDps, ResidualDecommissionDps: ddv.ResidualDecommissionDps, + DiskDisable: ddv.DiskDisable, } }