fix(master): only show running and failed decommisison disk when executing queryAllDecommissionDisk

Signed-off-by: chihe <chihe@oppo.com>
This commit is contained in:
chihe 2024-07-17 16:14:34 +08:00 committed by AmazingChi
parent 65acca7fdd
commit a2da2149c3
6 changed files with 35 additions and 49 deletions

View File

@ -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 {

View File

@ -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
}

View File

@ -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 {

View File

@ -144,6 +144,7 @@ const (
autoDpMetaRepairKey = "autoDpMetaRepair"
autoDpMetaRepairParallelCntKey = "autoDpMetaRepairParallelCnt"
dpTimeoutKey = "dpTimeout"
ShowAll = "showAll"
)
const (

View File

@ -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) {

View File

@ -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,
}
}