fix(master): set the repairingStatus of dp via raft.

close:#1000158166 #1000158134

Signed-off-by: shuqiang-zheng <zhengshuqiang@oppo.com>
This commit is contained in:
shuqiang-zheng 2025-06-04 20:21:21 +08:00 committed by zhumingze1108
parent 40ddcaafe0
commit 94ff433984
10 changed files with 186 additions and 64 deletions

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

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

View File

@ -877,6 +877,7 @@ type DataPartitionReport struct {
ForbidWriteOpOfProtoVer0 bool
ReadOnlyReasons uint32
IsMissingTinyExtent bool
IsRepairing bool
}
type DataNodeQosResponse struct {

View File

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