From 981b298ca30fc1a0bf16edfe022d07e8ee3cff3b Mon Sep 17 00:00:00 2001 From: leonrayang Date: Wed, 19 Jun 2024 10:00:36 +0800 Subject: [PATCH] refactor(master): refactor the code of zone moudule Signed-off-by: leonrayang --- master/api_service.go | 36 ++--- master/cluster.go | 143 +++++++++--------- master/cluster_stat.go | 4 +- master/data_node.go | 44 ++++-- master/gapi_cluster.go | 4 +- master/meta_node.go | 35 ++++- master/monitor_metrics.go | 8 +- master/node_selector.go | 192 +++++++----------------- master/nodeset_selector.go | 4 +- master/topology.go | 297 +++++++++++++------------------------ master/topology_test.go | 45 +++++- master/vol.go | 5 +- 12 files changed, 367 insertions(+), 450 deletions(-) diff --git a/master/api_service.go b/master/api_service.go index 5a7125e3d..889bcd431 100644 --- a/master/api_service.go +++ b/master/api_service.go @@ -368,7 +368,7 @@ func (m *Server) getTopology(w http.ResponseWriter, r *http.Request) { dataNode := value.(*DataNode) nsView.DataNodes = append(nsView.DataNodes, proto.NodeView{ ID: dataNode.ID, Addr: dataNode.Addr, - DomainAddr: dataNode.DomainAddr, IsActive: dataNode.isActive, IsWritable: dataNode.isWriteAble(), + DomainAddr: dataNode.DomainAddr, IsActive: dataNode.isActive, IsWritable: dataNode.IsWriteAble(), }) return true }) @@ -376,7 +376,7 @@ func (m *Server) getTopology(w http.ResponseWriter, r *http.Request) { metaNode := value.(*MetaNode) nsView.MetaNodes = append(nsView.MetaNodes, proto.NodeView{ ID: metaNode.ID, Addr: metaNode.Addr, - DomainAddr: metaNode.DomainAddr, IsActive: metaNode.IsActive, IsWritable: metaNode.isWritable(), + DomainAddr: metaNode.DomainAddr, IsActive: metaNode.IsActive, IsWritable: metaNode.IsWriteAble(), }) return true }) @@ -480,8 +480,8 @@ func (m *Server) listNodeSets(w http.ResponseWriter, r *http.Request) { ID: ns.ID, Capacity: ns.Capacity, Zone: zone.name, - CanAllocDataNodeCnt: ns.canAllocDataNodeCnt(), - CanAllocMetaNodeCnt: ns.canAllocMetaNodeCnt(), + CanAllocDataNodeCnt: ns.calcNodesForAlloc(ns.dataNodes), + CanAllocMetaNodeCnt: ns.calcNodesForAlloc(ns.metaNodes), DataNodeNum: ns.dataNodeLen(), MetaNodeNum: ns.metaNodeLen(), } @@ -521,8 +521,8 @@ func (m *Server) getNodeSet(w http.ResponseWriter, r *http.Request) { ID: ns.ID, Capacity: ns.Capacity, Zone: ns.zoneName, - CanAllocDataNodeCnt: ns.canAllocDataNodeCnt(), - CanAllocMetaNodeCnt: ns.canAllocMetaNodeCnt(), + CanAllocDataNodeCnt: ns.calcNodesForAlloc(ns.dataNodes), + CanAllocMetaNodeCnt: ns.calcNodesForAlloc(ns.metaNodes), DataNodeSelector: ns.GetDataNodeSelector(), MetaNodeSelector: ns.GetMetaNodeSelector(), } @@ -533,7 +533,7 @@ func (m *Server) getNodeSet(w http.ResponseWriter, r *http.Request) { Status: dn.isActive, DomainAddr: dn.DomainAddr, ID: dn.ID, - IsWritable: dn.isWriteAble(), + IsWritable: dn.IsWriteAble(), Total: dn.Total, Used: dn.Used, Avail: dn.Total - dn.Used, @@ -547,7 +547,7 @@ func (m *Server) getNodeSet(w http.ResponseWriter, r *http.Request) { Status: mn.IsActive, DomainAddr: mn.DomainAddr, ID: mn.ID, - IsWritable: mn.isWritable(), + IsWritable: mn.IsWriteAble(), Total: mn.Total, Used: mn.Used, Avail: mn.Total - mn.Used, @@ -2827,7 +2827,7 @@ func (m *Server) getDataNode(w http.ResponseWriter, r *http.Request) { ReportTime: dataNode.ReportTime, IsActive: dataNode.isActive, ToBeOffline: dataNode.ToBeOffline, - IsWriteAble: dataNode.isWriteAble(), + IsWriteAble: dataNode.IsWriteAble(), UsageRatio: dataNode.UsageRatio, SelectedTimes: dataNode.SelectedTimes, DataPartitionReports: dataNode.DataPartitionReports, @@ -2837,7 +2837,7 @@ func (m *Server) getDataNode(w http.ResponseWriter, r *http.Request) { BadDisks: dataNode.BadDisks, RdOnly: dataNode.RdOnly, CanAllocPartition: dataNode.canAlloc() && dataNode.canAllocDp(), - MaxDpCntLimit: dataNode.GetDpCntLimit(), + MaxDpCntLimit: dataNode.GetPartitionLimitCnt(), CpuUtil: dataNode.CpuUtil.Load(), IoUtils: dataNode.GetIoUtils(), DecommissionedDisk: dataNode.getDecommissionedDisks(), @@ -2927,7 +2927,7 @@ func (m *Server) migrateDataNodeHandler(w http.ResponseWriter, r *http.Request) return } - if !targetNode.isWriteAble() || !targetNode.dpCntInLimit() { + if !targetNode.IsWriteAble() || !targetNode.PartitionCntLimited() { err = fmt.Errorf("[%s] is not writable, can't used as target addr for migrate", targetAddr) sendErrReply(w, r, newErrHTTPReply(err)) return @@ -3472,7 +3472,7 @@ func (m *Server) buildNodeSetGrpInfo(nsg *nodeSetGroup) *proto.SimpleNodeSetGrpI nsg.nodeSets[i].dataNodes.Range(func(key, value interface{}) bool { node := value.(*DataNode) nsStat.DataTotal += node.Total - if node.isWriteAble() { + if node.IsWriteAble() { nsStat.DataUsed += node.Used } else { nsStat.DataUsed += node.Total @@ -3489,7 +3489,7 @@ func (m *Server) buildNodeSetGrpInfo(nsg *nodeSetGroup) *proto.SimpleNodeSetGrpI Addr: node.Addr, ReportTime: node.ReportTime, IsActive: node.isActive, - IsWriteAble: node.isWriteAble(), + IsWriteAble: node.IsWriteAble(), UsageRatio: node.UsageRatio, SelectedTimes: node.SelectedTimes, DataPartitionCount: node.DataPartitionCount, @@ -3512,7 +3512,7 @@ func (m *Server) buildNodeSetGrpInfo(nsg *nodeSetGroup) *proto.SimpleNodeSetGrpI ID: node.ID, Addr: node.Addr, IsActive: node.IsActive, - IsWriteAble: node.isWritable(), + IsWriteAble: node.IsWriteAble(), ZoneName: node.ZoneName, MaxMemAvailWeight: node.MaxMemAvailWeight, Total: node.Total, @@ -4466,7 +4466,7 @@ func (m *Server) getMetaNode(w http.ResponseWriter, r *http.Request) { Addr: metaNode.Addr, DomainAddr: metaNode.DomainAddr, IsActive: metaNode.IsActive, - IsWriteAble: metaNode.isWritable(), + IsWriteAble: metaNode.IsWriteAble(), ZoneName: metaNode.ZoneName, MaxMemAvailWeight: metaNode.MaxMemAvailWeight, Total: metaNode.Total, @@ -4478,8 +4478,8 @@ func (m *Server) getMetaNode(w http.ResponseWriter, r *http.Request) { MetaPartitionCount: metaNode.MetaPartitionCount, NodeSetID: metaNode.NodeSetID, PersistenceMetaPartitions: metaNode.PersistenceMetaPartitions, - CanAllowPartition: metaNode.isWritable() && metaNode.mpCntInLimit(), - MaxMpCntLimit: metaNode.GetMpCntLimit(), + CanAllowPartition: metaNode.IsWriteAble() && metaNode.PartitionCntLimited(), + MaxMpCntLimit: metaNode.GetPartitionLimitCnt(), CpuUtil: metaNode.CpuUtil.Load(), } sendOkReply(w, r, newSuccessHTTPReply(metaNodeInfo)) @@ -4636,7 +4636,7 @@ func (m *Server) migrateMetaNodeHandler(w http.ResponseWriter, r *http.Request) return } - if !targetNode.isWritable() || !targetNode.mpCntInLimit() { + if !targetNode.IsWriteAble() || !targetNode.PartitionCntLimited() { err = fmt.Errorf("[%s] is not writable, can't used as target addr for migrate", targetAddr) sendErrReply(w, r, newErrHTTPReply(err)) return diff --git a/master/cluster.go b/master/cluster.go index 922053af3..fd32a5d4f 100644 --- a/master/cluster.go +++ b/master/cluster.go @@ -1530,7 +1530,7 @@ func (c *Cluster) createDataPartition(volName string, preload *DataPartitionPreL goto errHandler } } else { - zoneNum := c.decideZoneNum(vol.crossZone) + zoneNum := c.decideZoneNum(vol) // zoneNum scope [1,3] if targetHosts, targetPeers, err = c.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, int(dpReplicaNum), zoneNum, zoneName); err != nil { goto errHandler @@ -1656,53 +1656,51 @@ func (c *Cluster) syncCreateMetaPartitionToMetaNode(host string, mp *MetaPartiti // if vol is not cross zone, return 1 // if vol enable cross zone and the zone number of cluster less than defaultReplicaNum return 2 // otherwise, return defaultReplicaNum -func (c *Cluster) decideZoneNum(crossZone bool) (zoneNum int) { - if !crossZone { +func (c *Cluster) decideZoneNum(vol *Vol) (zoneNum int) { + if !vol.crossZone { return 1 } - + specificZoneList := strings.Split(vol.zoneName, ",") var zoneLen int if c.FaultDomain { zoneLen = len(c.t.domainExcludeZones) } else { - zoneLen = c.t.zoneLen() + zoneLen = 1 + if len(specificZoneList) > 1 { + zoneLen = len(specificZoneList) + } else if vol.crossZone { + zoneLen = 2 + } } if zoneLen < defaultReplicaNum { zoneNum = 2 - } else { + } + if zoneLen > defaultReplicaNum { zoneNum = defaultReplicaNum } return zoneNum } -func (c *Cluster) chooseZone2Plus1(zones []*Zone, excludeNodeSets []uint64, excludeHosts []string, +func (c *Cluster) chooseZone2Plus1(rsMgr *rsManager, zones []*Zone, excludeNodeSets []uint64, excludeHosts []string, nodeType uint32, replicaNum int) (hosts []string, peers []proto.Peer, err error, ) { - if replicaNum < 2 || replicaNum > 3 { + if replicaNum < 2 || replicaNum > defaultReplicaNum { return nil, nil, fmt.Errorf("action[chooseZone2Plus1] replicaNum [%v]", replicaNum) } - zoneList := make([]*Zone, 2) - if zones[0].getSpaceLeft(nodeType) < zones[1].getSpaceLeft(nodeType) { - zoneList[0] = zones[0] - zoneList[1] = zones[1] - } else { - zoneList[0] = zones[1] - zoneList[1] = zones[0] - } - - for i := 2; i < len(zones); i++ { - spaceLeft := zones[i].getSpaceLeft(nodeType) - if spaceLeft > zoneList[0].getSpaceLeft(nodeType) { - if spaceLeft > zoneList[1].getSpaceLeft(nodeType) { - zoneList[1] = zones[i] - } else { - zoneList[0] = zones[i] - } + for i := range []int{1, 2} { + if rsMgr.zoneIndexForNode >= len(zones) { + rsMgr.zoneIndexForNode = 0 } + zoneList[i] = zones[rsMgr.zoneIndexForNode] + rsMgr.zoneIndexForNode++ } + sort.Slice(zoneList, func(i, j int) bool { + return zoneList[i].getSpaceLeft(nodeType) < zoneList[j].getSpaceLeft(nodeType) + }) + log.LogInfof("action[chooseZone2Plus1] type [%v] after check,zone0 [%v] left [%v] zone1 [%v] left [%v]", nodeType, zoneList[0].name, zoneList[0].getSpaceLeft(nodeType), zoneList[1].name, zoneList[1].getSpaceLeft(nodeType)) @@ -1750,46 +1748,57 @@ func (c *Cluster) chooseZoneNormal(zones []*Zone, excludeNodeSets []uint64, excl return } -func (c *Cluster) getHostFromNormalZone(nodeType uint32, excludeZones []string, excludeNodeSets []uint64, - excludeHosts []string, replicaNum int, - zoneNum int, specifiedZone string) (hosts []string, peers []proto.Peer, err error, -) { - var zones []*Zone - zones = make([]*Zone, 0) - if replicaNum <= zoneNum { - zoneNum = replicaNum - } +func (c *Cluster) getSpecificZoneList(specifiedZone string) (zones []*Zone, err error) { // when creating vol,user specified a zone,we reset zoneNum to 1,to be created partition with specified zone, // if specified zone is not writable,we choose a zone randomly - if specifiedZone != "" { - if err = c.checkNormalZoneName(specifiedZone); err != nil { + + if err = c.checkNormalZoneName(specifiedZone); err != nil { + Warn(c.Name, fmt.Sprintf("cluster[%v],specified zone[%v]is found", c.Name, specifiedZone)) + return + } + zoneList := strings.Split(specifiedZone, ",") + for i := 0; i < len(zoneList); i++ { + var zone *Zone + if zone, err = c.t.getZone(zoneList[i]); err != nil { Warn(c.Name, fmt.Sprintf("cluster[%v],specified zone[%v]is found", c.Name, specifiedZone)) return } - zoneList := strings.Split(specifiedZone, ",") - for i := 0; i < len(zoneList); i++ { - var zone *Zone - if zone, err = c.t.getZone(zoneList[i]); err != nil { - Warn(c.Name, fmt.Sprintf("cluster[%v],specified zone[%v]is found", c.Name, specifiedZone)) - return - } - zones = append(zones, zone) - } - } else { - if nodeType == TypeDataPartition { - if zones, err = c.t.allocZonesForDataNode(zoneNum, replicaNum, excludeZones); err != nil { - return - } - } else { - if zones, err = c.t.allocZonesForMetaNode(zoneNum, replicaNum, excludeZones); err != nil { - return - } - } + zones = append(zones, zone) + } + return +} + +func (c *Cluster) getHostFromNormalZone(nodeType uint32, excludeZones []string, excludeNodeSets []uint64, + excludeHosts []string, replicaNum int, + zoneNumNeed int, specifiedZoneName string) (hosts []string, peers []proto.Peer, err error, +) { + var zonesQualified []*Zone + zonesQualified = make([]*Zone, 0) + if replicaNum <= zoneNumNeed { + zoneNumNeed = replicaNum } - if len(zones) == 1 { - log.LogInfof("action[getHostFromNormalZone] zones [%v]", zones[0].name) - if hosts, peers, err = zones[0].getAvailNodeHosts(nodeType, excludeNodeSets, excludeHosts, replicaNum); err != nil { + var specifiedZones []*Zone + var rsMgr *rsManager + if specifiedZoneName != "" { + if specifiedZones, err = c.getSpecificZoneList(specifiedZoneName); err != nil { + return + } + } + if nodeType == TypeDataPartition { + rsMgr = &c.t.dataTopology + } else { + rsMgr = &c.t.metaTopology + } + + // get all zones that qualified + if zonesQualified, err = c.t.allocZonesForNode(rsMgr, zoneNumNeed, replicaNum, excludeZones, specifiedZones); err != nil { + return + } + + if len(zonesQualified) == 1 { + log.LogInfof("action[getHostFromNormalZone] zones [%v]", zonesQualified[0].name) + if hosts, peers, err = zonesQualified[0].getAvailNodeHosts(nodeType, excludeNodeSets, excludeHosts, replicaNum); err != nil { log.LogErrorf("action[getHostFromNormalZone],err[%v]", err) return } @@ -1799,23 +1808,23 @@ func (c *Cluster) getHostFromNormalZone(nodeType uint32, excludeZones []string, if excludeHosts == nil { excludeHosts = make([]string, 0) } - - if c.cfg.DefaultNormalZoneCnt == defaultNormalCrossZoneCnt && len(zones) >= defaultNormalCrossZoneCnt { - if hosts, peers, err = c.chooseZoneNormal(zones, excludeNodeSets, excludeHosts, nodeType, replicaNum); err != nil { + // The upper process tries to go and get the dedicated zones, and the latter tries to choose the right zones as possible. + if c.cfg.DefaultNormalZoneCnt == defaultNormalCrossZoneCnt && len(zonesQualified) >= defaultNormalCrossZoneCnt { + if hosts, peers, err = c.chooseZoneNormal(zonesQualified, excludeNodeSets, excludeHosts, nodeType, replicaNum); err != nil { return } } else { - if hosts, peers, err = c.chooseZone2Plus1(zones, excludeNodeSets, excludeHosts, nodeType, replicaNum); err != nil { + if hosts, peers, err = c.chooseZone2Plus1(rsMgr, zonesQualified, excludeNodeSets, excludeHosts, nodeType, replicaNum); err != nil { return } } result: - log.LogInfof("action[getHostFromNormalZone] replicaNum[%v],zoneNum[%v],selectedZones[%v],hosts[%v]", replicaNum, zoneNum, len(zones), hosts) + log.LogInfof("action[getHostFromNormalZone] replicaNum[%v],zoneNum[%v],selectedZones[%v],hosts[%v]", replicaNum, zoneNumNeed, len(zonesQualified), hosts) if len(hosts) != replicaNum { - log.LogErrorf("action[getHostFromNormalZone] replicaNum[%v],zoneNum[%v],selectedZones[%v],hosts[%v]", replicaNum, zoneNum, len(zones), hosts) + log.LogErrorf("action[getHostFromNormalZone] replicaNum[%v],zoneNum[%v],selectedZones[%v],hosts[%v]", replicaNum, zoneNumNeed, len(zonesQualified), hosts) return nil, nil, errors.Trace(proto.ErrNoDataNodeToCreateDataPartition, "hosts len[%v],replicaNum[%v],zoneNum[%v],selectedZones[%v]", - len(hosts), replicaNum, zoneNum, len(zones)) + len(hosts), replicaNum, zoneNumNeed, len(zonesQualified)) } return @@ -3438,7 +3447,7 @@ func (c *Cluster) allDataNodes() (dataNodes []proto.NodeView) { dataNode := node.(*DataNode) dataNodes = append(dataNodes, proto.NodeView{ Addr: dataNode.Addr, DomainAddr: dataNode.DomainAddr, - IsActive: dataNode.isActive, ID: dataNode.ID, IsWritable: dataNode.isWriteAble(), + IsActive: dataNode.isActive, ID: dataNode.ID, IsWritable: dataNode.IsWriteAble(), }) return true }) @@ -3451,7 +3460,7 @@ func (c *Cluster) allMetaNodes() (metaNodes []proto.NodeView) { metaNode := node.(*MetaNode) metaNodes = append(metaNodes, proto.NodeView{ ID: metaNode.ID, Addr: metaNode.Addr, DomainAddr: metaNode.DomainAddr, - IsActive: metaNode.IsActive, IsWritable: metaNode.isWritable(), + IsActive: metaNode.IsActive, IsWritable: metaNode.IsWriteAble(), }) return true }) diff --git a/master/cluster_stat.go b/master/cluster_stat.go index 297202da6..01ec427d1 100644 --- a/master/cluster_stat.go +++ b/master/cluster_stat.go @@ -73,7 +73,7 @@ func (c *Cluster) updateZoneStatInfo() { zone.dataNodes.Range(func(key, value interface{}) bool { zs.DataNodeStat.TotalNodes++ node := value.(*DataNode) - if node.isActive && node.isWriteAble() { + if node.isActive && node.IsWriteAble() { zs.DataNodeStat.WritableNodes++ } zs.DataNodeStat.Total += float64(node.Total) / float64(util.GB) @@ -90,7 +90,7 @@ func (c *Cluster) updateZoneStatInfo() { zone.metaNodes.Range(func(key, value interface{}) bool { zs.MetaNodeStat.TotalNodes++ node := value.(*MetaNode) - if node.IsActive && node.isWritable() { + if node.IsActive && node.IsWriteAble() { zs.MetaNodeStat.WritableNodes++ } zs.MetaNodeStat.Total += float64(node.Total) / float64(util.GB) diff --git a/master/data_node.go b/master/data_node.go index e2a0f7307..1acd33cd1 100644 --- a/master/data_node.go +++ b/master/data_node.go @@ -91,6 +91,10 @@ func newDataNode(addr, zoneName, clusterID string) (dataNode *DataNode) { return } +func (dataNode *DataNode) IsActiveNode() bool { + return dataNode.isActive +} + func (dataNode *DataNode) GetIoUtils() map[string]float64 { return dataNode.ioUtils.Load().(map[string]float64) } @@ -186,7 +190,7 @@ func (dataNode *DataNode) canAlloc() bool { return overSoldCap(dataNode.Total) >= dataNode.TotalPartitionSize } -func (dataNode *DataNode) isWriteAble() (ok bool) { +func (dataNode *DataNode) IsWriteAble() (ok bool) { dataNode.RLock() defer dataNode.RUnlock() @@ -203,7 +207,7 @@ func (dataNode *DataNode) availableDiskCount() (cnt int) { } func (dataNode *DataNode) canAllocDp() bool { - if !dataNode.isWriteAble() { + if !dataNode.IsWriteAble() { return false } @@ -217,24 +221,34 @@ func (dataNode *DataNode) canAllocDp() bool { return false } - if !dataNode.dpCntInLimit() { + if !dataNode.PartitionCntLimited() { return false } return true } -func (dataNode *DataNode) GetDpCntLimit() uint32 { +func (dataNode *DataNode) GetPartitionLimitCnt() uint32 { return uint32(dataNode.DpCntLimit.GetCntLimit()) } -func (dataNode *DataNode) dpCntInLimit() bool { - inLimit := dataNode.DataPartitionCount <= dataNode.GetDpCntLimit() - if !inLimit { +func (dataNode *DataNode) GetAvailableSpace() uint64 { + return dataNode.AvailableSpace +} + +func (dataNode *DataNode) PartitionCntLimited() bool { + limited := dataNode.DataPartitionCount <= dataNode.GetPartitionLimitCnt() + if !limited { log.LogInfof("dpCntInLimit: dp count is already over limit for node %s, cnt %d, limit %d", - dataNode.Addr, dataNode.DataPartitionCount, dataNode.GetDpCntLimit()) + dataNode.Addr, dataNode.DataPartitionCount, dataNode.GetPartitionLimitCnt()) } - return inLimit + return limited +} + +func (dataNode *DataNode) GetStorageInfo() string { + return fmt.Sprintf("data node(%v) cannot alloc dp, total space(%v) avaliable space(%v) used space(%v), offline(%v), avaliable disk cnt(%v), dp count(%v), over sold(%v))", + dataNode.GetAddr(), dataNode.GetTotal(), dataNode.GetTotal()-dataNode.GetUsed(), dataNode.GetUsed(), + dataNode.ToBeOffline, dataNode.availableDiskCount(), dataNode.DataPartitionCount, !dataNode.canAlloc()) } func (dataNode *DataNode) isWriteAbleWithSizeNoLock(size uint64) (ok bool) { @@ -253,6 +267,18 @@ func (dataNode *DataNode) isWriteAbleWithSize(size uint64) (ok bool) { return dataNode.isWriteAbleWithSizeNoLock(size) } +func (dataNode *DataNode) GetUsed() uint64 { + dataNode.RLock() + defer dataNode.RUnlock() + return dataNode.Used +} + +func (dataNode *DataNode) GetTotal() uint64 { + dataNode.RLock() + defer dataNode.RUnlock() + return dataNode.Total +} + func (dataNode *DataNode) GetID() uint64 { dataNode.RLock() defer dataNode.RUnlock() diff --git a/master/gapi_cluster.go b/master/gapi_cluster.go index 6ebc7e613..2ce4524a2 100644 --- a/master/gapi_cluster.go +++ b/master/gapi_cluster.go @@ -300,12 +300,12 @@ func (m *ClusterService) getTopology(ctx context.Context, args struct{}) (*proto cv.NodeSet[ns.ID] = nsView ns.dataNodes.Range(func(key, value interface{}) bool { dataNode := value.(*DataNode) - nsView.DataNodes = append(nsView.DataNodes, proto.NodeView{ID: dataNode.ID, Addr: dataNode.Addr, IsActive: dataNode.isActive, IsWritable: dataNode.isWriteAble()}) + nsView.DataNodes = append(nsView.DataNodes, proto.NodeView{ID: dataNode.ID, Addr: dataNode.Addr, IsActive: dataNode.isActive, IsWritable: dataNode.IsWriteAble()}) return true }) ns.metaNodes.Range(func(key, value interface{}) bool { metaNode := value.(*MetaNode) - nsView.MetaNodes = append(nsView.MetaNodes, proto.NodeView{ID: metaNode.ID, Addr: metaNode.Addr, IsActive: metaNode.IsActive, IsWritable: metaNode.isWritable()}) + nsView.MetaNodes = append(nsView.MetaNodes, proto.NodeView{ID: metaNode.ID, Addr: metaNode.Addr, IsActive: metaNode.IsActive, IsWritable: metaNode.IsWriteAble()}) return true }) } diff --git a/master/meta_node.go b/master/meta_node.go index 67ccb2d7e..bc58a36bf 100644 --- a/master/meta_node.go +++ b/master/meta_node.go @@ -15,6 +15,7 @@ package master import ( + "fmt" "sync" "time" @@ -61,10 +62,36 @@ func newMetaNode(addr, zoneName, clusterID string) (node *MetaNode) { return } +func (metaNode *MetaNode) IsActiveNode() bool { + return metaNode.IsActive +} + func (metaNode *MetaNode) clean() { metaNode.Sender.exitCh <- struct{}{} } +func (metaNode *MetaNode) GetStorageInfo() string { + return fmt.Sprintf("meta node(%v) cannot alloc dp, total space(%v) avaliable space(%v) used space(%v), offline(%v), mp count(%v)", + metaNode.GetAddr(), metaNode.GetTotal(), metaNode.GetTotal()-metaNode.GetUsed(), metaNode.GetUsed(), + metaNode.ToBeOffline, metaNode.MetaPartitionCount) +} + +func (metaNode *MetaNode) GetTotal() uint64 { + metaNode.RLock() + defer metaNode.RUnlock() + return metaNode.Total +} + +func (metaNode *MetaNode) GetUsed() uint64 { + metaNode.RLock() + defer metaNode.RUnlock() + return metaNode.Used +} + +func (metaNode *MetaNode) GetAvailableSpace() uint64 { + return metaNode.Total - metaNode.Used +} + func (metaNode *MetaNode) GetID() uint64 { metaNode.RLock() defer metaNode.RUnlock() @@ -84,7 +111,7 @@ func (metaNode *MetaNode) SelectNodeForWrite() { metaNode.SelectCount++ } -func (metaNode *MetaNode) isWritable() (ok bool) { +func (metaNode *MetaNode) IsWriteAble() (ok bool) { metaNode.RLock() defer metaNode.RUnlock() if metaNode.IsActive && metaNode.MaxMemAvailWeight > gConfig.metaNodeReservedMem && @@ -163,13 +190,13 @@ func (metaNode *MetaNode) checkHeartbeat() { } } -func (metaNode *MetaNode) GetMpCntLimit() (limit uint32) { +func (metaNode *MetaNode) GetPartitionLimitCnt() (limit uint32) { limit = uint32(metaNode.MpCntLimit.GetCntLimit()) return } -func (metaNode *MetaNode) mpCntInLimit() bool { - return uint32(metaNode.MetaPartitionCount) <= metaNode.GetMpCntLimit() +func (metaNode *MetaNode) PartitionCntLimited() bool { + return uint32(metaNode.MetaPartitionCount) <= metaNode.GetPartitionLimitCnt() } // LeaderMetaNode define the leader metaPartitions in meta node diff --git a/master/monitor_metrics.go b/master/monitor_metrics.go index d6ad46a06..a57e1959b 100644 --- a/master/monitor_metrics.go +++ b/master/monitor_metrics.go @@ -820,7 +820,7 @@ func (mm *monitorMetrics) updateMetaNodesStat() { mm.nodeStat.SetWithLabelValues(float64(metaNode.Used), MetricRoleMetaNode, metaNode.Addr, "memUsed") mm.nodeStat.SetWithLabelValues(float64(metaNode.MetaPartitionCount), MetricRoleMetaNode, metaNode.Addr, "mpCount") mm.nodeStat.SetWithLabelValues(float64(metaNode.Threshold), MetricRoleMetaNode, metaNode.Addr, "threshold") - mm.nodeStat.SetBoolWithLabelValues(metaNode.isWritable(), MetricRoleMetaNode, metaNode.Addr, "writable") + mm.nodeStat.SetBoolWithLabelValues(metaNode.IsWriteAble(), MetricRoleMetaNode, metaNode.Addr, "writable") mm.nodeStat.SetBoolWithLabelValues(metaNode.IsActive, MetricRoleMetaNode, metaNode.Addr, "active") return true @@ -871,7 +871,7 @@ func (mm *monitorMetrics) updateDataNodesStat() { mm.nodeStat.SetWithLabelValues(dataNode.UsageRatio, MetricRoleDataNode, dataNode.Addr, "usageRatio") mm.nodeStat.SetWithLabelValues(float64(len(dataNode.BadDisks)), MetricRoleDataNode, dataNode.Addr, "badDiskCount") mm.nodeStat.SetBoolWithLabelValues(dataNode.isActive, MetricRoleDataNode, dataNode.Addr, "active") - mm.nodeStat.SetBoolWithLabelValues(dataNode.isWriteAble(), MetricRoleDataNode, dataNode.Addr, "writable") + mm.nodeStat.SetBoolWithLabelValues(dataNode.IsWriteAble(), MetricRoleDataNode, dataNode.Addr, "writable") return true }) mm.dataNodesInactive.Set(float64(inactiveDataNodesCount)) @@ -897,7 +897,7 @@ func (mm *monitorMetrics) setNotWritableMetaNodesCount() { if !ok { return true } - if !metaNode.isWritable() { + if !metaNode.IsWriteAble() { notWritabelMetaNodesCount++ } return true @@ -912,7 +912,7 @@ func (mm *monitorMetrics) setNotWritableDataNodesCount() { if !ok { return true } - if !dataNode.isWriteAble() { + if !dataNode.IsWriteAble() { notWritabelDataNodesCount++ } return true diff --git a/master/node_selector.go b/master/node_selector.go index 0353a3726..231445156 100644 --- a/master/node_selector.go +++ b/master/node_selector.go @@ -65,6 +65,14 @@ type Node interface { SelectNodeForWrite() GetID() uint64 GetAddr() string + PartitionCntLimited() bool + IsActiveNode() bool + IsWriteAble() bool + GetPartitionLimitCnt() uint32 + GetTotal() uint64 + GetUsed() uint64 + GetAvailableSpace() uint64 + GetStorageInfo() string } // SortedWeightedNodes defines an array sorted by carry @@ -82,17 +90,8 @@ func (nodes SortedWeightedNodes) Swap(i, j int) { nodes[i], nodes[j] = nodes[j], nodes[i] } -func canAllocPartition(node interface{}, nodeType NodeType) bool { - switch nodeType { - case DataNodeType: - dataNode := node.(*DataNode) - return dataNode.canAlloc() && dataNode.canAllocDp() - case MetaNodeType: - metaNode := node.(*MetaNode) - return metaNode.isWritable() && metaNode.mpCntInLimit() - default: - panic("unknown node type") - } +func canAllocPartition(node Node) bool { + return node.IsWriteAble() && node.PartitionCntLimited() } func asNodeWrap(node interface{}, nodeType NodeType) Node { @@ -118,130 +117,63 @@ func (s *CarryWeightNodeSelector) GetName() string { return CarryWeightNodeSelectorName } -func (s *CarryWeightNodeSelector) prepareCarryForDataNodes(nodes *sync.Map, total uint64) { - nodes.Range(func(key, value interface{}) bool { - dataNode := value.(*DataNode) - if _, ok := s.carry[dataNode.ID]; !ok { - // use available space to calculate initial weight - s.carry[dataNode.ID] = float64(dataNode.AvailableSpace) / float64(total) - } - return true - }) -} - -func (s *CarryWeightNodeSelector) prepareCarryForMetaNodes(nodes *sync.Map, total uint64) { - nodes.Range(func(key, value interface{}) bool { - metaNode := value.(*MetaNode) - if _, ok := s.carry[metaNode.ID]; !ok { - // use available space to calculate initial weight - s.carry[metaNode.ID] = float64(metaNode.Total-metaNode.Used) / float64(total) - } - return true - }) -} - func (s *CarryWeightNodeSelector) prepareCarry(nodes *sync.Map, total uint64) { - switch s.nodeType { - case DataNodeType: - s.prepareCarryForDataNodes(nodes, total) - case MetaNodeType: - s.prepareCarryForMetaNodes(nodes, total) - default: - } -} - -func (s *CarryWeightNodeSelector) getTotalMaxForDataNodes(nodes *sync.Map) (total uint64) { nodes.Range(func(key, value interface{}) bool { - dataNode := value.(*DataNode) - if dataNode.Total > total { - total = dataNode.Total + node := value.(Node) + if _, ok := s.carry[node.GetID()]; !ok { + // use available space to calculate initial weight + s.carry[node.GetID()] = float64(node.GetAvailableSpace()) / float64(total) } return true }) - return -} - -func (s *CarryWeightNodeSelector) getTotalMaxForMetaNodes(nodes *sync.Map) (total uint64) { - nodes.Range(func(key, value interface{}) bool { - metaNode := value.(*MetaNode) - if metaNode.Total > total { - total = metaNode.Total - } - return true - }) - return } func (s *CarryWeightNodeSelector) getTotalMax(nodes *sync.Map) (total uint64) { - switch s.nodeType { - case DataNodeType: - total = s.getTotalMaxForDataNodes(nodes) - case MetaNodeType: - total = s.getTotalMaxForMetaNodes(nodes) - default: - } - return -} - -func (s *CarryWeightNodeSelector) getCarryDataNodes(maxTotal uint64, excludeHosts []string, dataNodes *sync.Map) (nodeTabs SortedWeightedNodes, availCount int) { - nodeTabs = make(SortedWeightedNodes, 0) - dataNodes.Range(func(key, value interface{}) bool { - dataNode := value.(*DataNode) - if contains(excludeHosts, dataNode.Addr) { - // log.LogDebugf("[getAvailCarryDataNodeTab] dataNode [%v] is excludeHosts", dataNode.Addr) - return true + nodes.Range(func(key, value interface{}) bool { + dataNode := value.(Node) + if dataNode.GetTotal() > total { + total = dataNode.GetTotal() } - if !canAllocPartition(dataNode, s.nodeType) { - log.LogWarnf("[getCarryDataNodes] data node(%v) cannot alloc dp, total space(%v) avaliable space(%v) used space(%v), offline(%v), avaliable disk cnt(%v), dp count(%v), over sold(%v), exclude hosts(%v)", dataNode.Addr, dataNode.Total, dataNode.AvailableSpace, dataNode.Used, dataNode.ToBeOffline, dataNode.availableDiskCount(), dataNode.DataPartitionCount, !dataNode.canAlloc(), excludeHosts) - return true - } - if s.carry[dataNode.ID] >= 1.0 { - availCount++ - } - - nt := new(weightedNode) - nt.Carry = s.carry[dataNode.ID] - nt.Weight = float64(dataNode.Total-dataNode.Used) / float64(maxTotal) - nt.Ptr = dataNode - nodeTabs = append(nodeTabs, nt) - return true - }) - return -} - -func (s *CarryWeightNodeSelector) getCarryMetaNodes(maxTotal uint64, excludeHosts []string, metaNodes *sync.Map) (nodes SortedWeightedNodes, availCount int) { - nodes = make(SortedWeightedNodes, 0) - metaNodes.Range(func(key, value interface{}) bool { - metaNode := value.(*MetaNode) - if contains(excludeHosts, metaNode.Addr) { - return true - } - if !canAllocPartition(metaNode, s.nodeType) { - log.LogWarnf("[getCarryMetaNodes] meta node(%v) cannot alloc mp, total space(%v) used space(%v), offline(%v), mp count(%v), exclude hosts(%v)", metaNode.Addr, metaNode.Total, metaNode.Used, metaNode.ToBeOffline, metaNode.MetaPartitionCount, excludeHosts) - return true - } - if s.carry[metaNode.ID] >= 1.0 { - availCount++ - } - nt := new(weightedNode) - nt.Carry = s.carry[metaNode.ID] - nt.Weight = (float64)(metaNode.Total-metaNode.Used) / (float64)(maxTotal) - nt.Ptr = metaNode - nodes = append(nodes, nt) return true }) return } func (s *CarryWeightNodeSelector) getCarryNodes(nset *nodeSet, maxTotal uint64, excludeHosts []string) (SortedWeightedNodes, int) { + var nodes *sync.Map switch s.nodeType { case DataNodeType: - return s.getCarryDataNodes(maxTotal, excludeHosts, nset.dataNodes) + nodes = nset.dataNodes case MetaNodeType: - return s.getCarryMetaNodes(maxTotal, excludeHosts, nset.metaNodes) + nodes = nset.metaNodes default: panic("unknown node type") } + + nodeTabs := make(SortedWeightedNodes, 0) + availCount := 0 + nodes.Range(func(key, value interface{}) bool { + node := value.(Node) + if contains(excludeHosts, node.GetAddr()) { + // log.LogDebugf("[getAvailCarryDataNodeTab] dataNode [%v] is excludeHosts", dataNode.Addr) + return true + } + if !canAllocPartition(node) { + log.LogWarnf("[getCarryDataNodes] nodeType (%v) storage info (%v) exclude hosts(%v)", s.nodeType, node.GetStorageInfo(), excludeHosts) + return true + } + if s.carry[node.GetID()] >= 1.0 { + availCount++ + } + + nt := new(weightedNode) + nt.Carry = s.carry[node.GetID()] + nt.Weight = float64(node.GetTotal()-node.GetUsed()) / float64(maxTotal) + nt.Ptr = node + nodeTabs = append(nodeTabs, nt) + return true + }) + return nodeTabs, availCount } func (s *CarryWeightNodeSelector) setNodeCarry(nodes SortedWeightedNodes, availCarryCount, replicaNum int) { @@ -322,16 +254,7 @@ type AvailableSpaceFirstNodeSelector struct { } func (s *AvailableSpaceFirstNodeSelector) getNodeAvailableSpace(node interface{}) uint64 { - switch s.nodeType { - case DataNodeType: - dataNode := node.(*DataNode) - return dataNode.AvailableSpace - case MetaNodeType: - metaNode := node.(*MetaNode) - return metaNode.Total - metaNode.Used - default: - panic("unkown node type") - } + return node.(Node).GetAvailableSpace() } func (s *AvailableSpaceFirstNodeSelector) GetName() string { @@ -349,7 +272,7 @@ func (s *AvailableSpaceFirstNodeSelector) Select(ns *nodeSet, excludeHosts []str nodes := ns.getNodes(s.nodeType) sortedNodes := make([]Node, 0) nodes.Range(func(key, value interface{}) bool { - sortedNodes = append(sortedNodes, asNodeWrap(value, s.nodeType)) + sortedNodes = append(sortedNodes, value.(Node)) return true }) // if we cannot get enough nodes, return error @@ -370,7 +293,7 @@ func (s *AvailableSpaceFirstNodeSelector) Select(ns *nodeSet, excludeHosts []str for nodeIndex < len(sortedNodes) { node := sortedNodes[nodeIndex] nodeIndex += 1 - if canAllocPartition(node, s.nodeType) { + if canAllocPartition(node) { if excludeHosts == nil || !contains(excludeHosts, node.GetAddr()) { selectedIndex = nodeIndex - 1 break @@ -428,7 +351,7 @@ func (s *RoundRobinNodeSelector) Select(ns *nodeSet, excludeHosts []string, repl nodes := ns.getNodes(s.nodeType) sortedNodes := make([]Node, 0) nodes.Range(func(key, value interface{}) bool { - sortedNodes = append(sortedNodes, asNodeWrap(value, s.nodeType)) + sortedNodes = append(sortedNodes, value.(Node)) return true }) // if we cannot get enough nodes, return error @@ -449,7 +372,7 @@ func (s *RoundRobinNodeSelector) Select(ns *nodeSet, excludeHosts []string, repl for nodeIndex < len(sortedNodes) { node := sortedNodes[(nodeIndex+s.index)%len(sortedNodes)] nodeIndex += 1 - if canAllocPartition(node, s.nodeType) { + if canAllocPartition(node) { if excludeHosts == nil || !contains(excludeHosts, node.GetAddr()) { selectedIndex = nodeIndex - 1 break @@ -503,16 +426,7 @@ func (s *StrawNodeSelector) GetName() string { } func (s *StrawNodeSelector) getWeight(node Node) float64 { - switch s.nodeType { - case DataNodeType: - dataNode := node.(*DataNode) - return float64(dataNode.AvailableSpace) / util.GB - case MetaNodeType: - metaNode := node.(*MetaNode) - return float64(metaNode.Total-metaNode.Used) / util.GB - default: - panic("unkown node type") - } + return float64(node.GetAvailableSpace()) / util.GB } func (s *StrawNodeSelector) selectOneNode(nodes []Node) (index int, maxNode Node) { @@ -549,7 +463,7 @@ func (s *StrawNodeSelector) Select(ns *nodeSet, excludeHosts []string, replicaNu nodes[0], nodes[index] = node, nodes[0] } nodes = nodes[1:] - if !canAllocPartition(node, s.nodeType) { + if !canAllocPartition(node) { continue } orderHosts = append(orderHosts, node.GetAddr()) diff --git a/master/nodeset_selector.go b/master/nodeset_selector.go index e199eef89..acc3d298d 100644 --- a/master/nodeset_selector.go +++ b/master/nodeset_selector.go @@ -77,9 +77,9 @@ func (ns *nodeSet) getMetaNodeTotalAvailableSpace() (space uint64) { func (ns *nodeSet) canWriteFor(nodeType NodeType, replica int) bool { switch nodeType { case DataNodeType: - return ns.canWriteForDataNode(replica) + return ns.canWriteForNode(ns.dataNodes, replica) case MetaNodeType: - return ns.canWriteForMetaNode(replica) + return ns.canWriteForNode(ns.metaNodes, replica) default: panic("unknow node type") } diff --git a/master/topology.go b/master/topology.go index 397db2999..e63c8181f 100644 --- a/master/topology.go +++ b/master/topology.go @@ -24,32 +24,46 @@ import ( "time" "github.com/cubefs/cubefs/proto" - "github.com/cubefs/cubefs/util" "github.com/cubefs/cubefs/util/auditlog" "github.com/cubefs/cubefs/util/errors" "github.com/cubefs/cubefs/util/log" ) +type rsManager struct { + nodeType NodeType + nodes *sync.Map + zoneIndexForNode int + zones []*Zone +} + +func (rsm *rsManager) clear() { + rsm.nodes = new(sync.Map) +} + type topology struct { - dataNodes *sync.Map - metaNodes *sync.Map - zoneMap *sync.Map - zoneIndexForDataNode int - zoneIndexForMetaNode int - zones []*Zone - domainExcludeZones []string // not domain zone, empty if domain disable. - zoneLock sync.RWMutex + zoneMap *sync.Map + zones []*Zone + domainExcludeZones []string // not domain zone, empty if domain disable. + metaTopology rsManager + dataTopology rsManager + zoneLock sync.RWMutex } func newTopology() (t *topology) { t = new(topology) t.zoneMap = new(sync.Map) - t.dataNodes = new(sync.Map) - t.metaNodes = new(sync.Map) + t.dataTopology.nodes = new(sync.Map) + t.metaTopology.nodes = new(sync.Map) t.zones = make([]*Zone, 0) return } +func (t *topology) getZoneLen() int { + t.zoneLock.RLock() + defer t.zoneLock.RUnlock() + return len(t.zones) +} + func (t *topology) zoneLen() int { t.zoneLock.RLock() defer t.zoneLock.RUnlock() @@ -57,14 +71,8 @@ func (t *topology) zoneLen() int { } func (t *topology) clear() { - t.dataNodes.Range(func(key, value interface{}) bool { - t.dataNodes.Delete(key) - return true - }) - t.metaNodes.Range(func(key, value interface{}) bool { - t.metaNodes.Delete(key) - return true - }) + t.metaTopology.clear() + t.dataTopology.clear() } func (t *topology) putZone(zone *Zone) (err error) { @@ -115,7 +123,7 @@ func (t *topology) getZone(name string) (zone *Zone, err error) { } func (t *topology) putDataNode(dataNode *DataNode) (err error) { - if _, ok := t.dataNodes.Load(dataNode.Addr); ok { + if _, ok := t.dataTopology.nodes.Load(dataNode.Addr); ok { return } zone, err := t.getZone(dataNode.ZoneName) @@ -129,7 +137,7 @@ func (t *topology) putDataNode(dataNode *DataNode) (err error) { } func (t *topology) putDataNodeToCache(dataNode *DataNode) { - t.dataNodes.Store(dataNode.Addr, dataNode) + t.dataTopology.nodes.Store(dataNode.Addr, dataNode) } func (t *topology) deleteDataNode(dataNode *DataNode) { @@ -138,11 +146,20 @@ func (t *topology) deleteDataNode(dataNode *DataNode) { return } zone.deleteDataNode(dataNode) - t.dataNodes.Delete(dataNode.Addr) + t.dataTopology.nodes.Delete(dataNode.Addr) +} + +func (t *topology) getZoneByDataNode(dataNode *DataNode) (zone *Zone, err error) { + _, ok := t.dataTopology.nodes.Load(dataNode.Addr) + if !ok { + return nil, errors.Trace(dataNodeNotFound(dataNode.Addr), "%v not found", dataNode.Addr) + } + + return t.getZone(dataNode.ZoneName) } func (t *topology) putMetaNode(metaNode *MetaNode) (err error) { - if _, ok := t.metaNodes.Load(metaNode.Addr); ok { + if _, ok := t.metaTopology.nodes.Load(metaNode.Addr); ok { return } zone, err := t.getZone(metaNode.ZoneName) @@ -155,7 +172,7 @@ func (t *topology) putMetaNode(metaNode *MetaNode) (err error) { } func (t *topology) deleteMetaNode(metaNode *MetaNode) { - t.metaNodes.Delete(metaNode.Addr) + t.metaTopology.nodes.Delete(metaNode.Addr) zone, err := t.getZone(metaNode.ZoneName) if err != nil { return @@ -164,7 +181,7 @@ func (t *topology) deleteMetaNode(metaNode *MetaNode) { } func (t *topology) putMetaNodeToCache(metaNode *MetaNode) { - t.metaNodes.Store(metaNode.Addr, metaNode) + t.metaTopology.nodes.Store(metaNode.Addr, metaNode) } type nodeSetCollection []*nodeSet @@ -305,13 +322,13 @@ func (nsgm *DomainManager) checkExcludeZoneState() { if zone.status == normalZone { log.LogInfof("action[checkExcludeZoneState] zone[%v] be set unavailableZone", zone.name) } - zone.status = unavailableZone + zone.setStatus(unavailableZone) } else { excludeNeedDomain = false if zone.status == unavailableZone { log.LogInfof("action[checkExcludeZoneState] zone[%v] be set normalZone", zone.name) } - zone.status = normalZone + zone.setStatus(normalZone) } } } @@ -359,7 +376,7 @@ func (nsgm *DomainManager) checkGrpState(domainGrpManager *DomainNodeSetGrpManag domainGrpManager.nodeSetGrpMap[i].nodeSets[j].dataNodes.Range(func(key, value interface{}) bool { node := value.(*DataNode) - if node.isWriteAble() { + if node.IsWriteAble() { used = used + node.Used } else { used = used + node.Total @@ -376,7 +393,7 @@ func (nsgm *DomainManager) checkGrpState(domainGrpManager *DomainNodeSetGrpManag } domainGrpManager.nodeSetGrpMap[i].nodeSets[j].metaNodes.Range(func(key, value interface{}) bool { node := value.(*MetaNode) - if node.isWritable() { + if node.IsWriteAble() { metaWorked = true log.LogInfof("action[checkGrpState] nodeset[%v] zonename[%v] used [%v] total [%v] threshold [%v] got available metanode", node.ID, node.ZoneName, node.Used, node.Total, node.Threshold) @@ -1032,39 +1049,11 @@ func (ns *nodeSet) deleteMetaNode(metaNode *MetaNode) { ns.metaNodes.Delete(metaNode.Addr) } -func (ns *nodeSet) canWriteForDataNode(replicaNum int) bool { +func (ns *nodeSet) canWriteForNode(nodes *sync.Map, replicaNum int) bool { var count int - ns.dataNodes.Range(func(key, value interface{}) bool { - node := value.(*DataNode) - if node.isWriteAble() && node.dpCntInLimit() { - count++ - } - if count >= replicaNum { - return false - } - return true - }) - log.LogInfof("canWriteForDataNode zone[%v], ns[%v],count[%v], replicaNum[%v]", - ns.zoneName, ns.ID, count, replicaNum) - return count >= replicaNum -} - -func (ns *nodeSet) canAllocDataNodeCnt() (cnt int) { - ns.dataNodes.Range(func(key, value interface{}) bool { - node := value.(*DataNode) - if node.isWriteAble() && node.dpCntInLimit() { - cnt++ - } - return true - }) - return -} - -func (ns *nodeSet) canWriteForMetaNode(replicaNum int) bool { - var count int - ns.metaNodes.Range(func(key, value interface{}) bool { - node := value.(*MetaNode) - if node.isWritable() && node.mpCntInLimit() { + nodes.Range(func(key, value interface{}) bool { + node := value.(Node) + if node.IsWriteAble() && node.PartitionCntLimited() { count++ } if count >= replicaNum { @@ -1077,10 +1066,10 @@ func (ns *nodeSet) canWriteForMetaNode(replicaNum int) bool { return count >= replicaNum } -func (ns *nodeSet) canAllocMetaNodeCnt() (cnt int) { - ns.metaNodes.Range(func(key, value interface{}) bool { - node := value.(*MetaNode) - if node.isWritable() && node.mpCntInLimit() { +func (ns *nodeSet) calcNodesForAlloc(nodes *sync.Map) (cnt int) { + nodes.Range(func(key, value interface{}) bool { + node := value.(Node) + if node.IsWriteAble() && node.PartitionCntLimited() { cnt++ } return true @@ -1286,23 +1275,30 @@ func (t *topology) getNodeSetByNodeSetId(nodeSetId uint64) (nodeSet *nodeSet, er return nil, errors.NewErrorf("set %v not found", nodeSetId) } -func calculateDemandWriteNodes(zoneNum int, replicaNum int) (demandWriteNodes int) { +func calculateDemandWriteNodes(zoneNum int, replicaNum int, isSpecialZoneName bool) (demandWriteNodesCntPerZone int) { + if isSpecialZoneName { + return 1 + } if zoneNum == 1 { - demandWriteNodes = replicaNum + demandWriteNodesCntPerZone = replicaNum } else { if replicaNum == 1 { - demandWriteNodes = 1 + demandWriteNodesCntPerZone = 1 } else { - demandWriteNodes = 2 + demandWriteNodesCntPerZone = 2 } } return } -func (t *topology) allocZonesForMetaNode(zoneNum, replicaNum int, excludeZone []string) (zones []*Zone, err error) { +// Choose the zone if it is writable and adapt to the rules for classifying zones +func (t *topology) allocZonesForNode(rsMgr *rsManager, zoneNumNeed, replicaNum int, excludeZone []string, specialZones []*Zone) (zones []*Zone, err error) { if len(t.domainExcludeZones) > 0 { zones = t.getDomainExcludeZones() log.LogInfof("action[allocZonesForMetaNode] getDomainExcludeZones zones [%v]", t.domainExcludeZones) + } else if specialZones != nil && len(specialZones) > 0 { + zones = specialZones + zoneNumNeed = len(specialZones) } else { // if domain enable, will not enter here zones = t.getAllZones() @@ -1314,31 +1310,25 @@ func (t *topology) allocZonesForMetaNode(zoneNum, replicaNum int, excludeZone [] excludeZone = make([]string, 0) } candidateZones := make([]*Zone, 0) - demandWriteNodes := calculateDemandWriteNodes(zoneNum, replicaNum) - for i := 0; i < len(zones); i++ { - if t.zoneIndexForMetaNode >= len(zones) { - t.zoneIndexForMetaNode = 0 - } - zone := zones[t.zoneIndexForMetaNode] - t.zoneIndexForMetaNode++ + demandWriteNodesCntPerZone := calculateDemandWriteNodes(zoneNumNeed, replicaNum, len(specialZones) > 1) + + for _, zone := range zones { if zone.status == unavailableZone { continue } if contains(excludeZone, zone.name) { continue } - if zone.canWriteForMetaNode(uint8(demandWriteNodes)) { + + if zone.canWriteForNode(rsMgr.nodes, uint8(demandWriteNodesCntPerZone)) { candidateZones = append(candidateZones, zone) } - if len(candidateZones) >= zoneNum { - break - } } // if across zone,candidateZones must be larger than or equal with 2,otherwise,must have a candidate zone - if (zoneNum >= 2 && len(candidateZones) < 2) || len(candidateZones) < 1 { - log.LogError(fmt.Sprintf("action[allocZonesForMetaNode],reqZoneNum[%v],candidateZones[%v],demandWriteNodes[%v],err:%v", - zoneNum, len(candidateZones), demandWriteNodes, proto.ErrNoZoneToCreateMetaPartition)) + if (zoneNumNeed >= 2 && len(candidateZones) < 2) || len(candidateZones) < 1 { + log.LogError(fmt.Sprintf("action[allocZonesForqa n Node],reqZoneNum[%v],candidateZones[%v],demandWriteNodes[%v],err:%v", + zoneNumNeed, len(candidateZones), demandWriteNodesCntPerZone, proto.ErrNoZoneToCreateMetaPartition)) return nil, proto.ErrNoZoneToCreateMetaPartition } zones = candidateZones @@ -1346,57 +1336,13 @@ func (t *topology) allocZonesForMetaNode(zoneNum, replicaNum int, excludeZone [] return } -func (t *topology) allocZonesForDataNode(zoneNum, replicaNum int, excludeZone []string) (zones []*Zone, err error) { - // domain enabled and have old zones to be used - if len(t.domainExcludeZones) > 0 { - zones = t.getDomainExcludeZones() - } else { - // if domain enable, will not enter here - zones = t.getAllZones() - } - - log.LogInfof("len(zones) = %v \n", len(zones)) - if t.isSingleZone() { - return zones, nil - } - if excludeZone == nil { - excludeZone = make([]string, 0) - } - - demandWriteNodes := calculateDemandWriteNodes(zoneNum, replicaNum) - candidateZones := make([]*Zone, 0) - - for i := 0; i < len(zones); i++ { - if t.zoneIndexForDataNode >= len(zones) { - t.zoneIndexForDataNode = 0 - } - - zone := zones[t.zoneIndexForDataNode] - t.zoneIndexForDataNode++ - - if zone.status == unavailableZone { - continue - } - if contains(excludeZone, zone.name) { - continue - } - if zone.canWriteForDataNode(uint8(demandWriteNodes)) { - candidateZones = append(candidateZones, zone) - } - if len(candidateZones) >= zoneNum { - break - } - } - - // if across zone,candidateZones must be larger than or equal with 2,otherwise,must have one candidate zone - if (zoneNum >= 2 && len(candidateZones) < 2) || len(candidateZones) < 1 { - log.LogError(fmt.Sprintf("action[allocZonesForDataNode],reqZoneNum[%v],candidateZones[%v],demandWriteNodes[%v],err:%v", - zoneNum, len(candidateZones), demandWriteNodes, proto.ErrNoZoneToCreateDataPartition)) - return nil, errors.NewError(proto.ErrNoZoneToCreateDataPartition) - } - zones = candidateZones - err = nil - return +func (ns *nodeSet) dataNodeCount() int { + var count int + ns.dataNodes.Range(func(key, value interface{}) bool { + count++ + return true + }) + return count } // Zone stores all the zone related information @@ -1430,7 +1376,7 @@ type zoneValue struct { func newZone(name string) (zone *Zone) { zone = &Zone{name: name} - zone.status = normalZone + zone.setStatus(normalZone) zone.dataNodes = new(sync.Map) zone.metaNodes = new(sync.Map) zone.nodeSetMap = make(map[uint64]*nodeSet) @@ -1705,16 +1651,16 @@ func (zone *Zone) allocNodeSetForMetaNode(excludeNodeSets []uint64, replicaNum u return ns, nil } -func (zone *Zone) canWriteForDataNode(replicaNum uint8) (can bool) { +func (zone *Zone) canWriteForNode(nodes *sync.Map, replicaNum uint8) (can bool) { zone.RLock() defer zone.RUnlock() var leastAlive uint8 - zone.dataNodes.Range(func(addr, value interface{}) bool { - dataNode := value.(*DataNode) - if !dataNode.dpCntInLimit() { + nodes.Range(func(addr, value interface{}) bool { + node := value.(Node) + if !node.PartitionCntLimited() { return true } - if dataNode.isActive && dataNode.isWriteAbleWithSize(30*util.GB) { + if node.IsActiveNode() && node.IsWriteAble() { leastAlive++ } if leastAlive >= replicaNum { @@ -1723,8 +1669,6 @@ func (zone *Zone) canWriteForDataNode(replicaNum uint8) (can bool) { } return true }) - log.LogInfof("action[canWriteForDataNode] zone name %v canWriteForDataNode leastAlive[%v],replicaNum[%v],count[%v]", - zone.name, leastAlive, replicaNum, zone.dataNodeCount()) return } @@ -1772,69 +1716,32 @@ func (zone *Zone) isUsedRatio(ratio float64) (can bool) { return false } -func (zone *Zone) getDataUsed() (dataNodeUsed uint64, dataNodeTotal uint64) { +func (zone *Zone) getUsed(dataType uint32) (dataNodeUsed uint64, dataNodeTotal uint64) { zone.RLock() defer zone.RUnlock() - zone.dataNodes.Range(func(addr, value interface{}) bool { - dataNode := value.(*DataNode) - if dataNode.isActive { - dataNodeUsed += dataNode.Used + var nodes *sync.Map + if dataType == uint32(MetaNodeType) { + nodes = zone.metaNodes + } else { + nodes = zone.dataNodes + } + nodes.Range(func(addr, value interface{}) bool { + dataNode := value.(Node) + if dataNode.IsActiveNode() { + dataNodeUsed += dataNode.GetUsed() } else { - dataNodeUsed += dataNode.Total + dataNodeUsed += dataNode.GetTotal() } - dataNodeTotal += dataNode.Total + dataNodeTotal += dataNode.GetTotal() return true }) return dataNodeUsed, dataNodeTotal } -func (zone *Zone) getMetaUsed() (metaNodeUsed uint64, metaNodeTotal uint64) { - zone.RLock() - defer zone.RUnlock() - - zone.metaNodes.Range(func(addr, value interface{}) bool { - metaNode := value.(*MetaNode) - if metaNode.IsActive && metaNode.isWritable() { - metaNodeUsed += metaNode.Used - } else { - metaNodeUsed += metaNode.Total - } - metaNodeTotal += metaNode.Total - return true - }) - return metaNodeUsed, metaNodeTotal -} - func (zone *Zone) getSpaceLeft(dataType uint32) (spaceLeft uint64) { - if dataType == TypeDataPartition { - dataNodeUsed, dataNodeTotal := zone.getDataUsed() - return dataNodeTotal - dataNodeUsed - } else { - metaNodeUsed, metaNodeTotal := zone.getMetaUsed() - return metaNodeTotal - metaNodeUsed - } -} - -func (zone *Zone) canWriteForMetaNode(replicaNum uint8) (can bool) { - zone.RLock() - defer zone.RUnlock() - var leastAlive uint8 - zone.metaNodes.Range(func(addr, value interface{}) bool { - metaNode := value.(*MetaNode) - if !metaNode.mpCntInLimit() { - return true - } - if metaNode.IsActive && metaNode.isWritable() { - leastAlive++ - } - if leastAlive >= replicaNum { - can = true - return false - } - return true - }) - return + dataNodeUsed, dataNodeTotal := zone.getUsed(dataType) + return dataNodeTotal - dataNodeUsed } func (zone *Zone) getAvailNodeHosts(nodeType uint32, excludeNodeSets []uint64, excludeHosts []string, replicaNum int) (newHosts []string, peers []proto.Peer, err error) { diff --git a/master/topology_test.go b/master/topology_test.go index e6eaaf2f2..7fd95036d 100644 --- a/master/topology_test.go +++ b/master/topology_test.go @@ -43,12 +43,12 @@ func TestSingleZone(t *testing.T) { // single zone exclude,if it is a single zone excludeZones don't take effect excludeZones := make([]string, 0) excludeZones = append(excludeZones, zoneName) - zones, err := topo.allocZonesForDataNode(replicaNum, replicaNum, excludeZones) + zones, err := topo.allocZonesForNode(&topo.metaTopology, replicaNum, replicaNum, excludeZones, []*Zone{}) require.NoError(t, err) require.EqualValues(t, 1, len(zones)) // single zone normal - zones, err = topo.allocZonesForDataNode(replicaNum, replicaNum, nil) + zones, err = topo.allocZonesForNode(&topo.dataTopology, replicaNum, replicaNum, nil, []*Zone{}) require.NoError(t, err) newHosts, _, err := zones[0].getAvailNodeHosts(TypeDataPartition, nil, nil, replicaNum) require.NoError(t, err) @@ -60,6 +60,36 @@ func TestAllocZones(t *testing.T) { topo := newTopology() c := new(Cluster) zoneCount := 3 + + hostZoneMap := make(map[string]string) + hostZoneMap[mds1Addr] = testZone1 + hostZoneMap[mds2Addr] = testZone1 + hostZoneMap[mds3Addr] = testZone2 + hostZoneMap[mds4Addr] = testZone2 + hostZoneMap[mds5Addr] = testZone3 + + zoneMap := make(map[string]bool) + zoneMap[testZone1] = false + zoneMap[testZone2] = false + zoneMap[testZone3] = false + + getZoneCntFunc := func(hosts []string) int { + for _, host := range hosts { + zoneNm := hostZoneMap[host] + zoneMap[zoneNm] = true + } + var zoneCnt int + for _, v := range zoneMap { + if v { + zoneCnt++ + } + } + for k := range zoneMap { + zoneMap[k] = false + } + return zoneCnt + } + // add three zones zoneName1 := testZone1 zone1 := newZone(zoneName1) @@ -91,9 +121,9 @@ func TestAllocZones(t *testing.T) { require.EqualValues(t, zoneCount, len(zones)) // only pass replica num replicaNum := 2 - zones, err := topo.allocZonesForDataNode(replicaNum, replicaNum, nil) + zones, err := topo.allocZonesForNode(&topo.dataTopology, replicaNum, replicaNum, nil, []*Zone{}) require.NoError(t, err) - require.EqualValues(t, replicaNum, len(zones)) + require.EqualValues(t, len(topo.getAllZones()), len(zones)) cluster := new(Cluster) cluster.t = topo @@ -109,12 +139,17 @@ func TestAllocZones(t *testing.T) { hosts, _, err = cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, replicaNum, 2, "") require.NoError(t, err) + // specific zone + hosts, _, err = cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, 3, 2, zoneName1+","+zoneName2) + require.NoError(t, err) + require.EqualValues(t, getZoneCntFunc(hosts), 2) + t.Logf("ChooseTargetDataHosts in multi zones,hosts[%v]", hosts) // after excluding zone3, alloc zones will be success excludeZones := make([]string, 0) excludeZones = append(excludeZones, zoneName3) - zones, err = topo.allocZonesForDataNode(2, replicaNum, excludeZones) + zones, err = topo.allocZonesForNode(&topo.dataTopology, 2, replicaNum, excludeZones, []*Zone{}) if err != nil { t.Logf("allocZonesForDataNode failed,err[%v]", err) } diff --git a/master/vol.go b/master/vol.go index 66c5e7d45..467808ac3 100644 --- a/master/vol.go +++ b/master/vol.go @@ -800,7 +800,7 @@ func (mp *MetaPartition) memUsedReachThreshold(clusterName, volName string) bool if !foundReadonlyReplica || readonlyReplica == nil { return false } - if readonlyReplica.metaNode.isWritable() { + if readonlyReplica.metaNode.IsWriteAble() { msg := fmt.Sprintf("action[checkSplitMetaPartition] vol[%v],max meta parition[%v] status is readonly\n", volName, mp.PartitionID) Warn(clusterName, msg) @@ -1483,8 +1483,7 @@ func (vol *Vol) doCreateMetaPartition(c *Cluster, start, end uint64) (mp *MetaPa } } else { var excludeZone []string - zoneNum := c.decideZoneNum(vol.crossZone) - + zoneNum := c.decideZoneNum(vol) if hosts, peers, err = c.getHostFromNormalZone(TypeMetaPartition, excludeZone, nil, nil, int(vol.mpReplicaNum), zoneNum, vol.zoneName); err != nil { log.LogErrorf("action[doCreateMetaPartition] getHostFromNormalZone err[%v]", err) return nil, errors.NewError(err)