fix(master): fix some issues found during decommission testing.

close:#1000059418

Signed-off-by: shuqiang-zheng <zhengshuqiang@oppo.com>
This commit is contained in:
shuqiang-zheng 2025-04-21 15:27:07 +08:00 committed by zhumingze1108
parent c0e6d5f0af
commit 3cbb46bd80
18 changed files with 933 additions and 769 deletions

View File

@ -88,6 +88,7 @@ func formatClusterView(cv *proto.ClusterView, cn *proto.ClusterNodeInfo, cp *pro
sb.WriteString(fmt.Sprintf(" EnableAutoDpMetaRepair : %v\n", cv.EnableAutoDpMetaRepair))
sb.WriteString(fmt.Sprintf(" AutoDpMetaRepairParallelCnt : %v\n", cv.AutoDpMetaRepairParallelCnt))
sb.WriteString(fmt.Sprintf(" MarkDiskBrokenThreshold : %v\n", strutil.FormatPercent(cv.MarkDiskBrokenThreshold)))
sb.WriteString(fmt.Sprintf(" DecommissionFirstHostDiskParallelLimit : %v\n", cv.DecommissionFirstHostDiskParallelLimit))
sb.WriteString(fmt.Sprintf(" DecommissionDpLimit : %v\n", cv.DecommissionLimit))
sb.WriteString(fmt.Sprintf(" DecommissionDiskLimit : %v\n", cv.DecommissionDiskLimit))
sb.WriteString(fmt.Sprintf(" DpBackupTimeout : %v\n", cv.DpBackupTimeout))

View File

@ -145,7 +145,7 @@ func (dp *DataPartition) repair(extentType uint8) {
dp.sendAllTinyExtentsToC(extentType, availableTinyExtents, brokenTinyExtents)
// error check
if dp.extentStore.AvailableTinyExtentCnt()+dp.extentStore.BrokenTinyExtentCnt() > storage.TinyExtentCount {
if dp.extentStore.AvailableTinyExtentCnt()+dp.extentStore.BrokenTinyExtentCnt() != storage.TinyExtentCount {
log.LogWarnf("action[repair] partition(%v) GoodTinyExtents(%v) "+
"BadTinyExtents(%v) finish cost[%v] extentType %v.", dp.partitionID, dp.extentStore.AvailableTinyExtentCnt(),
dp.extentStore.BrokenTinyExtentCnt(), time.Since(start).String(), extentType)

View File

@ -24,6 +24,8 @@ import (
"sync/atomic"
"time"
"github.com/cubefs/cubefs/datanode/storage"
syslog "log"
"github.com/cubefs/cubefs/proto"
@ -755,6 +757,7 @@ func (s *DataNode) buildHeartBeatResponse(response *proto.DataNodeHeartbeatRespo
TriggerDiskError: atomic.LoadUint64(&partition.diskErrCnt) > 0,
ForbidWriteOpOfProtoVer0: dpForbid,
ReadOnlyReasons: partition.ReadOnlyReasons(),
IsMissingTinyExtent: partition.extentStore.AvailableTinyExtentCnt()+partition.extentStore.BrokenTinyExtentCnt() < storage.TinyExtentCount,
}
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

@ -2090,7 +2090,7 @@ func sendErrReply(w http.ResponseWriter, r *http.Request, httpReply *proto.HTTPR
}
}
func parseRequestToUpdateDecommissionFirstHostTokenLimit(r *http.Request) (addr string, limit uint64, err error) {
func parseRequestToUpdateDecommissionFirstHostParallelLimit(r *http.Request) (addr string, limit uint64, err error) {
if err = r.ParseForm(); err != nil {
return
}
@ -2101,8 +2101,8 @@ func parseRequestToUpdateDecommissionFirstHostTokenLimit(r *http.Request) (addr
}
var value string
if value = r.FormValue(decommissionFirstHostTokenLimit); value == "" {
err = keyNotFound(decommissionFirstHostTokenLimit)
if value = r.FormValue(decommissionFirstHostParallelLimit); value == "" {
err = keyNotFound(decommissionFirstHostParallelLimit)
return
}
@ -2114,14 +2114,14 @@ func parseRequestToUpdateDecommissionFirstHostTokenLimit(r *http.Request) (addr
return
}
func parseRequestToUpdateDecommissionFirstHostDiskTokenLimit(r *http.Request) (limit uint64, err error) {
func parseRequestToUpdateDecommissionFirstHostDiskParallelLimit(r *http.Request) (limit uint64, err error) {
if err = r.ParseForm(); err != nil {
return
}
var value string
if value = r.FormValue(decommissionFirstHostDiskTokenLimit); value == "" {
err = keyNotFound(decommissionFirstHostDiskTokenLimit)
if value = r.FormValue(decommissionFirstHostDiskParallelLimit); value == "" {
err = keyNotFound(decommissionFirstHostDiskParallelLimit)
return
}

View File

@ -912,41 +912,41 @@ func (m *Server) getCluster(w http.ResponseWriter, r *http.Request) {
}
cv := &proto.ClusterView{
Name: m.cluster.Name,
CreateTime: time.Unix(m.cluster.CreateTime, 0).Format(proto.TimeFormat),
LeaderAddr: m.leaderInfo.addr,
DisableAutoAlloc: m.cluster.DisableAutoAllocate,
ForbidMpDecommission: m.cluster.ForbidMpDecommission,
MetaNodeThreshold: m.cluster.cfg.MetaNodeThreshold,
Applied: m.fsm.applied,
MaxDataPartitionID: m.cluster.idAlloc.dataPartitionID,
MaxMetaNodeID: m.cluster.idAlloc.commonID,
MaxMetaPartitionID: m.cluster.idAlloc.metaPartitionID,
VolDeletionDelayTimeHour: m.cluster.cfg.volDelayDeleteTimeHour,
MetaNodeGOGC: m.cluster.cfg.metaNodeGOGC,
DataNodeGOGC: m.cluster.cfg.dataNodeGOGC,
DpRepairTimeout: m.cluster.GetDecommissionDataPartitionRecoverTimeOut().String(),
DpBackupTimeout: m.cluster.GetDecommissionDataPartitionBackupTimeOut().String(),
MarkDiskBrokenThreshold: m.cluster.getMarkDiskBrokenThreshold(),
EnableAutoDpMetaRepair: m.cluster.getEnableAutoDpMetaRepair(),
AutoDpMetaRepairParallelCnt: m.cluster.GetAutoDpMetaRepairParallelCnt(),
EnableAutoDecommission: m.cluster.AutoDecommissionDiskIsEnabled(),
AutoDecommissionDiskInterval: m.cluster.GetAutoDecommissionDiskInterval().String(),
DecommissionLimit: atomic.LoadUint64(&m.cluster.DecommissionLimit),
DiskToRepairDpLimit: atomic.LoadUint64(&m.cluster.DecommissionFirstHostDiskTokenLimit),
DecommissionDiskLimit: m.cluster.GetDecommissionDiskLimit(),
DpTimeout: (time.Duration(m.cluster.getDataPartitionTimeoutSec()) * time.Second).String(),
MpTimeout: (time.Duration(m.cluster.getMetaPartitionTimeoutSec()) * time.Second).String(),
MasterNodes: make([]proto.NodeView, 0),
MetaNodes: make([]proto.NodeView, 0),
DataNodes: make([]proto.NodeView, 0),
VolStatInfo: make([]*proto.VolStatInfo, 0),
StatOfStorageClass: make([]*proto.StatOfStorageClass, 0),
StatMigrateStorageClass: make([]*proto.StatOfStorageClass, 0),
BadPartitionIDs: make([]proto.BadPartitionView, 0),
BadMetaPartitionIDs: make([]proto.BadPartitionView, 0),
ForbidWriteOpOfProtoVer0: m.cluster.cfg.forbidWriteOpOfProtoVer0,
LegacyDataMediaType: m.cluster.legacyDataMediaType,
Name: m.cluster.Name,
CreateTime: time.Unix(m.cluster.CreateTime, 0).Format(proto.TimeFormat),
LeaderAddr: m.leaderInfo.addr,
DisableAutoAlloc: m.cluster.DisableAutoAllocate,
ForbidMpDecommission: m.cluster.ForbidMpDecommission,
MetaNodeThreshold: m.cluster.cfg.MetaNodeThreshold,
Applied: m.fsm.applied,
MaxDataPartitionID: m.cluster.idAlloc.dataPartitionID,
MaxMetaNodeID: m.cluster.idAlloc.commonID,
MaxMetaPartitionID: m.cluster.idAlloc.metaPartitionID,
VolDeletionDelayTimeHour: m.cluster.cfg.volDelayDeleteTimeHour,
MetaNodeGOGC: m.cluster.cfg.metaNodeGOGC,
DataNodeGOGC: m.cluster.cfg.dataNodeGOGC,
DpRepairTimeout: m.cluster.GetDecommissionDataPartitionRecoverTimeOut().String(),
DpBackupTimeout: m.cluster.GetDecommissionDataPartitionBackupTimeOut().String(),
MarkDiskBrokenThreshold: m.cluster.getMarkDiskBrokenThreshold(),
EnableAutoDpMetaRepair: m.cluster.getEnableAutoDpMetaRepair(),
AutoDpMetaRepairParallelCnt: m.cluster.GetAutoDpMetaRepairParallelCnt(),
EnableAutoDecommission: m.cluster.AutoDecommissionDiskIsEnabled(),
AutoDecommissionDiskInterval: m.cluster.GetAutoDecommissionDiskInterval().String(),
DecommissionLimit: atomic.LoadUint64(&m.cluster.DecommissionLimit),
DecommissionFirstHostDiskParallelLimit: atomic.LoadUint64(&m.cluster.DecommissionFirstHostDiskParallelLimit),
DecommissionDiskLimit: m.cluster.GetDecommissionDiskLimit(),
DpTimeout: (time.Duration(m.cluster.getDataPartitionTimeoutSec()) * time.Second).String(),
MpTimeout: (time.Duration(m.cluster.getMetaPartitionTimeoutSec()) * time.Second).String(),
MasterNodes: make([]proto.NodeView, 0),
MetaNodes: make([]proto.NodeView, 0),
DataNodes: make([]proto.NodeView, 0),
VolStatInfo: make([]*proto.VolStatInfo, 0),
StatOfStorageClass: make([]*proto.StatOfStorageClass, 0),
StatMigrateStorageClass: make([]*proto.StatOfStorageClass, 0),
BadPartitionIDs: make([]proto.BadPartitionView, 0),
BadMetaPartitionIDs: make([]proto.BadPartitionView, 0),
ForbidWriteOpOfProtoVer0: m.cluster.cfg.forbidWriteOpOfProtoVer0,
LegacyDataMediaType: m.cluster.legacyDataMediaType,
RaftPartitionCanUsingDifferentPortEnabled: m.cluster.RaftPartitionCanUsingDifferentPortEnabled(),
FlashNodes: make([]proto.NodeView, 0),
FlashNodeHandleReadTimeout: m.cluster.cfg.flashNodeHandleReadTimeout,
@ -3448,10 +3448,7 @@ func (m *Server) cancelDecommissionDataNode(w http.ResponseWriter, r *http.Reque
rstMsg string
offLineAddr string
err error
// dps []uint64
node *DataNode
//zone *Zone
//ns *nodeSet
node *DataNode
)
metric := exporter.NewTPCnt(apiToMetricsName(proto.PauseDecommissionDataNode))
@ -3468,23 +3465,6 @@ func (m *Server) cancelDecommissionDataNode(w http.ResponseWriter, r *http.Reque
sendErrReply(w, r, newErrHTTPReply(err))
return
}
/*
if zone, err = m.cluster.t.getZone(node.ZoneName); err != nil {
ret := fmt.Sprintf("action[queryDataNodeDecoProgress] find datanode[%s] zone failed[%v]",
node.Addr, err.Error())
log.LogWarnf("%v", ret)
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: ret})
return
}
if ns, err = zone.getNodeSet(node.NodeSetID); err != nil {
ret := fmt.Sprintf("action[queryDataNodeDecoProgress] find datanode[%s] nodeset[%v] failed[%v]",
node.Addr, node.NodeSetID, err.Error())
log.LogWarnf("%v", ret)
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: ret})
return
}
*/
// can alloc dp
node.ToBeOffline = false
@ -6292,39 +6272,39 @@ func (m *Server) associateVolWithUser(userID, volName string) error {
return nil
}
func (m *Server) updateDecommissionFirstHostTokenLimit(w http.ResponseWriter, r *http.Request) {
func (m *Server) updateDecommissionFirstHostParallelLimit(w http.ResponseWriter, r *http.Request) {
var (
addr string
limit uint64
err error
)
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminUpdateDecommissionFirstHostTokenLimit))
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminUpdateDecommissionFirstHostParallelLimit))
defer func() {
doStatAndMetric(proto.AdminUpdateDecommissionFirstHostTokenLimit, metric, err, nil)
doStatAndMetric(proto.AdminUpdateDecommissionFirstHostParallelLimit, metric, err, nil)
}()
if addr, limit, err = parseRequestToUpdateDecommissionFirstHostTokenLimit(r); err != nil {
if addr, limit, err = parseRequestToUpdateDecommissionFirstHostParallelLimit(r); err != nil {
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()})
return
}
if err = m.cluster.setDecommissionFirstHostTokenLimit(addr, limit); err != nil {
if err = m.cluster.setDecommissionFirstHostParallelLimit(addr, limit); err != nil {
sendErrReply(w, r, newErrHTTPReply(fmt.Errorf("set master not worked %v", err)))
return
}
rstMsg := fmt.Sprintf("set decommission first host token limit to %v successfully", limit)
log.LogDebugf("action[updateDecommissionFirstHostTokenLimit] %v", rstMsg)
rstMsg := fmt.Sprintf("set decommission first host parallel limit to %v successfully", limit)
log.LogDebugf("action[updateDecommissionFirstHostParallelLimit] %v", rstMsg)
sendOkReply(w, r, newSuccessHTTPReply(rstMsg))
}
func (m *Server) queryDecommissionFirstHostTokenLimit(w http.ResponseWriter, r *http.Request) {
func (m *Server) queryDecommissionFirstHostParallelLimit(w http.ResponseWriter, r *http.Request) {
var (
addr string
dataNode *DataNode
err error
)
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminQueryDecommissionFirstHostTokenLimit))
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminQueryDecommissionFirstHostParallelLimit))
defer func() {
doStatAndMetric(proto.AdminQueryDecommissionFirstHostTokenLimit, metric, nil, nil)
doStatAndMetric(proto.AdminQueryDecommissionFirstHostParallelLimit, metric, nil, nil)
}()
if err = r.ParseForm(); err != nil {
@ -6343,62 +6323,62 @@ func (m *Server) queryDecommissionFirstHostTokenLimit(w http.ResponseWriter, r *
return
}
limit := atomic.LoadUint64(&dataNode.DecommissionFirstHostTokenLimit)
rstMsg := fmt.Sprintf("dataNode(%v) decommission first host token limit is %v", addr, limit)
log.LogDebugf("action[queryDecommissionFirstHostTokenTokenLimit] %v", rstMsg)
limit := atomic.LoadUint64(&dataNode.DecommissionFirstHostParallelLimit)
rstMsg := fmt.Sprintf("dataNode(%v) decommission first host parallel limit is %v", addr, limit)
log.LogDebugf("action[queryDecommissionFirstHostTokenParallelLimit] %v", rstMsg)
sendOkReply(w, r, newSuccessHTTPReply(rstMsg))
}
func (m *Server) updateDecommissionFirstHostDiskTokenLimit(w http.ResponseWriter, r *http.Request) {
func (m *Server) updateDecommissionFirstHostDiskParallelLimit(w http.ResponseWriter, r *http.Request) {
var (
limit uint64
err error
)
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminUpdateDecommissionFirstHostDiskTokenLimit))
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminUpdateDecommissionFirstHostDiskParallelLimit))
defer func() {
doStatAndMetric(proto.AdminUpdateDecommissionFirstHostDiskTokenLimit, metric, err, nil)
doStatAndMetric(proto.AdminUpdateDecommissionFirstHostDiskParallelLimit, metric, err, nil)
}()
if limit, err = parseRequestToUpdateDecommissionFirstHostDiskTokenLimit(r); err != nil {
if limit, err = parseRequestToUpdateDecommissionFirstHostDiskParallelLimit(r); err != nil {
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()})
return
}
if err = m.cluster.setDecommissionFirstHostDiskTokenLimit(limit); err != nil {
if err = m.cluster.setDecommissionFirstHostDiskParallelLimit(limit); err != nil {
sendErrReply(w, r, newErrHTTPReply(fmt.Errorf("set master not worked %v", err)))
return
}
rstMsg := fmt.Sprintf("set decommission first host disk token limit to %v successfully", limit)
log.LogDebugf("action[updateDecommissionFirstHostDiskTokenLimit] %v", rstMsg)
rstMsg := fmt.Sprintf("set decommission first host disk parallel limit to %v successfully", limit)
log.LogDebugf("action[updateDecommissionFirstHostDiskParallelLimit] %v", rstMsg)
sendOkReply(w, r, newSuccessHTTPReply(rstMsg))
}
func (m *Server) queryDecommissionFirstHostTokenDiskTokenLimit(w http.ResponseWriter, r *http.Request) {
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminQueryDecommissionFirstHostDiskTokenLimit))
func (m *Server) queryDecommissionFirstHostDiskParallelLimit(w http.ResponseWriter, r *http.Request) {
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminQueryDecommissionFirstHostDiskParallelLimit))
defer func() {
doStatAndMetric(proto.AdminQueryDecommissionFirstHostDiskTokenLimit, metric, nil, nil)
doStatAndMetric(proto.AdminQueryDecommissionFirstHostDiskParallelLimit, metric, nil, nil)
}()
limit := atomic.LoadUint64(&m.cluster.DecommissionFirstHostDiskTokenLimit)
rstMsg := fmt.Sprintf("decommission first host disk token limit is %v", limit)
log.LogDebugf("action[queryDecommissionFirstHostTokenDiskTokenLimit] %v", rstMsg)
limit := atomic.LoadUint64(&m.cluster.DecommissionFirstHostDiskParallelLimit)
rstMsg := fmt.Sprintf("decommission first host disk parallel limit is %v", limit)
log.LogDebugf("action[queryDecommissionFirstHostDiskParallelLimit] %v", rstMsg)
sendOkReply(w, r, newSuccessHTTPReply(rstMsg))
}
func (m *Server) queryDecommissionFirstHostTokenInfo(w http.ResponseWriter, r *http.Request) {
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminQueryDecommissionFirstHostTokenInfo))
func (m *Server) queryDecommissionFirstHostParallelInfo(w http.ResponseWriter, r *http.Request) {
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminQueryDecommissionFirstHostParallelInfo))
defer func() {
doStatAndMetric(proto.AdminQueryDecommissionFirstHostTokenInfo, metric, nil, nil)
doStatAndMetric(proto.AdminQueryDecommissionFirstHostParallelInfo, metric, nil, nil)
}()
infos := make([]*DataNodeToDecommissionRepairDpInfo, 0)
m.cluster.DataNodeToDecommissionRepairDpMap.Range(func(key, value interface{}) bool {
info := value.(*DataNodeToDecommissionRepairDpInfo)
if atomic.LoadUint64(&info.curParallel) != 0 {
if atomic.LoadUint64(&info.CurParallel) != 0 {
infos = append(infos, info)
}
return true
})
log.LogDebugf("action[queryDiskToRepairDpInfo] %v", infos)
log.LogDebugf("action[queryDecommissionFirstHostParallelInfo] %v", infos)
sendOkReply(w, r, newSuccessHTTPReply(infos))
}
@ -8068,8 +8048,6 @@ func (m *Server) cancelDecommissionDisk(w http.ResponseWriter, r *http.Request)
offLineAddr, diskPath string
err error
dataNode *DataNode
//zone *Zone
//ns *nodeSet
)
metric := exporter.NewTPCnt("req_cancelDecommissionDisk")
@ -8096,23 +8074,6 @@ func (m *Server) cancelDecommissionDisk(w http.ResponseWriter, r *http.Request)
return
}
/*
if zone, err = m.cluster.t.getZone(dataNode.ZoneName); err != nil {
ret := fmt.Sprintf("action[cancelDecommissionDisk] find datanode[%s] zone failed[%v]",
dataNode.Addr, err.Error())
log.LogWarnf("%v", ret)
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: ret})
return
}
if ns, err = zone.getNodeSet(dataNode.NodeSetID); err != nil {
ret := fmt.Sprintf("action[cancelDecommissionDisk] find datanode[%s] nodeset[%v] failed[%v]",
dataNode.Addr, dataNode.NodeSetID, err.Error())
log.LogWarnf("%v", ret)
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: ret})
return
}
*/
disk := value.(*DecommissionDisk)
status := disk.GetDecommissionStatus()
if status == DecommissionSuccess || status == DecommissionFail {

View File

@ -497,6 +497,40 @@ func TestDisk(t *testing.T) {
addr := mds5Addr
disk := "/cfs"
decommissionDisk(addr, disk, t)
cancelDecommissionDisk(addr, disk, t)
}
func cancelDecommissionDisk(addr, path string, t *testing.T) {
reqURL := fmt.Sprintf("%v%v?addr=%v&disk=%v",
hostAddr, proto.CancelDecommissionDisk, addr, path)
mocktest.Log(t, reqURL)
resp, err := http.Get(reqURL)
if err != nil {
t.Errorf("err is %v", err)
return
}
mocktest.Println(resp.StatusCode)
defer resp.Body.Close()
body, err := io.ReadAll(resp.Body)
if err != nil {
t.Errorf("err is %v", err)
return
}
mocktest.Println(string(body))
if resp.StatusCode != http.StatusOK {
t.Errorf("status code[%v]", resp.StatusCode)
return
}
reply := &proto.HTTPReply{}
if err = json.Unmarshal(body, reply); err != nil {
t.Error(err)
return
}
key := fmt.Sprintf("%s_%s", addr, path)
_, ok := server.cluster.DecommissionDisks.Load(key)
if ok {
t.Errorf("disk should be removed from DecommissionDisks")
}
}
func decommissionDisk(addr, path string, t *testing.T) {
@ -817,10 +851,11 @@ func TestDataPartitionDecommission(t *testing.T) {
time.Sleep(5 * time.Second)
partition := vol.dataPartitions.partitions[0]
offlineAddr := partition.Hosts[0]
reqURL := fmt.Sprintf("%v%v?name=%v&id=%v&addr=%v",
hostAddr, proto.AdminDecommissionDataPartition, vol.Name, partition.PartitionID, offlineAddr)
reqURL := fmt.Sprintf("%v%v?name=%v&id=%v&addr=%v&weight=%v",
hostAddr, proto.AdminDecommissionDataPartition, vol.Name, partition.PartitionID, offlineAddr, highPriorityDecommissionWeight)
process(reqURL, t)
require.EqualValues(t, markDecommission, partition.GetDecommissionStatus())
require.EqualValues(t, highPriorityDecommissionWeight, partition.DecommissionWeight)
}
// func TestGetAllVols(t *testing.T) {
@ -1779,6 +1814,31 @@ func TestSetDiscardDp(t *testing.T) {
require.False(t, dp.IsDiscard)
}
func TestUpdateDecommissionFirstHostParallelLimit(t *testing.T) {
reqUrl := fmt.Sprintf("%v%v", hostAddr, proto.AdminUpdateDecommissionFirstHostParallelLimit)
dataNode, _ := server.cluster.dataNode(mds1Addr)
oldVal := dataNode.DecommissionFirstHostParallelLimit
setVal := oldVal + 1
setUrl := fmt.Sprintf("%v?%v=%v&%v=%v", reqUrl, addrKey, mds1Addr, decommissionFirstHostParallelLimit, setVal)
unsetUrl := fmt.Sprintf("%v?%v=%v&%v=%v", reqUrl, addrKey, mds1Addr, decommissionFirstHostParallelLimit, oldVal)
process(setUrl, t)
require.EqualValues(t, setVal, dataNode.DecommissionFirstHostParallelLimit)
process(unsetUrl, t)
require.EqualValues(t, oldVal, dataNode.DecommissionFirstHostParallelLimit)
}
func TestUpdateDecommissionFirstHostDiskParallelLimit(t *testing.T) {
reqUrl := fmt.Sprintf("%v%v", hostAddr, proto.AdminUpdateDecommissionFirstHostDiskParallelLimit)
oldVal := server.cluster.DecommissionFirstHostDiskParallelLimit
setVal := oldVal + 1
setUrl := fmt.Sprintf("%v?%v=%v", reqUrl, decommissionFirstHostDiskParallelLimit, setVal)
unsetUrl := fmt.Sprintf("%v?%v=%v", reqUrl, decommissionFirstHostDiskParallelLimit, oldVal)
process(setUrl, t)
require.EqualValues(t, setVal, server.cluster.DecommissionFirstHostDiskParallelLimit)
process(unsetUrl, t)
require.EqualValues(t, oldVal, server.cluster.DecommissionFirstHostDiskParallelLimit)
}
func TestSetDecommissionDiskLimit(t *testing.T) {
oldVal := server.cluster.GetDecommissionDiskLimit()
reqUrl := fmt.Sprintf("%v%v", hostAddr, proto.AdminUpdateDecommissionDiskLimit)

View File

@ -89,29 +89,29 @@ type ClusterTopoSubItem struct {
type DataNodeToDecommissionRepairDpInfo struct {
mu sync.Mutex
curParallel uint64
addr string
diskToDecommissionRepairDpMap map[string]*DiskToDecommissionRepairDpInfo
CurParallel uint64
Addr string
DiskToDecommissionRepairDpMap map[string]*DiskToDecommissionRepairDpInfo
}
type DiskToDecommissionRepairDpInfo struct {
curParallel uint64
diskPath string
repairingDps map[uint64]struct{}
CurParallel uint64
DiskPath string
RepairingDps map[uint64]struct{}
}
// nolint: structcheck
type ClusterDecommission struct {
BadDataPartitionIds *sync.Map
BadMetaPartitionIds *sync.Map
DecommissionDisks sync.Map
DataNodeToDecommissionRepairDpMap sync.Map
DecommissionFirstHostDiskTokenLimit uint64
DecommissionLimit uint64
AutoDecommissionDiskMux sync.Mutex
DecommissionDiskLimit uint32
MarkDiskBrokenThreshold atomicutil.Float64
badPartitionMutex sync.RWMutex // BadDataPartitionIds and BadMetaPartitionIds operate mutex
BadDataPartitionIds *sync.Map
BadMetaPartitionIds *sync.Map
DecommissionDisks sync.Map
DataNodeToDecommissionRepairDpMap sync.Map
DecommissionFirstHostDiskParallelLimit uint64
DecommissionLimit uint64
AutoDecommissionDiskMux sync.Mutex
DecommissionDiskLimit uint32
MarkDiskBrokenThreshold atomicutil.Float64
badPartitionMutex sync.RWMutex // BadDataPartitionIds and BadMetaPartitionIds operate mutex
ForbidMpDecommission bool
EnableAutoDpMetaRepair atomicutil.Bool
@ -456,7 +456,7 @@ func newCluster(name string, leaderInfo *LeaderInfo, fsm *MetadataFsm, partition
c.QosAcceptLimit = rate.NewLimiter(rate.Limit(c.cfg.QosMasterAcceptLimit), proto.QosDefaultBurst)
c.apiLimiter = newApiLimiter()
c.DecommissionLimit = defaultDecommissionParallelLimit
c.DecommissionFirstHostDiskTokenLimit = defaultDecommissionFirstHostDiskTokenLimit
c.DecommissionFirstHostDiskParallelLimit = defaultDecommissionFirstHostDiskParallelLimit
c.checkAutoCreateDataPartition = false
c.masterClient = masterSDK.NewMasterClient(nil, false)
c.masterClient.SetTransport(proto.GetHttpTransporter(&proto.HttpCfg{
@ -2594,25 +2594,6 @@ func (c *Cluster) decommissionSingleDp(dp *DataPartition, newAddr, offlineAddr s
dp.PartitionID, newReplica.Addr, newReplica.Status)
}
if newReplica.isRepairing() { // wait for repair
masterNode, _ := dp.getReplica(dp.Hosts[0])
duration := time.Unix(masterNode.ReportTime, 0).Sub(time.Unix(newReplica.ReportTime, 0))
diskErrReplicas := dp.getAllDiskErrorReplica()
if isReplicasContainsHost(diskErrReplicas, dp.Hosts[0]) {
err = fmt.Errorf("action[decommissionSingleDp] dp %v host[0] %v is unavailable",
dp.PartitionID, dp.Hosts[0])
dp.DecommissionNeedRollback = false
newReplica.Status = proto.Unavailable // remove from data partition check
log.LogWarnf("action[decommissionSingleDp] dp %v err:%v", dp.PartitionID, err)
goto ERR
}
if math.Abs(duration.Minutes()) > 10 {
err = fmt.Errorf("action[decommissionSingleDp] dp %v host[0] %v is down",
dp.PartitionID, masterNode.Addr)
dp.DecommissionNeedRollback = false
newReplica.Status = proto.Unavailable // remove from data partition check
log.LogWarnf("action[decommissionSingleDp] dp %v err:%v", dp.PartitionID, err)
goto ERR
}
if time.Since(dp.RecoverStartTime) > c.GetDecommissionDataPartitionRecoverTimeOut() {
err = fmt.Errorf("action[decommissionSingleDp] dp %v new replica %v repair time out:%v",
dp.PartitionID, newAddr, time.Since(dp.RecoverStartTime))
@ -2629,6 +2610,22 @@ func (c *Cluster) decommissionSingleDp(dp *DataPartition, newAddr, offlineAddr s
break
}
} else {
masterNode, _ := dp.getReplica(dp.Hosts[0])
duration := time.Unix(masterNode.ReportTime, 0).Sub(time.Unix(newReplica.ReportTime, 0))
diskErrReplicas := dp.getAllDiskErrorReplica()
if isReplicasContainsHost(diskErrReplicas, dp.Hosts[0]) || math.Abs(duration.Minutes()) > 10 {
if isReplicasContainsHost(diskErrReplicas, dp.Hosts[0]) {
err = fmt.Errorf("action[decommissionSingleDp] dp %v host[0] %v is unavailable",
dp.PartitionID, dp.Hosts[0])
} else {
err = fmt.Errorf("action[decommissionSingleDp] dp %v host[0] %v is down",
dp.PartitionID, masterNode.Addr)
}
dp.DecommissionNeedRollback = true
newReplica.Status = proto.Unavailable // remove from data partition check
log.LogWarnf("action[decommissionSingleDp] dp %v err:%v", dp.PartitionID, err)
goto ERR
}
// newReplica repair failed or encounter bad disk ,need rollback
if newReplica.isUnavailable() {
err = fmt.Errorf("action[decommissionSingleDp] dp %v new replica %v is Unavailable",
@ -3167,7 +3164,6 @@ func (c *Cluster) addDataPartitionRaftMember(dp *DataPartition, addPeer proto.Pe
dp.Hosts = append(dp.Hosts, addPeer.Addr)
dp.Peers = append(dp.Peers, addPeer)
dp.Unlock()
// send task to leader addr first,if need to retry,then send to other addr
@ -5298,7 +5294,7 @@ func (c *Cluster) TryDecommissionDisk(disk *DecommissionDisk) {
disk.decommissionInfo(), dp.PartitionID)
IgnoreDecommissionDps = append(IgnoreDecommissionDps, proto.IgnoreDecommissionDP{
PartitionID: dp.PartitionID,
ErrMsg: proto.ErrPerformingDecommission.Error(),
ErrMsg: proto.ErrWaitForAutoAddReplica.Error(),
})
continue
} else {
@ -5898,25 +5894,25 @@ func (c *Cluster) setDecommissionDiskLimit(limit uint32) (err error) {
return
}
func (c *Cluster) setDecommissionFirstHostTokenLimit(addr string, limit uint64) (err error) {
func (c *Cluster) setDecommissionFirstHostParallelLimit(addr string, limit uint64) (err error) {
dataNode, err := c.dataNode(addr)
if err != nil {
log.LogErrorf("[setDecommissionFirstHostTokenLimit] failed , err(%v)", err)
log.LogErrorf("[setDecommissionFirstHostParallelLimit] failed , err(%v)", err)
return
}
atomic.StoreUint64(&dataNode.DecommissionFirstHostTokenLimit, limit)
atomic.StoreUint64(&dataNode.DecommissionFirstHostParallelLimit, limit)
if err = c.syncUpdateDataNode(dataNode); err != nil {
log.LogErrorf("[setDecommissionFirstHostTokenLimit] failed to set DecommissionFirstHostTokenLimit , err(%v)", err)
log.LogErrorf("[setDecommissionFirstHostParallelLimit] failed to set DecommissionFirstHostParallelLimit , err(%v)", err)
err = proto.ErrPersistenceByRaft
return
}
return
}
func (c *Cluster) setDecommissionFirstHostDiskTokenLimit(limit uint64) (err error) {
atomic.StoreUint64(&c.DecommissionFirstHostDiskTokenLimit, limit)
func (c *Cluster) setDecommissionFirstHostDiskParallelLimit(limit uint64) (err error) {
atomic.StoreUint64(&c.DecommissionFirstHostDiskParallelLimit, limit)
if err = c.syncPutCluster(); err != nil {
log.LogErrorf("[setDecommissionFirstHostDiskTokenLimit] failed to set DecommissionFirstHostDiskTokenLimit, err(%v)", err)
log.LogErrorf("[setDecommissionFirstHostDiskParallelLimit] failed to set DecommissionFirstHostDiskParallelLimit, err(%v)", err)
err = proto.ErrPersistenceByRaft
return
}

View File

@ -52,111 +52,111 @@ const (
forbiddenKey = "forbidden"
deleteVolKey = "delete"
forceDelVolKey = "forceDelVol"
ebsBlkSizeKey = "ebsBlkSize"
cacheThresholdKey = "cacheThreshold"
clientVersion = "version"
domainIdKey = "domainId"
volOwnerKey = "owner"
volAuthKey = "authKey"
replicaNumKey = "replicaNum"
followerReadKey = "followerRead"
authenticateKey = "authenticate"
akKey = "ak"
keywordsKey = "keywords"
zoneNameKey = "zoneName"
nodesetIdKey = "nodesetId"
crossZoneKey = "crossZone"
normalZonesFirstKey = "normalZonesFirst"
userKey = "user"
nodeDeleteBatchCountKey = "batchCount"
nodeMarkDeleteRateKey = "markDeleteRate"
nodeDeleteWorkerSleepMs = "deleteWorkerSleepMs"
nodeAutoRepairRateKey = "autoRepairRate"
nodeDpRepairTimeOutKey = "dpRepairTimeOut"
nodeDpBackupKey = "dpBackupTimeout"
nodeDpMaxRepairErrCntKey = "dpMaxRepairErrCnt"
clusterLoadFactorKey = "loadFactor"
maxDpCntLimitKey = "maxDpCntLimit"
maxMpCntLimitKey = "maxMpCntLimit"
clusterCreateTimeKey = "clusterCreateTime"
descriptionKey = "description"
dpSelectorNameKey = "dpSelectorName"
dpSelectorParmKey = "dpSelectorParm"
nodeTypeKey = "nodeType"
ratio = "ratio"
rdOnlyKey = "rdOnly"
srcAddrKey = "srcAddr"
targetAddrKey = "targetAddr"
forceKey = "force"
raftForceDelKey = "raftForceDel"
weightKey = "weight"
enablePosixAclKey = "enablePosixAcl"
enableTxMaskKey = "enableTxMask"
txTimeoutKey = "txTimeout"
txConflictRetryNumKey = "txConflictRetryNum"
txConflictRetryIntervalKey = "txConflictRetryInterval"
txOpLimitKey = "txOpLimit"
txForceResetKey = "txForceReset"
QosEnableKey = "qosEnable"
DiskEnableKey = "diskenable"
IopsWKey = "iopsWKey"
IopsRKey = "iopsRKey"
FlowWKey = "flowWKey"
FlowRKey = "flowRKey"
ClientReqPeriod = "reqPeriod"
ClientTriggerCnt = "triggerCnt"
QosMasterLimit = "qosLimit"
decommissionFirstHostDiskTokenLimit = "decommissionFirstHostDiskTokenLimit"
decommissionFirstHostTokenLimit = "decommissionFirstHostTokenLimit"
decommissionLimit = "decommissionLimit"
DiskDisableKey = "diskDisable"
Limit = "limit"
TimeOut = "timeout"
CountByMeta = "countByMeta"
dpReadOnlyWhenVolFull = "dpReadOnlyWhenVolFull"
PeriodicKey = "periodic"
IPKey = "ip"
OperateKey = "op"
UIDKey = "uid"
CapacityKey = "capacity"
configKey = "config"
MaxFilesKey = "maxFiles"
MaxBytesKey = "maxBytes"
quotaKey = "quotaId"
enableQuota = "enableQuota"
dpDiscardKey = "dpDiscard"
ignoreDiscardKey = "ignoreDiscard"
TrashIntervalKey = "trashInterval"
ClientIDKey = "clientIDKey"
verSeqKey = "verSeq"
Periodic = "periodic"
DecommissionType = "decommissionType"
decommissionDiskLimit = "decommissionDiskLimit"
dpRepairBlockSizeKey = "dpRepairBlockSize"
markDiskBrokenThresholdKey = "markDiskBrokenThreshold"
decommissionTypeKey = "decommissionType"
autoDecommissionDiskKey = "autoDecommissionDisk"
autoDecommissionDiskIntervalKey = "autoDecommissionDiskInterval"
autoDpMetaRepairKey = "autoDpMetaRepair"
autoDpMetaRepairParallelCntKey = "autoDpMetaRepairParallelCnt"
dpTimeoutKey = "dpTimeout"
mpTimeoutKey = "mpTimeout"
ShowAll = "showAll"
trashIntervalKey = "trashInterval"
accessTimeIntervalKey = "accessTimeValidInterval"
enablePersistAccessTimeKey = "enablePersistAccessTime"
mediaTypeKey = "mediaType"
allowedStorageClassKey = "allowedStorageClass"
volStorageClassKey = "volStorageClass"
opLogDimensionKey = "opLogDimension"
volNameKey = "volName"
dpIdKey = "dpId"
diskNameKey = "diskName"
forbidWriteOpOfProtoVersion0 = "forbidWriteOpOfProtoVersion0"
quotaClass = "quotaClass"
quotaOfClass = "quotaOfStorageClass"
dataMediaTypeKey = "dataMediaType"
forceDelVolKey = "forceDelVol"
ebsBlkSizeKey = "ebsBlkSize"
cacheThresholdKey = "cacheThreshold"
clientVersion = "version"
domainIdKey = "domainId"
volOwnerKey = "owner"
volAuthKey = "authKey"
replicaNumKey = "replicaNum"
followerReadKey = "followerRead"
authenticateKey = "authenticate"
akKey = "ak"
keywordsKey = "keywords"
zoneNameKey = "zoneName"
nodesetIdKey = "nodesetId"
crossZoneKey = "crossZone"
normalZonesFirstKey = "normalZonesFirst"
userKey = "user"
nodeDeleteBatchCountKey = "batchCount"
nodeMarkDeleteRateKey = "markDeleteRate"
nodeDeleteWorkerSleepMs = "deleteWorkerSleepMs"
nodeAutoRepairRateKey = "autoRepairRate"
nodeDpRepairTimeOutKey = "dpRepairTimeOut"
nodeDpBackupKey = "dpBackupTimeout"
nodeDpMaxRepairErrCntKey = "dpMaxRepairErrCnt"
clusterLoadFactorKey = "loadFactor"
maxDpCntLimitKey = "maxDpCntLimit"
maxMpCntLimitKey = "maxMpCntLimit"
clusterCreateTimeKey = "clusterCreateTime"
descriptionKey = "description"
dpSelectorNameKey = "dpSelectorName"
dpSelectorParmKey = "dpSelectorParm"
nodeTypeKey = "nodeType"
ratio = "ratio"
rdOnlyKey = "rdOnly"
srcAddrKey = "srcAddr"
targetAddrKey = "targetAddr"
forceKey = "force"
raftForceDelKey = "raftForceDel"
weightKey = "weight"
enablePosixAclKey = "enablePosixAcl"
enableTxMaskKey = "enableTxMask"
txTimeoutKey = "txTimeout"
txConflictRetryNumKey = "txConflictRetryNum"
txConflictRetryIntervalKey = "txConflictRetryInterval"
txOpLimitKey = "txOpLimit"
txForceResetKey = "txForceReset"
QosEnableKey = "qosEnable"
DiskEnableKey = "diskenable"
IopsWKey = "iopsWKey"
IopsRKey = "iopsRKey"
FlowWKey = "flowWKey"
FlowRKey = "flowRKey"
ClientReqPeriod = "reqPeriod"
ClientTriggerCnt = "triggerCnt"
QosMasterLimit = "qosLimit"
decommissionFirstHostDiskParallelLimit = "decommissionFirstHostDiskParallelLimit"
decommissionFirstHostParallelLimit = "decommissionFirstHostParallelLimit"
decommissionLimit = "decommissionLimit"
DiskDisableKey = "diskDisable"
Limit = "limit"
TimeOut = "timeout"
CountByMeta = "countByMeta"
dpReadOnlyWhenVolFull = "dpReadOnlyWhenVolFull"
PeriodicKey = "periodic"
IPKey = "ip"
OperateKey = "op"
UIDKey = "uid"
CapacityKey = "capacity"
configKey = "config"
MaxFilesKey = "maxFiles"
MaxBytesKey = "maxBytes"
quotaKey = "quotaId"
enableQuota = "enableQuota"
dpDiscardKey = "dpDiscard"
ignoreDiscardKey = "ignoreDiscard"
TrashIntervalKey = "trashInterval"
ClientIDKey = "clientIDKey"
verSeqKey = "verSeq"
Periodic = "periodic"
DecommissionType = "decommissionType"
decommissionDiskLimit = "decommissionDiskLimit"
dpRepairBlockSizeKey = "dpRepairBlockSize"
markDiskBrokenThresholdKey = "markDiskBrokenThreshold"
decommissionTypeKey = "decommissionType"
autoDecommissionDiskKey = "autoDecommissionDisk"
autoDecommissionDiskIntervalKey = "autoDecommissionDiskInterval"
autoDpMetaRepairKey = "autoDpMetaRepair"
autoDpMetaRepairParallelCntKey = "autoDpMetaRepairParallelCnt"
dpTimeoutKey = "dpTimeout"
mpTimeoutKey = "mpTimeout"
ShowAll = "showAll"
trashIntervalKey = "trashInterval"
accessTimeIntervalKey = "accessTimeValidInterval"
enablePersistAccessTimeKey = "enablePersistAccessTime"
mediaTypeKey = "mediaType"
allowedStorageClassKey = "allowedStorageClass"
volStorageClassKey = "volStorageClass"
opLogDimensionKey = "opLogDimension"
volNameKey = "volName"
dpIdKey = "dpId"
diskNameKey = "diskName"
forbidWriteOpOfProtoVersion0 = "forbidWriteOpOfProtoVersion0"
quotaClass = "quotaClass"
quotaOfClass = "quotaOfStorageClass"
dataMediaTypeKey = "dataMediaType"
remoteCacheEnable = "remoteCacheEnable"
remoteCacheAutoPrepare = "remoteCacheAutoPrepare"

View File

@ -34,59 +34,59 @@ import (
// DataNode stores all the information about a data node
type DataNode struct {
Total uint64 `json:"TotalWeight"`
Used uint64 `json:"UsedWeight"`
AvailableSpace uint64
ID uint64
ZoneName string `json:"Zone"`
Addr string
HeartbeatPort string `json:"HeartbeatPort"`
ReplicaPort string `json:"ReplicaPort"`
DomainAddr string
ReportTime time.Time
StartTime int64
LastUpdateTime time.Time
isActive bool
sync.RWMutex `graphql:"-"`
UsageRatio float64 // used / total space
SelectedTimes uint64 // number times that this datanode has been selected as the location for a data partition.
TaskManager *AdminTaskManager `graphql:"-"`
DataPartitionReports []*proto.DataPartitionReport
DataPartitionCount uint32
TotalPartitionSize uint64
NodeSetID uint64
PersistenceDataPartitions []uint64
BadDisks []string // Keep this old field for compatibility
DiskStats []proto.DiskStat // key:
BadDiskStats []proto.BadDiskStat // key: disk path
LostDisks []string
DecommissionedDisks sync.Map `json:"-"` // NOTE: the disks that already be decommissioned
AllDisks []string // TODO: remove me when merge to github master
ToBeOffline bool
RdOnly bool
MigrateLock sync.RWMutex
QosIopsRLimit uint64
QosIopsWLimit uint64
QosFlowRLimit uint64
QosFlowWLimit uint64
DecommissionStatus uint32
DecommissionDstAddr string
DecommissionRaftForce bool
DecommissionLimit int
DecommissionWeight int
DecommissionFirstHostTokenLimit uint64
DecommissionCompleteTime int64
DpCntLimit uint64 `json:"-"` // max count of data partition in a data node
CpuUtil atomicutil.Float64 `json:"-"`
ioUtils atomic.Value `json:"-"`
DecommissionDiskList []string // NOTE: the disks that running decommission
DecommissionDpTotal int
DecommissionSyncMutex sync.Mutex
BackupDataPartitions []proto.BackupDataPartitionInfo
MediaType uint32
ReceivedForbidWriteOpOfProtoVer0 bool
DiskOpLogs []proto.OpLog
DpOpLogs []proto.OpLog
Total uint64 `json:"TotalWeight"`
Used uint64 `json:"UsedWeight"`
AvailableSpace uint64
ID uint64
ZoneName string `json:"Zone"`
Addr string
HeartbeatPort string `json:"HeartbeatPort"`
ReplicaPort string `json:"ReplicaPort"`
DomainAddr string
ReportTime time.Time
StartTime int64
LastUpdateTime time.Time
isActive bool
sync.RWMutex `graphql:"-"`
UsageRatio float64 // used / total space
SelectedTimes uint64 // number times that this datanode has been selected as the location for a data partition.
TaskManager *AdminTaskManager `graphql:"-"`
DataPartitionReports []*proto.DataPartitionReport
DataPartitionCount uint32
TotalPartitionSize uint64
NodeSetID uint64
PersistenceDataPartitions []uint64
BadDisks []string // Keep this old field for compatibility
DiskStats []proto.DiskStat // key:
BadDiskStats []proto.BadDiskStat // key: disk path
LostDisks []string
DecommissionedDisks sync.Map `json:"-"` // NOTE: the disks that already be decommissioned
AllDisks []string // TODO: remove me when merge to github master
ToBeOffline bool
RdOnly bool
MigrateLock sync.RWMutex
QosIopsRLimit uint64
QosIopsWLimit uint64
QosFlowRLimit uint64
QosFlowWLimit uint64
DecommissionStatus uint32
DecommissionDstAddr string
DecommissionRaftForce bool
DecommissionLimit int
DecommissionWeight int
DecommissionFirstHostParallelLimit uint64
DecommissionCompleteTime int64
DpCntLimit uint64 `json:"-"` // max count of data partition in a data node
CpuUtil atomicutil.Float64 `json:"-"`
ioUtils atomic.Value `json:"-"`
DecommissionDiskList []string // NOTE: the disks that running decommission
DecommissionDpTotal int
DecommissionSyncMutex sync.Mutex
BackupDataPartitions []proto.BackupDataPartitionInfo
MediaType uint32
ReceivedForbidWriteOpOfProtoVer0 bool
DiskOpLogs []proto.OpLog
DpOpLogs []proto.OpLog
}
func newDataNode(addr, raftHeartbeatPort, raftReplicaPort, zoneName, clusterID string, mediaType uint32) (dataNode *DataNode) {
@ -103,7 +103,7 @@ func newDataNode(addr, raftHeartbeatPort, raftReplicaPort, zoneName, clusterID s
dataNode.LastUpdateTime = time.Now().Add(-time.Minute)
dataNode.TaskManager = newAdminTaskManager(dataNode.Addr, clusterID)
dataNode.DecommissionStatus = DecommissionInitial
dataNode.DecommissionFirstHostTokenLimit = defaultDecommissionFirstHostTokenLimit
dataNode.DecommissionFirstHostParallelLimit = defaultDecommissionFirstHostParallelLimit
dataNode.CpuUtil.Store(0)
dataNode.SetIoUtils(make(map[string]float64))
dataNode.AllDisks = make([]string, 0)

View File

@ -728,6 +728,7 @@ func (partition *DataPartition) updateMetric(vr *proto.DataPartitionReport, data
replica.IsLeader = vr.IsLeader
replica.ForbidWriteOpOfProtoVer0 = vr.ForbidWriteOpOfProtoVer0
replica.ReadOnlyReasons = vr.ReadOnlyReasons
replica.IsMissingTinyExtent = vr.IsMissingTinyExtent
partition.setForbidWriteOpOfProtoVer0()
if replica.IsLeader {
partition.LeaderReportTime = time.Now().Unix()
@ -1044,12 +1045,12 @@ const (
const InvalidDecommissionDpCnt = -1
const (
defaultDecommissionParallelLimit = 10
defaultDecommissionRetryLimit = 5
defaultDecommissionRollbackLimit = 3
defaultSetRestoreReplicaStatusLimit = 300
defaultDecommissionFirstHostDiskTokenLimit = 0
defaultDecommissionFirstHostTokenLimit = 0
defaultDecommissionParallelLimit = 10
defaultDecommissionRetryLimit = 5
defaultDecommissionRollbackLimit = 3
defaultSetRestoreReplicaStatusLimit = 300
defaultDecommissionFirstHostDiskParallelLimit = 10
defaultDecommissionFirstHostParallelLimit = 0
)
func GetDecommissionStatusMessage(status uint32) string {
@ -1136,31 +1137,30 @@ func (partition *DataPartition) ReleaseDecommissionFirstHostToken(c *Cluster) {
addr := keySlice[0]
diskPath := keySlice[1]
value, ok := c.DataNodeToDecommissionRepairDpMap.Load(addr)
if ok {
dataNodeToRepairDpInfo := value.(*DataNodeToDecommissionRepairDpInfo)
dataNodeToRepairDpInfo.mu.Lock()
defer dataNodeToRepairDpInfo.mu.Unlock()
diskToRepairDpInfo, found := dataNodeToRepairDpInfo.diskToDecommissionRepairDpMap[diskPath]
if !found {
return
}
if !ok {
return
}
dataNodeToRepairDpInfo := value.(*DataNodeToDecommissionRepairDpInfo)
dataNodeToRepairDpInfo.mu.Lock()
defer dataNodeToRepairDpInfo.mu.Unlock()
diskToRepairDpInfo, found := dataNodeToRepairDpInfo.DiskToDecommissionRepairDpMap[diskPath]
if !found {
return
}
if _, isExist := diskToRepairDpInfo.repairingDps[partition.PartitionID]; !isExist {
return
}
delete(diskToRepairDpInfo.repairingDps, partition.PartitionID)
if len(diskToRepairDpInfo.repairingDps) == 0 {
delete(dataNodeToRepairDpInfo.diskToDecommissionRepairDpMap, diskPath)
} else {
atomic.StoreUint64(&diskToRepairDpInfo.curParallel, uint64(len(diskToRepairDpInfo.repairingDps)))
dataNodeToRepairDpInfo.diskToDecommissionRepairDpMap[diskPath] = diskToRepairDpInfo
}
if _, isExist := diskToRepairDpInfo.RepairingDps[partition.PartitionID]; !isExist {
return
}
delete(diskToRepairDpInfo.RepairingDps, partition.PartitionID)
if len(diskToRepairDpInfo.RepairingDps) == 0 {
delete(dataNodeToRepairDpInfo.DiskToDecommissionRepairDpMap, diskPath)
} else {
atomic.StoreUint64(&diskToRepairDpInfo.CurParallel, uint64(len(diskToRepairDpInfo.RepairingDps)))
dataNodeToRepairDpInfo.DiskToDecommissionRepairDpMap[diskPath] = diskToRepairDpInfo
}
if atomic.LoadUint64(&dataNodeToRepairDpInfo.curParallel) > 0 {
atomic.AddUint64(&dataNodeToRepairDpInfo.curParallel, ^uint64(0))
}
c.DataNodeToDecommissionRepairDpMap.Store(addr, dataNodeToRepairDpInfo)
if atomic.LoadUint64(&dataNodeToRepairDpInfo.CurParallel) > 0 {
atomic.AddUint64(&dataNodeToRepairDpInfo.CurParallel, ^uint64(0))
}
}
@ -1183,40 +1183,40 @@ func (partition *DataPartition) AcquireDecommissionFirstHostToken(c *Cluster) bo
value, _ := c.DataNodeToDecommissionRepairDpMap.LoadOrStore(firstReplica.Addr, &DataNodeToDecommissionRepairDpInfo{
mu: sync.Mutex{},
curParallel: 0,
addr: firstReplica.Addr,
diskToDecommissionRepairDpMap: make(map[string]*DiskToDecommissionRepairDpInfo),
CurParallel: 0,
Addr: firstReplica.Addr,
DiskToDecommissionRepairDpMap: make(map[string]*DiskToDecommissionRepairDpInfo),
})
dataNodeToRepairDpInfo := value.(DataNodeToDecommissionRepairDpInfo)
dataNodeToRepairDpInfo := value.(*DataNodeToDecommissionRepairDpInfo)
dataNode, err := c.dataNode(firstReplica.Addr)
if err != nil {
log.LogErrorf("action[AcquireDecommissionFirstHostToken] failed, dp(%v) err(%v)", partition.PartitionID, err.Error())
return false
}
if atomic.LoadUint64(&dataNode.DecommissionFirstHostTokenLimit) != 0 &&
atomic.LoadUint64(&dataNodeToRepairDpInfo.curParallel) >= atomic.LoadUint64(&dataNode.DecommissionFirstHostTokenLimit) {
if atomic.LoadUint64(&dataNode.DecommissionFirstHostParallelLimit) != 0 &&
atomic.LoadUint64(&dataNodeToRepairDpInfo.CurParallel) >= atomic.LoadUint64(&dataNode.DecommissionFirstHostParallelLimit) {
return false
}
dataNodeToRepairDpInfo.mu.Lock()
defer dataNodeToRepairDpInfo.mu.Unlock()
diskToRepairDpInfo, found := dataNodeToRepairDpInfo.diskToDecommissionRepairDpMap[firstReplica.DiskPath]
diskToRepairDpInfo, found := dataNodeToRepairDpInfo.DiskToDecommissionRepairDpMap[firstReplica.DiskPath]
if !found {
diskToRepairDpInfo = &DiskToDecommissionRepairDpInfo{
curParallel: 0,
diskPath: firstReplica.DiskPath,
repairingDps: make(map[uint64]struct{}),
CurParallel: 0,
DiskPath: firstReplica.DiskPath,
RepairingDps: make(map[uint64]struct{}),
}
}
if atomic.LoadUint64(&c.DecommissionFirstHostDiskTokenLimit) != 0 &&
atomic.LoadUint64(&diskToRepairDpInfo.curParallel) >= atomic.LoadUint64(&c.DecommissionFirstHostDiskTokenLimit) {
if atomic.LoadUint64(&c.DecommissionFirstHostDiskParallelLimit) != 0 &&
atomic.LoadUint64(&diskToRepairDpInfo.CurParallel) >= atomic.LoadUint64(&c.DecommissionFirstHostDiskParallelLimit) {
return false
}
diskToRepairDpInfo.repairingDps[partition.PartitionID] = struct{}{}
atomic.StoreUint64(&diskToRepairDpInfo.curParallel, uint64(len(diskToRepairDpInfo.repairingDps)))
dataNodeToRepairDpInfo.diskToDecommissionRepairDpMap[firstReplica.DiskPath] = diskToRepairDpInfo
atomic.AddUint64(&dataNodeToRepairDpInfo.curParallel, 1)
diskToRepairDpInfo.RepairingDps[partition.PartitionID] = struct{}{}
atomic.StoreUint64(&diskToRepairDpInfo.CurParallel, uint64(len(diskToRepairDpInfo.RepairingDps)))
dataNodeToRepairDpInfo.DiskToDecommissionRepairDpMap[firstReplica.DiskPath] = diskToRepairDpInfo
atomic.AddUint64(&dataNodeToRepairDpInfo.CurParallel, 1)
c.DataNodeToDecommissionRepairDpMap.Store(firstReplica.Addr, dataNodeToRepairDpInfo)
key := fmt.Sprintf("%v_%v", firstReplica.Addr, firstReplica.DiskPath)
partition.DecommissionFirstHostDiskTokenKey = key
@ -1296,7 +1296,7 @@ func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk
if partition.ReplicaNum == 3 && len(partition.Hosts) == 3 {
diskErrReplicas := partition.getAllDiskErrorReplica()
if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) && isReplicasContainsHost(diskErrReplicas, partition.Hosts[1]) {
//raftForce delete host0 and host1
// raftForce delete host0 and host1
toDeleteHosts := partition.Hosts[:2]
for _, toDeleteHost := range toDeleteHosts {
if err = c.removeDataReplica(partition, toDeleteHost, false, true); err != nil {
@ -1308,7 +1308,7 @@ func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk
return
}
}
//decommission success, reset status
// decommission success, reset status
partition.ResetDecommissionStatus()
partition.setRestoreReplicaStop()
msg := fmt.Sprintf("dp(%v) replicaNum(%v) mark decommission found host0(%v) and host1(%v) unavailable, raftForce delete them",
@ -1318,29 +1318,6 @@ func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk
}
}
if partition.ReplicaNum == 2 && len(partition.Hosts) == 2 {
diskErrReplicas := partition.getAllDiskErrorReplica()
if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) {
//raftForce delete host0
toDeleteHost := partition.Hosts[0]
if err = c.removeDataReplica(partition, toDeleteHost, false, true); err != nil {
log.LogWarnf("action[MarkDecommissionStatus] dp[%v] replicaNum[%v] remove first data replica[%v] failed, err: %v",
partition.PartitionID, partition.ReplicaNum, toDeleteHost, err)
msg := fmt.Sprintf("dp(%v) replicaNum(%v) mark decommission found host0(%v) unavailable, raftForce delete it",
partition.decommissionInfo(), partition.ReplicaNum, toDeleteHost)
auditlog.LogMasterOp("DataPartitionDecommission", msg, err)
return
}
//decommission success, reset status
partition.ResetDecommissionStatus()
partition.setRestoreReplicaStop()
msg := fmt.Sprintf("dp(%v) replicaNum(%v) mark decommission found host0(%v) unavailable, raftForce delete it",
partition.decommissionInfo(), partition.ReplicaNum, toDeleteHost)
auditlog.LogMasterOp("DataPartitionDecommission", msg, nil)
return
}
}
raftForce = true
diskErrReplica := partition.getDiskErrorReplica()
if diskErrReplica != nil {
@ -1358,10 +1335,15 @@ func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk
// in the case of autoDecommission and dp no leader :
// 1. for three replicas dp with two diskErr replicas and one normal replica should be set to the highest priority.
// 2. for three replicas dp with one replica missing, one diskErr replica and one normal replica should be set to the highest priority.
// 3. for two replicas dp with one diskErr replica should be set to high priority.
// 3. for three replicas dp with one diskErr replica and two normal replicas should be set to the high priority.
// 4. for two replicas dp with one diskErr replica should be set to high priority.
if partition.ReplicaNum == 3 {
weight = highestPriorityDecommissionWeight
} else {
if (diskErrReplicaNum == 2 && len(partition.Hosts) == 3) || (diskErrReplicaNum == 1 && len(partition.Hosts) == 2) {
weight = highestPriorityDecommissionWeight
} else {
weight = highPriorityDecommissionWeight
}
} else if partition.ReplicaNum == 2 {
weight = highPriorityDecommissionWeight
}
}
@ -1379,29 +1361,56 @@ func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk
if diskErrReplica.Addr != srcAddr {
srcAddr = diskErrReplica.Addr
srcDisk = diskErrReplica.DiskPath
weight = highPriorityDecommissionWeight
log.LogWarnf("action[MarkDecommissionStatus] dp[%v] decommission bad replica %v_%v first",
partition.PartitionID, diskErrReplica.Addr, diskErrReplica.DiskPath)
err = proto.ErrDecommissionDiskErrDPFirst
}
weight = highPriorityDecommissionWeight
}
}
}
} else {
if partition.lostLeader(c) {
if partition.getReplicaDiskErrorNum() == partition.ReplicaNum {
// auto add replica may be skipped, so check with ReplicaNum or Peers
diskErrReplicaNum := partition.getReplicaDiskErrorNum()
if diskErrReplicaNum == partition.ReplicaNum || diskErrReplicaNum == uint8(len(partition.Peers)) {
log.LogWarnf("action[MarkDecommissionStatus] dp[%v] all replica is unavaliable, cannot handle in manual decommission mode",
partition.PartitionID)
return proto.ErrAllReplicaUnavailable
}
}
if migrateType == ManualDecommission && partition.ReplicaNum == 2 && len(partition.Hosts) >= 1 {
if migrateType == ManualDecommission && partition.ReplicaNum == 3 && len(partition.Hosts) >= 2 {
diskErrReplicas := partition.getAllDiskErrorReplica()
if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) {
// mark decommission failed
log.LogWarnf("action[MarkDecommissionStatus] dp[%v] replicaNum[%v] host0[%v] is unavaliable, cannot handle in manual decommission mode",
partition.PartitionID, partition.ReplicaNum, partition.Replicas[0].Addr)
return proto.ErrFirstHostUnavailable
if raftForce {
if (isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) && isReplicasContainsHost(diskErrReplicas, partition.Hosts[1])) ||
(isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) && srcAddr != partition.Hosts[0]) ||
(isReplicasContainsHost(diskErrReplicas, partition.Hosts[1]) && srcAddr == partition.Hosts[0]) {
// mark decommission failed
log.LogWarnf("action[MarkDecommissionStatus] dp[%v] replicaNum[%v] raftForce[%v] host0 other than the srcAddr is unavaliable, cannot handle in manual decommission mode",
partition.PartitionID, partition.ReplicaNum, raftForce)
return proto.ErrFirstHostUnavailable
}
}
}
if migrateType == ManualDecommission && partition.ReplicaNum == 2 && len(partition.Hosts) == 2 {
diskErrReplicas := partition.getAllDiskErrorReplica()
if raftForce {
if (isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) && srcAddr == partition.Hosts[1]) ||
(isReplicasContainsHost(diskErrReplicas, partition.Hosts[1]) && srcAddr == partition.Hosts[0]) {
// mark decommission failed
log.LogWarnf("action[MarkDecommissionStatus] dp[%v] replicaNum[%v] raftForce(%v) host0 other than the srcAddr is unavaliable, cannot handle in manual decommission mode",
partition.PartitionID, partition.ReplicaNum, raftForce)
return proto.ErrFirstHostUnavailable
}
} else {
if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) {
// mark decommission failed
log.LogWarnf("action[MarkDecommissionStatus] dp[%v] replicaNum[%v] host0[%v] is unavaliable, cannot handle in manual decommission mode",
partition.PartitionID, partition.ReplicaNum, partition.Hosts[0])
return proto.ErrFirstHostUnavailable
}
}
}
// in the case of manualDecommission :
@ -1568,7 +1577,7 @@ func (partition *DataPartition) Decommission(c *Cluster) bool {
srcReplica *DataReplica
resetDecommissionDst = true
begin = time.Now()
finalHosts []string
finalHosts = make([]string, len(partition.Hosts))
)
if partition.GetDecommissionStatus() == DecommissionInitial {
@ -1586,7 +1595,8 @@ func (partition *DataPartition) Decommission(c *Cluster) bool {
}
partition.RLock()
finalHosts = append(partition.Hosts, targetAddr) // add new one
copy(finalHosts, partition.Hosts)
finalHosts = append(finalHosts, targetAddr) // add new one
partition.RUnlock()
for i, host := range finalHosts {
if host == srcAddr {

View File

@ -2,9 +2,12 @@ package master
import (
"fmt"
"sync"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/cubefs/cubefs/proto"
"github.com/cubefs/cubefs/util"
)
@ -88,3 +91,121 @@ func loadDataPartitionTest(dp *DataPartition, t *testing.T) {
dp.validateCRC(server.cluster.Name)
dp.setToNormal()
}
func TestAcquireDecommissionFirstHostToken(t *testing.T) {
partition := &DataPartition{PartitionID: 1, Hosts: []string{"host0", "host1", "host2"}, ReplicaNum: 3}
partition.Replicas = []*DataReplica{
{DataReplica: proto.DataReplica{Addr: "host0", DiskPath: "/disk0"}},
{DataReplica: proto.DataReplica{Addr: "host1", DiskPath: "/disk1"}},
{DataReplica: proto.DataReplica{Addr: "host2", DiskPath: "/disk2"}},
}
partition.DecommissionSrcAddr = "host2"
partition.DecommissionType = ManualDecommission
cluster := &Cluster{
ClusterDecommission: ClusterDecommission{DecommissionFirstHostDiskParallelLimit: 0},
}
dataNode := &DataNode{
DecommissionFirstHostParallelLimit: 1,
}
cluster.dataNodes.Store("host0", dataNode)
dataNodeInfo := &DataNodeToDecommissionRepairDpInfo{
mu: sync.Mutex{},
Addr: "host0",
CurParallel: 1,
}
cluster.DataNodeToDecommissionRepairDpMap.Store("host0", dataNodeInfo)
assert.False(t, partition.AcquireDecommissionFirstHostToken(cluster))
cluster.DecommissionFirstHostDiskParallelLimit = 1
dataNode.DecommissionFirstHostParallelLimit = 2
dataNodeInfo = &DataNodeToDecommissionRepairDpInfo{
mu: sync.Mutex{},
Addr: "host0",
CurParallel: 1,
DiskToDecommissionRepairDpMap: map[string]*DiskToDecommissionRepairDpInfo{
"/disk0": {CurParallel: 1, DiskPath: "/disk0"},
},
}
cluster.DataNodeToDecommissionRepairDpMap.Store("host0", dataNodeInfo)
assert.False(t, partition.AcquireDecommissionFirstHostToken(cluster))
cluster.DecommissionFirstHostDiskParallelLimit = 2
dataNode.DecommissionFirstHostParallelLimit = 2
dataNodeInfo = &DataNodeToDecommissionRepairDpInfo{
mu: sync.Mutex{},
Addr: "host0",
CurParallel: 1,
DiskToDecommissionRepairDpMap: map[string]*DiskToDecommissionRepairDpInfo{
"/disk0": {
CurParallel: 1,
DiskPath: "/disk0",
RepairingDps: map[uint64]struct{}{
0: {},
},
},
},
}
cluster.DataNodeToDecommissionRepairDpMap.Store("host0", dataNodeInfo)
assert.True(t, partition.AcquireDecommissionFirstHostToken(cluster))
}
func TestReleaseDecommissionFirstHostToken(t *testing.T) {
partition := &DataPartition{PartitionID: 1, Hosts: []string{"host0", "host1", "host2"}, ReplicaNum: 3}
partition.Replicas = []*DataReplica{
{DataReplica: proto.DataReplica{Addr: "host0", DiskPath: "/disk0"}},
{DataReplica: proto.DataReplica{Addr: "host1", DiskPath: "/disk1"}},
{DataReplica: proto.DataReplica{Addr: "host2", DiskPath: "/disk2"}},
}
partition.DecommissionSrcAddr = "host2"
partition.DecommissionType = ManualDecommission
partition.DecommissionFirstHostDiskTokenKey = "host0_/disk0"
cluster := &Cluster{
ClusterDecommission: ClusterDecommission{DecommissionFirstHostDiskParallelLimit: 2},
}
dataNode := &DataNode{
DecommissionFirstHostParallelLimit: 2,
}
cluster.dataNodes.Store("host0", dataNode)
dataNodeInfo := &DataNodeToDecommissionRepairDpInfo{
mu: sync.Mutex{},
Addr: "host0",
CurParallel: 2,
DiskToDecommissionRepairDpMap: map[string]*DiskToDecommissionRepairDpInfo{
"/disk0": {
CurParallel: 2,
DiskPath: "/disk0",
RepairingDps: map[uint64]struct{}{
0: {},
1: {},
},
},
},
}
cluster.DataNodeToDecommissionRepairDpMap.Store("host0", dataNodeInfo)
partition.ReleaseDecommissionFirstHostToken(cluster)
value, ok := cluster.DataNodeToDecommissionRepairDpMap.Load("host0")
if !ok {
t.Errorf("dataNode should not be removed")
}
dataNodeInfoAfter := value.(*DataNodeToDecommissionRepairDpInfo)
diskInfo, ok := dataNodeInfoAfter.DiskToDecommissionRepairDpMap["/disk0"]
if !ok {
t.Errorf("disk should not be removed")
}
if len(diskInfo.RepairingDps) != 1 {
t.Errorf("repairingDps should have one dp left %v", diskInfo.RepairingDps)
return
}
if diskInfo.CurParallel != 1 {
t.Errorf("disk curParallel should be updated to 1 %v", diskInfo.CurParallel)
return
}
if dataNodeInfoAfter.CurParallel != 1 {
t.Errorf("datanode curParallel should be updated to 1 %v", dataNodeInfoAfter.CurParallel)
return
}
}

View File

@ -121,31 +121,19 @@ func (c *Cluster) checkDiskRecoveryProgress() {
masterNode, _ := partition.getReplica(partition.Hosts[0])
duration := time.Unix(masterNode.ReportTime, 0).Sub(time.Unix(newReplica.ReportTime, 0))
diskErrReplicas := partition.getAllDiskErrorReplica()
if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) {
if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) || math.Abs(duration.Minutes()) > 10 {
if partition.DecommissionType == ManualAddReplica {
partition.resetForManualAddReplica()
} else {
partition.markRollbackFailed(false)
partition.markRollbackFailed(true)
}
partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v cannot finish recover"+
" for host[0] %v is unavailable", partition.DecommissionDstAddr, partition.Hosts[0])
Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] %v",
c.Name, partitionID, partition.DecommissionErrorMessage))
partition.RLock()
err = c.syncUpdateDataPartition(partition)
if err != nil {
log.LogErrorf("[checkDiskRecoveryProgress] update dp(%v) fail, err(%v)", partitionID, err)
}
partition.RUnlock()
continue
} else if math.Abs(duration.Minutes()) > 10 {
if partition.DecommissionType == ManualAddReplica {
partition.resetForManualAddReplica()
if isReplicasContainsHost(diskErrReplicas, partition.Hosts[0]) {
partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v cannot finish recover"+
" for host[0] %v is unavailable", partition.DecommissionDstAddr, partition.Hosts[0])
} else {
partition.markRollbackFailed(false)
partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v cannot finish recover"+
" for host[0] %v is down ", partition.DecommissionDstAddr, masterNode.Addr)
}
partition.DecommissionErrorMessage = fmt.Sprintf("Decommission target node %v cannot finish recover"+
" for host[0] %v is down ", partition.DecommissionDstAddr, masterNode.Addr)
Warn(c.Name, fmt.Sprintf("action[checkDiskRecoveryProgress]clusterID[%v],partitionID[%v] %v",
c.Name, partitionID, partition.DecommissionErrorMessage))
partition.RLock()
@ -564,7 +552,6 @@ func (dd *DecommissionDisk) cancelDecommission(cluster *Cluster) (err error) {
if dp.GetDecommissionStatus() == DecommissionSuccess || dp.IsRollbackFailed() {
continue
}
if dp.DecommissionDstAddr != "" {
ns, _, err = getTargetNodeset(dp.DecommissionDstAddr, cluster)
if err != nil {
@ -575,17 +562,7 @@ func (dd *DecommissionDisk) cancelDecommission(cluster *Cluster) (err error) {
}
if ns.HasDecommissionToken(dp.PartitionID) {
if dp.isSpecialReplicaCnt() {
/*
if dp.IsMarkDecommission() || dp.IsDecommissionPrepare() {
//wait to SpecialDecommissionWaitAddRes or decommissionFailed
for {
break
}
}
*/
if (dp.IsDecommissionRunning() && dp.GetSpecialReplicaDecommissionStep() == SpecialDecommissionWaitAddRes) || dp.IsDecommissionFailed() {
dp.SpecialReplicaDecommissionStop <- false // todo how to gracefully exit a dp that is going offline
log.LogDebugf("action[CancelDataPartitionDecommission] try delete dp[%v] replica %v",
dp.PartitionID, dp.DecommissionDstAddr)
// delete it from BadDataPartitionIds
@ -609,22 +586,6 @@ func (dd *DecommissionDisk) cancelDecommission(cluster *Cluster) (err error) {
continue
}
} else {
/*
waitTimes := 0
for dp.IsMarkDecommission() || dp.IsDecommissionPrepare() {
//wait to decommissionRunning or decommissionFailed
waitTimes++
if waitTimes >= 300 {
log.LogWarnf("wait for dp(%v) decommission status to success or failure timeout", dp.PartitionID)
failedDpIds = append(failedDpIds, dp.PartitionID)
break
}
}
if waitTimes >= 300 {
continue
}
*/
if dp.IsDecommissionRunning() || dp.IsDecommissionFailed() {
log.LogDebugf("action[CancelDataPartitionDecommission] try delete dp[%v] replica %v",
dp.PartitionID, dp.DecommissionDstAddr)
@ -641,17 +602,15 @@ func (dd *DecommissionDisk) cancelDecommission(cluster *Cluster) (err error) {
}
}
}
dp.ReleaseDecommissionToken(cluster)
dp.ReleaseDecommissionFirstHostToken(cluster)
}
msg := fmt.Sprintf("dp(%v) cancel decommission", dp.decommissionInfo())
dp.ResetDecommissionStatus()
dp.setRestoreReplicaStop()
cluster.syncUpdateDataPartition(dp)
auditlog.LogMasterOp("CancelDataPartitionDecommission", msg, nil)
}
msg := fmt.Sprintf("dp(%v) cancel decommission", dp.decommissionInfo())
dp.ResetDecommissionStatus()
dp.setRestoreReplicaStop()
cluster.syncUpdateDataPartition(dp)
auditlog.LogMasterOp("CancelDataPartitionDecommission", msg, nil)
}
dd.SetDecommissionStatus(DecommissionCancel)

View File

@ -333,20 +333,20 @@ func (m *Server) registerAPIRoutes(router *mux.Router) {
Path(proto.AdminGetConfig).
HandlerFunc(m.getConfigHandler)
router.NewRoute().Methods(http.MethodGet, http.MethodPost).
Path(proto.AdminUpdateDecommissionFirstHostDiskTokenLimit).
HandlerFunc(m.updateDecommissionFirstHostDiskTokenLimit)
Path(proto.AdminUpdateDecommissionFirstHostDiskParallelLimit).
HandlerFunc(m.updateDecommissionFirstHostDiskParallelLimit)
router.NewRoute().Methods(http.MethodGet, http.MethodPost).
Path(proto.AdminQueryDecommissionFirstHostDiskTokenLimit).
HandlerFunc(m.queryDecommissionFirstHostTokenDiskTokenLimit)
Path(proto.AdminQueryDecommissionFirstHostDiskParallelLimit).
HandlerFunc(m.queryDecommissionFirstHostDiskParallelLimit)
router.NewRoute().Methods(http.MethodGet, http.MethodPost).
Path(proto.AdminUpdateDecommissionFirstHostTokenLimit).
HandlerFunc(m.updateDecommissionFirstHostTokenLimit)
Path(proto.AdminUpdateDecommissionFirstHostParallelLimit).
HandlerFunc(m.updateDecommissionFirstHostParallelLimit)
router.NewRoute().Methods(http.MethodGet, http.MethodPost).
Path(proto.AdminQueryDecommissionFirstHostTokenLimit).
HandlerFunc(m.updateDecommissionFirstHostTokenLimit)
Path(proto.AdminQueryDecommissionFirstHostParallelLimit).
HandlerFunc(m.queryDecommissionFirstHostParallelLimit)
router.NewRoute().Methods(http.MethodGet, http.MethodPost).
Path(proto.AdminQueryDecommissionFirstHostTokenInfo).
HandlerFunc(m.queryDecommissionFirstHostTokenInfo)
Path(proto.AdminQueryDecommissionFirstHostParallelInfo).
HandlerFunc(m.queryDecommissionFirstHostParallelInfo)
router.NewRoute().Methods(http.MethodGet, http.MethodPost).
Path(proto.AdminUpdateDecommissionLimit).
HandlerFunc(m.updateDecommissionLimit)

View File

@ -35,103 +35,104 @@ import (
transferred over the network. */
type clusterValue struct {
Name string
CreateTime int64
Threshold float32
LoadFactor float32
DisableAutoAllocate bool
ForbidMpDecommission bool
DataNodeDeleteLimitRate uint64
MetaNodeDeleteBatchCount uint64
MetaNodeDeleteWorkerSleepMs uint64
DataNodeAutoRepairLimitRate uint64
MaxDpCntLimit uint64
MaxMpCntLimit uint64
FaultDomain bool
DiskQosEnable bool
QosLimitUpload uint64
DirChildrenNumLimit uint32
DecommissionLimit uint64
DecommissionFirstHostDiskTokenLimit uint64
CheckDataReplicasEnable bool
FileStatsEnable bool
FileStatsThresholds []uint64
ClusterUuid string
ClusterUuidEnable bool
MetaPartitionInodeIdStep uint64
MaxConcurrentLcNodes uint64
DpMaxRepairErrCnt uint64
DpRepairTimeOut uint64
DpBackupTimeOut uint64
EnableAutoDecommissionDisk bool
AutoDecommissionDiskInterval int64
DecommissionDiskLimit uint32
VolDeletionDelayTimeHour int64
MetaNodeGOGC int
DataNodeGOGC int
MarkDiskBrokenThreshold float64
EnableAutoDpMetaRepair bool
AutoDpMetaRepairParallelCnt uint32
DataPartitionTimeoutSec int64
MetaPartitionTimeoutSec int64
ForbidWriteOpOfProtoVer0 bool
LegacyDataMediaType uint32
RaftPartitionAlreadyUseDifferentPort bool
MetaNodeMemoryHighPer float64
MetaNodeMemoryLowPer float64
AutoMpMigrate bool
FlashNodeHandleReadTimeout int
FlashNodeReadDataNodeTimeout int
Name string
CreateTime int64
Threshold float32
LoadFactor float32
DisableAutoAllocate bool
ForbidMpDecommission bool
DataNodeDeleteLimitRate uint64
MetaNodeDeleteBatchCount uint64
MetaNodeDeleteWorkerSleepMs uint64
DataNodeAutoRepairLimitRate uint64
MaxDpCntLimit uint64
MaxMpCntLimit uint64
FaultDomain bool
DiskQosEnable bool
QosLimitUpload uint64
DirChildrenNumLimit uint32
DecommissionLimit uint64
DecommissionFirstHostDiskParallelLimit uint64
CheckDataReplicasEnable bool
FileStatsEnable bool
FileStatsThresholds []uint64
ClusterUuid string
ClusterUuidEnable bool
MetaPartitionInodeIdStep uint64
MaxConcurrentLcNodes uint64
DpMaxRepairErrCnt uint64
DpRepairTimeOut uint64
DpBackupTimeOut uint64
EnableAutoDecommissionDisk bool
AutoDecommissionDiskInterval int64
DecommissionDiskLimit uint32
VolDeletionDelayTimeHour int64
MetaNodeGOGC int
DataNodeGOGC int
MarkDiskBrokenThreshold float64
EnableAutoDpMetaRepair bool
AutoDpMetaRepairParallelCnt uint32
DataPartitionTimeoutSec int64
MetaPartitionTimeoutSec int64
ForbidWriteOpOfProtoVer0 bool
LegacyDataMediaType uint32
RaftPartitionAlreadyUseDifferentPort bool
MetaNodeMemoryHighPer float64
MetaNodeMemoryLowPer float64
AutoMpMigrate bool
FlashNodeHandleReadTimeout int
FlashNodeReadDataNodeTimeout int
}
func newClusterValue(c *Cluster) (cv *clusterValue) {
cv = &clusterValue{
Name: c.Name,
CreateTime: c.CreateTime,
LoadFactor: c.cfg.ClusterLoadFactor,
Threshold: c.cfg.MetaNodeThreshold,
DataNodeDeleteLimitRate: c.cfg.DataNodeDeleteLimitRate,
MetaNodeDeleteBatchCount: c.cfg.MetaNodeDeleteBatchCount,
MetaNodeDeleteWorkerSleepMs: c.cfg.MetaNodeDeleteWorkerSleepMs,
DataNodeAutoRepairLimitRate: c.cfg.DataNodeAutoRepairLimitRate,
DisableAutoAllocate: c.DisableAutoAllocate,
ForbidMpDecommission: c.ForbidMpDecommission,
MaxDpCntLimit: c.getMaxDpCntLimit(),
MaxMpCntLimit: c.getMaxMpCntLimit(),
FaultDomain: c.FaultDomain,
DiskQosEnable: c.diskQosEnable,
QosLimitUpload: uint64(c.QosAcceptLimit.Limit()),
DirChildrenNumLimit: c.cfg.DirChildrenNumLimit,
DecommissionLimit: c.DecommissionLimit,
CheckDataReplicasEnable: c.checkDataReplicasEnable,
FileStatsEnable: c.fileStatsEnable,
FileStatsThresholds: c.fileStatsThresholds,
ClusterUuid: c.clusterUuid,
ClusterUuidEnable: c.clusterUuidEnable,
MetaPartitionInodeIdStep: c.cfg.MetaPartitionInodeIdStep,
MaxConcurrentLcNodes: c.cfg.MaxConcurrentLcNodes,
DpMaxRepairErrCnt: c.cfg.DpMaxRepairErrCnt,
DpRepairTimeOut: c.cfg.DpRepairTimeOut,
DpBackupTimeOut: c.cfg.DpBackupTimeOut,
EnableAutoDecommissionDisk: c.EnableAutoDecommissionDisk.Load(),
AutoDecommissionDiskInterval: c.AutoDecommissionInterval.Load(),
DecommissionDiskLimit: c.GetDecommissionDiskLimit(),
VolDeletionDelayTimeHour: c.cfg.volDelayDeleteTimeHour,
MetaNodeGOGC: c.cfg.metaNodeGOGC,
DataNodeGOGC: c.cfg.dataNodeGOGC,
MarkDiskBrokenThreshold: c.getMarkDiskBrokenThreshold(),
EnableAutoDpMetaRepair: c.getEnableAutoDpMetaRepair(),
AutoDpMetaRepairParallelCnt: c.AutoDpMetaRepairParallelCnt.Load(),
DataPartitionTimeoutSec: c.getDataPartitionTimeoutSec(),
MetaPartitionTimeoutSec: c.getMetaPartitionTimeoutSec(),
ForbidWriteOpOfProtoVer0: c.cfg.forbidWriteOpOfProtoVer0,
LegacyDataMediaType: c.legacyDataMediaType,
RaftPartitionAlreadyUseDifferentPort: c.cfg.raftPartitionAlreadyUseDifferentPort.Load(),
MetaNodeMemoryHighPer: c.cfg.metaNodeMemHighPer,
MetaNodeMemoryLowPer: c.cfg.metaNodeMemLowPer,
AutoMpMigrate: c.cfg.AutoMpMigrate,
FlashNodeHandleReadTimeout: c.cfg.flashNodeHandleReadTimeout,
FlashNodeReadDataNodeTimeout: c.cfg.flashNodeReadDataNodeTimeout,
Name: c.Name,
CreateTime: c.CreateTime,
LoadFactor: c.cfg.ClusterLoadFactor,
Threshold: c.cfg.MetaNodeThreshold,
DataNodeDeleteLimitRate: c.cfg.DataNodeDeleteLimitRate,
MetaNodeDeleteBatchCount: c.cfg.MetaNodeDeleteBatchCount,
MetaNodeDeleteWorkerSleepMs: c.cfg.MetaNodeDeleteWorkerSleepMs,
DataNodeAutoRepairLimitRate: c.cfg.DataNodeAutoRepairLimitRate,
DisableAutoAllocate: c.DisableAutoAllocate,
ForbidMpDecommission: c.ForbidMpDecommission,
MaxDpCntLimit: c.getMaxDpCntLimit(),
MaxMpCntLimit: c.getMaxMpCntLimit(),
FaultDomain: c.FaultDomain,
DiskQosEnable: c.diskQosEnable,
QosLimitUpload: uint64(c.QosAcceptLimit.Limit()),
DirChildrenNumLimit: c.cfg.DirChildrenNumLimit,
DecommissionFirstHostDiskParallelLimit: c.DecommissionFirstHostDiskParallelLimit,
DecommissionLimit: c.DecommissionLimit,
CheckDataReplicasEnable: c.checkDataReplicasEnable,
FileStatsEnable: c.fileStatsEnable,
FileStatsThresholds: c.fileStatsThresholds,
ClusterUuid: c.clusterUuid,
ClusterUuidEnable: c.clusterUuidEnable,
MetaPartitionInodeIdStep: c.cfg.MetaPartitionInodeIdStep,
MaxConcurrentLcNodes: c.cfg.MaxConcurrentLcNodes,
DpMaxRepairErrCnt: c.cfg.DpMaxRepairErrCnt,
DpRepairTimeOut: c.cfg.DpRepairTimeOut,
DpBackupTimeOut: c.cfg.DpBackupTimeOut,
EnableAutoDecommissionDisk: c.EnableAutoDecommissionDisk.Load(),
AutoDecommissionDiskInterval: c.AutoDecommissionInterval.Load(),
DecommissionDiskLimit: c.GetDecommissionDiskLimit(),
VolDeletionDelayTimeHour: c.cfg.volDelayDeleteTimeHour,
MetaNodeGOGC: c.cfg.metaNodeGOGC,
DataNodeGOGC: c.cfg.dataNodeGOGC,
MarkDiskBrokenThreshold: c.getMarkDiskBrokenThreshold(),
EnableAutoDpMetaRepair: c.getEnableAutoDpMetaRepair(),
AutoDpMetaRepairParallelCnt: c.AutoDpMetaRepairParallelCnt.Load(),
DataPartitionTimeoutSec: c.getDataPartitionTimeoutSec(),
MetaPartitionTimeoutSec: c.getMetaPartitionTimeoutSec(),
ForbidWriteOpOfProtoVer0: c.cfg.forbidWriteOpOfProtoVer0,
LegacyDataMediaType: c.legacyDataMediaType,
RaftPartitionAlreadyUseDifferentPort: c.cfg.raftPartitionAlreadyUseDifferentPort.Load(),
MetaNodeMemoryHighPer: c.cfg.metaNodeMemHighPer,
MetaNodeMemoryLowPer: c.cfg.metaNodeMemLowPer,
AutoMpMigrate: c.cfg.AutoMpMigrate,
FlashNodeHandleReadTimeout: c.cfg.flashNodeHandleReadTimeout,
FlashNodeReadDataNodeTimeout: c.cfg.flashNodeReadDataNodeTimeout,
}
return cv
}
@ -486,54 +487,54 @@ func newVolValueFromBytes(raw []byte) (*volValue, error) {
}
type dataNodeValue struct {
ID uint64
NodeSetID uint64
Addr string
HeartbeatPort string
ReplicaPort string
ZoneName string
RdOnly bool
DecommissionedDisks []string
DecommissionStatus uint32
DecommissionDstAddr string
DecommissionRaftForce bool
DecommissionLimit int
DecommissionWeight int
DecommissionFirstHostTokenLimit uint64
DecommissionCompleteTime int64
ToBeOffline bool
DecommissionDiskList []string
DecommissionDpTotal int
BadDisks []string
AllDisks []string
MediaType uint32
MaxDpCntLimit uint64
ID uint64
NodeSetID uint64
Addr string
HeartbeatPort string
ReplicaPort string
ZoneName string
RdOnly bool
DecommissionedDisks []string
DecommissionStatus uint32
DecommissionDstAddr string
DecommissionRaftForce bool
DecommissionLimit int
DecommissionWeight int
DecommissionFirstHostParallelLimit uint64
DecommissionCompleteTime int64
ToBeOffline bool
DecommissionDiskList []string
DecommissionDpTotal int
BadDisks []string
AllDisks []string
MediaType uint32
MaxDpCntLimit uint64
}
func newDataNodeValue(dataNode *DataNode) *dataNodeValue {
return &dataNodeValue{
ID: dataNode.ID,
NodeSetID: dataNode.NodeSetID,
Addr: dataNode.Addr,
HeartbeatPort: dataNode.HeartbeatPort,
ReplicaPort: dataNode.ReplicaPort,
ZoneName: dataNode.ZoneName,
RdOnly: dataNode.RdOnly,
DecommissionedDisks: dataNode.getDecommissionedDisks(),
DecommissionStatus: atomic.LoadUint32(&dataNode.DecommissionStatus),
DecommissionDstAddr: dataNode.DecommissionDstAddr,
DecommissionRaftForce: dataNode.DecommissionRaftForce,
DecommissionLimit: dataNode.DecommissionLimit,
DecommissionWeight: dataNode.DecommissionWeight,
DecommissionFirstHostTokenLimit: dataNode.DecommissionFirstHostTokenLimit,
DecommissionCompleteTime: dataNode.DecommissionCompleteTime,
ToBeOffline: dataNode.ToBeOffline,
DecommissionDiskList: dataNode.DecommissionDiskList,
DecommissionDpTotal: dataNode.DecommissionDpTotal,
AllDisks: dataNode.AllDisks,
BadDisks: dataNode.BadDisks,
MediaType: dataNode.MediaType,
MaxDpCntLimit: dataNode.DpCntLimit,
ID: dataNode.ID,
NodeSetID: dataNode.NodeSetID,
Addr: dataNode.Addr,
HeartbeatPort: dataNode.HeartbeatPort,
ReplicaPort: dataNode.ReplicaPort,
ZoneName: dataNode.ZoneName,
RdOnly: dataNode.RdOnly,
DecommissionedDisks: dataNode.getDecommissionedDisks(),
DecommissionStatus: atomic.LoadUint32(&dataNode.DecommissionStatus),
DecommissionDstAddr: dataNode.DecommissionDstAddr,
DecommissionRaftForce: dataNode.DecommissionRaftForce,
DecommissionLimit: dataNode.DecommissionLimit,
DecommissionWeight: dataNode.DecommissionWeight,
DecommissionFirstHostParallelLimit: dataNode.DecommissionFirstHostParallelLimit,
DecommissionCompleteTime: dataNode.DecommissionCompleteTime,
ToBeOffline: dataNode.ToBeOffline,
DecommissionDiskList: dataNode.DecommissionDiskList,
DecommissionDpTotal: dataNode.DecommissionDpTotal,
AllDisks: dataNode.AllDisks,
BadDisks: dataNode.BadDisks,
MediaType: dataNode.MediaType,
MaxDpCntLimit: dataNode.DpCntLimit,
}
}
@ -1316,7 +1317,7 @@ func (c *Cluster) loadClusterValue() (err error) {
c.diskQosEnable = cv.DiskQosEnable
c.cfg.QosMasterAcceptLimit = cv.QosLimitUpload
c.DecommissionLimit = cv.DecommissionLimit // dont update nodesets limit for nodesets are not loaded
c.DecommissionFirstHostDiskTokenLimit = cv.DecommissionFirstHostDiskTokenLimit
c.DecommissionFirstHostDiskParallelLimit = cv.DecommissionFirstHostDiskParallelLimit
c.fileStatsEnable = cv.FileStatsEnable
c.fileStatsThresholds = cv.FileStatsThresholds
c.clusterUuid = cv.ClusterUuid
@ -1627,7 +1628,7 @@ func (c *Cluster) loadDataNodes() (err error) {
dataNode.DecommissionRaftForce = dnv.DecommissionRaftForce
dataNode.DecommissionLimit = dnv.DecommissionLimit
dataNode.DecommissionWeight = dnv.DecommissionWeight
dataNode.DecommissionFirstHostTokenLimit = dnv.DecommissionFirstHostTokenLimit
dataNode.DecommissionFirstHostParallelLimit = dnv.DecommissionFirstHostParallelLimit
dataNode.DecommissionCompleteTime = dnv.DecommissionCompleteTime
dataNode.ToBeOffline = dnv.ToBeOffline
dataNode.DecommissionDiskList = dnv.DecommissionDiskList
@ -1645,10 +1646,10 @@ func (c *Cluster) loadDataNodes() (err error) {
c.dataNodes.Store(dataNode.Addr, dataNode)
log.LogInfof("action[loadDataNodes],dataNode[%v],dataNodeID[%v],MediaType[%v],zone[%v],ns[%v] DecommissionStatus [%v] "+
"DecommissionDstAddr[%v] DecommissionRaftForce[%v] DecommissionDpTotal[%v] DecommissionLimit[%v] DecommissionWeight[%v] DecommissionFirstHostTokenLimit[%v] DpCntLimit[%v]"+
"DecommissionDstAddr[%v] DecommissionRaftForce[%v] DecommissionDpTotal[%v] DecommissionLimit[%v] DecommissionWeight[%v] DecommissionFirstHostParallelLimit[%v] DpCntLimit[%v]"+
"DecommissionCompleteTime [%v] ToBeOffline[%v]",
dataNode.Addr, dataNode.ID, dataNode.MediaType, dnv.ZoneName, dnv.NodeSetID, dataNode.DecommissionStatus,
dataNode.DecommissionDstAddr, dataNode.DecommissionRaftForce, dataNode.DecommissionDpTotal, dataNode.DecommissionLimit, dataNode.DecommissionWeight, dataNode.DecommissionFirstHostTokenLimit,
dataNode.DecommissionDstAddr, dataNode.DecommissionRaftForce, dataNode.DecommissionDpTotal, dataNode.DecommissionLimit, dataNode.DecommissionWeight, dataNode.DecommissionFirstHostParallelLimit,
dataNode.DpCntLimit, time.Unix(dataNode.DecommissionCompleteTime, 0).Format("2006-01-02 15:04:05"), dataNode.ToBeOffline)
log.LogInfof("action[loadDataNodes],dataNode[%v],dataNodeID[%v],zone[%v],ns[%v],MediaType[%v]",

View File

@ -31,42 +31,43 @@ import (
// metrics
const (
StatPeriod = time.Minute * time.Duration(1)
MetricDataNodesUsedGB = "dataNodes_used_GB"
MetricDataNodesTotalGB = "dataNodes_total_GB"
MetricDataNodesStat = "dataNodes_stats"
MetricDataNodesIncreasedGB = "dataNodes_increased_GB"
MetricMetaNodesUsedGB = "metaNodes_used_GB"
MetricMetaNodesTotalGB = "metaNodes_total_GB"
MetricMetaNodesIncreasedGB = "metaNodes_increased_GB"
MetricDataNodesCount = "dataNodes_count"
MetricMetaNodesCount = "metaNodes_count"
MetricNodeStat = "node_stat"
MetricVolCount = "vol_count"
MetricVolTotalGB = "vol_total_GB"
MetricVolUsedGB = "vol_used_GB"
MetricVolUsageGB = "vol_usage_ratio"
MetricVolStats = "vol_stats"
MetricVolMetaCount = "vol_meta_count"
MetricBadMpCount = "bad_mp_count"
MetricBadDpCount = "bad_dp_count"
MetricDiskError = "disk_error"
MetricFlashNodesDiskError = "flashNodes_disk_error"
MetricDiskLost = "disk_lost"
MetricDpUnableDecommission = "dp_unable_decommission"
MetricDataNodesInactive = "dataNodes_inactive"
MetricInactiveDataNodeInfo = "inactive_dataNodes_info"
MetricMetaNodesInactive = "metaNodes_inactive"
MetricDataNodesNotWritable = "dataNodes_not_writable"
MetricDataNodesAllocable = "dataNodes_allocable"
MetricMetaNodesNotWritable = "metaNodes_not_writable"
MetricInactiveMetaNodeInfo = "inactive_metaNodes_info"
MetricMetaInconsistent = "mp_inconsistent"
MetricMasterNoLeader = "master_no_leader"
MetricMasterNoCache = "master_no_cache"
MetricMasterSnapshot = "master_snapshot"
MetricMastersInactive = "masters_inactive"
MetricInactiveMasterInfo = "inactive_masters_info"
StatPeriod = time.Minute * time.Duration(1)
MetricDataNodesUsedGB = "dataNodes_used_GB"
MetricDataNodesTotalGB = "dataNodes_total_GB"
MetricDataNodesStat = "dataNodes_stats"
MetricDataNodesIncreasedGB = "dataNodes_increased_GB"
MetricMetaNodesUsedGB = "metaNodes_used_GB"
MetricMetaNodesTotalGB = "metaNodes_total_GB"
MetricMetaNodesIncreasedGB = "metaNodes_increased_GB"
MetricDataNodesCount = "dataNodes_count"
MetricMetaNodesCount = "metaNodes_count"
MetricNodeStat = "node_stat"
MetricVolCount = "vol_count"
MetricVolTotalGB = "vol_total_GB"
MetricVolUsedGB = "vol_used_GB"
MetricVolUsageGB = "vol_usage_ratio"
MetricVolStats = "vol_stats"
MetricVolMetaCount = "vol_meta_count"
MetricBadMpCount = "bad_mp_count"
MetricBadDpCount = "bad_dp_count"
MetricDiskError = "disk_error"
MetricFlashNodesDiskError = "flashNodes_disk_error"
MetricDiskLost = "disk_lost"
MetricDpUnableDecommissionCount = "dp_unable_decommission_count"
MetricDpMissingTinyExtent = "dp_missing_tinyExtent"
MetricDataNodesInactive = "dataNodes_inactive"
MetricInactiveDataNodeInfo = "inactive_dataNodes_info"
MetricMetaNodesInactive = "metaNodes_inactive"
MetricDataNodesNotWritable = "dataNodes_not_writable"
MetricDataNodesAllocable = "dataNodes_allocable"
MetricMetaNodesNotWritable = "metaNodes_not_writable"
MetricInactiveMetaNodeInfo = "inactive_metaNodes_info"
MetricMetaInconsistent = "mp_inconsistent"
MetricMasterNoLeader = "master_no_leader"
MetricMasterNoCache = "master_no_cache"
MetricMasterSnapshot = "master_snapshot"
MetricMastersInactive = "masters_inactive"
MetricInactiveMasterInfo = "inactive_masters_info"
MetricMissingDp = "missing_dp"
MetricDpNoLeader = "dp_no_leader"
@ -105,56 +106,57 @@ const (
var WarnMetrics *warningMetrics
type monitorMetrics struct {
cluster *Cluster
dataNodesCount *exporter.Gauge
metaNodesCount *exporter.Gauge
volCount *exporter.Gauge
dataNodesTotal *exporter.Gauge
dataNodesUsed *exporter.Gauge
dataNodeStat *exporter.GaugeVec
dataNodeIncreased *exporter.Gauge
metaNodesTotal *exporter.Gauge
metaNodesUsed *exporter.Gauge
metaNodesIncreased *exporter.Gauge
volTotalSpace *exporter.GaugeVec
volUsedSpace *exporter.GaugeVec
volUsage *exporter.GaugeVec
volMetaCount *exporter.GaugeVec
volStats *exporter.GaugeVec
badMpCount *exporter.Gauge
badDpCount *exporter.Gauge
diskError *exporter.GaugeVec
flashNodesDiskError *exporter.GaugeVec
diskLost *exporter.GaugeVec
dpUnableDecommission *exporter.GaugeVec
dataNodesNotWritable *exporter.Gauge
dataNodesAllocable *exporter.Gauge
metaNodesNotWritable *exporter.Gauge
dataNodesInactive *exporter.Gauge
InactiveDataNodeInfo *exporter.GaugeVec
metaNodesInactive *exporter.Gauge
InactiveMetaNodeInfo *exporter.GaugeVec
mastersInactive *exporter.Gauge
InactiveMasterInfo *exporter.GaugeVec
ReplicaMissingDPCount *exporter.GaugeVec
DpMissingLeaderCount *exporter.GaugeVec
MpMissingLeaderCount *exporter.Gauge
MpMissingReplicaCount *exporter.Gauge
dataNodesetInactiveCount *exporter.GaugeVec
metaNodesetInactiveCount *exporter.GaugeVec
metaEqualCheckFail *exporter.GaugeVec
masterNoLeader *exporter.Gauge
masterNoCache *exporter.GaugeVec
masterSnapshot *exporter.Gauge
nodesetMetaTotal *exporter.GaugeVec
nodesetMetaUsed *exporter.GaugeVec
nodesetMetaUsageRatio *exporter.GaugeVec
nodesetDataTotal *exporter.GaugeVec
nodesetDataUsed *exporter.GaugeVec
nodesetDataUsageRatio *exporter.GaugeVec
nodesetMpReplicaCount *exporter.GaugeVec
nodesetDpReplicaCount *exporter.GaugeVec
nodeStat *exporter.GaugeVec
cluster *Cluster
dataNodesCount *exporter.Gauge
metaNodesCount *exporter.Gauge
volCount *exporter.Gauge
dataNodesTotal *exporter.Gauge
dataNodesUsed *exporter.Gauge
dataNodeStat *exporter.GaugeVec
dataNodeIncreased *exporter.Gauge
metaNodesTotal *exporter.Gauge
metaNodesUsed *exporter.Gauge
metaNodesIncreased *exporter.Gauge
volTotalSpace *exporter.GaugeVec
volUsedSpace *exporter.GaugeVec
volUsage *exporter.GaugeVec
volMetaCount *exporter.GaugeVec
volStats *exporter.GaugeVec
badMpCount *exporter.Gauge
badDpCount *exporter.Gauge
diskError *exporter.GaugeVec
flashNodesDiskError *exporter.GaugeVec
diskLost *exporter.GaugeVec
dpUnableDecommissionCount *exporter.Gauge
dpMissingTinyExtent *exporter.GaugeVec
dataNodesNotWritable *exporter.Gauge
dataNodesAllocable *exporter.Gauge
metaNodesNotWritable *exporter.Gauge
dataNodesInactive *exporter.Gauge
InactiveDataNodeInfo *exporter.GaugeVec
metaNodesInactive *exporter.Gauge
InactiveMetaNodeInfo *exporter.GaugeVec
mastersInactive *exporter.Gauge
InactiveMasterInfo *exporter.GaugeVec
ReplicaMissingDPCount *exporter.GaugeVec
DpMissingLeaderCount *exporter.GaugeVec
MpMissingLeaderCount *exporter.Gauge
MpMissingReplicaCount *exporter.Gauge
dataNodesetInactiveCount *exporter.GaugeVec
metaNodesetInactiveCount *exporter.GaugeVec
metaEqualCheckFail *exporter.GaugeVec
masterNoLeader *exporter.Gauge
masterNoCache *exporter.GaugeVec
masterSnapshot *exporter.Gauge
nodesetMetaTotal *exporter.GaugeVec
nodesetMetaUsed *exporter.GaugeVec
nodesetMetaUsageRatio *exporter.GaugeVec
nodesetDataTotal *exporter.GaugeVec
nodesetDataUsed *exporter.GaugeVec
nodesetDataUsageRatio *exporter.GaugeVec
nodesetMpReplicaCount *exporter.GaugeVec
nodesetDpReplicaCount *exporter.GaugeVec
nodeStat *exporter.GaugeVec
volNames map[string]struct{}
badDisks map[string]string
@ -513,7 +515,8 @@ func (mm *monitorMetrics) start() {
mm.diskError = exporter.NewGaugeVec(MetricDiskError, "", []string{"addr", "path"})
mm.flashNodesDiskError = exporter.NewGaugeVec(MetricFlashNodesDiskError, "", []string{"addr", "path"})
mm.diskLost = exporter.NewGaugeVec(MetricDiskLost, "", []string{"addr", "path"})
mm.dpUnableDecommission = exporter.NewGaugeVec(MetricDpUnableDecommission, "", []string{"dpId"})
mm.dpUnableDecommissionCount = exporter.NewGauge(MetricDpUnableDecommissionCount)
mm.dpMissingTinyExtent = exporter.NewGaugeVec(MetricDpMissingTinyExtent, "", []string{"dpId", "addr"})
mm.nodeStat = exporter.NewGaugeVec(MetricNodeStat, "", []string{"type", "addr", "stat"})
mm.dataNodesInactive = exporter.NewGauge(MetricDataNodesInactive)
mm.InactiveDataNodeInfo = exporter.NewGaugeVec(MetricInactiveDataNodeInfo, "", []string{"clusterName", "addr"})
@ -617,6 +620,7 @@ func (mm *monitorMetrics) doStat() {
mm.setDiskLostMetric()
mm.setFlashNodesDiskErrorMetric()
mm.setDpUnableDecommissionMetric()
mm.setDpMissingTinyExtentMetric()
mm.setNotWritableDataNodesCount()
mm.setNotWritableMetaNodesCount()
mm.setMpInconsistentErrorMetric()
@ -924,15 +928,32 @@ func (mm *monitorMetrics) setDiskLostMetric() {
}
func (mm *monitorMetrics) setDpUnableDecommissionMetric() {
mm.dpUnableDecommission.Reset()
dpUnableDecommissionCount := 0
vols := mm.cluster.allVols()
for _, vol := range vols {
partitions := vol.dataPartitions.clonePartitions()
for _, dp := range partitions {
if dp.GetDecommissionStatus() == DecommissionFail &&
strings.Contains(dp.DecommissionErrorMessage, proto.ErrAllReplicaUnavailable.Error()) && !dp.IsDiscard {
dpUnableDecommissionCount++
}
}
}
mm.dpUnableDecommissionCount.Set(float64(dpUnableDecommissionCount))
}
func (mm *monitorMetrics) setDpMissingTinyExtentMetric() {
mm.dpMissingTinyExtent.Reset()
vols := mm.cluster.allVols()
for _, vol := range vols {
partitions := vol.dataPartitions.clonePartitions()
for _, dp := range partitions {
if dp.GetDecommissionStatus() == DecommissionFail && strings.Contains(dp.DecommissionErrorMessage, proto.ErrAllReplicaUnavailable.Error()) {
idStr := strconv.FormatUint(dp.PartitionID, 10)
mm.dpUnableDecommission.SetWithLabelValues(1, idStr)
for _, replica := range dp.Replicas {
if replica.IsMissingTinyExtent {
idStr := strconv.FormatUint(dp.PartitionID, 10)
mm.dpMissingTinyExtent.SetWithLabelValues(1, idStr, replica.Addr)
}
}
}
}
@ -1320,7 +1341,8 @@ func (mm *monitorMetrics) resetAllLeaderMetrics() {
mm.metaNodesIncreased.Set(0)
// mm.diskError.Set(0)
mm.diskLost.Reset()
mm.dpUnableDecommission.Reset()
mm.dpUnableDecommissionCount.Set(0)
mm.dpMissingTinyExtent.Reset()
mm.diskDecommissioned.Reset()
mm.dataNodesInactive.Set(0)
mm.metaNodesInactive.Set(0)

View File

@ -2328,6 +2328,34 @@ func (l *DecommissionDataPartitionList) traverse(c *Cluster) {
return
}
allDecommissionDP := l.GetAllDecommissionDataPartitions()
for _, dp := range allDecommissionDP {
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())
err := proto.ErrAllReplicaUnavailable
dp.DecommissionErrorMessage = err.Error()
dp.markRollbackFailed(false)
continue
}
if dp.DecommissionType == AutoDecommission {
diskErrReplicas := dp.getAllDiskErrorReplica()
if isReplicasContainsHost(diskErrReplicas, dp.DecommissionSrcAddr) {
if dp.ReplicaNum == 3 {
if (diskErrReplicaNum == 2 && len(dp.Hosts) == 3) || (diskErrReplicaNum == 1 && len(dp.Hosts) == 2) {
dp.DecommissionWeight = highestPriorityDecommissionWeight
} else if diskErrReplicaNum == 1 && len(dp.Hosts) == 3 {
dp.DecommissionWeight = highPriorityDecommissionWeight
}
} else if dp.ReplicaNum == 2 {
if diskErrReplicaNum == 1 && len(dp.Hosts) == 2 {
dp.DecommissionWeight = highPriorityDecommissionWeight
}
}
}
}
}
sort.Slice(allDecommissionDP, func(i, j int) bool {
return allDecommissionDP[i].DecommissionWeight > allDecommissionDP[j].DecommissionWeight
})
@ -2388,7 +2416,7 @@ func (l *DecommissionDataPartitionList) traverse(c *Cluster) {
l.Remove(dp)
dp.ResetDecommissionStatus()
c.syncUpdateDataPartition(dp)
} else if dp.IsMarkDecommission() { //&& dp.TryAcquireDecommissionToken(c) {
} else if dp.IsMarkDecommission() {
if dp.AcquireDecommissionFirstHostToken(c) {
if dp.TryAcquireDecommissionToken(c) {
go func(dp *DataPartition) {

View File

@ -32,75 +32,75 @@ 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"
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"
AdminUpdateDecommissionFirstHostDiskTokenLimit = "/admin/updateDecommissionFirstHostDiskTokenLimit"
AdminQueryDecommissionFirstHostDiskTokenLimit = "/admin/queryDecommissionFirstHostDiskTokenLimit"
AdminUpdateDecommissionFirstHostTokenLimit = "/admin/updateDecommissionFirstHostTokenLimit"
AdminQueryDecommissionFirstHostTokenLimit = "/admin/queryDecommissionFirstHostTokenLimit"
AdminQueryDecommissionFirstHostTokenInfo = "/admin/queryDecommissionFirstHostTokenInfo"
AdminUpdateDecommissionLimit = "/admin/updateDecommissionLimit"
AdminQueryDecommissionLimit = "/admin/queryDecommissionLimit"
AdminQueryDecommissionFailedDisk = "/admin/queryDecommissionFailedDisk"
AdminAbortDecommissionDisk = "/admin/abortDecommissionDisk"
AdminResetDataPartitionRestoreStatus = "/admin/resetDataPartitionRestoreStatus"
AdminGetOpLog = "/admin/getOpLog"
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"
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"
// #nosec G101
AdminQueryDecommissionToken = "/admin/queryDecommissionToken"
@ -860,6 +860,7 @@ type DataPartitionReport struct {
TriggerDiskError bool
ForbidWriteOpOfProtoVer0 bool
ReadOnlyReasons uint32
IsMissingTinyExtent bool
}
type DataNodeQosResponse struct {

View File

@ -149,7 +149,7 @@ type ClusterView struct {
AutoDpMetaRepairParallelCnt int
EnableAutoDecommission bool
AutoDecommissionDiskInterval string
DiskToRepairDpLimit uint64
DecommissionFirstHostDiskParallelLimit uint64
DecommissionLimit uint64
DecommissionDiskLimit uint32
DpRepairTimeout string
@ -377,6 +377,7 @@ type DataReplica struct {
TriggerDiskError bool
ForbidWriteOpOfProtoVer0 bool
ReadOnlyReasons uint32
IsMissingTinyExtent bool
}
// data partition diagnosis represents the inactive data nodes, corrupt data partitions, and data partitions lack of replicas