feat(master): add status update records to dp decommission process for querying.

close: #1000329528

Signed-off-by: shuqiang-zheng <zhengshuqiang@oppo.com>
This commit is contained in:
shuqiang-zheng 2025-08-11 19:25:02 +08:00 committed by chihe
parent eab8a04023
commit 17b991b351
14 changed files with 413 additions and 264 deletions

View File

@ -29,6 +29,7 @@ const (
CliOpDecommission = "decommission"
CliOpRecommission = "recommission"
CliOpQueryProgress = "query-progress"
CliOpQueryStatusUpdateRecords = "query-status-update-records"
CliOpQueryDiskStat = "query-disk-stat"
CliOpQueryNodeStat = "query-node-stat"
CliOpAbortDecommission = "abort-decommission"

View File

@ -46,22 +46,24 @@ func newDataPartitionCmd(client *master.MasterClient) *cobra.Command {
newDataPartitionResetRestoreStatusCmd(client),
newDataPartitionQueryDiskDecommissionInfoStat(client),
newDataPartitionQueryDataNodeDecommissionInfoStat(client),
newDataPartitionQueryDecommissionStatusUpdateRecords(client),
)
return cmd
}
const (
cmdDataPartitionGetShort = "Display detail information of a data partition"
cmdCheckCorruptDataPartitionShort = "Check and list unhealthy data partitions"
cmdDataPartitionDecommissionShort = "Decommission a replication of the data partition to a new address"
cmdDataPartitionReplicateShort = "Add a replication of the data partition on a new address"
cmdDataPartitionDeleteReplicaShort = "Delete a replication of the data partition on a fixed address"
cmdDataPartitionGetDiscardShort = "Display all discard data partitions"
cmdDataPartitionSetDiscardShort = "Set discard flag for data partition"
cmdDataPartitionQueryDecommissionProgressShort = "Query data partition decommission progress"
cmdDataPartitionResetRestoreStatusShort = "Reset data partition restore status"
cmdDataPartitionQueryDiskDecommissionInfoStatShort = "Query data partition disk decommission info stat"
cmdDataPartitionQueryDataNodeDecommissionInfoStatShort = "Query data partition datanode decommission info stat"
cmdDataPartitionGetShort = "Display detail information of a data partition"
cmdCheckCorruptDataPartitionShort = "Check and list unhealthy data partitions"
cmdDataPartitionDecommissionShort = "Decommission a replication of the data partition to a new address"
cmdDataPartitionReplicateShort = "Add a replication of the data partition on a new address"
cmdDataPartitionDeleteReplicaShort = "Delete a replication of the data partition on a fixed address"
cmdDataPartitionGetDiscardShort = "Display all discard data partitions"
cmdDataPartitionSetDiscardShort = "Set discard flag for data partition"
cmdDataPartitionQueryDecommissionProgressShort = "Query data partition decommission progress"
cmdDataPartitionResetRestoreStatusShort = "Reset data partition restore status"
cmdDataPartitionQueryDecommissionStatusUpdateRecordsShort = "Query data partition decommission status update records"
cmdDataPartitionQueryDiskDecommissionInfoStatShort = "Query data partition disk decommission info stat"
cmdDataPartitionQueryDataNodeDecommissionInfoStatShort = "Query data partition datanode decommission info stat"
)
func newDataPartitionGetCmd(client *master.MasterClient) *cobra.Command {
@ -680,6 +682,40 @@ func newDataPartitionResetRestoreStatusCmd(client *master.MasterClient) *cobra.C
return cmd
}
func newDataPartitionQueryDecommissionStatusUpdateRecords(client *master.MasterClient) *cobra.Command {
cmd := &cobra.Command{
Use: CliOpQueryStatusUpdateRecords + " [DATA PARTITION ID]",
Short: cmdDataPartitionQueryDecommissionStatusUpdateRecordsShort,
Args: cobra.MinimumNArgs(1),
Run: func(cmd *cobra.Command, args []string) {
var (
err error
dpId uint64
)
defer func() {
errout(err)
}()
dpId, err = strconv.ParseUint(args[0], 10, 64)
if err != nil {
return
}
records, err := client.AdminAPI().QueryDataPartitionDecommissionStatusUpdateRecords(dpId)
if err != nil {
return
}
if len(records) == 0 {
stdout("decommission status update records is empty, dp %v may not be in decommissioning\n", dpId)
return
}
stdout("%v", formatDataPartitionDecommissionStatusUpdateRecords(records))
},
}
return cmd
}
func newDataPartitionQueryDiskDecommissionInfoStat(client *master.MasterClient) *cobra.Command {
cmd := &cobra.Command{
Use: CliOpQueryDiskStat,

View File

@ -1406,6 +1406,21 @@ func formatDataPartitionDecommissionProgress(info *proto.DecommissionDataPartiti
return sb.String()
}
func formatDataPartitionDecommissionStatusUpdateRecords(records []*proto.DecommissionStatusRecord) string {
sb := strings.Builder{}
if len(records) != 0 {
sb.WriteString("decommission status update records: \n")
for _, record := range records {
sb.WriteString(fmt.Sprintf(" Condition : %v\n", record.Condition))
sb.WriteString(fmt.Sprintf(" Status : %v\n", record.Status))
sb.WriteString(fmt.Sprintf(" Time : %v\n", record.Time))
sb.WriteString(fmt.Sprintf(" ErrMessage : %v\n", record.ErrMessage))
sb.WriteString("\n")
}
}
return sb.String()
}
func formatDataPartitionDecommissionInfoStat(infos []*proto.DecommissionInfoStat) string {
sb := strings.Builder{}
if len(infos) != 0 {

View File

@ -1843,7 +1843,7 @@ func (m *Server) addDataReplica(w http.ResponseWriter, r *http.Request) {
dp.DecommissionType = ManualAddReplica
dp.RecoverStartTime = time.Now()
dp.RecoverUpdateTime = time.Now()
dp.SetDecommissionStatus(DecommissionRunning)
dp.SetDecommissionStatus(DecommissionRunning, "manualAddReplica", "")
var newReplica *DataReplica
if newReplica, err = dp.getReplica(addr); err != nil {
@ -2157,7 +2157,8 @@ func (m *Server) decommissionDataPartition(w http.ResponseWriter, r *http.Reques
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: rstMsg})
return
}
err = m.cluster.markDecommissionDataPartition(dp, node, dstNodeSet, raftForce, uint32(decommissionType), weight)
triggerCondition := fmt.Sprintf("manualDecommission_dp(%v)", dp.PartitionID)
err = m.cluster.markDecommissionDataPartition(dp, node, dstNodeSet, raftForce, uint32(decommissionType), weight, triggerCondition)
if err != nil {
sendErrReply(w, r, newErrHTTPReply(err))
return
@ -2296,6 +2297,27 @@ func (m *Server) resetDataPartitionDecommissionStatus(w http.ResponseWriter, r *
sendOkReply(w, r, newSuccessHTTPReply(msg))
}
func (m *Server) queryDataPartitionDecommissionStatusUpdateRecords(w http.ResponseWriter, r *http.Request) {
var (
dp *DataPartition
partitionID uint64
err error
records []*proto.DecommissionStatusRecord
)
if partitionID, err = parseRequestToLoadDataPartition(r); err != nil {
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()})
return
}
if dp, err = m.cluster.getDataPartitionByID(partitionID); err != nil {
sendErrReply(w, r, newErrHTTPReply(proto.ErrDataPartitionNotExists))
return
}
records = dp.cloneDecommissionStatusRecords()
sendOkReply(w, r, newSuccessHTTPReply(records))
}
func (m *Server) queryDataPartitionDecommissionStatus(w http.ResponseWriter, r *http.Request) {
var (
dp *DataPartition

View File

@ -516,8 +516,7 @@ func (c *Cluster) scheduleTask() {
c.scheduleStartBalanceTask()
c.scheduleToUpdateFlashGroupSlots()
c.scheduleToCheckDataPartitionRepairingStatus()
c.scheduleToCheckDataPartitionDecommissionDiskRetryMap()
c.scheduleToBalanceDataNode()
c.scheduleToCheckDataPartitionDecommissionInfoRecords()
}
func (c *Cluster) masterAddr() (addr string) {
@ -2569,7 +2568,7 @@ func (c *Cluster) decommissionSingleDp(dp *DataPartition, newAddr, offlineAddr s
}
// if addDataReplica is success, can add to BadDataPartitionIds
dp.SetSpecialReplicaDecommissionStep(SpecialDecommissionWaitAddRes)
dp.SetDecommissionStatus(DecommissionRunning)
dp.SetDecommissionStatus(DecommissionRunning, "decommission_singleDp_addNewReplica", "")
dp.isRecover = true
dp.Status = proto.ReadOnly
dp.RecoverUpdateTime = time.Now()
@ -2587,7 +2586,7 @@ func (c *Cluster) decommissionSingleDp(dp *DataPartition, newAddr, offlineAddr s
case decommContinue = <-dp.SpecialReplicaDecommissionStop: //
if !decommContinue {
err = fmt.Errorf("action[decommissionSingleDp] dp %v wait addDataReplica is stopped", dp.PartitionID)
dp.SetDecommissionStatus(DecommissionPause)
dp.SetDecommissionStatus(DecommissionPause, "decommission_singleDp_waitForRepair", err.Error())
log.LogWarnf("action[decommissionSingleDp] dp %v err:%v", dp.PartitionID, err)
goto ERR
}
@ -2712,7 +2711,7 @@ func (c *Cluster) decommissionSingleDp(dp *DataPartition, newAddr, offlineAddr s
case decommContinue = <-dp.SpecialReplicaDecommissionStop:
if !decommContinue {
err = fmt.Errorf("action[decommissionSingleDp] dp %v wait for leader selection is stopped", dp.PartitionID)
dp.SetDecommissionStatus(DecommissionPause)
dp.SetDecommissionStatus(DecommissionPause, "decommission_singleDp_waitForLeader", err.Error())
goto ERR
}
}
@ -2728,7 +2727,7 @@ func (c *Cluster) decommissionSingleDp(dp *DataPartition, newAddr, offlineAddr s
goto ERR
}
dp.SetSpecialReplicaDecommissionStep(SpecialDecommissionInitial)
dp.SetDecommissionStatus(DecommissionSuccess)
dp.SetDecommissionStatus(DecommissionSuccess, "decommission_singleDp_deleteOfflineReplica_complete", "")
// dp may not add into decommission list when master restart or leader change
dp.setRestoreReplicaStop()
c.syncUpdateDataPartition(dp)
@ -5299,8 +5298,9 @@ func (c *Cluster) handleDataNodeBadDisk(dataNode *DataNode) {
log.LogInfof("[handleDataNodeBadDisk] data node(%v) not found in dp(%v) maybe decommissioned?", dataNode.Addr, dpId)
continue
}
err = c.markDecommissionDataPartition(dp, dataNode, 0, false, AutoDecommission, highPriorityDecommissionWeight)
if err != nil {
triggerCondition := fmt.Sprintf("autoDecommission_diskErrDp(%v)", dp.PartitionID)
err = c.markDecommissionDataPartition(dp, dataNode, 0, false, AutoDecommission, highPriorityDecommissionWeight, triggerCondition)
if err != nil && !strings.Contains(err.Error(), proto.ErrPerformingDecommission.Error()) {
log.LogErrorf("[handleDataNodeBadDisk] failed to decommssion dp(%v) on data node(%v) disk(%v), err(%v)", dataNode.Addr, disk.DiskPath, dp.PartitionID, err)
continue
}
@ -5332,7 +5332,9 @@ func (c *Cluster) TryDecommissionRunningDiskIgnoreDps(disk *DecommissionDisk) {
ignorePartitions = make([]*DataPartition, 0)
)
defer func() {
auditlog.LogMasterOp("RunningDiskDecommissionIgnoredDps", rstMsg, err)
if len(ignorePartitionIds) != 0 || err != nil {
auditlog.LogMasterOp("RunningDiskDecommissionIgnoredDps", rstMsg, err)
}
}()
for _, ignoreDecommissionDpInfo := range disk.IgnoreDecommissionDps {
@ -5348,9 +5350,8 @@ func (c *Cluster) TryDecommissionRunningDiskIgnoreDps(disk *DecommissionDisk) {
disk.SrcAddr, disk.DiskPath, len(ignorePartitionIds), ignorePartitionIds)
ignorePartitionIds = ignorePartitionIds[:0]
if len(ignorePartitions) == 0 {
log.LogInfof("action[TryDecommissionRunningDiskIgnoreDps] no any ignore partitions on disk[%v_%v]",
disk.SrcAddr, disk.DiskPath)
rstMsg = fmt.Sprintf("no any ignore partitions on disk[%v]", disk.decommissionInfo())
log.LogInfof("action[TryDecommissionRunningDiskIgnoreDps] %v", rstMsg)
return
}
if node, err = c.dataNode(disk.SrcAddr); err != nil {
@ -5372,7 +5373,7 @@ func (c *Cluster) TryDecommissionRunningDiskIgnoreDps(disk *DecommissionDisk) {
for _, ignoreDp := range ignorePartitions {
triggerCondition := fmt.Sprintf("disk(%v)_%v_dp(%v)", disk.SrcAddr+"_"+disk.DiskPath, disk.Type, ignoreDp.PartitionID)
if err = ignoreDp.MarkDecommissionStatus(node.Addr, disk.DstAddr, disk.DiskPath, 0, disk.DecommissionRaftForce,
disk.DecommissionTerm, disk.Type, disk.DecommissionWeight, c, ns); err != nil {
disk.DecommissionTerm, disk.Type, disk.DecommissionWeight, c, ns, triggerCondition); err != nil {
if strings.Contains(err.Error(), proto.ErrDecommissionDiskErrDPFirst.Error()) {
c.syncUpdateDataPartition(ignoreDp)
// still decommission dp but not involved in the calculation of the decommission progress.
@ -5533,8 +5534,9 @@ func (c *Cluster) TryDecommissionDisk(disk *DecommissionDisk) {
ignoreIDs = append(ignoreIDs, dp.PartitionID)
continue
}
triggerCondition := fmt.Sprintf("disk(%v)_%v_dp(%v)", disk.SrcAddr+"_"+disk.DiskPath, disk.Type, dp.PartitionID)
if err = dp.MarkDecommissionStatus(node.Addr, disk.DstAddr, disk.DiskPath, 0, disk.DecommissionRaftForce,
disk.DecommissionTerm, disk.Type, disk.DecommissionWeight, c, ns); err != nil {
disk.DecommissionTerm, disk.Type, disk.DecommissionWeight, c, ns, triggerCondition); err != nil {
if strings.Contains(err.Error(), proto.ErrDecommissionDiskErrDPFirst.Error()) {
c.syncUpdateDataPartition(dp)
// still decommission dp but not involved in the calculation of the decommission progress.
@ -5575,7 +5577,7 @@ func (c *Cluster) TryDecommissionDisk(disk *DecommissionDisk) {
// mark as failed and set decommission src, make sure it can be included in the calculation of progress
dp.DecommissionSrcAddr = node.Addr
dp.DecommissionSrcDiskPath = disk.DiskPath
dp.markRollbackFailed(false)
dp.markRollbackFailed(false, triggerCondition, err.Error())
dp.DecommissionErrorMessage = err.Error()
dp.DecommissionTerm = disk.DecommissionTerm
dp.addRetryTimesByDiskPath(dp.DecommissionSrcAddr + "_" + dp.DecommissionSrcDiskPath)
@ -6350,7 +6352,7 @@ func (c *Cluster) rangeAllParitions(f func(d *DataPartition) bool) {
}
}
func (c *Cluster) markDecommissionDataPartition(dp *DataPartition, src *DataNode, dstNodeSetID uint64, raftForce bool, migrateType uint32, weight int) (err error) {
func (c *Cluster) markDecommissionDataPartition(dp *DataPartition, src *DataNode, dstNodeSetID uint64, raftForce bool, migrateType uint32, weight int, triggerCondition string) (err error) {
addr := src.Addr
replica, err := dp.getReplica(addr)
if err != nil {
@ -6368,9 +6370,10 @@ func (c *Cluster) markDecommissionDataPartition(dp *DataPartition, src *DataNode
return
}
if err = dp.MarkDecommissionStatus(addr, "", replica.DiskPath, dstNodeSetID, raftForce, uint64(time.Now().Unix()), migrateType, weight, c, ns); err != nil {
if !strings.Contains(err.Error(), proto.ErrDecommissionDiskErrDPFirst.Error()) {
dp.markRollbackFailed(false)
if err = dp.MarkDecommissionStatus(addr, "", replica.DiskPath, dstNodeSetID, raftForce, uint64(time.Now().Unix()), migrateType, weight, c, ns, triggerCondition); err != nil {
if !strings.Contains(err.Error(), proto.ErrDecommissionDiskErrDPFirst.Error()) && !strings.Contains(err.Error(), proto.ErrPerformingDecommission.Error()) &&
!strings.Contains(err.Error(), proto.ErrWaitForAutoAddReplica.Error()) {
dp.markRollbackFailed(false, triggerCondition, err.Error())
dp.DecommissionErrorMessage = err.Error()
c.syncUpdateDataPartition(dp)
return
@ -6378,9 +6381,9 @@ func (c *Cluster) markDecommissionDataPartition(dp *DataPartition, src *DataNode
}
// TODO: handle error
err = c.syncUpdateDataPartition(dp)
if err != nil {
return
updateErr := c.syncUpdateDataPartition(dp)
if updateErr != nil {
return errors.NewErrorf("dp(%v) mark decommission status failed, err(%v), updateErr(%v)", dp.PartitionID, err, updateErr)
}
if dp.GetDecommissionStatus() == markDecommission {
@ -6460,25 +6463,28 @@ func (c *Cluster) syncRecoverBackupDataPartitionReplica(host, disk string, dp *D
return
}
func (c *Cluster) scheduleToCheckDataPartitionDecommissionDiskRetryMap() {
func (c *Cluster) scheduleToCheckDataPartitionDecommissionInfoRecords() {
c.runTask(&cTask{
tickTime: time.Second * time.Duration(c.cfg.IntervalToCheckDataPartition),
name: "scheduleToCheckDataPartitionDecommissionDiskRetryMap",
name: "scheduleToCheckDataPartitionDecommissionInfoRecords",
function: func() (fin bool) {
if c.partition != nil && c.partition.IsRaftLeader() {
c.checkDataPartitionDecommissionDiskRetryMap()
c.checkDataPartitionDecommissionInfoRecords()
}
return
},
})
}
func (c *Cluster) checkDataPartitionDecommissionDiskRetryMap() {
func (c *Cluster) checkDataPartitionDecommissionInfoRecords() {
vols := c.allVols()
for _, vol := range vols {
partitions := vol.dataPartitions.clonePartitions()
for _, dp := range partitions {
dp.deleteInvalidRetryTimesRecord()
if dp.GetDecommissionStatus() == DecommissionInitial {
dp.clearDecommissionStatusRecords()
}
}
}
}

View File

@ -678,7 +678,7 @@ func (dataNode *DataNode) GetLatestDecommissionDataPartition(c *Cluster) (remain
if dd.GetDecommissionStatus() == markDecommission {
remainingDpCnt += dd.GetDecommissionTotalDpCnt(c)
} else {
remainingDpCnt += len(partitions)
remainingDpCnt += len(dps)
}
}
}

View File

@ -59,8 +59,9 @@ type DataPartition struct {
RdOnly bool
addReplicaMutex sync.RWMutex
DecommissionDiskRetryMapMutex sync.RWMutex
DecommissionInfoRecordMutex sync.RWMutex // used for decommissionDiskRetryMap and decommissionStatusUpdateRecords
DecommissionDiskRetryMap map[string]int
DecommissionStatusUpdateRecords []*proto.DecommissionStatusRecord
DecommissionRetry int
DecommissionStatus uint32
DecommissionSrcAddr string
@ -77,17 +78,18 @@ type DataPartition struct {
DecommissionWeight int
SpecialReplicaDecommissionStop chan bool // used for stop
SpecialReplicaDecommissionStep uint32
IsDiscard bool
VerSeq uint64
RecoverStartTime time.Time
RecoverUpdateTime time.Time
RecoverLastConsumeTime time.Duration
DecommissionRetryTime time.Time
RepairBlockSize uint64
DecommissionType uint32
RestoreReplica uint32
MediaType uint32
ForbidWriteOpOfProtoVer0 bool
proto.DecommissionInfoStat
IsDiscard bool
VerSeq uint64
RecoverStartTime time.Time
RecoverUpdateTime time.Time
RecoverLastConsumeTime time.Duration
DecommissionRetryTime time.Time
RepairBlockSize uint64
DecommissionType uint32
RestoreReplica uint32
MediaType uint32
ForbidWriteOpOfProtoVer0 bool
}
func newDataPartition(ID uint64, replicaNum uint8, volName string, volID uint64,
@ -103,6 +105,7 @@ func newDataPartition(ID uint64, replicaNum uint8, volName string, volID uint64,
partition.FilesWithMissingReplica = make(map[string]int64)
partition.MissingNodes = make(map[string]int64)
partition.DecommissionDiskRetryMap = make(map[string]int)
partition.DecommissionStatusUpdateRecords = make([]*proto.DecommissionStatusRecord, 0)
partition.Status = proto.ReadOnly
partition.VolName = volName
@ -1225,6 +1228,7 @@ func (partition *DataPartition) AcquireDecommissionFirstHostToken(c *Cluster) bo
diskToRepairDpInfo *DiskToDecommissionRepairDpInfo
dataNodeParallel uint64
)
defer c.syncUpdateDataPartition(partition)
for _, host := range partition.Hosts {
// for AutoAddReplica , firstHost does not need to consider the decommission source address since only adding and not deleting replica
@ -1289,7 +1293,8 @@ func (partition *DataPartition) AcquireDecommissionFirstHostToken(c *Cluster) bo
log.LogInfof("action[AcquireDecommissionFirstHostToken] dp(%v) acquire first host token(%v) success", partition.PartitionID, partition.DecommissionFirstHostDiskTokenKey)
return true
errHandle:
partition.markRollbackFailed(false)
triggerCondition := fmt.Sprintf("acquireFirsthostToken_firstHost(%v)", firstHost)
partition.markRollbackFailed(false, triggerCondition, err.Error())
partition.DecommissionErrorMessage = err.Error()
log.LogWarnf("action[AcquireDecommissionFirstHostToken] clusterID[%v] vol[%v] partitionID[%v]"+
" retry [%v] status [%v] DecommissionDstAddrSpecify [%v] DecommissionDstAddr [%v] DecommissionDstNodeSet [%v] failed",
@ -1298,6 +1303,34 @@ errHandle:
return false
}
func (partition *DataPartition) recordDecommissionStatus(condition string, errMsg string) {
partition.DecommissionInfoRecordMutex.Lock()
defer partition.DecommissionInfoRecordMutex.Unlock()
record := &proto.DecommissionStatusRecord{
Condition: condition,
Status: GetDecommissionStatusMessage(partition.DecommissionStatus),
Time: time.Now().Format("2006-01-02 15:04:05"),
ErrMessage: errMsg,
}
partition.DecommissionStatusUpdateRecords = append(partition.DecommissionStatusUpdateRecords, record)
}
func (partition *DataPartition) cloneDecommissionStatusRecords() []*proto.DecommissionStatusRecord {
partition.DecommissionInfoRecordMutex.RLock()
defer partition.DecommissionInfoRecordMutex.RUnlock()
records := make([]*proto.DecommissionStatusRecord, 0)
records = append(records, partition.DecommissionStatusUpdateRecords...)
return records
}
func (partition *DataPartition) clearDecommissionStatusRecords() {
partition.DecommissionInfoRecordMutex.Lock()
defer partition.DecommissionInfoRecordMutex.Unlock()
if len(partition.DecommissionStatusUpdateRecords) != 0 {
partition.DecommissionStatusUpdateRecords = make([]*proto.DecommissionStatusRecord, 0)
}
}
func isReplicasContainsHost(replicas []*DataReplica, host string) bool {
for _, replica := range replicas {
if replica.Addr == host {
@ -1308,7 +1341,7 @@ func isReplicasContainsHost(replicas []*DataReplica, host string) bool {
}
func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk string, dstNodeSetID uint64, raftForce bool, term uint64,
migrateType uint32, weight int, c *Cluster, ns *nodeSet,
migrateType uint32, weight int, c *Cluster, ns *nodeSet, triggerCondition string,
) (err error) {
defer func() {
if err != nil {
@ -1547,7 +1580,7 @@ directly:
break
}
partition.DecommissionRetry = 0
partition.SetDecommissionStatus(markDecommission)
partition.SetDecommissionStatus(markDecommission, triggerCondition, "")
// update decommissionTerm for next time query
partition.DecommissionTerm = term
partition.DecommissionWeight = weight
@ -1577,7 +1610,7 @@ directly:
// initial or failed restart
partition.ResetDecommissionStatus()
partition.DecommissionType = migrateType
partition.SetDecommissionStatus(markDecommission)
partition.SetDecommissionStatus(markDecommission, triggerCondition, "")
partition.DecommissionSrcAddr = srcAddr
partition.DecommissionDstAddr = dstAddr
partition.DecommissionSrcDiskPath = srcDisk
@ -1604,9 +1637,10 @@ directly:
return
}
func (partition *DataPartition) SetDecommissionStatus(status uint32) {
func (partition *DataPartition) SetDecommissionStatus(status uint32, triggerCondition string, errMsg string) {
log.LogDebugf("[SetDecommissionStatus] set dp(%v) decommission status to status(%v)", partition.PartitionID, status)
atomic.StoreUint32(&partition.DecommissionStatus, status)
partition.recordDecommissionStatus(triggerCondition, errMsg)
}
func (partition *DataPartition) SetSpecialReplicaDecommissionStep(step uint32) {
@ -1656,8 +1690,8 @@ func (partition *DataPartition) IsDoingDecommission() bool {
}
func (partition *DataPartition) cloneDecommissionDiskRetryMap() (result map[string]int) {
partition.DecommissionDiskRetryMapMutex.RLock()
defer partition.DecommissionDiskRetryMapMutex.RUnlock()
partition.DecommissionInfoRecordMutex.RLock()
defer partition.DecommissionInfoRecordMutex.RUnlock()
result = make(map[string]int)
for disk, retryTimes := range partition.DecommissionDiskRetryMap {
result[disk] = retryTimes
@ -1666,8 +1700,8 @@ func (partition *DataPartition) cloneDecommissionDiskRetryMap() (result map[stri
}
func (partition *DataPartition) addRetryTimesByDiskPath(diskPath string) {
partition.DecommissionDiskRetryMapMutex.Lock()
defer partition.DecommissionDiskRetryMapMutex.Unlock()
partition.DecommissionInfoRecordMutex.Lock()
defer partition.DecommissionInfoRecordMutex.Unlock()
if partition.DecommissionDiskRetryMap[diskPath] >= math.MaxInt {
partition.DecommissionDiskRetryMap[diskPath] = 0
} else {
@ -1676,29 +1710,29 @@ func (partition *DataPartition) addRetryTimesByDiskPath(diskPath string) {
}
func (partition *DataPartition) deleteRetryTimesRecordByDiskPath(diskPath string) {
partition.DecommissionDiskRetryMapMutex.Lock()
defer partition.DecommissionDiskRetryMapMutex.Unlock()
partition.DecommissionInfoRecordMutex.Lock()
defer partition.DecommissionInfoRecordMutex.Unlock()
delete(partition.DecommissionDiskRetryMap, diskPath)
}
func (partition *DataPartition) getRetryTimesRecordByDiskPath(diskPath string) (retryTimes int) {
partition.DecommissionDiskRetryMapMutex.RLock()
defer partition.DecommissionDiskRetryMapMutex.RUnlock()
partition.DecommissionInfoRecordMutex.RLock()
defer partition.DecommissionInfoRecordMutex.RUnlock()
retryTimes = partition.DecommissionDiskRetryMap[diskPath]
return retryTimes
}
func (partition *DataPartition) deleteInvalidRetryTimesRecord() {
partition.DecommissionDiskRetryMapMutex.RLock()
partition.DecommissionInfoRecordMutex.RLock()
if len(partition.DecommissionDiskRetryMap) == 0 {
partition.DecommissionDiskRetryMapMutex.RUnlock()
partition.DecommissionInfoRecordMutex.RUnlock()
return
}
diskRetryMap := make(map[string]int)
for disk, retryTimes := range partition.DecommissionDiskRetryMap {
diskRetryMap[disk] = retryTimes
}
partition.DecommissionDiskRetryMapMutex.RUnlock()
partition.DecommissionInfoRecordMutex.RUnlock()
for key := range diskRetryMap {
arr := strings.Split(key, "_")
if len(arr) == 2 {
@ -1708,9 +1742,9 @@ func (partition *DataPartition) deleteInvalidRetryTimesRecord() {
continue
}
}
partition.DecommissionDiskRetryMapMutex.Lock()
partition.DecommissionInfoRecordMutex.Lock()
delete(partition.DecommissionDiskRetryMap, key)
partition.DecommissionDiskRetryMapMutex.Unlock()
partition.DecommissionInfoRecordMutex.Unlock()
}
}
@ -1733,6 +1767,7 @@ func (partition *DataPartition) Decommission(c *Cluster) bool {
srcAddr = partition.DecommissionSrcAddr
targetAddr = partition.DecommissionDstAddr
srcReplica *DataReplica
triggerCondition string
resetDecommissionDst = true
begin = time.Now()
finalHosts = make([]string, len(partition.Hosts))
@ -1741,14 +1776,14 @@ func (partition *DataPartition) Decommission(c *Cluster) bool {
if partition.GetDecommissionStatus() == DecommissionInitial {
log.LogWarnf("action[decommissionDataPartition] dp [%v] may be cancel", partition.decommissionInfo())
partition.DecommissionErrorMessage = "cancel decommission"
partition.markRollbackFailed(false)
partition.markRollbackFailed(false, "decommission_statusInitial", "cancel decommission")
return false
}
if !c.AutoDecommissionDiskIsEnabled() && partition.DecommissionType == AutoDecommission {
log.LogWarnf("action[decommissionDataPartition] dp [%v] decommission is disable", partition.decommissionInfo())
partition.DecommissionErrorMessage = "disable auto " +
" decommission"
partition.markRollbackFailed(false)
partition.markRollbackFailed(false, "decommission_autoDecommissionCheck", "disable auto decommission")
return false
}
@ -1769,11 +1804,11 @@ func (partition *DataPartition) Decommission(c *Cluster) bool {
if partition.ReplicaNum == 1 && partition.DecommissionRaftForce {
log.LogWarnf("action[decommissionDataPartition] dp [%v] single replica does not support raftForce deletion", partition.decommissionInfo())
partition.DecommissionErrorMessage = "single replica does not support raftForce deletion"
partition.markRollbackFailed(false)
partition.markRollbackFailed(false, "decommission_raftForceCheck", "single replica does not support raftForce deletion")
return false
}
partition.SetDecommissionStatus(DecommissionPrepare)
partition.SetDecommissionStatus(DecommissionPrepare, "decommission_prepare", "")
err = c.syncUpdateDataPartition(partition)
if err != nil {
log.LogWarnf("action[decommissionDataPartition] dp [%v] update to prepare failed", partition.PartitionID)
@ -1789,7 +1824,7 @@ func (partition *DataPartition) Decommission(c *Cluster) bool {
// log.LogWarnf("[decommissionDataPartition] delete dp(%v) discard(%v)", partition.PartitionID, partition.IsDiscard)
// vol.deleteDataPartition(c, partition)
// }
partition.SetDecommissionStatus(DecommissionSuccess)
partition.SetDecommissionStatus(DecommissionSuccess, "decommission_discardCheck", "")
log.LogWarnf("action[decommissionDataPartition] skip dp(%v) discard(%v)", partition.PartitionID, partition.IsDiscard)
return true
}
@ -1801,7 +1836,8 @@ func (partition *DataPartition) Decommission(c *Cluster) bool {
srcReplica, _ = partition.getReplica(partition.DecommissionSrcAddr)
if len(partition.Replicas) == int(partition.ReplicaNum) && srcReplica == nil {
partition.SetDecommissionStatus(DecommissionSuccess)
triggerCondition = fmt.Sprintf("decommission_srcReplica(%v)_hasBeenDeleted", partition.DecommissionSrcAddr)
partition.SetDecommissionStatus(DecommissionSuccess, triggerCondition, "")
log.LogWarnf("action[decommissionDataPartition]dp(%v) status(%v) is already decommissioned",
partition.PartitionID, partition.Status)
return true
@ -1821,7 +1857,7 @@ func (partition *DataPartition) Decommission(c *Cluster) bool {
log.LogWarnf("action[decommissionDataPartition] %s", msg)
auditlog.LogMasterOp("DataPartitionDecommission", msg, nil)
partition.DecommissionErrorMessage = msg
partition.markRollbackFailed(false)
partition.markRollbackFailed(false, "decommission_raftForceCheck", msg)
return false
}
}
@ -1866,7 +1902,7 @@ func (partition *DataPartition) Decommission(c *Cluster) bool {
newReplica.Status = proto.Recovering // in case heartbeat response is not arrived
partition.isRecover = true
partition.Status = proto.ReadOnly
partition.SetDecommissionStatus(DecommissionRunning)
partition.SetDecommissionStatus(DecommissionRunning, "decommission_waitForRepair", "")
partition.RecoverUpdateTime = time.Now()
partition.RecoverStartTime = time.Now()
c.putBadDataPartitionIDsByDiskPath(partition.DecommissionSrcDiskPath, partition.DecommissionSrcAddr, partition.PartitionID)
@ -1897,17 +1933,18 @@ errHandler:
// if need rollback, set to fail
// do not reset DecommissionDstAddr outside the rollback operation, as it may cause rollback failure
if partition.DecommissionNeedRollback {
partition.SetDecommissionStatus(DecommissionFail)
partition.SetDecommissionStatus(DecommissionFail, "decommission_needRollBack", err.Error())
} else {
// The maximum number of retries for the DP error has been reached,
// and a rollback is still required, even if the rollback conditions have not been triggered.
if partition.DecommissionRetry >= defaultDecommissionRetryLimit {
partition.markRollbackFailed(true)
triggerCondition = fmt.Sprintf("decommission_retryOverLimit_count(%v)", partition.DecommissionRetry)
partition.markRollbackFailed(true, triggerCondition, err.Error())
} else {
// remove dp from BadDataPartitionIDs, preventing errors caused by disk manager not finding the replica
err := c.removeDPFromBadDataPartitionIDs(partition.DecommissionSrcAddr, partition.DecommissionSrcDiskPath, partition.PartitionID)
if err != nil {
log.LogWarnf("action[decommissionDataPartition] del dp[%v] from bad dataPartitionIDs failed:%v", partition.PartitionID, err)
removeErr := c.removeDPFromBadDataPartitionIDs(partition.DecommissionSrcAddr, partition.DecommissionSrcDiskPath, partition.PartitionID)
if removeErr != nil {
log.LogWarnf("action[decommissionDataPartition] del dp[%v] from bad dataPartitionIDs failed:%v", partition.PartitionID, removeErr)
}
partition.ReleaseDecommissionToken(c)
partition.ReleaseDecommissionFirstHostToken(c)
@ -1916,7 +1953,8 @@ errHandler:
partition.DecommissionDstAddr = ""
log.LogWarnf("action[decommissionDataPartition] partitionID:%v reset DecommissionDstAddr", partition.PartitionID)
}
partition.SetDecommissionStatus(markDecommission)
triggerCondition = fmt.Sprintf("decommission_retry_count(%v)", partition.DecommissionRetry)
partition.SetDecommissionStatus(markDecommission, triggerCondition, err.Error())
}
}
msg = fmt.Sprintf("clusterID[%v] info[%v] offline failed:%v consume[%v]seconds",
@ -1940,7 +1978,7 @@ func (partition *DataPartition) PauseDecommission(c *Cluster) bool {
partition.PartitionID, partition.GetDecommissionStatus())
if status == markDecommission {
partition.SetDecommissionStatus(DecommissionPause)
partition.SetDecommissionStatus(DecommissionPause, "pauseDecommission", "")
return true
}
if partition.isSpecialReplicaCnt() {
@ -1962,7 +2000,7 @@ func (partition *DataPartition) PauseDecommission(c *Cluster) bool {
partition.PartitionID, partition.GetDecommissionStatus())
}
}
partition.SetDecommissionStatus(DecommissionPause)
partition.SetDecommissionStatus(DecommissionPause, "pauseDecommission", "")
partition.isRecover = false
return true
}
@ -1980,13 +2018,14 @@ func (partition *DataPartition) ResetDecommissionStatus() {
partition.DecommissionDstNodeSet = 0
partition.DecommissionNeedRollback = false
atomic.StoreUint32(&partition.DecommissionNeedRollbackTimes, 0)
partition.SetDecommissionStatus(DecommissionInitial)
partition.SetDecommissionStatus(DecommissionInitial, "resetDecommissionStatus", "")
partition.SetSpecialReplicaDecommissionStep(SpecialDecommissionInitial)
partition.DecommissionErrorMessage = ""
partition.DecommissionType = InitialDecommission
partition.RecoverStartTime = time.Time{}
partition.RecoverUpdateTime = time.Time{}
partition.DecommissionRetryTime = time.Time{}
partition.clearDecommissionStatusRecords()
}
func (partition *DataPartition) resetRestoreMeta(expected uint32) (ok bool) {
@ -2030,7 +2069,7 @@ func (partition *DataPartition) rollback(c *Cluster) {
partition.isRecover = false
partition.DecommissionNeedRollback = false
partition.DecommissionErrorMessage = ""
partition.SetDecommissionStatus(markDecommission)
partition.SetDecommissionStatus(markDecommission, "rollback_complete", "")
partition.SetSpecialReplicaDecommissionStep(SpecialDecommissionInitial)
// specify dst addr do not need rollback
// keep DecommissionSrcAddr to prevent allocate DecommissionSrcAddr data node during acquire token
@ -2413,7 +2452,8 @@ errHandler:
partition.DecommissionRetry++
partition.DecommissionRetryTime = time.Now()
if partition.DecommissionRetry >= defaultDecommissionRetryLimit {
partition.markRollbackFailed(false)
triggerCondition := "acquireNsDecommissionToken"
partition.markRollbackFailed(false, triggerCondition, err.Error())
}
partition.DecommissionErrorMessage = err.Error()
log.LogWarnf("action[TryAcquireDecommissionToken] clusterID[%v] vol[%v] partitionID[%v]"+
@ -2502,8 +2542,8 @@ func (partition *DataPartition) needRollback(c *Cluster) bool {
return true
}
func (partition *DataPartition) markRollbackFailed(needRollback bool) {
partition.SetDecommissionStatus(DecommissionFail)
func (partition *DataPartition) markRollbackFailed(needRollback bool, triggerCondition string, errMsg string) {
partition.SetDecommissionStatus(DecommissionFail, triggerCondition, errMsg)
partition.DecommissionNeedRollbackTimes = defaultDecommissionRollbackLimit
partition.DecommissionNeedRollback = needRollback
}
@ -2875,7 +2915,8 @@ func (partition *DataPartition) checkReplicaMeta(c *Cluster) (err error) {
partition.PartitionID, addr)
return nil
}
err = c.markDecommissionDataPartition(partition, node, 0, false, AutoAddReplica, highPriorityDecommissionWeight)
triggerCondition := fmt.Sprintf("autoAddReplica_dp(%v)", partition.PartitionID)
err = c.markDecommissionDataPartition(partition, node, 0, false, AutoAddReplica, highPriorityDecommissionWeight, triggerCondition)
auditMsg = fmt.Sprintf("dp(%v) ReplicaNum %v hostsNum %v auto add replica",
partition.PartitionID, partition.ReplicaNum, len(partition.Hosts))
log.LogDebugf("action[checkReplicaMeta]%v: err %v", auditMsg, err)
@ -3000,11 +3041,11 @@ func (partition *DataPartition) removeHostByForce(c *Cluster, peerAddr string) {
}
}
func (partition *DataPartition) resetForManualAddReplica() {
func (partition *DataPartition) resetForManualAddReplica(triggerCondition string, errMsg string) {
partition.DecommissionDstAddr = ""
partition.DecommissionType = InitialDecommission
partition.isRecover = false
partition.SetDecommissionStatus(DecommissionInitial)
partition.SetDecommissionStatus(DecommissionInitial, triggerCondition, errMsg)
partition.setRestoreReplicaStop()
}

View File

@ -77,12 +77,12 @@ func (c *Cluster) checkDiskRecoveryProgress() {
log.LogInfof("action[checkDiskRecoveryProgress] dp %v isSpec %v replicas %v conf replicas num %v status(%v)",
partition.decommissionInfo(), partition.isSpecialReplicaCnt(), len(partition.Replicas), int(partition.ReplicaNum), partition.GetDecommissionStatus())
if len(partition.Replicas) == 0 {
partition.SetDecommissionStatus(DecommissionSuccess)
partition.SetDecommissionStatus(DecommissionSuccess, "checkDiskRecoveryProgress_dpMaybeDeleted", "")
log.LogWarnf("action[checkDiskRecoveryProgress] dp %v maybe deleted", partition.PartitionID)
continue
}
if partition.IsDiscard {
partition.SetDecommissionStatus(DecommissionSuccess)
partition.SetDecommissionStatus(DecommissionSuccess, "checkDiskRecoveryProgress_discardCheck", "")
log.LogWarnf("[checkDiskRecoveryProgress] dp(%v) is discard, decommission successfully", partition.PartitionID)
continue
}
@ -97,15 +97,16 @@ func (c *Cluster) checkDiskRecoveryProgress() {
newReplica, _ := partition.getReplica(partition.DecommissionDstAddr)
if newReplica == nil {
errMsg := fmt.Sprintf("Decommission target node %v not found", partition.DecommissionDstAddr)
log.LogWarnf("action[checkDiskRecoveryProgress] dp %v cannot find replica %v", partition.PartitionID,
partition.DecommissionDstAddr)
if partition.DecommissionType == ManualAddReplica {
partition.resetForManualAddReplica()
partition.resetForManualAddReplica("checkDiskRecoveryProgress", errMsg)
} else {
partition.DecommissionNeedRollback = true
partition.SetDecommissionStatus(DecommissionFail)
partition.SetDecommissionStatus(DecommissionFail, "checkDiskRecoveryProgress", errMsg)
}
partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v not found", partition.DecommissionDstAddr)
partition.DecommissionErrorMessage = errMsg
partition.RLock()
err = c.syncUpdateDataPartition(partition)
if err != nil {
@ -122,18 +123,20 @@ func (c *Cluster) checkDiskRecoveryProgress() {
duration := time.Unix(masterNode.ReportTime, 0).Sub(time.Unix(newReplica.ReportTime, 0))
diskErrReplicas := partition.getAllDiskErrorReplica()
if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) || math.Abs(duration.Minutes()) > 10 {
if partition.DecommissionType == ManualAddReplica {
partition.resetForManualAddReplica()
} else {
partition.markRollbackFailed(true)
}
var errMsg string
if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) {
partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v cannot finish recover"+
errMsg = fmt.Sprintf("Decommission target node %v cannot finish recover"+
" for host[0] %v is unavailable", partition.DecommissionDstAddr, partition.Hosts[0])
} else {
partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v cannot finish recover"+
errMsg = fmt.Sprintf("Decommission target node %v cannot finish recover"+
" for host[0] %v is down ", partition.DecommissionDstAddr, masterNode.Addr)
}
if partition.DecommissionType == ManualAddReplica {
partition.resetForManualAddReplica("checkDiskRecoveryProgress", errMsg)
} else {
partition.markRollbackFailed(true, "checkDiskRecoveryProgress", errMsg)
}
partition.DecommissionErrorMessage = errMsg
Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] %v",
c.Name, partitionID, partition.DecommissionErrorMessage))
partition.RLock()
@ -144,13 +147,14 @@ func (c *Cluster) checkDiskRecoveryProgress() {
partition.RUnlock()
continue
} else if time.Since(partition.RecoverUpdateTime) > c.GetDecommissionDataPartitionRecoverTimeOut() {
errMsg := fmt.Sprintf("Decommission target node %v repair timeout", partition.DecommissionDstAddr)
if partition.DecommissionType == ManualAddReplica {
partition.resetForManualAddReplica()
partition.resetForManualAddReplica("checkDiskRecoveryProgress", errMsg)
} else {
partition.DecommissionNeedRollback = true
partition.SetDecommissionStatus(DecommissionFail)
partition.SetDecommissionStatus(DecommissionFail, "checkDiskRecoveryProgress", errMsg)
}
partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v repair timeout", partition.DecommissionDstAddr)
partition.DecommissionErrorMessage = errMsg
Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] replica %v_%v recovered timeout,recoverUpdateTime %s",
c.Name, partitionID, newReplica.Addr, newReplica.DiskPath, time.Since(partition.RecoverUpdateTime)))
partition.RLock()
@ -165,8 +169,10 @@ func (c *Cluster) checkDiskRecoveryProgress() {
newBadDpIds = append(newBadDpIds, partitionID)
} else {
if partition.DecommissionType == ManualAddReplica {
var errMsg string
if newReplica.isUnavailable() {
partition.DecommissionErrorMessage = fmt.Sprintf("New replica %v is unavailable", partition.DecommissionDstAddr)
errMsg = fmt.Sprintf("New replica %v is unavailable", partition.DecommissionDstAddr)
partition.DecommissionErrorMessage = errMsg
Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] replica %v has recovered failed",
c.Name, partitionID, partition.DecommissionDstAddr))
} else {
@ -174,7 +180,11 @@ func (c *Cluster) checkDiskRecoveryProgress() {
Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] replica %v has recovered success",
c.Name, partitionID, partition.DecommissionDstAddr))
}
partition.resetForManualAddReplica()
partition.resetForManualAddReplica("checkDiskRecoveryProgress", errMsg)
if errMsg == "" {
partition.clearDecommissionStatusRecords()
}
log.LogInfof("[checkDiskRecoveryProgress] dp(%v) manual add new replica addr %v status(%v)",
partitionID, newReplica.Addr, newReplica.Status)
partition.RLock()
@ -192,14 +202,15 @@ func (c *Cluster) checkDiskRecoveryProgress() {
}
// do not add to BadDataPartitionIds
if newReplica.isUnavailable() {
errMsg := fmt.Sprintf("New replica %v is unavailable", partition.DecommissionDstAddr)
partition.DecommissionNeedRollback = true
partition.SetDecommissionStatus(DecommissionFail)
partition.DecommissionErrorMessage = fmt.Sprintf("New replica %v is unavailable", partition.DecommissionDstAddr)
partition.SetDecommissionStatus(DecommissionFail, "checkDiskRecoveryProgress", errMsg)
partition.DecommissionErrorMessage = errMsg
Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] replica %v has recovered failed",
c.Name, partitionID, partition.DecommissionDstAddr))
} else {
partition.DecommissionErrorMessage = ""
partition.SetDecommissionStatus(DecommissionSuccess) // can be readonly or readwrite
partition.SetDecommissionStatus(DecommissionSuccess, "checkDiskRecoveryProgress", "") // can be readonly or readwrite
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()))

View File

@ -634,6 +634,9 @@ func (m *Server) registerAPIRoutes(router *mux.Router) {
router.NewRoute().Methods(http.MethodGet).
Path(proto.AdminQueryDataPartitionDecommissionStatus).
HandlerFunc(m.queryDataPartitionDecommissionStatus)
router.NewRoute().Methods(http.MethodGet).
Path(proto.AdminQueryDataPartitionDecommissionStatusUpdateRecords).
HandlerFunc(m.queryDataPartitionDecommissionStatusUpdateRecords)
router.NewRoute().Methods(http.MethodGet).
Path(proto.AdminCheckReplicaMeta).
HandlerFunc(m.checkReplicaMeta)

View File

@ -179,42 +179,43 @@ func newMetaPartitionValue(mp *MetaPartition) (mpv *metaPartitionValue) {
}
type dataPartitionValue struct {
PartitionID uint64
ReplicaNum uint8
Hosts string
Peers []proto.Peer
Status int8
VolID uint64
VolName string
OfflinePeerID uint64
Replicas []*replicaValue
IsRecover bool
PartitionType int
RdOnly bool
IsDiscard bool
DecommissionDiskRetryMap map[string]int
DecommissionRetry int
DecommissionStatus uint32
DecommissionSrcAddr string
DecommissionDstAddr string
DecommissionRaftForce bool
DecommissionSrcDiskPath string
DecommissionTerm uint64
DecommissionWeight int
SpecialReplicaDecommissionStep uint32
DecommissionDstAddrSpecify bool
DecommissionDstNodeSet uint64
DecommissionNeedRollback bool
RecoverStartTime int64
RecoverUpdateTime int64
RecoverLastConsumeTime float64
DecommissionRetryTime int64
Forbidden bool
DecommissionErrorMessage string
DecommissionNeedRollbackTimes uint32
DecommissionType uint32
RestoreReplica uint32
MediaType uint32
PartitionID uint64
ReplicaNum uint8
Hosts string
Peers []proto.Peer
Status int8
VolID uint64
VolName string
OfflinePeerID uint64
Replicas []*replicaValue
IsRecover bool
PartitionType int
RdOnly bool
IsDiscard bool
DecommissionDiskRetryMap map[string]int
DecommissionStatusUpdateRecords []*proto.DecommissionStatusRecord
DecommissionRetry int
DecommissionStatus uint32
DecommissionSrcAddr string
DecommissionDstAddr string
DecommissionRaftForce bool
DecommissionSrcDiskPath string
DecommissionTerm uint64
DecommissionWeight int
SpecialReplicaDecommissionStep uint32
DecommissionDstAddrSpecify bool
DecommissionDstNodeSet uint64
DecommissionNeedRollback bool
RecoverStartTime int64
RecoverUpdateTime int64
RecoverLastConsumeTime float64
DecommissionRetryTime int64
Forbidden bool
DecommissionErrorMessage string
DecommissionNeedRollbackTimes uint32
DecommissionType uint32
RestoreReplica uint32
MediaType uint32
}
func (dpv *dataPartitionValue) Restore(c *Cluster) (dp *DataPartition) {
@ -268,6 +269,7 @@ func (dpv *dataPartitionValue) Restore(c *Cluster) (dp *DataPartition) {
for disk, retryTimes := range dpv.DecommissionDiskRetryMap {
dp.DecommissionDiskRetryMap[disk] = retryTimes
}
dp.DecommissionStatusUpdateRecords = append(dp.DecommissionStatusUpdateRecords, dpv.DecommissionStatusUpdateRecords...)
return dp
}
@ -278,50 +280,47 @@ type replicaValue struct {
func newDataPartitionValue(dp *DataPartition) (dpv *dataPartitionValue) {
dpv = &dataPartitionValue{
PartitionID: dp.PartitionID,
ReplicaNum: dp.ReplicaNum,
Hosts: dp.hostsToString(),
Peers: dp.Peers,
Status: dp.Status,
VolID: dp.VolID,
VolName: dp.VolName,
OfflinePeerID: dp.OfflinePeerID,
Replicas: make([]*replicaValue, 0),
IsRecover: dp.isRecover,
PartitionType: dp.PartitionType,
RdOnly: dp.RdOnly,
IsDiscard: dp.IsDiscard,
DecommissionDiskRetryMap: make(map[string]int),
DecommissionRetry: dp.DecommissionRetry,
DecommissionStatus: atomic.LoadUint32(&dp.DecommissionStatus),
DecommissionSrcAddr: dp.DecommissionSrcAddr,
DecommissionDstAddr: dp.DecommissionDstAddr,
DecommissionRaftForce: dp.DecommissionRaftForce,
DecommissionSrcDiskPath: dp.DecommissionSrcDiskPath,
DecommissionTerm: dp.DecommissionTerm,
DecommissionWeight: dp.DecommissionWeight,
SpecialReplicaDecommissionStep: dp.SpecialReplicaDecommissionStep,
DecommissionDstAddrSpecify: dp.DecommissionDstAddrSpecify,
DecommissionDstNodeSet: dp.DecommissionDstNodeSet,
DecommissionNeedRollback: dp.DecommissionNeedRollback,
RecoverStartTime: dp.RecoverStartTime.Unix(),
RecoverUpdateTime: dp.RecoverUpdateTime.Unix(),
RecoverLastConsumeTime: dp.RecoverLastConsumeTime.Seconds(),
DecommissionRetryTime: dp.DecommissionRetryTime.Unix(),
DecommissionErrorMessage: dp.DecommissionErrorMessage,
DecommissionNeedRollbackTimes: dp.DecommissionNeedRollbackTimes,
DecommissionType: dp.DecommissionType,
RestoreReplica: atomic.LoadUint32(&dp.RestoreReplica),
MediaType: dp.MediaType,
PartitionID: dp.PartitionID,
ReplicaNum: dp.ReplicaNum,
Hosts: dp.hostsToString(),
Peers: dp.Peers,
Status: dp.Status,
VolID: dp.VolID,
VolName: dp.VolName,
OfflinePeerID: dp.OfflinePeerID,
Replicas: make([]*replicaValue, 0),
IsRecover: dp.isRecover,
PartitionType: dp.PartitionType,
RdOnly: dp.RdOnly,
IsDiscard: dp.IsDiscard,
DecommissionDiskRetryMap: dp.cloneDecommissionDiskRetryMap(),
DecommissionStatusUpdateRecords: dp.cloneDecommissionStatusRecords(),
DecommissionRetry: dp.DecommissionRetry,
DecommissionStatus: atomic.LoadUint32(&dp.DecommissionStatus),
DecommissionSrcAddr: dp.DecommissionSrcAddr,
DecommissionDstAddr: dp.DecommissionDstAddr,
DecommissionRaftForce: dp.DecommissionRaftForce,
DecommissionSrcDiskPath: dp.DecommissionSrcDiskPath,
DecommissionTerm: dp.DecommissionTerm,
DecommissionWeight: dp.DecommissionWeight,
SpecialReplicaDecommissionStep: dp.SpecialReplicaDecommissionStep,
DecommissionDstAddrSpecify: dp.DecommissionDstAddrSpecify,
DecommissionDstNodeSet: dp.DecommissionDstNodeSet,
DecommissionNeedRollback: dp.DecommissionNeedRollback,
RecoverStartTime: dp.RecoverStartTime.Unix(),
RecoverUpdateTime: dp.RecoverUpdateTime.Unix(),
RecoverLastConsumeTime: dp.RecoverLastConsumeTime.Seconds(),
DecommissionRetryTime: dp.DecommissionRetryTime.Unix(),
DecommissionErrorMessage: dp.DecommissionErrorMessage,
DecommissionNeedRollbackTimes: dp.DecommissionNeedRollbackTimes,
DecommissionType: dp.DecommissionType,
RestoreReplica: atomic.LoadUint32(&dp.RestoreReplica),
MediaType: dp.MediaType,
}
for _, replica := range dp.Replicas {
rv := &replicaValue{Addr: replica.Addr, DiskPath: replica.DiskPath}
dpv.Replicas = append(dpv.Replicas, rv)
}
retryTimesMap := dp.cloneDecommissionDiskRetryMap()
for disk, retryTimes := range retryTimesMap {
dpv.DecommissionDiskRetryMap[disk] = retryTimes
}
return
}

View File

@ -2197,7 +2197,7 @@ func (l *DecommissionDataPartitionList) Put(id uint64, value *DataPartition, c *
}
// prepare status reset to mark status to retry again
if value.GetDecommissionStatus() == DecommissionPrepare {
value.SetDecommissionStatus(markDecommission)
value.SetDecommissionStatus(markDecommission, "leaderChange_updatePrepareToMark", "")
}
l.mu.Lock()
if _, ok := l.cacheMap[value.PartitionID]; ok {
@ -2223,7 +2223,7 @@ func (l *DecommissionDataPartitionList) Put(id uint64, value *DataPartition, c *
// restore special replica decommission progress
if value.isSpecialReplicaCnt() && value.GetDecommissionStatus() == DecommissionRunning && !value.DecommissionRaftForce {
value.SetDecommissionStatus(markDecommission)
value.SetDecommissionStatus(markDecommission, "leaderChange_updateRunningToMark", "")
value.isRecover = false // can pass decommission validate check
log.LogInfof("action[DecommissionDataPartitionListPut] ns[%v] dp[%v] set status from DecommissionRunning to markDecommission",
id, value.PartitionID)
@ -2323,16 +2323,16 @@ func (l *DecommissionDataPartitionList) startTraverse() {
func updateDecommissionWeight(dps []*DataPartition, c *Cluster) {
for _, dp := range dps {
if dp.IsDiscard {
dp.SetDecommissionStatus(DecommissionSuccess)
dp.SetDecommissionStatus(DecommissionSuccess, "traverDecommissionList_updateDecommissionWeight_discardCheck", "")
log.LogWarnf("action[DecommissionListTraverse] skip dp(%v) discard(%v)", dp.PartitionID, dp.IsDiscard)
continue
}
diskErrReplicaNum := dp.getReplicaDiskErrorNum()
if diskErrReplicaNum == dp.ReplicaNum || diskErrReplicaNum == uint8(len(dp.Peers)) {
log.LogWarnf("action[DecommissionListTraverse] dp[%v] all live replica is unavaliable", dp.decommissionInfo())
log.LogWarnf("action[DecommissionListTraverse] dp[%v] all live replica is unavailable", dp.decommissionInfo())
err := proto.ErrAllReplicaUnavailable
dp.DecommissionErrorMessage = err.Error()
dp.markRollbackFailed(false)
dp.markRollbackFailed(false, "traverDecommissionList_updateDecommissionWeight_diskErrReplicaNumCheck", err.Error())
continue
}
if dp.DecommissionType == AutoDecommission && dp.IsMarkDecommission() {

View File

@ -34,77 +34,77 @@ type ContextUserKey string
// api
const (
// Admin APIs
AdminGetMasterApiList = "/admin/getMasterApiList"
AdminSetApiQpsLimit = "/admin/setApiQpsLimit"
AdminGetApiQpsLimit = "/admin/getApiQpsLimit"
AdminRemoveApiQpsLimit = "/admin/rmApiQpsLimit"
AdminGetCluster = "/admin/getCluster"
AdminSetClusterInfo = "/admin/setClusterInfo"
AdminGetMonitorPushAddr = "/admin/getMonitorPushAddr"
AdminGetClusterDataNodes = "/admin/cluster/getAllDataNodes"
AdminGetClusterMetaNodes = "/admin/cluster/getAllMetaNodes"
AdminGetDataPartition = "/dataPartition/get"
AdminLoadDataPartition = "/dataPartition/load"
AdminCreateDataPartition = "/dataPartition/create"
AdminCreatePreLoadDataPartition = "/dataPartition/createPreLoad"
AdminDecommissionDataPartition = "/dataPartition/decommission"
AdminDiagnoseDataPartition = "/dataPartition/diagnose"
AdminResetDataPartitionDecommissionStatus = "/dataPartition/resetDecommissionStatus"
AdminQueryDataPartitionDecommissionStatus = "/dataPartition/queryDecommissionStatus"
AdminCheckReplicaMeta = "/dataPartition/checkReplicaMeta"
AdminRecoverReplicaMeta = "/dataPartition/recoverReplicaMeta"
AdminRecoverBackupDataReplica = "/dataPartition/recoverBackupDataReplica"
AdminDeleteDataReplica = "/dataReplica/delete"
AdminAddDataReplica = "/dataReplica/add"
AdminDeleteVol = "/vol/delete"
AdminUpdateVol = "/vol/update"
AdminVolShrink = "/vol/shrink"
AdminVolExpand = "/vol/expand"
AdminVolForbidden = "/vol/forbidden"
AdminVolEnableAuditLog = "/vol/auditlog"
AdminVolSetDpRepairBlockSize = "/vol/setDpRepairBlockSize"
AdminCreateVol = "/admin/createVol"
AdminGetVol = "/admin/getVol"
AdminClusterFreeze = "/cluster/freeze"
AdminClusterForbidMpDecommission = "/cluster/forbidMetaPartitionDecommission"
AdminClusterStat = "/cluster/stat"
AdminSetCheckDataReplicasEnable = "/cluster/setCheckDataReplicasEnable"
AdminGetIP = "/admin/getIp"
AdminCreateMetaPartition = "/metaPartition/create"
AdminSetMetaNodeThreshold = "/threshold/set"
AdminSetMasterVolDeletionDelayTime = "/volDeletionDelayTime/set"
AdminSetMetaNodeGOGC = "/metaNodeGOGC/set"
AdminSetDataNodeGOGC = "/dataNodeGOGC/set"
AdminListVols = "/vol/list"
AdminSetNodeInfo = "/admin/setNodeInfo"
AdminGetNodeInfo = "/admin/getNodeInfo"
AdminGetAllNodeSetGrpInfo = "/admin/getDomainInfo"
AdminGetNodeSetGrpInfo = "/admin/getDomainNodeSetGrpInfo"
AdminGetIsDomainOn = "/admin/getIsDomainOn"
AdminUpdateNodeSetCapcity = "/admin/updateNodeSetCapcity"
AdminUpdateNodeSetId = "/admin/updateNodeSetId"
AdminUpdateNodeSetNodeSelector = "/admin/updateNodeSetNodeSelector"
AdminUpdateDomainDataUseRatio = "/admin/updateDomainDataRatio"
AdminUpdateZoneExcludeRatio = "/admin/updateZoneExcludeRatio"
AdminSetNodeRdOnly = "/admin/setNodeRdOnly"
AdminSetDpRdOnly = "/admin/setDpRdOnly"
AdminSetConfig = "/admin/setConfig"
AdminGetConfig = "/admin/getConfig"
AdminDataPartitionChangeLeader = "/dataPartition/changeleader"
AdminChangeMasterLeader = "/master/changeleader"
AdminOpFollowerPartitionsRead = "/master/opFollowerPartitionRead"
AdminUpdateDecommissionFirstHostDiskParallelLimit = "/admin/updateDecommissionFirstHostDiskParallelLimit"
AdminQueryDecommissionFirstHostDiskParallelLimit = "/admin/queryDecommissionFirstHostDiskParallelLimit"
AdminUpdateDecommissionFirstHostParallelLimit = "/admin/updateDecommissionFirstHostParallelLimit"
AdminQueryDecommissionFirstHostParallelLimit = "/admin/queryDecommissionFirstHostParallelLimit"
AdminQueryDecommissionFirstHostParallelInfo = "/admin/queryDecommissionFirstHostParallelInfo"
AdminUpdateDecommissionLimit = "/admin/updateDecommissionLimit"
AdminQueryDecommissionLimit = "/admin/queryDecommissionLimit"
AdminQueryDecommissionFailedDisk = "/admin/queryDecommissionFailedDisk"
AdminAbortDecommissionDisk = "/admin/abortDecommissionDisk"
AdminResetDataPartitionRestoreStatus = "/admin/resetDataPartitionRestoreStatus"
AdminGetOpLog = "/admin/getOpLog"
AdminGetRemoteCacheConfig = "/admin/getRemoteCacheConfig"
AdminGetMasterApiList = "/admin/getMasterApiList"
AdminSetApiQpsLimit = "/admin/setApiQpsLimit"
AdminGetApiQpsLimit = "/admin/getApiQpsLimit"
AdminRemoveApiQpsLimit = "/admin/rmApiQpsLimit"
AdminGetCluster = "/admin/getCluster"
AdminSetClusterInfo = "/admin/setClusterInfo"
AdminGetMonitorPushAddr = "/admin/getMonitorPushAddr"
AdminGetClusterDataNodes = "/admin/cluster/getAllDataNodes"
AdminGetClusterMetaNodes = "/admin/cluster/getAllMetaNodes"
AdminGetDataPartition = "/dataPartition/get"
AdminLoadDataPartition = "/dataPartition/load"
AdminCreateDataPartition = "/dataPartition/create"
AdminDecommissionDataPartition = "/dataPartition/decommission"
AdminDiagnoseDataPartition = "/dataPartition/diagnose"
AdminResetDataPartitionDecommissionStatus = "/dataPartition/resetDecommissionStatus"
AdminQueryDataPartitionDecommissionStatus = "/dataPartition/queryDecommissionStatus"
AdminQueryDataPartitionDecommissionStatusUpdateRecords = "/dataPartition/queryDecommissionStatusUpdateRecords"
AdminCheckReplicaMeta = "/dataPartition/checkReplicaMeta"
AdminRecoverReplicaMeta = "/dataPartition/recoverReplicaMeta"
AdminRecoverBackupDataReplica = "/dataPartition/recoverBackupDataReplica"
AdminDeleteDataReplica = "/dataReplica/delete"
AdminAddDataReplica = "/dataReplica/add"
AdminDeleteVol = "/vol/delete"
AdminUpdateVol = "/vol/update"
AdminVolShrink = "/vol/shrink"
AdminVolExpand = "/vol/expand"
AdminVolForbidden = "/vol/forbidden"
AdminVolEnableAuditLog = "/vol/auditlog"
AdminVolSetDpRepairBlockSize = "/vol/setDpRepairBlockSize"
AdminCreateVol = "/admin/createVol"
AdminGetVol = "/admin/getVol"
AdminClusterFreeze = "/cluster/freeze"
AdminClusterForbidMpDecommission = "/cluster/forbidMetaPartitionDecommission"
AdminClusterStat = "/cluster/stat"
AdminSetCheckDataReplicasEnable = "/cluster/setCheckDataReplicasEnable"
AdminGetIP = "/admin/getIp"
AdminCreateMetaPartition = "/metaPartition/create"
AdminSetMetaNodeThreshold = "/threshold/set"
AdminSetMasterVolDeletionDelayTime = "/volDeletionDelayTime/set"
AdminSetMetaNodeGOGC = "/metaNodeGOGC/set"
AdminSetDataNodeGOGC = "/dataNodeGOGC/set"
AdminListVols = "/vol/list"
AdminSetNodeInfo = "/admin/setNodeInfo"
AdminGetNodeInfo = "/admin/getNodeInfo"
AdminGetAllNodeSetGrpInfo = "/admin/getDomainInfo"
AdminGetNodeSetGrpInfo = "/admin/getDomainNodeSetGrpInfo"
AdminGetIsDomainOn = "/admin/getIsDomainOn"
AdminUpdateNodeSetCapcity = "/admin/updateNodeSetCapcity"
AdminUpdateNodeSetId = "/admin/updateNodeSetId"
AdminUpdateNodeSetNodeSelector = "/admin/updateNodeSetNodeSelector"
AdminUpdateDomainDataUseRatio = "/admin/updateDomainDataRatio"
AdminUpdateZoneExcludeRatio = "/admin/updateZoneExcludeRatio"
AdminSetNodeRdOnly = "/admin/setNodeRdOnly"
AdminSetDpRdOnly = "/admin/setDpRdOnly"
AdminSetConfig = "/admin/setConfig"
AdminGetConfig = "/admin/getConfig"
AdminDataPartitionChangeLeader = "/dataPartition/changeleader"
AdminChangeMasterLeader = "/master/changeleader"
AdminOpFollowerPartitionsRead = "/master/opFollowerPartitionRead"
AdminUpdateDecommissionFirstHostDiskParallelLimit = "/admin/updateDecommissionFirstHostDiskParallelLimit"
AdminQueryDecommissionFirstHostDiskParallelLimit = "/admin/queryDecommissionFirstHostDiskParallelLimit"
AdminUpdateDecommissionFirstHostParallelLimit = "/admin/updateDecommissionFirstHostParallelLimit"
AdminQueryDecommissionFirstHostParallelLimit = "/admin/queryDecommissionFirstHostParallelLimit"
AdminQueryDecommissionFirstHostParallelInfo = "/admin/queryDecommissionFirstHostParallelInfo"
AdminUpdateDecommissionLimit = "/admin/updateDecommissionLimit"
AdminQueryDecommissionLimit = "/admin/queryDecommissionLimit"
AdminQueryDecommissionFailedDisk = "/admin/queryDecommissionFailedDisk"
AdminAbortDecommissionDisk = "/admin/abortDecommissionDisk"
AdminResetDataPartitionRestoreStatus = "/admin/resetDataPartitionRestoreStatus"
AdminGetOpLog = "/admin/getOpLog"
AdminGetRemoteCacheConfig = "/admin/getRemoteCacheConfig"
// #nosec G101
AdminQueryDecommissionToken = "/admin/queryDecommissionToken"

View File

@ -496,6 +496,13 @@ type DiscardDataPartitionInfos struct {
DiscardDps []DataPartitionInfo
}
type DecommissionStatusRecord struct {
Condition string
Status string
Time string
ErrMessage string
}
type DecommissionInfoStat struct {
Key string
RepairSourceDp []uint64

View File

@ -253,6 +253,14 @@ func (api *AdminAPI) AddMetaReplica(metaPartitionID uint64, nodeAddr string, cli
return
}
func (api *AdminAPI) QueryDataPartitionDecommissionStatusUpdateRecords(partitionId uint64) (records []*proto.DecommissionStatusRecord, err error) {
request := newRequest(get, proto.AdminQueryDataPartitionDecommissionStatusUpdateRecords).Header(api.h)
request.addParam("id", strconv.FormatUint(partitionId, 10))
records = make([]*proto.DecommissionStatusRecord, 0)
err = api.mc.requestWith(&records, request)
return
}
func (api *AdminAPI) QueryDataPartitionDecommissionStatus(partitionId uint64) (info *proto.DecommissionDataPartitionInfo, err error) {
request := newRequest(get, proto.AdminQueryDataPartitionDecommissionStatus).Header(api.h)
request.addParam("id", strconv.FormatUint(partitionId, 10))