diff --git a/cli/cmd/fmt.go b/cli/cmd/fmt.go index 6f111604e..d3870884c 100644 --- a/cli/cmd/fmt.go +++ b/cli/cmd/fmt.go @@ -879,7 +879,7 @@ func formatTimeToString(t time.Time) string { var dataReplicaTableRowPattern = "%-65v %-12v %-12v %-12v %-12v %-12v %-12v %-12v %-18v %-10v" func formatDataReplicaTableHeader() string { - return fmt.Sprintf(dataReplicaTableRowPattern, "ADDR", "USEDSIZE", "TOTALSIZE", "ISLEADER", "FILECOUNT", "HASLOADRESPONSE", "NEEDSTOCOMPARE", "STATUS", "DISKPATH", "REPORT TIME") + return fmt.Sprintf(dataReplicaTableRowPattern, "ADDR", "USEDSIZE", "TOTALSIZE", "ISLEADER", "FILECOUNT", "HASLOADRESPONSE", "NEEDSTOCOMPARE", "ISREPAIRING", "STATUS", "DISKPATH", "REPORT TIME") } var dataFileInCoreTableRowPattern = "%-12v %-12v %-10v %-10v" @@ -919,7 +919,7 @@ func formatDataReplica(index int, replica *proto.DataReplica, rowTable bool) str if rowTable { return fmt.Sprintf(dataReplicaTableRowPattern, formatAddr(replica.Addr, replica.DomainAddr), formatSize(replica.Used), formatSize(replica.Total), replica.IsLeader, replica.FileCount, - replica.HasLoadResponse, replica.NeedsToCompare, formatDataPartitionStatus(replica.Status), + replica.HasLoadResponse, replica.NeedsToCompare, replica.IsRepairing, formatDataPartitionStatus(replica.Status), replica.DiskPath, formatTime(replica.ReportTime)) } return alignColumnIndex(index, @@ -930,6 +930,7 @@ func formatDataReplica(index int, replica *proto.DataReplica, rowTable bool) str arow("FileCount", replica.FileCount), arow("HasLoadResponse", replica.HasLoadResponse), arow("NeedsToCompare", replica.NeedsToCompare), + arow("IsRepairing", replica.IsRepairing), arow("Status", formatDataPartitionStatus(replica.Status)), arow("DiskPath", replica.DiskPath), arow("ReportTime", formatTime(replica.ReportTime)), diff --git a/datanode/partition.go b/datanode/partition.go index ba7dcc9e8..4dfc86665 100644 --- a/datanode/partition.go +++ b/datanode/partition.go @@ -501,6 +501,42 @@ func newDataPartition(dpCfg *dataPartitionCfg, disk *Disk, isCreate bool) (dp *D return } +func (partition *DataPartition) HandleSetRepairingStatusOp(req *proto.SetDataPartitionRepairingStatusRequest) (err error) { + var ( + reqData []byte + pItem *RaftCmdItem + ) + if reqData, err = json.Marshal(req); err != nil { + return + } + pItem = &RaftCmdItem{ + Op: uint32(proto.OpSetRepairingStatus), + K: []byte("setRepairingStatus"), + V: reqData, + } + data, _ := MarshalRaftCmd(pItem) + _, err = partition.Submit(data) + return +} + +func (partition *DataPartition) fsmSetRepairingStatusOp(opItem *RaftCmdItem) (err error) { + req := new(proto.SetDataPartitionRepairingStatusRequest) + if err = json.Unmarshal(opItem.V, req); err != nil { + log.LogErrorf("action[fsmSetRepairingStatusOp] dp[%v] op item %v", partition.partitionID, opItem) + return + } + + oldStatus := partition.isRepairing + partition.isRepairing = req.RepairingStatus + if err = partition.PersistMetadata(); err != nil { + log.LogErrorf("action[fsmSetRepairingStatusOp] persist dp %v metadata failed, err: %v", partition.partitionID, err) + partition.isRepairing = oldStatus + return + } + log.LogInfof("action[fsmSetRepairingStatusOp] %v set repairingStatus %v success", partition.partitionID, req.RepairingStatus) + return +} + func (partition *DataPartition) HandleVersionOp(req *proto.MultiVersionOpRequest) (err error) { var ( verData []byte diff --git a/datanode/partition_raftfsm.go b/datanode/partition_raftfsm.go index 5f236fb41..e17365e03 100644 --- a/datanode/partition_raftfsm.go +++ b/datanode/partition_raftfsm.go @@ -51,6 +51,9 @@ func (dp *DataPartition) Apply(command []byte, index uint64) (resp interface{}, if opItem.Op == uint32(proto.OpVersionOp) { dp.fsmVersionOp(opItem) return + } else if opItem.Op == uint32(proto.OpSetRepairingStatus) { + dp.fsmSetRepairingStatusOp(opItem) + return } return } diff --git a/datanode/space_manager.go b/datanode/space_manager.go index b4c069b30..e61b421b9 100644 --- a/datanode/space_manager.go +++ b/datanode/space_manager.go @@ -760,6 +760,7 @@ func (s *DataNode) buildHeartBeatResponse(response *proto.DataNodeHeartbeatRespo ForbidWriteOpOfProtoVer0: dpForbid, ReadOnlyReasons: partition.ReadOnlyReasons(), IsMissingTinyExtent: partition.extentStore.AvailableTinyExtentCnt()+partition.extentStore.BrokenTinyExtentCnt() < storage.TinyExtentCount, + IsRepairing: partition.isRepairing, } log.LogDebugf("action[Heartbeats] dpid(%v), status(%v) total(%v) used(%v) leader(%v) isLeader(%v) "+ "TriggerDiskError(%v) reqId(%v) testID(%v) cost(%v).", diff --git a/datanode/wrap_operator.go b/datanode/wrap_operator.go index 6294495a8..f76b901f7 100644 --- a/datanode/wrap_operator.go +++ b/datanode/wrap_operator.go @@ -1806,14 +1806,15 @@ func (s *DataNode) handlePacketToSetRepairingStatus(p *repl.Packet) { log.LogWarnf("action[handlePacketToSetRepairStatus] cannot find dp %v", request.PartitionId) return } - oldStatus := dp.isRepairing - dp.isRepairing = request.RepairingStatus - if err = dp.PersistMetadata(); err != nil { - log.LogErrorf("action[handlePacketToSetRepairingStatus] persist dp %v metadata failed, err: %v", dp.partitionID, err) - dp.isRepairing = oldStatus + + _, isLeader := dp.IsRaftLeader() + if !isLeader { + err = raft.ErrNotLeader return } - log.LogInfof("action[handlePacketToSetRepairingStatus] %v set repairingStatus %v success", request.PartitionId, request.RepairingStatus) + + err = dp.HandleSetRepairingStatusOp(request) + log.LogInfof("action[handlePacketToSetRepairStatus] opcode %v dpid %v after raft submit err %v resultCode %v", p.Opcode, p.PartitionID, err) } func (s *DataNode) handlePacketToStopDataPartitionRepair(p *repl.Packet) { diff --git a/master/cluster.go b/master/cluster.go index f420bffc6..245251d9d 100644 --- a/master/cluster.go +++ b/master/cluster.go @@ -507,6 +507,7 @@ func (c *Cluster) scheduleTask() { c.scheduleToUpdateFlashGroupRespCache() c.scheduleStartBalanceTask() c.scheduleToUpdateFlashGroupSlots() + c.scheduleToCheckDataPartitionRepairingStatus() } func (c *Cluster) masterAddr() (addr string) { @@ -2545,7 +2546,7 @@ func (c *Cluster) decommissionSingleDp(dp *DataPartition, newAddr, offlineAddr s err = fmt.Errorf("action[decommissionSingleDp] dp %v addDataReplica %v fail err %v", dp.PartitionID, newAddr, err) goto ERR } - if err = c.setAllReplicasRepairingStatus(dp, true, true); err != nil { + if err = c.setDpRepairingStatus(dp, true, true); err != nil { err = fmt.Errorf("action[decommissionSingleDp] dp %v set all replicas repairingStatus to true fail err %v", dp.PartitionID, err) goto ERR } @@ -3077,21 +3078,6 @@ func (c *Cluster) addDataReplica(dp *DataPartition, addr string, needRollBack, i return } -func (c *Cluster) setAllReplicasRepairingStatus(dp *DataPartition, repairingStatus bool, needRollBack bool) (err error) { - for _, host := range dp.Hosts { - if !dp.setReplicaRepairingStatus(host, repairingStatus, c) { - err = fmt.Errorf("dp %v addr %v setRepairingStatus failed", dp.PartitionID, host) - log.LogWarnf("action[setAllReplicasRepairingStatus] dp %v addr %v setRepairingStatus failed", dp.PartitionID, host) - if needRollBack { - dp.DecommissionNeedRollback = true - c.syncUpdateDataPartition(dp) - } - return - } - } - return -} - // update datanode size with to replica size func (c *Cluster) updateDataNodeSize(addr string, dp *DataPartition) error { if len(dp.Replicas) == 0 { @@ -3140,6 +3126,87 @@ func (c *Cluster) returnDataSize(addr string, dp *DataPartition) { dataNode.AvailableSpace += leaderSize } +func (c *Cluster) buildSetDpRepairStatusTaskAndSyncSendTask(dp *DataPartition, repairingStatus bool, leaderAddr string) (resp *proto.Packet, err error) { + log.LogInfof("action[buildSetDpRepairStatusTaskAndSyncSendTask] dp[%v] repairStatus[%v] start", dp.PartitionID, repairingStatus) + defer func() { + var resultCode uint8 + if resp != nil { + resultCode = resp.ResultCode + } + if err != nil { + log.LogErrorf("vol[%v],data partition[%v],leader addr[%v],resultCode[%v],err[%v]", dp.VolName, dp.PartitionID, leaderAddr, resultCode, err) + } else { + log.LogWarnf("vol[%v],data partition[%v],leader addr[%v],resultCode[%v],err[%v]", dp.VolName, dp.PartitionID, leaderAddr, resultCode, err) + } + }() + task, err := dp.createTaskToSetRepairingStatus(leaderAddr, repairingStatus) + if err != nil { + return + } + leaderDataNode, err := c.dataNode(leaderAddr) + if err != nil { + return + } + if resp, err = leaderDataNode.TaskManager.syncSendAdminTask(task); err != nil { + return + } + log.LogInfof("action[buildSetDpRepairStatusTaskAndSyncSendTask] dp[%v] repairStatus[%v] finished", dp.PartitionID, repairingStatus) + return +} + +func (c *Cluster) setDpRepairingStatus(dp *DataPartition, repairingStatus bool, needRollBack bool) (err error) { + var ( + candidateAddrs []string + leaderAddr string + ) + + defer func() { + if err != nil && needRollBack { + dp.DecommissionNeedRollback = true + c.syncUpdateDataPartition(dp) + } + }() + + dp.RLock() + candidateAddrs = make([]string, 0, len(dp.Hosts)) + leaderAddr = dp.getLeaderAddr() + if leaderAddr != "" && contains(dp.Hosts, leaderAddr) { + candidateAddrs = append(candidateAddrs, leaderAddr) + } else { + leaderAddr = "" + } + for _, host := range dp.Hosts { + if host == leaderAddr { + continue + } + candidateAddrs = append(candidateAddrs, host) + } + dp.RUnlock() + + // send task to leader addr first,if need to retry,then send to other addr + for index, host := range candidateAddrs { + if leaderAddr == "" && len(candidateAddrs) < int(dp.ReplicaNum) { + time.Sleep(retrySendSyncTaskInternal) + } + _, err = c.buildSetDpRepairStatusTaskAndSyncSendTask(dp, repairingStatus, 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 { + return err + } + } + + if index < len(candidateAddrs)-1 { + time.Sleep(retrySendSyncTaskInternal) + } + } + return +} + 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() { @@ -6225,6 +6292,38 @@ func (c *Cluster) syncRecoverBackupDataPartitionReplica(host, disk string, dp *D return } +func (c *Cluster) scheduleToCheckDataPartitionRepairingStatus() { + c.runTask(&cTask{ + tickTime: time.Second * time.Duration(c.cfg.IntervalToCheckDataPartition), + name: "scheduleToCheckDataPartitionRepairingStatus", + function: func() (fin bool) { + if c.partition != nil && c.partition.IsRaftLeader() { + c.checkDataPartitionRepairingStatus() + } + return + }, + }) +} + +func (c *Cluster) checkDataPartitionRepairingStatus() { + vols := c.allVols() + for _, vol := range vols { + partitions := vol.dataPartitions.clonePartitions() + for _, dp := range partitions { + if !c.processDataPartitionDecommission(dp.PartitionID) { + for _, replica := range dp.Replicas { + if replica.IsRepairing { + if err := c.setDpRepairingStatus(dp, false, false); err != nil { + log.LogWarnf("action[checkDataPartitionRepairingStatus] dp(%v) set repairingStatus to false failed, err(%v)", dp.PartitionID, err) + } + break + } + } + } + } + } +} + func (c *Cluster) scheduleToCheckDataReplicaMeta() { c.runTask(&cTask{ tickTime: time.Second * time.Duration(c.cfg.IntervalToCheckDataPartition), diff --git a/master/data_partition.go b/master/data_partition.go index 129ce7564..743f232b6 100644 --- a/master/data_partition.go +++ b/master/data_partition.go @@ -217,6 +217,12 @@ func (partition *DataPartition) createTaskToTryToChangeLeader(addr string) (task return } +func (partition *DataPartition) createTaskToSetRepairingStatus(addr string, repairingStatus bool) (task *proto.AdminTask, err error) { + task = proto.NewAdminTask(proto.OpSetRepairingStatus, addr, newSetRepairingStatusRequest(partition.PartitionID, repairingStatus)) + partition.resetTaskID(task) + return +} + func (partition *DataPartition) createTaskToAddRaftMember(addPeer proto.Peer, leaderAddr string) (task *proto.AdminTask, err error) { task = proto.NewAdminTask(proto.OpAddDataPartitionRaftMember, leaderAddr, newAddDataPartitionRaftMemberRequest(partition.PartitionID, addPeer)) partition.resetTaskID(task) @@ -732,6 +738,7 @@ func (partition *DataPartition) updateMetric(vr *proto.DataPartitionReport, data replica.ForbidWriteOpOfProtoVer0 = vr.ForbidWriteOpOfProtoVer0 replica.ReadOnlyReasons = vr.ReadOnlyReasons replica.IsMissingTinyExtent = vr.IsMissingTinyExtent + replica.IsRepairing = vr.IsRepairing partition.setForbidWriteOpOfProtoVer0() if replica.IsLeader { partition.LeaderReportTime = time.Now().Unix() @@ -1727,7 +1734,7 @@ func (partition *DataPartition) Decommission(c *Cluster) bool { if err = c.addDataReplica(partition, targetAddr, true, false); err != nil { goto errHandler } - if err = c.setAllReplicasRepairingStatus(partition, true, true); err != nil { + if err = c.setDpRepairingStatus(partition, true, true); err != nil { goto errHandler } newReplica, _ := partition.getReplica(targetAddr) @@ -1883,7 +1890,7 @@ func (partition *DataPartition) rollback(c *Cluster) { partition.DecommissionErrorMessage = fmt.Sprintf("rollback failed:%v", err.Error()) return } - err = c.setAllReplicasRepairingStatus(partition, false, false) + err = c.setDpRepairingStatus(partition, false, false) if err != nil { // keep decommission status to failed for rollback log.LogWarnf("action[rollback] dp[%v] rollback to set all replicas repairingStatus to false failed:%v", @@ -2002,40 +2009,6 @@ func (partition *DataPartition) IsRollbackFailed() bool { atomic.LoadUint32(&partition.DecommissionNeedRollbackTimes) >= defaultDecommissionRollbackLimit } -func (partition *DataPartition) setReplicaRepairingStatus(replicaAddr string, repairingStatus bool, c *Cluster) bool { - const RetryMax = 5 - var ( - dataNode *DataNode - err error - retry = 0 - ) - for retry <= RetryMax { - if dataNode, err = c.dataNode(replicaAddr); err != nil { - retry++ - time.Sleep(time.Second) - log.LogWarnf("action[setReplicaRepairingStatus] dp[%v] can't find dataNode %v", partition.PartitionID, replicaAddr) - continue - } - task := partition.createTaskToSetRepairingStatus(replicaAddr, repairingStatus) - packet, err := dataNode.TaskManager.syncSendAdminTask(task) - if err != nil { - retry++ - time.Sleep(time.Second) - log.LogWarnf("action[setReplicaRepairingStatus] dp[%v] send repairingStatus set task failed %v", partition.PartitionID, err.Error()) - continue - } - log.LogDebugf("action[setReplicaRepairingStatus] dp[%v] send repairingStatus set task to replica %v packet %v", partition.PartitionID, replicaAddr, packet) - return true - } - return false -} - -func (partition *DataPartition) createTaskToSetRepairingStatus(addr string, repairingStatus bool) (task *proto.AdminTask) { - task = proto.NewAdminTask(proto.OpSetRepairingStatus, addr, newSetRepairingStatusRequest(partition.PartitionID, repairingStatus)) - partition.resetTaskID(task) - return -} - func (partition *DataPartition) pauseReplicaRepair(replicaAddr string, stop bool, c *Cluster) bool { index := partition.findReplica(replicaAddr) if index == -1 { diff --git a/master/topology.go b/master/topology.go index 9749a3a78..78fd68c9e 100644 --- a/master/topology.go +++ b/master/topology.go @@ -2376,8 +2376,8 @@ func (l *DecommissionDataPartitionList) traverse(c *Cluster) { } log.LogDebugf("[DecommissionListTraverse]ns %v(%p) traverse dp(%v)", l.nsId, l, dp.decommissionInfo()) if dp.IsDecommissionSuccess() { - if err := c.setAllReplicasRepairingStatus(dp, false, false); err != nil { - continue + if err := c.setDpRepairingStatus(dp, false, false); err != nil { + log.LogWarnf("action[DecommissionListTraverse]ns %v(%p) dp[%v] set repairStatus to false failed, err %v", l.nsId, l, dp.decommissionInfo(), err) } l.Remove(dp) dp.ReleaseDecommissionToken(c) @@ -2400,8 +2400,8 @@ func (l *DecommissionDataPartitionList) traverse(c *Cluster) { if !dp.tryRollback(c) { log.LogDebugf("action[DecommissionListTraverse]ns %v(%p) Remove dp[%v] for fail", l.nsId, l, dp.PartitionID) - if err := c.setAllReplicasRepairingStatus(dp, false, false); err != nil { - continue + if err := c.setDpRepairingStatus(dp, false, false); err != nil { + log.LogWarnf("action[DecommissionListTraverse]ns %v(%p) dp[%v] set repairStatus to false failed, err %v", l.nsId, l, dp.decommissionInfo(), err) } l.Remove(dp) // if dp is not removed from decommission list, do not reset RestoreReplica @@ -2419,10 +2419,16 @@ func (l *DecommissionDataPartitionList) traverse(c *Cluster) { l.nsId, l, dp.PartitionID) dp.ReleaseDecommissionToken(c) dp.ReleaseDecommissionFirstHostToken(c) + if err := c.setDpRepairingStatus(dp, false, false); err != nil { + log.LogWarnf("action[DecommissionListTraverse]ns %v(%p) dp[%v] set repairStatus to false failed, err %v", l.nsId, l, dp.decommissionInfo(), err) + } l.Remove(dp) dp.setRestoreReplicaStop() c.syncUpdateDataPartition(dp) } else if dp.IsDecommissionInitial() { // fixed done ,not release token + if err := c.setDpRepairingStatus(dp, false, false); err != nil { + log.LogWarnf("action[DecommissionListTraverse]ns %v(%p) dp[%v] set repairStatus to false failed, err %v", l.nsId, l, dp.decommissionInfo(), err) + } l.Remove(dp) dp.ResetDecommissionStatus() c.syncUpdateDataPartition(dp) diff --git a/proto/admin_proto.go b/proto/admin_proto.go index 8fecd6cb7..5991b0c3b 100644 --- a/proto/admin_proto.go +++ b/proto/admin_proto.go @@ -877,6 +877,7 @@ type DataPartitionReport struct { ForbidWriteOpOfProtoVer0 bool ReadOnlyReasons uint32 IsMissingTinyExtent bool + IsRepairing bool } type DataNodeQosResponse struct { diff --git a/proto/model.go b/proto/model.go index a5746e893..e587ebd23 100644 --- a/proto/model.go +++ b/proto/model.go @@ -381,6 +381,7 @@ type DataReplica struct { ForbidWriteOpOfProtoVer0 bool ReadOnlyReasons uint32 IsMissingTinyExtent bool + IsRepairing bool } // data partition diagnosis represents the inactive data nodes, corrupt data partitions, and data partitions lack of replicas