mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
feat(master): Add disk decommission success metric. #1000012562
Signed-off-by: zhumingze <zhumingze@oppo.com>
This commit is contained in:
parent
3cbb46bd80
commit
02fda728d7
@ -16,47 +16,48 @@ package cmd
|
||||
|
||||
const (
|
||||
// List of operation name for cli
|
||||
CliOpGet = "get"
|
||||
CliOpList = "list"
|
||||
CliOpStatus = "stat"
|
||||
CliOpCreate = "create"
|
||||
CliOpDelete = "delete"
|
||||
CliOpRemove = "remove"
|
||||
CliOpInfo = "info"
|
||||
CliOpAdd = "add"
|
||||
CliOpSet = "set"
|
||||
CliOpUpdate = "update"
|
||||
CliOpDecommission = "decommission"
|
||||
CliOpRecommission = "recommission"
|
||||
CliOpQueryProgress = "query-progress"
|
||||
CliOpAbortDecommission = "abort-decommission"
|
||||
CliOpMigrate = "migrate"
|
||||
CliOpDownloadZip = "load"
|
||||
CliOpMetaCompatibility = "meta"
|
||||
CliOpFreeze = "freeze"
|
||||
CliOpSetThreshold = "threshold"
|
||||
CliOpSetVolDeletionDelayTime = "volDeletionDelayTime"
|
||||
CliOpSetCluster = "set"
|
||||
CliOpCheck = "check"
|
||||
CliOpReset = "reset"
|
||||
CliOpReplicate = "add-replica"
|
||||
CliOpDelReplica = "del-replica"
|
||||
CliOpExpand = "expand"
|
||||
CliOpShrink = "shrink"
|
||||
CliOpGetDiscard = "get-discard"
|
||||
CliOpSetDiscard = "set-discard"
|
||||
CliOpForbidMpDecommission = "forbid-mp-decommission"
|
||||
CliOpQueryDecommissionedDisk = "query-decommissioned-disk"
|
||||
CliOpEnableAutoDecommission = "enable-auto-decommission"
|
||||
CliOpQueryDecommissionFailedDisk = "query-decommission-failed-disk"
|
||||
CliOpSetDecommissionDiskLimit = "set-decommission-disk-limit"
|
||||
CliOpResetRestoreStatus = "reset-restore-status"
|
||||
CliOpCancelDecommission = "cancel-decommission"
|
||||
CliOpDiskOp = "diskop"
|
||||
CliOpDpOp = "dpop"
|
||||
CliOpDataNodeOp = "datanodeop"
|
||||
CliOpVolOp = "volop"
|
||||
CliOpToLeader = "to-leader"
|
||||
CliOpGet = "get"
|
||||
CliOpList = "list"
|
||||
CliOpStatus = "stat"
|
||||
CliOpCreate = "create"
|
||||
CliOpDelete = "delete"
|
||||
CliOpRemove = "remove"
|
||||
CliOpInfo = "info"
|
||||
CliOpAdd = "add"
|
||||
CliOpSet = "set"
|
||||
CliOpUpdate = "update"
|
||||
CliOpDecommission = "decommission"
|
||||
CliOpRecommission = "recommission"
|
||||
CliOpQueryProgress = "query-progress"
|
||||
CliOpAbortDecommission = "abort-decommission"
|
||||
CliOpMigrate = "migrate"
|
||||
CliOpDownloadZip = "load"
|
||||
CliOpMetaCompatibility = "meta"
|
||||
CliOpFreeze = "freeze"
|
||||
CliOpSetThreshold = "threshold"
|
||||
CliOpSetVolDeletionDelayTime = "volDeletionDelayTime"
|
||||
CliOpSetCluster = "set"
|
||||
CliOpCheck = "check"
|
||||
CliOpReset = "reset"
|
||||
CliOpReplicate = "add-replica"
|
||||
CliOpDelReplica = "del-replica"
|
||||
CliOpExpand = "expand"
|
||||
CliOpShrink = "shrink"
|
||||
CliOpGetDiscard = "get-discard"
|
||||
CliOpSetDiscard = "set-discard"
|
||||
CliOpForbidMpDecommission = "forbid-mp-decommission"
|
||||
CliOpQueryDecommissionedDisk = "query-decommissioned-disk"
|
||||
CliOpQueryDecommissionSuccessDisk = "query-decommissionSuccess-disk"
|
||||
CliOpEnableAutoDecommission = "enable-auto-decommission"
|
||||
CliOpQueryDecommissionFailedDisk = "query-decommission-failed-disk"
|
||||
CliOpSetDecommissionDiskLimit = "set-decommission-disk-limit"
|
||||
CliOpResetRestoreStatus = "reset-restore-status"
|
||||
CliOpCancelDecommission = "cancel-decommission"
|
||||
CliOpDiskOp = "diskop"
|
||||
CliOpDpOp = "dpop"
|
||||
CliOpDataNodeOp = "datanodeop"
|
||||
CliOpVolOp = "volop"
|
||||
CliOpToLeader = "to-leader"
|
||||
|
||||
CliOpSetDecommissionLimit = "set-decommission-limit"
|
||||
CliOpQueryDecommissionStatus = "query-decommission-status"
|
||||
|
||||
@ -42,6 +42,7 @@ func newDataNodeCmd(client *master.MasterClient) *cobra.Command {
|
||||
newDataNodeMigrateCmd(client),
|
||||
newDataNodeQueryDecommissionProgress(client),
|
||||
newDataNodeQueryDecommissionedDisk(client),
|
||||
newDataNodeQueryDecommissionSuccessDisk(client),
|
||||
newDataNodeCancelDecommissionCmd(client),
|
||||
// newDataNodeDiskOpCmd(client),
|
||||
// newDataNodeDpOpCmd(client),
|
||||
@ -50,12 +51,13 @@ func newDataNodeCmd(client *master.MasterClient) *cobra.Command {
|
||||
}
|
||||
|
||||
const (
|
||||
cmdDataNodeListShort = "List information of data nodes"
|
||||
cmdDataNodeInfoShort = "Show information of a data node"
|
||||
cmdDataNodeDecommissionInfoShort = "decommission partitions in a data node to others"
|
||||
cmdDataNodeQueryDecommissionedDisksShort = "query datanode decommissioned disks"
|
||||
cmdDataNodeCancelDecommissionedDisksShort = "cancel decommission progress for datanode"
|
||||
cmdDataNodeQueryDecommissionProgress = "query datanode decommission progress"
|
||||
cmdDataNodeListShort = "List information of data nodes"
|
||||
cmdDataNodeInfoShort = "Show information of a data node"
|
||||
cmdDataNodeDecommissionInfoShort = "decommission partitions in a data node to others"
|
||||
cmdDataNodeQueryDecommissionedDisksShort = "query datanode decommissioned disks"
|
||||
cmdDataNodeQueryDecommissionSuccessDisksShort = "query datanode decommissionSuccess disks"
|
||||
cmdDataNodeCancelDecommissionedDisksShort = "cancel decommission progress for datanode"
|
||||
cmdDataNodeQueryDecommissionProgress = "query datanode decommission progress"
|
||||
// cmdDataNodeDiskOpShort = "Show Disk_op information of a data node"
|
||||
// cmdDataNodeDpOpShort = "Show Dp_op information of a data node"
|
||||
)
|
||||
@ -229,6 +231,27 @@ func newDataNodeQueryDecommissionedDisk(client *master.MasterClient) *cobra.Comm
|
||||
return cmd
|
||||
}
|
||||
|
||||
func newDataNodeQueryDecommissionSuccessDisk(client *master.MasterClient) *cobra.Command {
|
||||
cmd := &cobra.Command{
|
||||
Use: CliOpQueryDecommissionSuccessDisk + " [{HOST}:{PORT}]",
|
||||
Short: cmdDataNodeQueryDecommissionSuccessDisksShort,
|
||||
Args: cobra.MinimumNArgs(1),
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
disks, err := client.NodeAPI().QueryDecommissionSuccessDisks(args[0])
|
||||
if err != nil {
|
||||
stdout("%v", err)
|
||||
return err
|
||||
}
|
||||
stdoutln("[DecommissionSuccess disks]")
|
||||
for _, disk := range disks.Disks {
|
||||
stdout("%v\n", disk)
|
||||
}
|
||||
return nil
|
||||
},
|
||||
}
|
||||
return cmd
|
||||
}
|
||||
|
||||
func newDataNodeCancelDecommissionCmd(client *master.MasterClient) *cobra.Command {
|
||||
cmd := &cobra.Command{
|
||||
Use: CliOpCancelDecommission + " [{HOST}:{PORT}]",
|
||||
|
||||
@ -999,30 +999,31 @@ func formatDataNodeDetail(dn *proto.DataNodeInfo, rowTable bool) string {
|
||||
formatSize(dn.Used), formatSize(dn.Total), formatNodeStatus(dn.IsActive), formatTimeToString(dn.ReportTime))
|
||||
}
|
||||
sb := strings.Builder{}
|
||||
sb.WriteString(fmt.Sprintf(" ID : %v\n", dn.ID))
|
||||
sb.WriteString(fmt.Sprintf(" Address : %v\n", formatAddr(dn.Addr, dn.DomainAddr)))
|
||||
sb.WriteString(fmt.Sprintf(" RaftHeartbeatPort : %v\n", dn.RaftHeartbeatPort))
|
||||
sb.WriteString(fmt.Sprintf(" RaftReplicaPort : %v\n", dn.RaftReplicaPort))
|
||||
sb.WriteString(fmt.Sprintf(" Allocated ratio : %v\n", dn.UsageRatio))
|
||||
sb.WriteString(fmt.Sprintf(" Allocated : %v\n", formatSize(dn.Used)))
|
||||
sb.WriteString(fmt.Sprintf(" Available : %v\n", formatSize(dn.AvailableSpace)))
|
||||
sb.WriteString(fmt.Sprintf(" Total : %v\n", formatSize(dn.Total)))
|
||||
sb.WriteString(fmt.Sprintf(" Zone : %v\n", dn.ZoneName))
|
||||
sb.WriteString(fmt.Sprintf(" Rdonly : %v\n", dn.RdOnly))
|
||||
sb.WriteString(fmt.Sprintf(" Status : %v\n", formatNodeStatus(dn.IsActive)))
|
||||
sb.WriteString(fmt.Sprintf(" MediaType : %v\n", proto.MediaTypeString(dn.MediaType)))
|
||||
sb.WriteString(fmt.Sprintf(" ToBeOffline : %v\n", formatNodeOfflineStatus(dn.ToBeOffline)))
|
||||
sb.WriteString(fmt.Sprintf(" Report time : %v\n", formatTimeToString(dn.ReportTime)))
|
||||
sb.WriteString(fmt.Sprintf(" Partition count : %v\n", dn.DataPartitionCount))
|
||||
sb.WriteString(fmt.Sprintf(" AllDisks : %v\n", dn.AllDisks))
|
||||
sb.WriteString(fmt.Sprintf(" Bad disks : %v\n", dn.BadDisks))
|
||||
sb.WriteString(fmt.Sprintf(" Lost disks : %v\n", dn.LostDisks))
|
||||
sb.WriteString(fmt.Sprintf(" Decommissioned disks: %v\n", dn.DecommissionedDisk))
|
||||
sb.WriteString(fmt.Sprintf(" Persist partitions : %v\n", dn.PersistenceDataPartitions))
|
||||
sb.WriteString(fmt.Sprintf(" Backup partitions : %v\n", dn.BackupDataPartitions))
|
||||
sb.WriteString(fmt.Sprintf(" Can alloc partition : %v\n", dn.CanAllocPartition))
|
||||
sb.WriteString(fmt.Sprintf(" Max partition count : %v\n", dn.MaxDpCntLimit))
|
||||
sb.WriteString(fmt.Sprintf(" CpuUtil : %.1f%%\n", dn.CpuUtil))
|
||||
sb.WriteString(fmt.Sprintf(" ID : %v\n", dn.ID))
|
||||
sb.WriteString(fmt.Sprintf(" Address : %v\n", formatAddr(dn.Addr, dn.DomainAddr)))
|
||||
sb.WriteString(fmt.Sprintf(" RaftHeartbeatPort : %v\n", dn.RaftHeartbeatPort))
|
||||
sb.WriteString(fmt.Sprintf(" RaftReplicaPort : %v\n", dn.RaftReplicaPort))
|
||||
sb.WriteString(fmt.Sprintf(" Allocated ratio : %v\n", dn.UsageRatio))
|
||||
sb.WriteString(fmt.Sprintf(" Allocated : %v\n", formatSize(dn.Used)))
|
||||
sb.WriteString(fmt.Sprintf(" Available : %v\n", formatSize(dn.AvailableSpace)))
|
||||
sb.WriteString(fmt.Sprintf(" Total : %v\n", formatSize(dn.Total)))
|
||||
sb.WriteString(fmt.Sprintf(" Zone : %v\n", dn.ZoneName))
|
||||
sb.WriteString(fmt.Sprintf(" Rdonly : %v\n", dn.RdOnly))
|
||||
sb.WriteString(fmt.Sprintf(" Status : %v\n", formatNodeStatus(dn.IsActive)))
|
||||
sb.WriteString(fmt.Sprintf(" MediaType : %v\n", proto.MediaTypeString(dn.MediaType)))
|
||||
sb.WriteString(fmt.Sprintf(" ToBeOffline : %v\n", formatNodeOfflineStatus(dn.ToBeOffline)))
|
||||
sb.WriteString(fmt.Sprintf(" Report time : %v\n", formatTimeToString(dn.ReportTime)))
|
||||
sb.WriteString(fmt.Sprintf(" Partition count : %v\n", dn.DataPartitionCount))
|
||||
sb.WriteString(fmt.Sprintf(" AllDisks : %v\n", dn.AllDisks))
|
||||
sb.WriteString(fmt.Sprintf(" Bad disks : %v\n", dn.BadDisks))
|
||||
sb.WriteString(fmt.Sprintf(" Lost disks : %v\n", dn.LostDisks))
|
||||
sb.WriteString(fmt.Sprintf(" Decommissioned disks : %v\n", dn.DecommissionedDisk))
|
||||
sb.WriteString(fmt.Sprintf(" DecommissionSuccess disks : %v\n", dn.DecommissionSuccessDisk))
|
||||
sb.WriteString(fmt.Sprintf(" Persist partitions : %v\n", dn.PersistenceDataPartitions))
|
||||
sb.WriteString(fmt.Sprintf(" Backup partitions : %v\n", dn.BackupDataPartitions))
|
||||
sb.WriteString(fmt.Sprintf(" Can alloc partition : %v\n", dn.CanAllocPartition))
|
||||
sb.WriteString(fmt.Sprintf(" Max partition count : %v\n", dn.MaxDpCntLimit))
|
||||
sb.WriteString(fmt.Sprintf(" CpuUtil : %.1f%%\n", dn.CpuUtil))
|
||||
sb.WriteString(" IoUtils :\n")
|
||||
for device, used := range dn.IoUtils {
|
||||
sb.WriteString(fmt.Sprintf(" %v:%.1f%%\n", device, used))
|
||||
|
||||
@ -3253,6 +3253,7 @@ func (m *Server) getDataNode(w http.ResponseWriter, r *http.Request) {
|
||||
CpuUtil: dataNode.CpuUtil.Load(),
|
||||
IoUtils: dataNode.GetIoUtils(),
|
||||
DecommissionedDisk: dataNode.getDecommissionedDisks(),
|
||||
DecommissionSuccessDisk: dataNode.getDecommissionSuccessDisks(),
|
||||
BackupDataPartitions: dataNode.getBackupDataPartitionIDs(),
|
||||
PersistenceDataPartitionsWithDiskPath: m.cluster.getAllDataPartitionWithDiskPathByDataNode(nodeAddr),
|
||||
MediaType: dataNode.MediaType,
|
||||
@ -4664,6 +4665,7 @@ func (m *Server) recommissionDisk(w http.ResponseWriter, r *http.Request) {
|
||||
node *DataNode
|
||||
rstMsg string
|
||||
onLineAddr, diskPath string
|
||||
recommissionType string
|
||||
err error
|
||||
)
|
||||
metric := exporter.NewTPCnt(apiToMetricsName(proto.RecommissionDisk))
|
||||
@ -4676,22 +4678,35 @@ func (m *Server) recommissionDisk(w http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
if node, err = m.cluster.dataNode(onLineAddr); err != nil {
|
||||
sendErrReply(w, r, newErrHTTPReply(errors.NewErrorf("disk %v on dataNode %v is bad disk", diskPath, onLineAddr)))
|
||||
if recommissionType, err = parseRecommissionType(r); err != nil {
|
||||
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
if node.isBadDisk(diskPath) {
|
||||
if node, err = m.cluster.dataNode(onLineAddr); err != nil {
|
||||
sendErrReply(w, r, newErrHTTPReply(proto.ErrDataNodeNotExists))
|
||||
return
|
||||
}
|
||||
if err = m.cluster.deleteAndSyncDecommissionedDisk(node, diskPath); err != nil {
|
||||
sendErrReply(w, r, newErrHTTPReply(err))
|
||||
return
|
||||
|
||||
if recommissionType == "decommissioned" {
|
||||
if node.isBadDisk(diskPath) {
|
||||
sendErrReply(w, r, newErrHTTPReply(errors.NewErrorf("disk %v on dataNode %v is bad disk", diskPath, onLineAddr)))
|
||||
return
|
||||
}
|
||||
if err = m.cluster.deleteAndSyncDecommissionedDisk(node, diskPath); err != nil {
|
||||
sendErrReply(w, r, newErrHTTPReply(err))
|
||||
return
|
||||
}
|
||||
}
|
||||
if recommissionType == "decommissionSuccess" {
|
||||
if err = m.cluster.deleteAndSyncDecommissionSuccessDisk(node, diskPath); err != nil {
|
||||
sendErrReply(w, r, newErrHTTPReply(err))
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
rstMsg = fmt.Sprintf("receive recommissionDisk node[%v] disk[%v], and recommission successfully",
|
||||
node.Addr, diskPath)
|
||||
rstMsg = fmt.Sprintf("receive recommissionDisk node[%v] disk[%v] recommissionType[%v], and recommission successfully",
|
||||
node.Addr, diskPath, recommissionType)
|
||||
|
||||
Warn(m.clusterName, rstMsg)
|
||||
sendOkReply(w, r, newSuccessHTTPReply(rstMsg))
|
||||
@ -5427,6 +5442,18 @@ func parseNodeAddrAndDisk(r *http.Request) (nodeAddr, diskPath string, err error
|
||||
return
|
||||
}
|
||||
|
||||
func parseRecommissionType(r *http.Request) (recommissionType string, err error) {
|
||||
if err = r.ParseForm(); err != nil {
|
||||
return
|
||||
}
|
||||
recommissionType = r.FormValue(RecommissionType)
|
||||
if recommissionType == "" {
|
||||
err = keyNotFound(addrKey)
|
||||
return
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func parseReqToQueryDecoDisk(r *http.Request) (decommissionType int, showAll bool, err error) {
|
||||
if err = r.ParseForm(); err != nil {
|
||||
return
|
||||
@ -6616,6 +6643,41 @@ func (m *Server) queryDisableDisk(w http.ResponseWriter, r *http.Request) {
|
||||
sendOkReply(w, r, newSuccessHTTPReply(disksInfo))
|
||||
}
|
||||
|
||||
func (m *Server) queryDecommissionSuccessDisk(w http.ResponseWriter, r *http.Request) {
|
||||
var (
|
||||
node *DataNode
|
||||
rstMsg string
|
||||
nodeAddr string
|
||||
err error
|
||||
)
|
||||
metric := exporter.NewTPCnt(apiToMetricsName(proto.QueryDecommissionSuccessDisk))
|
||||
defer func() {
|
||||
doStatAndMetric(proto.QueryDecommissionSuccessDisk, metric, err, nil)
|
||||
}()
|
||||
|
||||
if nodeAddr, err = parseAndExtractNodeAddr(r); err != nil {
|
||||
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
if node, err = m.cluster.dataNode(nodeAddr); err != nil {
|
||||
sendErrReply(w, r, newErrHTTPReply(proto.ErrDataNodeNotExists))
|
||||
return
|
||||
}
|
||||
|
||||
disks := node.getDecommissionSuccessDisks()
|
||||
|
||||
disksInfo := &proto.DecommissionedDisks{
|
||||
Node: nodeAddr,
|
||||
Disks: disks,
|
||||
}
|
||||
rstMsg = fmt.Sprintf("datanode[%v] decommission success disk[%v]",
|
||||
nodeAddr, disks)
|
||||
|
||||
Warn(m.clusterName, rstMsg)
|
||||
sendOkReply(w, r, newSuccessHTTPReply(disksInfo))
|
||||
}
|
||||
|
||||
func parseReqToDecoDataNodeProgress(r *http.Request) (nodeAddr string, err error) {
|
||||
if err = r.ParseForm(); err != nil {
|
||||
return
|
||||
@ -8525,6 +8587,10 @@ func (m *Server) resetDecommissionDataNodeStatus(w http.ResponseWriter, r *http.
|
||||
sendErrReply(w, r, newErrHTTPReply(err))
|
||||
return
|
||||
}
|
||||
if err = m.cluster.deleteAndSyncDecommissionSuccessDisk(dn, disk); err != nil {
|
||||
sendErrReply(w, r, newErrHTTPReply(err))
|
||||
return
|
||||
}
|
||||
}
|
||||
dn.resetDecommissionStatus()
|
||||
m.cluster.syncUpdateDataNode(dn)
|
||||
|
||||
@ -5043,6 +5043,7 @@ func (c *Cluster) checkDecommissionDisk() {
|
||||
// keep failed decommission disk in list for preventing the reuse of a
|
||||
// term in future decommissioning operations
|
||||
if status == DecommissionSuccess {
|
||||
c.addAndSyncDecommissionSuccessDisk(disk.SrcAddr, disk.DiskPath)
|
||||
if time.Since(time.Unix(disk.DecommissionCompleteTime, 0)) > (120 * time.Hour) {
|
||||
if err := c.syncDeleteDecommissionDisk(disk); err != nil {
|
||||
msg := fmt.Sprintf("action[checkDecommissionDisk],clusterID[%v] node[%v] disk[%v],"+
|
||||
|
||||
@ -133,6 +133,7 @@ const (
|
||||
Periodic = "periodic"
|
||||
DecommissionType = "decommissionType"
|
||||
decommissionDiskLimit = "decommissionDiskLimit"
|
||||
RecommissionType = "recommissionType"
|
||||
dpRepairBlockSizeKey = "dpRepairBlockSize"
|
||||
markDiskBrokenThresholdKey = "markDiskBrokenThreshold"
|
||||
decommissionTypeKey = "decommissionType"
|
||||
|
||||
@ -60,7 +60,8 @@ type DataNode struct {
|
||||
DiskStats []proto.DiskStat // key:
|
||||
BadDiskStats []proto.BadDiskStat // key: disk path
|
||||
LostDisks []string
|
||||
DecommissionedDisks sync.Map `json:"-"` // NOTE: the disks that already be decommissioned
|
||||
DecommissionedDisks sync.Map `json:"-"` // NOTE: the disks that already be executed decommission and in disable status
|
||||
DecommissionSuccessDisks sync.Map `json:"-"` // NOTE: the disks that already complete decommission
|
||||
AllDisks []string // TODO: remove me when merge to github master
|
||||
ToBeOffline bool
|
||||
RdOnly bool
|
||||
@ -520,6 +521,28 @@ func (dataNode *DataNode) checkDecommissionedDisks(d string) (ok bool) {
|
||||
return
|
||||
}
|
||||
|
||||
func (dataNode *DataNode) addDecommissionSuccessDisk(diskPath string) (exist bool) {
|
||||
_, exist = dataNode.DecommissionSuccessDisks.LoadOrStore(diskPath, struct{}{})
|
||||
log.LogInfof("action[addDecommissionSuccessDisk] finish, exist[%v], decommissionSuccess disk[%v], dataNode[%v]", exist, diskPath, dataNode.Addr)
|
||||
return
|
||||
}
|
||||
|
||||
func (dataNode *DataNode) deleteDecommissionSuccessDisk(diskPath string) (exist bool) {
|
||||
_, exist = dataNode.DecommissionSuccessDisks.LoadAndDelete(diskPath)
|
||||
log.LogInfof("action[deleteDecommissionSuccessDisk] finish, exist[%v], decommissionSuccess disk[%v], dataNode[%v]", exist, diskPath, dataNode.Addr)
|
||||
return
|
||||
}
|
||||
|
||||
func (dataNode *DataNode) getDecommissionSuccessDisks() (decommissionSuccessDisks []string) {
|
||||
dataNode.DecommissionSuccessDisks.Range(func(key, value interface{}) bool {
|
||||
if diskPath, ok := key.(string); ok {
|
||||
decommissionSuccessDisks = append(decommissionSuccessDisks, diskPath)
|
||||
}
|
||||
return true
|
||||
})
|
||||
return
|
||||
}
|
||||
|
||||
func (dataNode *DataNode) updateDecommissionStatus(c *Cluster, debug, persist bool) (uint32, float64) {
|
||||
var (
|
||||
totalDisk = len(dataNode.DecommissionDiskList)
|
||||
|
||||
@ -253,6 +253,40 @@ func (c *Cluster) deleteAndSyncDecommissionedDisk(dataNode *DataNode, diskPath s
|
||||
return
|
||||
}
|
||||
|
||||
func (c *Cluster) addAndSyncDecommissionSuccessDisk(addr string, diskPath string) (err error) {
|
||||
var dataNode *DataNode
|
||||
if dataNode, err = c.dataNode(addr); err != nil {
|
||||
log.LogWarnf("action[addAndSyncDecommissionSuccessDisk] cannot find dataNode[%s]", addr)
|
||||
return
|
||||
}
|
||||
|
||||
if exist := dataNode.addDecommissionSuccessDisk(diskPath); exist {
|
||||
return
|
||||
}
|
||||
if err = c.syncUpdateDataNode(dataNode); err != nil {
|
||||
dataNode.deleteDecommissionSuccessDisk(diskPath)
|
||||
log.LogWarnf("action[addAndSyncDecommissionSuccessDisk]submit raft failed: %v, delete disks[%v], dataNode[%v]",
|
||||
err, diskPath, dataNode.Addr)
|
||||
return
|
||||
}
|
||||
log.LogInfof("action[addAndSyncDecommissionSuccessDisk] finish, remaining decommissionSuccess disks[%v], dataNode[%v]", dataNode.getDecommissionSuccessDisks(), dataNode.Addr)
|
||||
return
|
||||
}
|
||||
|
||||
func (c *Cluster) deleteAndSyncDecommissionSuccessDisk(dataNode *DataNode, diskPath string) (err error) {
|
||||
if exist := dataNode.deleteDecommissionSuccessDisk(diskPath); !exist {
|
||||
return
|
||||
}
|
||||
if err = c.syncUpdateDataNode(dataNode); err != nil {
|
||||
dataNode.addDecommissionSuccessDisk(diskPath)
|
||||
log.LogWarnf("action[deleteAndSyncDecommissionSuccessDisk]submit raft failed: %v, delete disks[%v], dataNode[%v]",
|
||||
err, diskPath, dataNode.Addr)
|
||||
return
|
||||
}
|
||||
log.LogInfof("action[deleteAndSyncDecommissionSuccessDisk] finish, remaining decommissionSuccess disks[%v], dataNode[%v]", dataNode.getDecommissionSuccessDisks(), dataNode.Addr)
|
||||
return
|
||||
}
|
||||
|
||||
func (c *Cluster) decommissionDisk(dataNode *DataNode, raftForce bool, badDiskPath string,
|
||||
badPartitions []*DataPartition, diskDisable bool,
|
||||
) (err error) {
|
||||
|
||||
@ -760,6 +760,9 @@ func (m *Server) registerAPIRoutes(router *mux.Router) {
|
||||
router.NewRoute().Methods(http.MethodGet, http.MethodPost).
|
||||
Path(proto.QueryDisableDisk).
|
||||
HandlerFunc(m.queryDisableDisk)
|
||||
router.NewRoute().Methods(http.MethodGet, http.MethodPost).
|
||||
Path(proto.QueryDecommissionSuccessDisk).
|
||||
HandlerFunc(m.queryDecommissionSuccessDisk)
|
||||
router.NewRoute().Methods(http.MethodGet, http.MethodPost).
|
||||
Path(proto.CancelDecommissionDisk).
|
||||
HandlerFunc(m.cancelDecommissionDisk)
|
||||
|
||||
@ -495,6 +495,7 @@ type dataNodeValue struct {
|
||||
ZoneName string
|
||||
RdOnly bool
|
||||
DecommissionedDisks []string
|
||||
DecommissionSuccessDisks []string
|
||||
DecommissionStatus uint32
|
||||
DecommissionDstAddr string
|
||||
DecommissionRaftForce bool
|
||||
@ -521,6 +522,7 @@ func newDataNodeValue(dataNode *DataNode) *dataNodeValue {
|
||||
ZoneName: dataNode.ZoneName,
|
||||
RdOnly: dataNode.RdOnly,
|
||||
DecommissionedDisks: dataNode.getDecommissionedDisks(),
|
||||
DecommissionSuccessDisks: dataNode.getDecommissionSuccessDisks(),
|
||||
DecommissionStatus: atomic.LoadUint32(&dataNode.DecommissionStatus),
|
||||
DecommissionDstAddr: dataNode.DecommissionDstAddr,
|
||||
DecommissionRaftForce: dataNode.DecommissionRaftForce,
|
||||
@ -1623,6 +1625,9 @@ func (c *Cluster) loadDataNodes() (err error) {
|
||||
for _, disk := range dnv.DecommissionedDisks {
|
||||
dataNode.addDecommissionedDisk(disk)
|
||||
}
|
||||
for _, disk := range dnv.DecommissionSuccessDisks {
|
||||
dataNode.addDecommissionSuccessDisk(disk)
|
||||
}
|
||||
dataNode.DecommissionStatus = dnv.DecommissionStatus
|
||||
dataNode.DecommissionDstAddr = dnv.DecommissionDstAddr
|
||||
dataNode.DecommissionRaftForce = dnv.DecommissionRaftForce
|
||||
|
||||
@ -175,7 +175,7 @@ type monitorMetrics struct {
|
||||
lcVolMigrateBytes *exporter.GaugeVec
|
||||
lcVolError *exporter.GaugeVec
|
||||
|
||||
diskDecommissioned *exporter.GaugeVec
|
||||
diskDecommissionSuccess *exporter.GaugeVec
|
||||
}
|
||||
|
||||
func newMonitorMetrics(c *Cluster) *monitorMetrics {
|
||||
@ -555,7 +555,7 @@ func (mm *monitorMetrics) start() {
|
||||
mm.lcVolMigrateBytes = exporter.NewGaugeVec(MetricLcVolMigrateBytes, "", []string{"id", "type"})
|
||||
mm.lcVolError = exporter.NewGaugeVec(MetricLcVolError, "", []string{"id", "type"})
|
||||
|
||||
mm.diskDecommissioned = exporter.NewGaugeVec(MetricDiskDecommissionSuccess, "", []string{"addr", "path"})
|
||||
mm.diskDecommissionSuccess = exporter.NewGaugeVec(MetricDiskDecommissionSuccess, "", []string{"addr", "path"})
|
||||
go mm.statMetrics()
|
||||
}
|
||||
|
||||
@ -960,17 +960,17 @@ func (mm *monitorMetrics) setDpMissingTinyExtentMetric() {
|
||||
}
|
||||
|
||||
func (mm *monitorMetrics) setDiskDecommissionedMetric() {
|
||||
mm.diskDecommissioned.Reset()
|
||||
mm.diskDecommissionSuccess.Reset()
|
||||
|
||||
mm.cluster.dataNodes.Range(func(addr, node interface{}) bool {
|
||||
dataNode, ok := node.(*DataNode)
|
||||
if !ok {
|
||||
return true
|
||||
}
|
||||
disks := dataNode.getDecommissionedDisks()
|
||||
disks := dataNode.getDecommissionSuccessDisks()
|
||||
for _, disk := range disks {
|
||||
key := fmt.Sprintf("%s_%s", dataNode.Addr, disk)
|
||||
mm.diskDecommissioned.SetWithLabelValues(1, dataNode.Addr, key)
|
||||
mm.diskDecommissionSuccess.SetWithLabelValues(1, dataNode.Addr, key)
|
||||
}
|
||||
return true
|
||||
})
|
||||
@ -1343,7 +1343,7 @@ func (mm *monitorMetrics) resetAllLeaderMetrics() {
|
||||
mm.diskLost.Reset()
|
||||
mm.dpUnableDecommissionCount.Set(0)
|
||||
mm.dpMissingTinyExtent.Reset()
|
||||
mm.diskDecommissioned.Reset()
|
||||
mm.diskDecommissionSuccess.Reset()
|
||||
mm.dataNodesInactive.Set(0)
|
||||
mm.metaNodesInactive.Set(0)
|
||||
mm.mastersInactive.Set(0)
|
||||
|
||||
@ -243,7 +243,8 @@ const (
|
||||
|
||||
AddLcNode = "/lcNode/add"
|
||||
|
||||
QueryDisableDisk = "/dataNode/queryDisableDisk"
|
||||
QueryDisableDisk = "/dataNode/queryDisableDisk"
|
||||
QueryDecommissionSuccessDisk = "/dataNode/queryDecommissionSuccessDisk"
|
||||
// Operation response
|
||||
GetMetaNodeTaskResponse = "/metaNode/response" // Method: 'POST', ContentType: 'application/json'
|
||||
GetDataNodeTaskResponse = "/dataNode/response" // Method: 'POST', ContentType: 'application/json'
|
||||
|
||||
@ -83,6 +83,7 @@ type DataNodeInfo struct {
|
||||
CpuUtil float64 `json:"cpuUtil"`
|
||||
IoUtils map[string]float64 `json:"ioUtil"`
|
||||
DecommissionedDisk []string
|
||||
DecommissionSuccessDisk []string
|
||||
BackupDataPartitions []uint64
|
||||
MediaType uint32
|
||||
DiskOpLogs []OpLog
|
||||
|
||||
@ -182,6 +182,12 @@ func (api *NodeAPI) QueryDecommissionedDisks(addr string) (disks *proto.Decommis
|
||||
return
|
||||
}
|
||||
|
||||
func (api *NodeAPI) QueryDecommissionSuccessDisks(addr string) (disks *proto.DecommissionedDisks, err error) {
|
||||
disks = &proto.DecommissionedDisks{}
|
||||
err = api.mc.requestWith(disks, newRequest(get, proto.QueryDecommissionSuccessDisk).Header(api.h).addParam("addr", addr))
|
||||
return
|
||||
}
|
||||
|
||||
func (api *NodeAPI) QueryCancelDecommissionedDataNode(addr string) (err error) {
|
||||
err = api.mc.request(newRequest(get, proto.CancelDecommissionDataNode).Header(api.h).addParam("addr", addr))
|
||||
return
|
||||
|
||||
Loading…
Reference in New Issue
Block a user