fix(master): add display of remainingDpCnt when querying disk or dataNode decommission progress.

close: #1000329505

Signed-off-by: shuqiang-zheng <zhengshuqiang@oppo.com>
(cherry picked from commit 83748a2a03)
This commit is contained in:
shuqiang-zheng 2025-07-23 17:01:22 +08:00 committed by chihe
parent 99b040a620
commit 1c0a800f93
5 changed files with 41 additions and 22 deletions

View File

@ -1341,6 +1341,7 @@ func formatDataNodeDecommissionProgress(progress *proto.DataDecommissionProgress
sb.WriteString(fmt.Sprintf("Status : %v\n", progress.StatusMessage))
sb.WriteString(fmt.Sprintf("Progress : %v\n", progress.Progress))
sb.WriteString(fmt.Sprintf("TotalDpCnt : %v\n", progress.TotalDpCnt))
sb.WriteString(fmt.Sprintf("RemainingDpCnt: %v\n", progress.RemainingDpCnt))
if len(progress.RunningDps) != 0 {
sb.WriteString("running Dps: \n")
for i, info := range progress.RunningDps {
@ -1362,6 +1363,7 @@ func formatDecommissionProgress(progress *proto.DecommissionProgress) string {
sb.WriteString(fmt.Sprintf("Status : %v\n", progress.StatusMessage))
sb.WriteString(fmt.Sprintf("Progress : %v\n", progress.Progress))
sb.WriteString(fmt.Sprintf("TotalDpCnt : %v\n", progress.TotalDpCnt))
sb.WriteString(fmt.Sprintf("RemainingDpCnt: %v\n", progress.RemainingDpCnt))
if len(progress.RunningDps) != 0 {
sb.WriteString("running Dps: \n")
for i, info := range progress.RunningDps {

View File

@ -4978,15 +4978,22 @@ func (m *Server) queryDiskDecoProgress(w http.ResponseWriter, r *http.Request) {
resp := &proto.DecommissionProgress{
Progress: fmt.Sprintf("%.2f%%", progress*float64(100)),
StatusMessage: GetDecommissionStatusMessage(status),
TotalDpCnt: disk.DecommissionDpTotal,
TotalDpCnt: disk.GetDecommissionTotalDpCnt(m.cluster),
IgnoreDps: disk.IgnoreDecommissionDps,
ResidualDps: disk.residualDecommissionDpsGetAll(),
StartTime: time.Unix(int64(disk.DecommissionTerm), 0).String(),
IsManualDecommissionDisk: disk.IsManualDecommissionDisk(),
}
failedDps, runningDps := disk.GetDecommissionFailedAndRunningDPByTerm(m.cluster)
resp.FailedDps = failedDps
resp.RunningDps = runningDps
if status == markDecommission {
resp.RemainingDpCnt = resp.TotalDpCnt
resp.FailedDps = nil
resp.RunningDps = nil
} else {
remainingDpCnt, failedDps, runningDps := disk.GetDecommissionFailedAndRunningDPByTerm(m.cluster)
resp.RemainingDpCnt = remainingDpCnt + len(resp.IgnoreDps) + len(resp.ResidualDps)
resp.FailedDps = failedDps
resp.RunningDps = runningDps
}
retryOverLimitDps := disk.GetDecommissionDiskRetryOverLimitDP(m.cluster)
resp.RetryOverLimitDps = retryOverLimitDps
sendOkReply(w, r, newSuccessHTTPReply(resp))
@ -5053,10 +5060,12 @@ func (m *Server) queryAllDecommissionDisk(w http.ResponseWriter, r *http.Request
IsManualDecommissionDisk: disk.IsManualDecommissionDisk(),
}
if status == markDecommission {
decommissionProgress.FailedDps = make([]proto.FailedDpInfo, 0)
decommissionProgress.RunningDps = make([]uint64, 0)
decommissionProgress.RemainingDpCnt = decommissionProgress.TotalDpCnt
decommissionProgress.FailedDps = nil
decommissionProgress.RunningDps = nil
} else {
failedDps, runningDps := disk.GetDecommissionFailedAndRunningDPByTerm(m.cluster)
remainingDpCnt, failedDps, runningDps := disk.GetDecommissionFailedAndRunningDPByTerm(m.cluster)
decommissionProgress.RemainingDpCnt = remainingDpCnt + len(decommissionProgress.IgnoreDps) + len(decommissionProgress.ResidualDps)
decommissionProgress.FailedDps = failedDps
decommissionProgress.RunningDps = runningDps
}
@ -6803,11 +6812,12 @@ func (m *Server) queryDataNodeDecoProgress(w http.ResponseWriter, r *http.Reques
StatusMessage: GetDecommissionStatusMessage(status),
TotalDpCnt: dn.DecommissionDpTotal,
}
failedDps, runningDps := dn.GetDecommissionFailedAndRunningDPByTerm(m.cluster)
remainingDpCnt, failedDps, runningDps := dn.GetDecommissionFailedAndRunningDPByTerm(m.cluster)
resp.FailedDps = failedDps
resp.RunningDps = runningDps
resp.IgnoreDps = dn.getIgnoreDecommissionDpList(m.cluster)
resp.ResidualDps = dn.getResidualDecommissionDpList(m.cluster)
resp.RemainingDpCnt = remainingDpCnt + len(resp.IgnoreDps) + len(resp.ResidualDps)
sendOkReply(w, r, newSuccessHTTPReply(resp))
}

View File

@ -660,7 +660,7 @@ func (dataNode *DataNode) updateDecommissionStatus(c *Cluster, debug, persist bo
return dataNode.GetDecommissionStatus(), progress / float64(totalDisk)
}
func (dataNode *DataNode) GetLatestDecommissionDataPartition(c *Cluster) (partitions []*DataPartition) {
func (dataNode *DataNode) GetLatestDecommissionDataPartition(c *Cluster) (remainingDpCnt int, partitions []*DataPartition) {
log.LogDebugf("action[GetLatestDecommissionDataPartition]dataNode %v diskList %v", dataNode.Addr, dataNode.DecommissionDiskList)
for _, disk := range dataNode.DecommissionDiskList {
key := fmt.Sprintf("%s_%s", dataNode.Addr, disk)
@ -675,6 +675,11 @@ func (dataNode *DataNode) GetLatestDecommissionDataPartition(c *Cluster) (partit
}
log.LogDebugf("action[GetLatestDecommissionDataPartition]dataNode %v disk %v dps[%v]",
dataNode.Addr, dd.DiskPath, dpIds)
if dd.GetDecommissionStatus() == markDecommission {
remainingDpCnt += dd.GetDecommissionTotalDpCnt(c)
} else {
remainingDpCnt += len(partitions)
}
}
}
return
@ -688,12 +693,12 @@ func (dataNode *DataNode) SetDecommissionStatus(status uint32) {
atomic.StoreUint32(&dataNode.DecommissionStatus, status)
}
func (dataNode *DataNode) GetDecommissionFailedAndRunningDPByTerm(c *Cluster) ([]proto.FailedDpInfo, []uint64) {
func (dataNode *DataNode) GetDecommissionFailedAndRunningDPByTerm(c *Cluster) (int, []proto.FailedDpInfo, []uint64) {
var (
failedDps []proto.FailedDpInfo
runningDps []uint64
)
partitions := dataNode.GetLatestDecommissionDataPartition(c)
remainingDpCnt, partitions := dataNode.GetLatestDecommissionDataPartition(c)
log.LogDebugf("action[GetDecommissionDataNodeFailedDP] partitions len %v", len(partitions))
for _, dp := range partitions {
if dp.IsRollbackFailed() {
@ -705,7 +710,7 @@ func (dataNode *DataNode) GetDecommissionFailedAndRunningDPByTerm(c *Cluster) ([
}
}
log.LogWarnf("action[GetDecommissionDataNodeFailedDP] failed dp list [%v]", failedDps)
return failedDps, runningDps
return remainingDpCnt, failedDps, runningDps
}
func (dataNode *DataNode) GetDecommissionFailedDP(c *Cluster) (error, []uint64) {

View File

@ -540,7 +540,7 @@ func (dd *DecommissionDisk) GetDecommissionDiskRetryOverLimitDP(c *Cluster) []ui
return retryOverLimitDps
}
func (dd *DecommissionDisk) GetDecommissionFailedAndRunningDPByTerm(c *Cluster) ([]proto.FailedDpInfo, []uint64) {
func (dd *DecommissionDisk) GetDecommissionFailedAndRunningDPByTerm(c *Cluster) (int, []proto.FailedDpInfo, []uint64) {
partitions := c.getAllDecommissionDataPartitionByDiskAndTerm(dd.SrcAddr, dd.DiskPath, dd.DecommissionTerm)
var (
failedDps []proto.FailedDpInfo
@ -557,7 +557,7 @@ func (dd *DecommissionDisk) GetDecommissionFailedAndRunningDPByTerm(c *Cluster)
}
}
log.LogWarnf("action[GetDecommissionFailedAndRunningDPByTerm] failed dp list [%v]", failedDps)
return failedDps, runningDps
return len(partitions), failedDps, runningDps
}
func (dd *DecommissionDisk) GetDecommissionFailedDP(c *Cluster) (error, []uint64) {

View File

@ -451,6 +451,7 @@ type DecommissionProgress struct {
StatusMessage string
Progress string
TotalDpCnt int
RemainingDpCnt int
RunningDps []uint64
FailedDps []FailedDpInfo
IgnoreDps []IgnoreDecommissionDP
@ -461,14 +462,15 @@ type DecommissionProgress struct {
}
type DataDecommissionProgress struct {
Status uint32
StatusMessage string
Progress string
TotalDpCnt int
RunningDps []uint64
FailedDps []FailedDpInfo
IgnoreDps []IgnoreDecommissionDP
ResidualDps []IgnoreDecommissionDP
Status uint32
StatusMessage string
Progress string
TotalDpCnt int
RemainingDpCnt int
RunningDps []uint64
FailedDps []FailedDpInfo
IgnoreDps []IgnoreDecommissionDP
ResidualDps []IgnoreDecommissionDP
}
type DiskInfo struct {