fix(master):if dp has emtpy replicas when execute decommission,return error

Signed-off-by: chihe <chihe@oppo.com>
This commit is contained in:
chihe 2024-06-07 19:16:26 +08:00 committed by longerfly
parent 83b51e236c
commit 55f8ce3b67
11 changed files with 251 additions and 24 deletions

View File

@ -1470,13 +1470,15 @@ func (dp *DataPartition) reload(s *SpaceManager) error {
disk := dp.disk
rootDir := dp.path
log.LogDebugf("data partition disk %v rootDir %v", disk, rootDir)
s.partitionMutex.Lock()
delete(s.partitions, dp.partitionID)
s.partitionMutex.Unlock()
s.DetachDataPartition(dp.partitionID)
dp.Stop()
dp.Disk().DetachDataPartition(dp)
log.LogDebugf("data partition %v is detached", dp.partitionID)
_, err := LoadDataPartition(rootDir, disk)
dp2, err := LoadDataPartition(rootDir, disk)
if err != nil {
return err
}
s.AttachPartition(dp2)
return err
}

View File

@ -767,3 +767,171 @@ func (s *DataNode) markDiskBroken(w http.ResponseWriter, r *http.Request) {
}
s.buildSuccessResp(w, "success")
}
func (s *DataNode) setDiskExtentReadLimitStatus(w http.ResponseWriter, r *http.Request) {
const (
paramStatus = "status"
)
if err := r.ParseForm(); err != nil {
err = fmt.Errorf("parse form fail: %v", err)
s.buildFailureResp(w, http.StatusBadRequest, err.Error())
return
}
status, err := strconv.ParseBool(r.FormValue(paramStatus))
if err != nil {
err = fmt.Errorf("parse param %v fail: %v", paramStatus, err)
s.buildFailureResp(w, http.StatusBadRequest, err.Error())
return
}
for _, disk := range s.space.disks {
disk.SetExtentRepairReadLimitStatus(status)
}
s.buildSuccessResp(w, "success")
}
type DiskExtentReadLimitInfo struct {
DiskPath string `json:"diskPath"`
ExtentReadLimitStatus bool `json:"extentReadLimitStatus"`
Dp uint64 `json:"dp"`
}
type DiskExtentReadLimitStatusResponse struct {
Infos []DiskExtentReadLimitInfo `json:"infos"`
}
func (s *DataNode) queryDiskExtentReadLimitStatus(w http.ResponseWriter, r *http.Request) {
resp := &DiskExtentReadLimitStatusResponse{}
for _, disk := range s.space.disks {
status, dp := disk.QueryExtentRepairReadLimitStatus()
resp.Infos = append(resp.Infos, DiskExtentReadLimitInfo{DiskPath: disk.Path, ExtentReadLimitStatus: status, Dp: dp})
}
s.buildSuccessResp(w, resp)
}
func (s *DataNode) detachDataPartition(w http.ResponseWriter, r *http.Request) {
const (
paramID = "id"
)
if err := r.ParseForm(); err != nil {
err = fmt.Errorf("parse form fail: %v", err)
s.buildFailureResp(w, http.StatusBadRequest, err.Error())
return
}
partitionID, err := strconv.ParseUint(r.FormValue(paramID), 10, 64)
if err != nil {
err = fmt.Errorf("parse param %v fail: %v", paramID, err)
s.buildFailureResp(w, http.StatusBadRequest, err.Error())
return
}
partition := s.space.Partition(partitionID)
if partition == nil {
s.buildFailureResp(w, http.StatusBadRequest, "partition not exist")
return
}
// store disk path and root of dp
disk := partition.disk
rootDir := partition.path
log.LogDebugf("data partition disk %v rootDir %v", disk, rootDir)
s.space.partitionMutex.Lock()
delete(s.space.partitions, partitionID)
s.space.partitionMutex.Unlock()
partition.Stop()
partition.Disk().DetachDataPartition(partition)
log.LogDebugf("data partition %v is detached", partitionID)
s.buildSuccessResp(w, "success")
}
func (s *DataNode) releaseDiskExtentReadLimitToken(w http.ResponseWriter, r *http.Request) {
const (
paramDisk = "disk"
)
if err := r.ParseForm(); err != nil {
err = fmt.Errorf("parse form fail: %v", err)
s.buildFailureResp(w, http.StatusBadRequest, err.Error())
return
}
diskPath := r.FormValue(paramDisk)
// store disk path and root of dp
disk, err := s.space.GetDisk(diskPath)
if err != nil {
log.LogErrorf("action[loadDataPartition] disk(%v) is not found err(%v).", diskPath, err)
s.buildFailureResp(w, http.StatusBadRequest, fmt.Sprintf("disk %v is not found", diskPath))
return
}
disk.ReleaseReadExtentToken()
s.buildSuccessResp(w, "success")
}
func (s *DataNode) loadDataPartition(w http.ResponseWriter, r *http.Request) {
const (
paramID = "id"
paramDisk = "disk"
)
if err := r.ParseForm(); err != nil {
err = fmt.Errorf("parse form fail: %v", err)
s.buildFailureResp(w, http.StatusBadRequest, err.Error())
return
}
partitionID, err := strconv.ParseUint(r.FormValue(paramID), 10, 64)
if err != nil {
err = fmt.Errorf("parse param %v fail: %v", paramID, err)
s.buildFailureResp(w, http.StatusBadRequest, err.Error())
return
}
partition := s.space.Partition(partitionID)
if partition != nil {
s.buildFailureResp(w, http.StatusBadRequest, "partition is already loaded")
return
}
diskPath := r.FormValue(paramDisk)
// store disk path and root of dp
disk, err := s.space.GetDisk(diskPath)
if err != nil {
log.LogErrorf("action[loadDataPartition] disk(%v) is not found err(%v).", diskPath, err)
s.buildFailureResp(w, http.StatusBadRequest, fmt.Sprintf("disk %v is not found", diskPath))
return
}
fileInfoList, err := os.ReadDir(disk.Path)
if err != nil {
log.LogErrorf("action[loadDataPartition] read dir(%v) err(%v).", disk.Path, err)
s.buildFailureResp(w, http.StatusBadRequest, fmt.Sprintf(" read dir(%v) err(%v)", disk.Path, err))
return
}
rootDir := ""
for _, fileInfo := range fileInfoList {
filename := fileInfo.Name()
if !disk.isPartitionDir(filename) {
if disk.isExpiredPartitionDir(filename) {
}
continue
}
if id, _, err := unmarshalPartitionName(filename); err != nil {
log.LogErrorf("action[RestorePartition] unmarshal partitionName(%v) from disk(%v) err(%v) ",
filename, disk.Path, err.Error())
continue
} else {
if id == partitionID {
rootDir = filename
}
}
}
if rootDir == "" {
log.LogErrorf("action[loadDataPartition] dp root not found in dir(%v) .", disk.Path)
s.buildFailureResp(w, http.StatusBadRequest, fmt.Sprintf("dp root not found in dir(%v)", disk.Path))
return
}
log.LogDebugf("data partition disk %v rootDir %v", disk, rootDir)
dp, err := LoadDataPartition(path.Join(diskPath, rootDir), disk)
if err != nil {
s.buildFailureResp(w, http.StatusBadRequest, err.Error())
} else {
s.space.AttachPartition(dp)
s.buildSuccessResp(w, "success")
}
}

