fix(master): fix new replica repair blocking due to host0 unavailability during dp decommission.

close:#1000059565

Signed-off-by: shuqiang-zheng <zhengshuqiang@oppo.com>
This commit is contained in:
shuqiang-zheng 2025-04-15 20:19:30 +08:00 committed by zhumingze1108
parent 40d975d7a9
commit fdf97e54ad
6 changed files with 150 additions and 16 deletions

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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