From c8d253aa39bd9bc628d3eac5d65975716d8d7227 Mon Sep 17 00:00:00 2001 From: zhumingze Date: Fri, 28 Mar 2025 15:09:50 +0800 Subject: [PATCH] feat(master): mp&dp replica missing check time can be adjusted. #1000044373 Signed-off-by: zhumingze --- cli/cmd/cluster.go | 17 ++++++++++++- cli/cmd/const.go | 1 + cli/cmd/fmt.go | 5 ++-- master/api_args_parse.go | 11 +++++++++ master/api_service.go | 10 ++++++++ master/api_service_test.go | 12 +++++++++ master/cluster.go | 24 ++++++++++++++++-- master/cluster_task.go | 2 +- master/config.go | 3 +++ master/const.go | 1 + master/meta_partition.go | 50 +++++++++++++++++++------------------- master/metadata_fsm_op.go | 7 ++++++ master/monitor_metrics.go | 4 +-- master/server.go | 7 ++++++ master/vol.go | 8 +++--- master/vol_test.go | 2 +- proto/model.go | 1 + sdk/master/api_admin.go | 5 +++- 18 files changed, 131 insertions(+), 39 deletions(-) diff --git a/cli/cmd/cluster.go b/cli/cmd/cluster.go index b96cd18e3..178989761 100644 --- a/cli/cmd/cluster.go +++ b/cli/cmd/cluster.go @@ -281,6 +281,7 @@ func newClusterSetParasCmd(client *master.MasterClient) *cobra.Command { opMaxMpCntLimit := "" dpRepairTimeout := "" dpTimeout := "" + mpTimeout := "" dpBackupTimeout := "" decommissionDpLimit := "" decommissionDiskLimit := "" @@ -361,6 +362,19 @@ func newClusterSetParasCmd(client *master.MasterClient) *cobra.Command { dpTimeout = strconv.FormatInt(int64(heartbeatTimeout.Seconds()), 10) } + if mpTimeout != "" { + var mpHeartbeatTimeout time.Duration + mpHeartbeatTimeout, err = time.ParseDuration(mpTimeout) + if err != nil { + return + } + if mpHeartbeatTimeout < time.Second { + err = fmt.Errorf("dp timeout %v smaller than 1s", mpHeartbeatTimeout) + return + } + + mpTimeout = strconv.FormatInt(int64(mpHeartbeatTimeout.Seconds()), 10) + } if dpBackupTimeout != "" { var backupTimeout time.Duration backupTimeout, err = time.ParseDuration(dpBackupTimeout) @@ -413,7 +427,7 @@ func newClusterSetParasCmd(client *master.MasterClient) *cobra.Command { optAutoRepairRate, optLoadFactor, opMaxDpCntLimit, opMaxMpCntLimit, clientIDKey, autoDecommissionDisk, autoDecommissionDiskInterval, autoDpMetaRepair, autoDpMetaRepairParallelCnt, - dpRepairTimeout, dpTimeout, dpBackupTimeout, decommissionDpLimit, decommissionDiskLimit, + dpRepairTimeout, dpTimeout, mpTimeout, dpBackupTimeout, decommissionDpLimit, decommissionDiskLimit, forbidWriteOpOfProtoVersion0, dataMediaType, handleTimeout, readDataNodeTimeout); err != nil { return } @@ -437,6 +451,7 @@ func newClusterSetParasCmd(client *master.MasterClient) *cobra.Command { cmd.Flags().StringVar(&autoDpMetaRepairParallelCnt, CliFlagAutoDpMetaRepairParallelCnt, "", "Parallel count of auto data partition meta repair") cmd.Flags().StringVar(&dpRepairTimeout, CliFlagDpRepairTimeout, "", "Data partition repair timeout(example: 1h)") cmd.Flags().StringVar(&dpTimeout, CliFlagDpTimeout, "", "Data partition heartbeat timeout(example: 10s)") + cmd.Flags().StringVar(&mpTimeout, CliFlagMpTimeout, "", "Meta partition heartbeat timeout(example: 10s)") cmd.Flags().StringVar(&autoDecommissionDisk, CliFlagAutoDecommissionDisk, "", "Enable or disable auto decommission disk") cmd.Flags().StringVar(&autoDecommissionDiskInterval, CliFlagAutoDecommissionDiskInterval, "", "Interval of auto decommission disk(example: 10s)") cmd.Flags().StringVar(&dpBackupTimeout, CliFlagDpBackupTimeout, "", "Data partition backup directory timeout(example: 1h)") diff --git a/cli/cmd/const.go b/cli/cmd/const.go index b34f5a94b..5752494fe 100644 --- a/cli/cmd/const.go +++ b/cli/cmd/const.go @@ -148,6 +148,7 @@ const ( CliFlagAutoDpMetaRepairParallelCnt = "autoDpMetaRepairParallelCnt" CliFlagDpRepairTimeout = "dpRepairTimeout" CliFlagDpTimeout = "dpHeartbeatTimeout" + CliFlagMpTimeout = "mpHeartbeatTimeout" CliFlagAutoDecommissionDisk = "autoDecommissionDisk" CliFlagAutoDecommissionDiskInterval = "autoDecommissionDiskInterval" CliFlagDpBackupTimeout = "dpBackupTimeout" diff --git a/cli/cmd/fmt.go b/cli/cmd/fmt.go index 38ee28ef9..9a0aaf897 100644 --- a/cli/cmd/fmt.go +++ b/cli/cmd/fmt.go @@ -79,9 +79,10 @@ func formatClusterView(cv *proto.ClusterView, cn *proto.ClusterNodeInfo, cp *pro sb.WriteString(fmt.Sprintf(" LoadFactor : %v\n", cn.LoadFactor)) sb.WriteString(fmt.Sprintf(" DpRepairTimeout : %v\n", cv.DpRepairTimeout)) sb.WriteString(fmt.Sprintf(" DataPartitionTimeout : %v\n", cv.DpTimeout)) + sb.WriteString(fmt.Sprintf(" MetaPartitionTimeout : %v\n", cv.MpTimeout)) sb.WriteString(fmt.Sprintf(" volDeletionDelayTime : %v h\n", cv.VolDeletionDelayTimeHour)) - sb.WriteString(fmt.Sprintf(" MetaNodeGOGC : %v\n", cv.MetaNodeGOGC)) - sb.WriteString(fmt.Sprintf(" DataNodeGOGC : %v\n", cv.DataNodeGOGC)) + sb.WriteString(fmt.Sprintf(" MetaNodeGOGC : %v\n", cv.MetaNodeGOGC)) + sb.WriteString(fmt.Sprintf(" DataNodeGOGC : %v\n", cv.DataNodeGOGC)) sb.WriteString(fmt.Sprintf(" EnableAutoDecommission : %v\n", cv.EnableAutoDecommission)) sb.WriteString(fmt.Sprintf(" AutoDecommissionDiskInterval : %v\n", cv.AutoDecommissionDiskInterval)) sb.WriteString(fmt.Sprintf(" EnableAutoDpMetaRepair : %v\n", cv.EnableAutoDpMetaRepair)) diff --git a/master/api_args_parse.go b/master/api_args_parse.go index 2c077b9ef..df8e38195 100644 --- a/master/api_args_parse.go +++ b/master/api_args_parse.go @@ -1721,6 +1721,17 @@ func parseAndExtractSetNodeInfoParams(r *http.Request) (params map[string]interf params[dpTimeoutKey] = val } + if value = r.FormValue(mpTimeoutKey); value != "" { + noParams = false + val := int64(0) + val, err = strconv.ParseInt(value, 10, 64) + if err != nil { + err = unmatchedKey(mpTimeoutKey) + return + } + params[mpTimeoutKey] = val + } + if value = r.FormValue(decommissionLimit); value != "" { noParams = false val := uint64(0) diff --git a/master/api_service.go b/master/api_service.go index 0fc828848..b5eed1256 100644 --- a/master/api_service.go +++ b/master/api_service.go @@ -935,6 +935,7 @@ func (m *Server) getCluster(w http.ResponseWriter, r *http.Request) { DecommissionLimit: atomic.LoadUint64(&m.cluster.DecommissionLimit), 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), @@ -3824,6 +3825,15 @@ func (m *Server) setNodeInfoHandler(w http.ResponseWriter, r *http.Request) { } } + if val, ok := params[mpTimeoutKey]; ok { + if mpTimeout, ok := val.(int64); ok { + if err = m.cluster.setMetaPartitionTimeout(mpTimeout); err != nil { + sendErrReply(w, r, newErrHTTPReply(err)) + return + } + } + } + if val, ok := params[decommissionDiskLimit]; ok { if diskLimit, ok := val.(uint64); ok { if err = m.cluster.setDecommissionDiskLimit(uint32(diskLimit)); err != nil { diff --git a/master/api_service_test.go b/master/api_service_test.go index a4804404f..e3ac9c87a 100644 --- a/master/api_service_test.go +++ b/master/api_service_test.go @@ -1760,6 +1760,18 @@ func TestSetDpTimeout(t *testing.T) { require.EqualValues(t, oldVal, server.cluster.getDataPartitionTimeoutSec()) } +func TestSetMpTimeout(t *testing.T) { + reqUrl := fmt.Sprintf("%v%v", hostAddr, proto.AdminSetNodeInfo) + oldVal := server.cluster.getMetaPartitionTimeoutSec() + setVal := int64(10) + setUrl := fmt.Sprintf("%v?%v=%v&dirSizeLimit=0", reqUrl, mpTimeoutKey, setVal) + unsetUrl := fmt.Sprintf("%v?%v=%v&dirSizeLimit=0", reqUrl, mpTimeoutKey, oldVal) + process(setUrl, t) + require.EqualValues(t, setVal, server.cluster.getMetaPartitionTimeoutSec()) + process(unsetUrl, t) + require.EqualValues(t, oldVal, server.cluster.getMetaPartitionTimeoutSec()) +} + func TestSetDpRepairTimeout(t *testing.T) { reqUrl := fmt.Sprintf("%v%v", hostAddr, proto.AdminSetNodeInfo) oldVal := server.cluster.cfg.DpRepairTimeOut diff --git a/master/cluster.go b/master/cluster.go index a734c5f8a..28369d8b3 100644 --- a/master/cluster.go +++ b/master/cluster.go @@ -4490,6 +4490,18 @@ func (c *Cluster) setDataPartitionTimeout(val int64) (err error) { return } +func (c *Cluster) setMetaPartitionTimeout(val int64) (err error) { + oldVal := atomic.LoadInt64(&c.cfg.MetaPartitionTimeOutSec) + atomic.StoreInt64(&c.cfg.MetaPartitionTimeOutSec, val) + if err = c.syncPutCluster(); err != nil { + log.LogErrorf("[setMetaPartitionTimeout] failed to set dp timeout, err(%v)", err) + atomic.StoreInt64(&c.cfg.MetaPartitionTimeOutSec, oldVal) + err = proto.ErrPersistenceByRaft + return + } + return +} + func (c *Cluster) setMetaNodeDeleteWorkerSleepMs(val uint64) (err error) { oldVal := atomic.LoadUint64(&c.cfg.MetaNodeDeleteWorkerSleepMs) atomic.StoreUint64(&c.cfg.MetaNodeDeleteWorkerSleepMs, val) @@ -4579,6 +4591,11 @@ func (c *Cluster) setEnableAutoDpMetaRepair(val bool) (err error) { return } +func (c *Cluster) getEnableAutoDpMetaRepair() (v bool) { + v = c.EnableAutoDpMetaRepair.Load() + return +} + func (c *Cluster) getDataPartitionTimeoutSec() (val int64) { val = atomic.LoadInt64(&c.cfg.DataPartitionTimeOutSec) if val == 0 { @@ -4587,8 +4604,11 @@ func (c *Cluster) getDataPartitionTimeoutSec() (val int64) { return } -func (c *Cluster) getEnableAutoDpMetaRepair() (v bool) { - v = c.EnableAutoDpMetaRepair.Load() +func (c *Cluster) getMetaPartitionTimeoutSec() (val int64) { + val = atomic.LoadInt64(&c.cfg.MetaPartitionTimeOutSec) + if val == 0 { + val = defaultMetaPartitionTimeOutSec + } return } diff --git a/master/cluster_task.go b/master/cluster_task.go index 29da10ff4..2373f3721 100644 --- a/master/cluster_task.go +++ b/master/cluster_task.go @@ -294,7 +294,7 @@ func (c *Cluster) checkReplicaMetaPartitions() ( vol.mpsLock.RLock() for _, mp := range vol.MetaPartitions { - if uint8(len(mp.Hosts)) < mp.ReplicaNum || uint8(len(mp.getActiveAddrs())) < mp.ReplicaNum { + if uint8(len(mp.Hosts)) < mp.ReplicaNum || uint8(len(mp.getActiveAddrs(defaultMetaPartitionTimeOutSec))) < mp.ReplicaNum { lackReplicaMetaPartitions = append(lackReplicaMetaPartitions, mp) } diff --git a/master/config.go b/master/config.go index 3a7176302..3a3634aa7 100644 --- a/master/config.go +++ b/master/config.go @@ -39,6 +39,7 @@ const ( cfgDpNoLeaderReportIntervalSec = "dpNoLeaderReportIntervalSec" cfgMpNoLeaderReportIntervalSec = "mpNoLeaderReportIntervalSec" dataPartitionTimeOutSec = "dataPartitionTimeOutSec" + metaPartitionTimeOutSec = "metaPartitionTimeOutSec" NumberOfDataPartitionsToLoad = "numberOfDataPartitionsToLoad" secondsToFreeDataPartitionAfterLoad = "secondsToFreeDataPartitionAfterLoad" nodeSetCapacity = "nodeSetCap" @@ -135,6 +136,7 @@ type clusterConfig struct { DpNoLeaderReportIntervalSec int64 MpNoLeaderReportIntervalSec int64 DataPartitionTimeOutSec int64 + MetaPartitionTimeOutSec int64 IntervalToAlarmMissingDataPartition int64 PeriodToLoadALLDataPartitions int64 metaNodeReservedMem uint64 @@ -217,6 +219,7 @@ func newClusterConfig() (cfg *clusterConfig) { cfg.DpNoLeaderReportIntervalSec = defaultDpNoLeaderReportIntervalSec cfg.MpNoLeaderReportIntervalSec = defaultMpNoLeaderReportIntervalSec cfg.DataPartitionTimeOutSec = defaultDataPartitionTimeOutSec + cfg.MetaPartitionTimeOutSec = defaultMetaPartitionTimeOutSec cfg.IntervalToCheckDataPartition = defaultIntervalToCheckDataPartition cfg.IntervalToCheckQos = defaultIntervalToCheckQos cfg.IntervalToAlarmMissingDataPartition = defaultIntervalToAlarmMissingDataPartition diff --git a/master/const.go b/master/const.go index bcb567636..101a164f8 100644 --- a/master/const.go +++ b/master/const.go @@ -144,6 +144,7 @@ const ( autoDpMetaRepairKey = "autoDpMetaRepair" autoDpMetaRepairParallelCntKey = "autoDpMetaRepairParallelCnt" dpTimeoutKey = "dpTimeout" + mpTimeoutKey = "mpTimeout" ShowAll = "showAll" trashIntervalKey = "trashInterval" accessTimeIntervalKey = "accessTimeValidInterval" diff --git a/master/meta_partition.go b/master/meta_partition.go index b9649d59a..501e604b6 100644 --- a/master/meta_partition.go +++ b/master/meta_partition.go @@ -280,11 +280,11 @@ func (mp *MetaPartition) isLeaderExist() bool { return false } -func (mp *MetaPartition) checkLeader(clusterID string) { +func (mp *MetaPartition) checkLeader(clusterID string, timeOutSec int64) { mp.Lock() defer mp.Unlock() for _, mr := range mp.Replicas { - if !mr.isActive() { + if !mr.isActive(timeOutSec) { mr.IsLeader = false } } @@ -298,7 +298,7 @@ func (mp *MetaPartition) checkLeader(clusterID string) { } } -func (mp *MetaPartition) checkStatus(clusterID string, writeLog bool, replicaNum int, maxPartitionID uint64, metaPartitionInodeIdStep uint64, forbiddenVol bool) (doSplit bool) { +func (mp *MetaPartition) checkStatus(clusterID string, writeLog bool, replicaNum int, maxPartitionID uint64, metaPartitionInodeIdStep uint64, forbiddenVol bool, timeOutSec int64) (doSplit bool) { if mp.IsFreeze { return } @@ -306,8 +306,8 @@ func (mp *MetaPartition) checkStatus(clusterID string, writeLog bool, replicaNum mp.Lock() defer mp.Unlock() - mp.checkReplicas() - liveReplicas := mp.getLiveReplicas() + mp.checkReplicas(timeOutSec) + liveReplicas := mp.getLiveReplicas(timeOutSec) if len(liveReplicas) <= replicaNum/2 { mp.Status = proto.Unavailable @@ -465,7 +465,7 @@ func (mp *MetaPartition) updateMetaPartition(mgr *proto.MetaPartitionReport, met } func (mp *MetaPartition) canBeOffline(nodeAddr string, replicaNum int) (err error) { - liveReplicas := mp.getLiveReplicas() + liveReplicas := mp.getLiveReplicas(defaultMetaPartitionTimeOutSec) if len(liveReplicas) < int(mp.ReplicaNum/2+1) { err = proto.ErrNoEnoughReplica return @@ -505,19 +505,19 @@ func (mp *MetaPartition) getLiveReplicasAddr(liveReplicas []*MetaReplica) (addrs return } -func (mp *MetaPartition) getLiveReplicas() (liveReplicas []*MetaReplica) { +func (mp *MetaPartition) getLiveReplicas(timeOutSec int64) (liveReplicas []*MetaReplica) { liveReplicas = make([]*MetaReplica, 0) for _, mr := range mp.Replicas { - if mr.isActive() { + if mr.isActive(timeOutSec) { liveReplicas = append(liveReplicas, mr) } } return } -func (mp *MetaPartition) checkReplicas() { +func (mp *MetaPartition) checkReplicas(timeOutSec int64) { for _, mr := range mp.Replicas { - if !mr.isActive() { + if !mr.isActive(timeOutSec) { mr.Status = proto.Unavailable mr.StatByStorageClass = make([]*proto.StatOfStorageClass, 0) mr.StatByMigrateStorageClass = make([]*proto.StatOfStorageClass, 0) @@ -544,18 +544,18 @@ func (mp *MetaPartition) persistToRocksDB(action, volName string, newHosts []str return } -func (mp *MetaPartition) getActiveAddrs() (liveAddrs []string) { +func (mp *MetaPartition) getActiveAddrs(timeOutSec int64) (liveAddrs []string) { liveAddrs = make([]string, 0) for _, mr := range mp.Replicas { - if mr.isActive() { + if mr.isActive(timeOutSec) { liveAddrs = append(liveAddrs, mr.Addr) } } return liveAddrs } -func (mp *MetaPartition) isMissingReplica(addr string) bool { - return !contains(mp.getActiveAddrs(), addr) +func (mp *MetaPartition) isMissingReplica(addr string, timeOutSec int64) bool { + return !contains(mp.getActiveAddrs(timeOutSec), addr) } func (mp *MetaPartition) shouldReportMissingReplica(addr string, interval int64) (isWarn bool) { @@ -571,12 +571,12 @@ func (mp *MetaPartition) shouldReportMissingReplica(addr string, interval int64) // return false } -func (mp *MetaPartition) reportMissingReplicas(clusterID, leaderAddr string, seconds, interval int64) { +func (mp *MetaPartition) reportMissingReplicas(clusterID, leaderAddr string, timeOutSec int64, interval int64) { mp.Lock() defer mp.Unlock() for _, replica := range mp.Replicas { // reduce the alarm frequency - if contains(mp.Hosts, replica.Addr) && replica.isMissing() { + if contains(mp.Hosts, replica.Addr) && replica.isMissing(timeOutSec) { if mp.shouldReportMissingReplica(replica.Addr, interval) { metaNode := replica.metaNode var lastReportTime time.Time @@ -587,7 +587,7 @@ func (mp *MetaPartition) reportMissingReplicas(clusterID, leaderAddr string, sec } msg := fmt.Sprintf("action[reportMissingReplicas], clusterID[%v] volName[%v] partition:%v on node:%v "+ "miss time > :%v vlocLastRepostTime:%v dnodeLastReportTime:%v nodeisActive:%v", - clusterID, mp.volName, mp.PartitionID, replica.Addr, seconds, replica.ReportTime, lastReportTime, isActive) + clusterID, mp.volName, mp.PartitionID, replica.Addr, timeOutSec, replica.ReportTime, lastReportTime, isActive) Warn(clusterID, msg) if WarnMetrics != nil { WarnMetrics.WarnMissingMp(clusterID, replica.Addr, mp.PartitionID, true) @@ -603,10 +603,10 @@ func (mp *MetaPartition) reportMissingReplicas(clusterID, leaderAddr string, sec WarnMetrics.CleanObsoleteMpMissing(clusterID, mp) } for _, addr := range mp.Hosts { - if mp.isMissingReplica(addr) && mp.shouldReportMissingReplica(addr, interval) { + if mp.isMissingReplica(addr, timeOutSec) && mp.shouldReportMissingReplica(addr, interval) { msg := fmt.Sprintf("action[reportMissingReplicas],clusterID[%v] volName[%v] partition:%v on node:%v "+ "miss time > %v ", - clusterID, mp.volName, mp.PartitionID, addr, defaultMetaPartitionTimeOutSec) + clusterID, mp.volName, mp.PartitionID, addr, timeOutSec) Warn(clusterID, msg) msg = fmt.Sprintf("decommissionMetaPartitionURL is http://%v/dataPartition/decommission?id=%v&addr=%v", leaderAddr, mp.PartitionID, addr) Warn(clusterID, msg) @@ -762,13 +762,13 @@ func (mr *MetaReplica) createTaskToLoadMetaPartition(partitionID uint64) (t *pro return } -func (mr *MetaReplica) isMissing() (miss bool) { - return time.Now().Unix()-mr.ReportTime > defaultMetaPartitionTimeOutSec +func (mr *MetaReplica) isMissing(timeOutSec int64) (miss bool) { + return time.Now().Unix()-mr.ReportTime > timeOutSec } -func (mr *MetaReplica) isActive() (active bool) { +func (mr *MetaReplica) isActive(timeOutSec int64) (active bool) { return mr.metaNode.IsActive && mr.Status != proto.Unavailable && - time.Now().Unix()-mr.ReportTime < defaultMetaPartitionTimeOutSec + time.Now().Unix()-mr.ReportTime < timeOutSec } func (mr *MetaReplica) setLastReportTime() { @@ -874,7 +874,7 @@ func (mp *MetaPartition) getMinusOfMaxInodeID() (minus float64) { func (mp *MetaPartition) activeMaxInodeSimilar() bool { minus := float64(0) var sentry float64 - replicas := mp.getLiveReplicas() + replicas := mp.getLiveReplicas(defaultMetaPartitionTimeOutSec) for index, replica := range replicas { if index == 0 { sentry = float64(replica.MaxInodeID) @@ -936,7 +936,7 @@ func (mp *MetaPartition) setDentryCount() { func (mp *MetaPartition) SetForbidWriteOpOfProtoVer0() { for _, r := range mp.Replicas { - if !r.isActive() { + if !r.isActive(defaultMetaPartitionTimeOutSec) { continue } if !r.ForbidWriteOpOfProtoVer0 { diff --git a/master/metadata_fsm_op.go b/master/metadata_fsm_op.go index a17bbd13b..de5b7e3ea 100644 --- a/master/metadata_fsm_op.go +++ b/master/metadata_fsm_op.go @@ -71,6 +71,7 @@ type clusterValue struct { EnableAutoDpMetaRepair bool AutoDpMetaRepairParallelCnt uint32 DataPartitionTimeoutSec int64 + MetaPartitionTimeoutSec int64 ForbidWriteOpOfProtoVer0 bool LegacyDataMediaType uint32 RaftPartitionAlreadyUseDifferentPort bool @@ -119,6 +120,7 @@ func newClusterValue(c *Cluster) (cv *clusterValue) { 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(), @@ -1141,6 +1143,10 @@ func (c *Cluster) updateDataPartitionTimeoutSec(val int64) { atomic.StoreInt64(&c.cfg.DataPartitionTimeOutSec, val) } +func (c *Cluster) updateMetaPartitionTimeoutSec(val int64) { + atomic.StoreInt64(&c.cfg.MetaPartitionTimeOutSec, val) +} + func (c *Cluster) updateDataNodeAutoRepairLimit(val uint64) { atomic.StoreUint64(&c.cfg.DataNodeAutoRepairLimitRate, val) } @@ -1367,6 +1373,7 @@ func (c *Cluster) loadClusterValue() (err error) { c.updateAutoDpMetaRepairParallelCnt(cv.AutoDpMetaRepairParallelCnt) c.updateDataPartitionTimeoutSec(cv.DataPartitionTimeoutSec) c.cfg.raftPartitionAlreadyUseDifferentPort.Store(cv.RaftPartitionAlreadyUseDifferentPort) + c.updateMetaPartitionTimeoutSec(cv.MetaPartitionTimeoutSec) c.cfg.forbidWriteOpOfProtoVer0 = cv.ForbidWriteOpOfProtoVer0 c.legacyDataMediaType = cv.LegacyDataMediaType if cv.MetaNodeMemoryHighPer <= 0.001 { diff --git a/master/monitor_metrics.go b/master/monitor_metrics.go index b955350f2..11962ff64 100644 --- a/master/monitor_metrics.go +++ b/master/monitor_metrics.go @@ -648,7 +648,7 @@ func (mm *monitorMetrics) setMpAndDpMetrics() { continue } - if replicaNum > uint8(len(dp.liveReplicas(defaultDataPartitionTimeOutSec))) { + if replicaNum > uint8(len(dp.liveReplicas(mm.cluster.getDataPartitionTimeoutSec()))) { dpMissingReplicaMap[uint64(replicaNum)]++ } if proto.IsNormalDp(dp.PartitionType) && dp.getLeaderAddr() == "" && time.Now().Unix()-dp.LeaderReportTime > mm.cluster.cfg.DpNoLeaderReportIntervalSec { @@ -660,7 +660,7 @@ func (mm *monitorMetrics) setMpAndDpMetrics() { if !mp.isLeaderExist() && time.Now().Unix()-mp.LeaderReportTime > mm.cluster.cfg.MpNoLeaderReportIntervalSec { mpMissingLeaderCount++ } - if len(mp.getActiveAddrs()) < int(mp.ReplicaNum) { + if len(mp.getActiveAddrs(mm.cluster.getMetaPartitionTimeoutSec())) < int(mp.ReplicaNum) { mpMissingReplicaCount++ } } diff --git a/master/server.go b/master/server.go index 6a940ceff..f97884599 100644 --- a/master/server.go +++ b/master/server.go @@ -351,6 +351,13 @@ func (m *Server) checkConfig(cfg *config.Config) (err error) { } } + metaPartitionTimeOutSec := cfg.GetString(metaPartitionTimeOutSec) + if metaPartitionTimeOutSec != "" { + if m.config.MetaPartitionTimeOutSec, err = strconv.ParseInt(metaPartitionTimeOutSec, 10, 0); err != nil { + return fmt.Errorf("%v,err:%v", proto.ErrInvalidCfg, err.Error()) + } + } + numberOfDataPartitionsToLoad := cfg.GetString(NumberOfDataPartitionsToLoad) if numberOfDataPartitionsToLoad != "" { if m.config.numberOfDataPartitionsToLoad, err = strconv.Atoi(numberOfDataPartitionsToLoad); err != nil { diff --git a/master/vol.go b/master/vol.go index 6dc0cc23b..3587c8c3e 100644 --- a/master/vol.go +++ b/master/vol.go @@ -975,7 +975,7 @@ func (vol *Vol) checkMetaPartitions(c *Cluster) { quotaByClass := vol.getQuotaByClass() for _, mp := range mps { - doSplit = mp.checkStatus(c.Name, true, int(vol.mpReplicaNum), maxPartitionID, metaPartitionInodeIdStep, vol.Forbidden) + doSplit = mp.checkStatus(c.Name, true, int(vol.mpReplicaNum), maxPartitionID, metaPartitionInodeIdStep, vol.Forbidden, c.getMetaPartitionTimeoutSec()) if doSplit && !c.cfg.DisableAutoCreate { nextStart := mp.MaxInodeID + metaPartitionInodeIdStep log.LogInfof(c.Name, fmt.Sprintf("cluster[%v],vol[%v],meta partition[%v] splits start[%v] maxinodeid:[%v] default step:[%v],nextStart[%v]", @@ -985,10 +985,10 @@ func (vol *Vol) checkMetaPartitions(c *Cluster) { } } - mp.checkLeader(c.Name) + mp.checkLeader(c.Name, c.getMetaPartitionTimeoutSec()) mp.checkReplicaNum(c, vol.Name, vol.mpReplicaNum) mp.checkEnd(c, maxPartitionID) - mp.reportMissingReplicas(c.Name, c.leaderInfo.addr, defaultMetaPartitionTimeOutSec, defaultIntervalToAlarmMissingMetaPartition) + mp.reportMissingReplicas(c.Name, c.leaderInfo.addr, c.getMetaPartitionTimeoutSec(), defaultIntervalToAlarmMissingMetaPartition) tasks = append(tasks, mp.replicaCreationTasks(c.Name, vol.Name)...) for _, mpStat := range mp.StatByStorageClass { @@ -1062,7 +1062,7 @@ func (vol *Vol) checkSplitMetaPartition(c *Cluster, metaPartitionInodeStep uint6 } func (mp *MetaPartition) memUsedReachThreshold(clusterName, volName string) bool { - liveReplicas := mp.getLiveReplicas() + liveReplicas := mp.getLiveReplicas(defaultMetaPartitionTimeOutSec) foundReadonlyReplica := false var readonlyReplica *MetaReplica for _, replica := range liveReplicas { diff --git a/master/vol_test.go b/master/vol_test.go index 820bb1e3a..3401cddc1 100644 --- a/master/vol_test.go +++ b/master/vol_test.go @@ -307,7 +307,7 @@ func checkMetaPartitionsWritableTest(vol *Vol, t *testing.T) { maxPartitionID := vol.maxMetaPartitionID() maxMp := vol.MetaPartitions[maxPartitionID] // after check meta partitions ,the status must be writable - maxMp.checkStatus(server.cluster.Name, false, int(vol.mpReplicaNum), maxPartitionID, 4194304, vol.Forbidden) + maxMp.checkStatus(server.cluster.Name, false, int(vol.mpReplicaNum), maxPartitionID, 4194304, vol.Forbidden, defaultMetaPartitionTimeOutSec) if maxMp.Status != proto.ReadWrite { t.Errorf("expect partition status[%v],real status[%v]\n", proto.ReadWrite, maxMp.Status) return diff --git a/proto/model.go b/proto/model.go index b54526a88..9d2bff6d4 100644 --- a/proto/model.go +++ b/proto/model.go @@ -153,6 +153,7 @@ type ClusterView struct { DpRepairTimeout string DpBackupTimeout string DpTimeout string + MpTimeout string DataNodeStatInfo *NodeStatInfo MetaNodeStatInfo *NodeStatInfo VolStatInfo []*VolStatInfo diff --git a/sdk/master/api_admin.go b/sdk/master/api_admin.go index f2c323ece..964a9f165 100644 --- a/sdk/master/api_admin.go +++ b/sdk/master/api_admin.go @@ -592,7 +592,7 @@ func (api *AdminAPI) SetMasterVolDeletionDelayTime(volDeletionDelayTimeHour int) func (api *AdminAPI) SetClusterParas(batchCount, markDeleteRate, deleteWorkerSleepMs, autoRepairRate, loadFactor, maxDpCntLimit, maxMpCntLimit, clientIDKey string, enableAutoDecommissionDisk string, autoDecommissionDiskInterval string, enableAutoDpMetaRepair string, autoDpMetaRepairParallelCnt string, - dpRepairTimeout string, dpTimeout string, dpBackupTimeout string, + dpRepairTimeout string, dpTimeout string, mpTimeout string, dpBackupTimeout string, decommissionDpLimit, decommissionDiskLimit, forbidWriteOpOfProtoVersion0 string, mediaType string, handleTimeout string, readDataNodeTimeout string, ) (err error) { @@ -631,6 +631,9 @@ func (api *AdminAPI) SetClusterParas(batchCount, markDeleteRate, deleteWorkerSle if dpTimeout != "" { request.addParam("dpTimeout", dpTimeout) } + if mpTimeout != "" { + request.addParam("mpTimeout", mpTimeout) + } if dpBackupTimeout != "" { request.addParam("dpBackupTimeout", dpBackupTimeout) }