View File

@ -1325,7 +1325,8 @@ func (s *DataNode) handlePacketToAddDataPartitionRaftMember(p *repl.Packet) {
return
}
log.LogInfof("action[handlePacketToAddDataPartitionRaftMember] %v, partition id %v", req.AddPeer, req.PartitionId)
log.LogInfof("action[handlePacketToAddDataPartitionRaftMember]req(%v) addPeer %v, partition id %v",
p.GetReqID(), req.AddPeer, req.PartitionId)
p.AddMesgLog(string(reqData))
dp := s.space.Partition(req.PartitionId)
@ -1348,11 +1349,11 @@ func (s *DataNode) handlePacketToAddDataPartitionRaftMember(p *repl.Packet) {
isRaftLeader, err = s.forwardToRaftLeader(dp, p, false)
if !isRaftLeader {
if err != nil {
log.LogWarnf("action[handlePacketToAddDataPartitionRaftMember]dp %v req %v forward to leader failed:%v",
dp.partitionID, p.GetReqID(), err)
log.LogWarnf("action[handlePacketToAddDataPartitionRaftMember]dp %v req %v addPeer %v forward to leader failed:%v",
dp.partitionID, p.GetReqID(), req.AddPeer, err)
} else {
log.LogWarnf("action[handlePacketToAddDataPartitionRaftMember]dp %v req %v forward to leader",
dp.partitionID, p.GetReqID())
log.LogWarnf("action[handlePacketToAddDataPartitionRaftMember]dp %v req %v addPeer %v forward to leader",
dp.partitionID, p.GetReqID(), req.AddPeer)
}
return
}
@ -1408,8 +1409,9 @@ func (s *DataNode) handlePacketToRemoveDataPartitionRaftMember(p *repl.Packet) {
p.GetReqID(), string(reqData), req.RemovePeer.Addr, dp.partitionID, dp.replicaNum, dp.config.Peers, dp.replicas)
p.PartitionID = req.PartitionId
// do not return error to keep decommission progress go forward
// do not check replica existence on leader for autoRemove enable, follower may be contains redundant peers
// do not return error to keep master decommission progress go forward
// do not check replica existence on leader for autoRemove enable, follower may be contains redundant peers, make sure
// leader can send remove wal logs to follower
if !dp.IsExistReplica(req.RemovePeer.Addr) && !req.Force && !req.AutoRemove {
log.LogWarnf("action[handlePacketToRemoveDataPartitionRaftMember]dp %v receive MasterCommand: req %v "+
"RemoveRaftPeer(%v) force(%v) autoRemove(%v) has not exist", dp.partitionID, p.GetReqID(), req.RemovePeer, req.Force, req.AutoRemove)
@ -1456,7 +1458,8 @@ func (s *DataNode) handlePacketToRemoveDataPartitionRaftMember(p *repl.Packet) {
}
}
}
if !found && !req.AutoRemove {
if !found && !req.AutoRemove && !req.Force {
err = errors.NewErrorf("cannot found peer(%v) in dp(%v) peers", req.RemovePeer.Addr, dp.partitionID)
log.LogWarnf("handlePacketToRemoveDataPartitionRaftMember:%v", err.Error())
return

View File

@ -211,6 +211,7 @@ func (s *raft) runApply() {
}
s.doStop()
s.resetApply()
log.LogWarnf("raft(%v) quit runApply", s.raftFsm.id)
}()
loopCount := 0
@ -271,6 +272,7 @@ func (s *raft) run() {
s.stopSnapping()
s.raftConfig.Storage.Close()
close(s.done)
log.LogWarnf("raft(%v) quit run", s.raftFsm.id)
}()
s.prevHardSt.Term = s.raftFsm.term

View File

@ -121,8 +121,8 @@ func (r *raftFsm) campaign(force bool, t CampaignType) {
}
li, lt := r.raftLog.lastIndexAndTerm()
if logger.IsEnableDebug() {
logger.Debug("[raft->campaign][%v,%v logterm: %d, index: %d] sent "+
"%v request to %v at term %d. raftFSM[%p]", msgType, r.id, r.config.ReplicateAddr, lt, li, id, r.term, r)
logger.Debug("[raft->campaign][%v,raft %v, term: %d] sent "+
"index %v request to %v at term %d. raftFSM[%p]", msgType, r.id, lt, li, id, r.term, r)
}
m := proto.GetMessage()

View File

@ -234,6 +234,8 @@ func (rs *RaftServer) Status(id uint64) (status *Status) {
if ok {
status = raft.status()
} else {
logger.Warn("raftServer cannot found, id:%d", id)
}
if status == nil {
status = &Status{

View File

@ -2473,6 +2473,9 @@ func (c *Cluster) addDataReplica(dp *DataPartition, addr string, ignoreDecommiss
// update datanode size with to replica size
func (c *Cluster) updateDataNodeSize(addr string, dp *DataPartition) error {
if len(dp.Replicas) == 0 {
return errors.NewErrorf("dp %v has empty replica", dp.decommissionInfo())
}
leaderSize := dp.Replicas[0].Used
dataNode, err := c.dataNode(addr)
if err != nil {
@ -2576,7 +2579,17 @@ func (c *Cluster) addDataPartitionRaftMember(dp *DataPartition, addPeer proto.Pe
_, err = c.buildAddDataPartitionRaftMemberTaskAndSyncSendTask(dp, addPeer, host)
if err == nil {
break
} else {
// if send to leader raise err, it may send to follower ,then follower forward
// this request to leader, return nil. so when leader encounter en error, should
// return err
if leaderAddr != "" && leaderAddr == host {
dp.Hosts = oldHosts
dp.Peers = oldPeers
return err
}
}
if index < len(candidateAddrs)-1 {
time.Sleep(retrySendSyncTaskInternal)
}

View File

@ -16,6 +16,7 @@ package master
import (
"fmt"
"github.com/cubefs/cubefs/util/auditlog"
"sync"
"sync/atomic"
"time"
@ -99,10 +100,12 @@ func (dataNode *DataNode) SetIoUtils(used map[string]float64) {
func (dataNode *DataNode) checkLiveness() {
dataNode.Lock()
defer dataNode.Unlock()
log.LogInfof("action[checkLiveness] datanode[%v] report time[%v],since report time[%v], need gap [%v]",
dataNode.Addr, dataNode.ReportTime, time.Since(dataNode.ReportTime), time.Second*time.Duration(defaultNodeTimeOutSec))
if time.Since(dataNode.ReportTime) > time.Second*time.Duration(defaultNodeTimeOutSec) {
dataNode.isActive = false
msg := fmt.Sprintf("datanode[%v] report time[%v],since report time[%v], need gap [%v]",
dataNode.Addr, dataNode.ReportTime, time.Since(dataNode.ReportTime), time.Second*time.Duration(defaultNodeTimeOutSec))
log.LogWarnf("action[checkLiveness] %v", msg)
auditlog.LogMasterOp("DataNodeLive", msg, nil)
}
}

View File

@ -616,9 +616,12 @@ func (partition *DataPartition) getLiveReplicasFromHosts(timeOutSec int64) (repl
if replica.isLive(partition.PartitionID, timeOutSec) {
replicas = append(replicas, replica)
} else {
msg := fmt.Sprintf("dp %v replica addr %v is unavailable, datanode active %v replica status %v and is active %v",
partition.PartitionID, replica.Addr, replica.dataNode.isActive, replica.Status, replica.isActive(timeOutSec))
replica.Status = proto.Unavailable
log.LogWarnf("action[getLiveReplicasFromHosts] vol %v dp %v replica %v is unavailable",
partition.VolName, partition.PartitionID, replica.Addr)
auditlog.LogMasterOp("DataPartitionReplicaStatus", msg, nil)
}
}
@ -2044,7 +2047,8 @@ func (partition *DataPartition) checkReplicaMeta(c *Cluster) (err error) {
}
// remove raft member
err = partition.createTaskToRemoveRaftMember(c, peer, force, true)
auditMsg = fmt.Sprintf("dp(%v) remove redundant peer %v force %v", partition.PartitionID, peer, force)
auditMsg = fmt.Sprintf("dp(%v) remove redundant peer %v force %v:to replica %v: LocalPeers%v",
partition.decommissionInfo(), peer, force, replica.Addr, replica.LocalPeers)
log.LogDebugf("action[checkReplicaMeta]%v, err %v", auditMsg, err)
auditlog.LogMasterOp("RestoreReplicaMeta", auditMsg, err)
if err != nil {
@ -2061,8 +2065,25 @@ func (partition *DataPartition) checkReplicaMeta(c *Cluster) (err error) {
redundantPeers := findPeersToDeleteByConfig(partition.Peers, replica.LocalPeers)
for _, peer := range redundantPeers {
err = c.removeHostMember(partition, peer)
auditMsg = fmt.Sprintf("dp(%v) remove redundant peer %v for master", partition.PartitionID, peer)
log.LogDebugf("action[checkReplicaMeta]%v: err %v", auditMsg, err)
auditMsg = fmt.Sprintf("dp(%v) remove redundant peer %v for master,base on replica %v,localPeers(%v) ",
partition.decommissionInfo(), peer, replica.Addr, replica.LocalPeers)
auditlog.LogMasterOp("RestoreReplicaMeta", auditMsg, err)
if err != nil {
return
}
// redundant peers on master may exist on dataNode, and the redundant replica will be
// added into partition.Replicas again by hear beat.
var dataNode *DataNode
dataNode, err = c.dataNode(peer.Addr)
auditMsg = fmt.Sprintf("dp(%v) cannot found datanode for replica %v ,base on replica %v,localPeers(%v) ",
partition.decommissionInfo(), peer.Addr, replica.Addr, replica.LocalPeers)
auditlog.LogMasterOp("RestoreReplicaMeta", auditMsg, err)
if err != nil {
return
}
err = c.deleteDataReplica(partition, dataNode)
auditMsg = fmt.Sprintf("dp(%v) remove redundant replica on %v for master,base on replica %v,localPeers(%v) ",
partition.decommissionInfo(), peer.Addr, replica.Addr, replica.LocalPeers)
auditlog.LogMasterOp("RestoreReplicaMeta", auditMsg, err)
if err != nil {
return
@ -2133,9 +2154,9 @@ func (partition *DataPartition) lostLeader(c *Cluster) bool {
}
func (partition *DataPartition) decommissionInfo() string {
return fmt.Sprintf("vol(%v)_dp(%v)_src(%v)_dst(%v)_hosts(%v)_retry(%v)_isRecover(%v)_status(%v)_specialStatus(%v)"+
return fmt.Sprintf("vol(%v)_dp(%v)_replicaNum(%v)_src(%v)_dst(%v)_hosts(%v)_retry(%v)_isRecover(%v)_status(%v)_specialStatus(%v)"+
"_needRollback(%v)_rollbackTimes(%v)_force(%v)_type(%v)_RestoreReplica(%v)_errMsg(%v)_discard(%v)_term(%v)",
partition.VolName, partition.PartitionID, partition.DecommissionSrcAddr, partition.DecommissionDstAddr,
partition.VolName, partition.PartitionID, partition.ReplicaNum, partition.DecommissionSrcAddr, partition.DecommissionDstAddr,
partition.Hosts, partition.DecommissionRetry, partition.isRecover, GetDecommissionStatusMessage(partition.GetDecommissionStatus()),
GetSpecialDecommissionStatusMessage(partition.GetSpecialReplicaDecommissionStep()), partition.DecommissionNeedRollback,
partition.DecommissionNeedRollbackTimes, partition.DecommissionRaftForce, GetDecommissionTypeMessage(partition.DecommissionType),
@ -2225,8 +2246,14 @@ func (partition *DataPartition) setRestoreReplicaStop() bool {
}
func (partition *DataPartition) tryRestoreReplicaMeta(c *Cluster, migrateType uint32) error {
// AutoAddReplica do not need to check meta for replica again
// AutoAddReplica do not need to check meta for replica again, only have to check
// dp is performing decommission
if migrateType == AutoAddReplica {
if partition.isPerformingDecommission(c) {
log.LogDebugf("action[checkReplicaMeta]dp(%v) is performing decommission, skip it",
partition.PartitionID)
return proto.ErrPerformingDecommission
}
return nil
}
//

View File

@ -191,8 +191,9 @@ func (c *Cluster) checkDiskRecoveryProgress() {
} else {
partition.DecommissionErrorMessage = ""
partition.SetDecommissionStatus(DecommissionSuccess) // can be readonly or readwrite
Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] replica %v has recovered success",
c.Name, partitionID, partition.DecommissionDstAddr))
Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] "+
"replica %v has recovered success,cost(%v)",
c.Name, partitionID, partition.DecommissionDstAddr, time.Since(partition.RecoverStartTime).String()))
}
partition.RLock()
err = c.syncUpdateDataPartition(partition)

View File

@ -17,6 +17,7 @@ package master
import (
"container/list"
"fmt"
"github.com/cubefs/cubefs/util/auditlog"
"sort"
"strings"
"sync"
@ -2191,6 +2192,8 @@ func (l *DecommissionDataPartitionList) traverse(c *Cluster) {
if dp.IsDecommissionSuccess() {
l.Remove(dp)
dp.ReleaseDecommissionToken(c)
msg := fmt.Sprintf("dp %v decommission success, cost %v",
dp.decommissionInfo(), time.Since(dp.RecoverStartTime))
dp.ResetDecommissionStatus()
dp.setRestoreReplicaStop()
err := c.syncUpdateDataPartition(dp)
@ -2201,6 +2204,7 @@ func (l *DecommissionDataPartitionList) traverse(c *Cluster) {
log.LogDebugf("action[DecommissionListTraverse]Remove dp[%v] for success",
dp.PartitionID)
}
auditlog.LogMasterOp("TraverseDataPartition", msg, err)
} else if dp.IsDecommissionFailed() {
if !dp.tryRollback(c) {
log.LogDebugf("action[DecommissionListTraverse]Remove dp[%v] for fail",
@ -2212,6 +2216,8 @@ func (l *DecommissionDataPartitionList) traverse(c *Cluster) {
// rollback fail/success need release token
dp.ReleaseDecommissionToken(c)
c.syncUpdateDataPartition(dp)
msg := fmt.Sprintf("dp %v decommission failed", dp.decommissionInfo())
auditlog.LogMasterOp("TraverseDataPartition", msg, nil)
} else if dp.IsDecommissionPaused() {
log.LogDebugf("action[DecommissionListTraverse]Remove dp[%v] for paused ",
dp.PartitionID)