feat(master): add retries for ignore dataPartitions during disk decommission.

close: #1000329518

Signed-off-by: shuqiang-zheng <zhengshuqiang@oppo.com>
This commit is contained in:
shuqiang-zheng 2025-07-24 16:35:39 +08:00 committed by 贺迟
parent 83748a2a03
commit e1165a242e
2 changed files with 150 additions and 0 deletions

View File

@ -5310,6 +5310,120 @@ func (c *Cluster) checkBadDisk() {
})
}
func (c *Cluster) TryDecommissionRunningDiskIgnoreDps(disk *DecommissionDisk) {
var (
dp *DataPartition
node *DataNode
zone *Zone
ns *nodeSet
err error
rstMsg string
ignorePartitionIds = make([]uint64, 0)
ignorePartitions = make([]*DataPartition, 0)
)
defer func() {
auditlog.LogMasterOp("RunningDiskDecommissionIgnoredDps", rstMsg, err)
}()
for _, ignoreDecommissionDpInfo := range disk.IgnoreDecommissionDps {
if dp, err = c.getDataPartitionByID(ignoreDecommissionDpInfo.PartitionID); err != nil {
log.LogWarnf("action[TryDecommissionRunningDiskIgnoreDps] find dataPartition[%v] failed[%v]",
ignoreDecommissionDpInfo.PartitionID, err.Error())
} else {
ignorePartitions = append(ignorePartitions, dp)
ignorePartitionIds = append(ignorePartitionIds, ignoreDecommissionDpInfo.PartitionID)
}
}
log.LogInfof("action[TryDecommissionRunningDiskIgnoreDps] disk[%v_%v] len(%v) ignorePartitionIds %v",
disk.SrcAddr, disk.DiskPath, len(ignorePartitionIds), ignorePartitionIds)
ignorePartitionIds = ignorePartitionIds[:0]
if len(ignorePartitions) == 0 {
log.LogInfof("action[TryDecommissionRunningDiskIgnoreDps] no any ignore partitions on disk[%v_%v]",
disk.SrcAddr, disk.DiskPath)
rstMsg = fmt.Sprintf("no any ignore partitions on disk[%v]", disk.decommissionInfo())
return
}
if node, err = c.dataNode(disk.SrcAddr); err != nil {
log.LogWarnf("action[TryDecommissionRunningDiskIgnoreDps] cannot find dataNode[%s]", disk.SrcAddr)
return
}
if zone, err = c.t.getZone(node.ZoneName); err != nil {
log.LogWarnf("action[TryDecommissionRunningDiskIgnoreDps] find datanode[%s] zone failed[%v]",
node.Addr, err.Error())
return
}
if ns, err = zone.getNodeSet(node.NodeSetID); err != nil {
log.LogWarnf("action[TryDecommissionRunningDiskIgnoreDps] find datanode[%s] nodeset[%v] failed[%v]",
node.Addr, node.NodeSetID, err.Error())
return
}
var ignoreAgainIDs []uint64
IgnoreAgainDecommissionDps := make([]proto.IgnoreDecommissionDP, 0)
for _, ignoreDp := range ignorePartitions {
triggerCondition := fmt.Sprintf("disk(%v)_%v_dp(%v)", disk.SrcAddr+"_"+disk.DiskPath, disk.Type, ignoreDp.PartitionID)
if err = ignoreDp.MarkDecommissionStatus(node.Addr, disk.DstAddr, disk.DiskPath, 0, disk.DecommissionRaftForce,
disk.DecommissionTerm, disk.Type, disk.DecommissionWeight, c, ns); err != nil {
if strings.Contains(err.Error(), proto.ErrDecommissionDiskErrDPFirst.Error()) {
c.syncUpdateDataPartition(ignoreDp)
// still decommission dp but not involved in the calculation of the decommission progress.
// disk.DecommissionDpTotal -= 1
ns.AddToDecommissionDataPartitionList(ignoreDp, c)
ignoreAgainIDs = append(ignoreAgainIDs, ignoreDp.PartitionID)
IgnoreAgainDecommissionDps = append(IgnoreAgainDecommissionDps, proto.IgnoreDecommissionDP{
PartitionID: ignoreDp.PartitionID,
ErrMsg: proto.ErrDecommissionDiskErrDPFirst.Error(),
})
continue
} else if strings.Contains(err.Error(), proto.ErrPerformingDecommission.Error()) {
if ignoreDp.DecommissionSrcAddr != node.Addr || ignoreDp.DecommissionType == AutoAddReplica {
// disk.DecommissionDpTotal -= 1
ignoreAgainIDs = append(ignoreAgainIDs, ignoreDp.PartitionID)
log.LogWarnf("action[TryDecommissionRunningDiskIgnoreDps] disk(%v) dp(%v) is decommissioning",
disk.decommissionInfo(), ignoreDp.PartitionID)
IgnoreAgainDecommissionDps = append(IgnoreAgainDecommissionDps, proto.IgnoreDecommissionDP{
PartitionID: ignoreDp.PartitionID,
ErrMsg: proto.ErrPerformingDecommission.Error(),
})
continue
} else {
log.LogDebugf("action[TryDecommissionRunningDiskIgnoreDps] disk(%v) dp(%v) may be mark decommission before leader change",
disk.decommissionInfo(), ignoreDp.PartitionID)
ns.AddToDecommissionDataPartitionList(ignoreDp, c)
}
} else if strings.Contains(err.Error(), proto.ErrWaitForAutoAddReplica.Error()) {
ignoreAgainIDs = append(ignoreAgainIDs, ignoreDp.PartitionID)
log.LogWarnf("action[TryDecommissionRunningDiskIgnoreDps] disk(%v) dp(%v) is auto add replica",
disk.decommissionInfo(), ignoreDp.PartitionID)
IgnoreAgainDecommissionDps = append(IgnoreAgainDecommissionDps, proto.IgnoreDecommissionDP{
PartitionID: ignoreDp.PartitionID,
ErrMsg: proto.ErrWaitForAutoAddReplica.Error(),
})
continue
} else {
// mark as failed and set decommission src, make sure it can be included in the calculation of progress
ignoreDp.DecommissionSrcAddr = node.Addr
ignoreDp.DecommissionSrcDiskPath = disk.DiskPath
ignoreDp.markRollbackFailed(false, triggerCondition, err.Error())
ignoreDp.DecommissionErrorMessage = err.Error()
ignoreDp.DecommissionTerm = disk.DecommissionTerm
ignoreDp.addRetryTimesByDiskPath(ignoreDp.DecommissionSrcAddr + "_" + ignoreDp.DecommissionSrcDiskPath)
log.LogWarnf("action[TryDecommissionRunningDiskIgnoreDps] disk(%v) set dp(%v) DecommissionTerm %v",
disk.decommissionInfo(), ignoreDp.PartitionID, disk.DecommissionTerm)
}
} else {
if ignoreDp.GetDecommissionStatus() == markDecommission {
ns.AddToDecommissionDataPartitionList(ignoreDp, c)
}
}
c.syncUpdateDataPartition(ignoreDp)
ignorePartitionIds = append(ignorePartitionIds, ignoreDp.PartitionID)
}
disk.IgnoreDecommissionDps = IgnoreAgainDecommissionDps
rstMsg = fmt.Sprintf("disk[%v] ignorePartitionIds %v offline successfully, ignoreAgain (%v) %v",
disk.decommissionInfo(), ignorePartitionIds, len(ignoreAgainIDs), ignoreAgainIDs)
log.LogInfof("action[TryDecommissionRunningDiskIgnoreDps] %s", rstMsg)
}
func (c *Cluster) TryDecommissionDisk(disk *DecommissionDisk) {
var (
node *DataNode

View File

@ -1202,6 +1202,25 @@ func (ns *nodeSet) removeAutoDecommissionDisk(dd *DecommissionDisk) {
ns.autoDecommissionDiskList.Remove(ns.ID, dd)
}
func (ns *nodeSet) findRunningDiskAndDecommissionIgnoreDps(c *Cluster) {
// try to decommission the dps of running disks that have been ignored in this round of decommissioning
allDecommissionRunningDisks := ns.manualDecommissionDiskList.PopDecommissionRunningDisk()
if c.AutoDecommissionDiskIsEnabled() {
allAutoDecommissionRunningDisks := ns.autoDecommissionDiskList.PopDecommissionRunningDisk()
allDecommissionRunningDisks = append(allDecommissionRunningDisks, allAutoDecommissionRunningDisks...)
}
sort.Slice(allDecommissionRunningDisks, func(i, j int) bool {
return allDecommissionRunningDisks[i].DecommissionWeight > allDecommissionRunningDisks[j].DecommissionWeight
})
log.LogDebugf("ns %v(%p) traverseDecommissionRunningDisk traverse allDecommissionRunningDiskCnt %v",
ns.ID, ns, len(allDecommissionRunningDisks))
if len(allDecommissionRunningDisks) > 0 {
for _, decommissionRunningDisk := range allDecommissionRunningDisks {
c.TryDecommissionRunningDiskIgnoreDps(decommissionRunningDisk)
}
}
}
func (ns *nodeSet) traverseDecommissionDisk(c *Cluster) {
t := time.NewTicker(DecommissionInterval)
// wait for loading all decommissionDisk when reload metadata
@ -1248,6 +1267,8 @@ func (ns *nodeSet) traverseDecommissionDisk(c *Cluster) {
return true
})
ns.findRunningDiskAndDecommissionIgnoreDps(c)
decommissionDiskCnt, allDecommissionDisks := ns.manualDecommissionDiskList.PopMarkDecommissionDisk(0)
if c.AutoDecommissionDiskIsEnabled() {
autoDecommissionDiskCnt, allAutoDecommissionDisks := ns.autoDecommissionDiskList.PopMarkDecommissionDisk(0)
@ -2579,6 +2600,21 @@ func (l *DecommissionDiskList) PopMarkDecommissionDisk(limit int) (count int, co
return count, collection
}
func (l *DecommissionDiskList) PopDecommissionRunningDisk() (collection []*DecommissionDisk) {
l.mu.Lock()
defer l.mu.Unlock()
collection = make([]*DecommissionDisk, 0)
for elm := l.decommissionList.Front(); elm != nil; elm = elm.Next() {
disk := elm.Value.(*DecommissionDisk)
if disk.GetDecommissionStatus() != DecommissionRunning {
continue
}
collection = append(collection, disk)
log.LogDebugf("action[PopDecommissionRunningDisk] pop disk[%v]", disk)
}
return collection
}
func (l *DecommissionDataPartitionList) Has(id uint64) bool {
l.mu.Lock()
_, ok := l.cacheMap[id]