fix(master): fix for dp no leader due to failure to add raft member during decommission.

close:#1000105754

Signed-off-by: shuqiang-zheng <zhengshuqiang@oppo.com>
This commit is contained in:
shuqiang-zheng 2025-05-08 17:06:46 +08:00 committed by zhumingze1108
parent a3cf843526
commit 89c1915676
3 changed files with 14 additions and 10 deletions

View File

@ -1822,7 +1822,7 @@ func (m *Server) addDataReplica(w http.ResponseWriter, r *http.Request) {
break
}
if err = m.cluster.addDataReplica(dp, addr, false); err != nil {
if err = m.cluster.addDataReplica(dp, addr, false, false); err != nil {
sendErrReply(w, r, newErrHTTPReply(err))
return
}
@ -8404,7 +8404,7 @@ func (m *Server) recoverBackupDataReplica(w http.ResponseWriter, r *http.Request
addPeer := proto.Peer{ID: dataNode.ID, Addr: addr, HeartbeatPort: dataNode.HeartbeatPort, ReplicaPort: dataNode.ReplicaPort}
log.LogInfof("action[recoverBackupDataReplica] dp %v dst addr %v try add raft member, node id %v", dp.PartitionID, addr, dataNode.ID)
if err = m.cluster.addDataPartitionRaftMember(dp, addPeer); err != nil {
if err = m.cluster.addDataPartitionRaftMember(dp, addPeer, false); err != nil {
log.LogWarnf("action[recoverBackupDataReplica] dp %v addr %v try add raft member err [%v]", dp.PartitionID, addr, err)
sendErrReply(w, r, newErrHTTPReply(err))
return

View File

@ -2537,7 +2537,7 @@ func (c *Cluster) decommissionSingleDp(dp *DataPartition, newAddr, offlineAddr s
}()
// 1. add new replica first
if dp.GetSpecialReplicaDecommissionStep() == SpecialDecommissionEnter {
if err = c.addDataReplica(dp, newAddr, false); err != nil {
if err = c.addDataReplica(dp, newAddr, true, false); err != nil {
err = fmt.Errorf("action[decommissionSingleDp] dp %v addDataReplica %v fail err %v", dp.PartitionID, newAddr, err)
goto ERR
}
@ -2934,7 +2934,7 @@ func (c *Cluster) migrateDataPartition(srcAddr, targetAddr string, dp *DataParti
if err = c.removeDataReplica(dp, srcAddr, false, raftForce); err != nil {
goto errHandler
}
if err = c.addDataReplica(dp, newAddr, false); err != nil {
if err = c.addDataReplica(dp, newAddr, false, false); err != nil {
goto errHandler
}
@ -3021,7 +3021,7 @@ func (c *Cluster) validateDecommissionDataPartition(dp *DataPartition, offlineAd
return
}
func (c *Cluster) addDataReplica(dp *DataPartition, addr string, ignoreDecommissionDisk bool) (err error) {
func (c *Cluster) addDataReplica(dp *DataPartition, addr string, needRollBack, ignoreDecommissionDisk bool) (err error) {
defer func() {
if err != nil {
log.LogErrorf("action[addDataReplica],vol[%v],dp %v ,err[%v]", dp.VolName, dp.PartitionID, err)
@ -3054,7 +3054,7 @@ func (c *Cluster) addDataReplica(dp *DataPartition, addr string, ignoreDecommiss
}
log.LogInfof("action[addDataReplica] dp %v dst addr %v try add raft member, node id %v", dp.PartitionID, addr, targetDataNode.ID)
if err = c.addDataPartitionRaftMember(dp, addPeer); err != nil {
if err = c.addDataPartitionRaftMember(dp, addPeer, needRollBack); err != nil {
log.LogWarnf("action[addDataReplica] dp %v addr %v try add raft member err [%v]", dp.PartitionID, addr, err)
return
}
@ -3117,7 +3117,7 @@ func (c *Cluster) returnDataSize(addr string, dp *DataPartition) {
dataNode.AvailableSpace += leaderSize
}
func (c *Cluster) buildAddDataPartitionRaftMemberTaskAndSyncSendTask(dp *DataPartition, addPeer proto.Peer, leaderAddr string) (resp *proto.Packet, err error) {
func (c *Cluster) buildAddDataPartitionRaftMemberTaskAndSyncSendTask(dp *DataPartition, addPeer proto.Peer, leaderAddr string, needRollBack bool) (resp *proto.Packet, err error) {
log.LogInfof("action[buildAddDataPartitionRaftMemberTaskAndSyncSendTask] add peer [%v] start", addPeer)
defer func() {
var resultCode uint8
@ -3139,13 +3139,17 @@ func (c *Cluster) buildAddDataPartitionRaftMemberTaskAndSyncSendTask(dp *DataPar
return
}
if resp, err = leaderDataNode.TaskManager.syncSendAdminTask(task); err != nil {
if needRollBack {
dp.DecommissionNeedRollback = true
c.syncUpdateDataPartition(dp)
}
return
}
log.LogInfof("action[buildAddDataPartitionRaftMemberTaskAndSyncSendTask] add peer [%v] finished", addPeer)
return
}
func (c *Cluster) addDataPartitionRaftMember(dp *DataPartition, addPeer proto.Peer) (err error) {
func (c *Cluster) addDataPartitionRaftMember(dp *DataPartition, addPeer proto.Peer, needRollBack bool) (err error) {
var (
candidateAddrs []string
leaderAddr string
@ -3172,7 +3176,7 @@ func (c *Cluster) addDataPartitionRaftMember(dp *DataPartition, addPeer proto.Pe
if leaderAddr == "" && len(candidateAddrs) < int(dp.ReplicaNum) {
time.Sleep(retrySendSyncTaskInternal)
}
_, err = c.buildAddDataPartitionRaftMemberTaskAndSyncSendTask(dp, addPeer, host)
_, err = c.buildAddDataPartitionRaftMemberTaskAndSyncSendTask(dp, addPeer, host, needRollBack)
if err == nil {
break
} else {

View File

@ -1691,7 +1691,7 @@ func (partition *DataPartition) Decommission(c *Cluster) bool {
if err = c.removeDataReplica(partition, srcAddr, false, partition.DecommissionRaftForce); err != nil {
goto errHandler
}
if err = c.addDataReplica(partition, targetAddr, false); err != nil {
if err = c.addDataReplica(partition, targetAddr, true, false); err != nil {
goto errHandler
}
newReplica, _ := partition.getReplica(targetAddr)