feat(master): consider rack info when create partition. #1000327851

Signed-off-by: Victor1319 <zengxuewei@oppo.com>
This commit is contained in:
Victor1319 2025-09-08 17:01:13 +08:00
parent 8a41591165
commit 1a4b6f53c2
11 changed files with 2274 additions and 104 deletions

View File

@ -4259,7 +4259,12 @@ func (m *Server) updateNodesetCapcity(zoneName string, nodesetId uint64, capcity
return
}
ns.Capacity = int(capcity)
ns.racksLock.Lock()
for _, rack := range ns.racks {
rack.Capacity = int(capcity) / 3
}
ns.racksLock.Unlock()
m.cluster.syncUpdateNodeSet(ns)
log.LogInfof("action[updateNodesetCapcity] zonename %v nodeset %v capcity %v", zoneName, nodesetId, capcity)
return

View File

@ -27,6 +27,7 @@ import (
"sync/atomic"
"time"
pt "github.com/cubefs/cubefs/proto"
"github.com/cubefs/cubefs/util/auditlog"
"github.com/google/uuid"
"golang.org/x/time/rate"
@ -521,6 +522,33 @@ func (c *Cluster) masterAddr() (addr string) {
return c.leaderInfo.addr
}
func (c *Cluster) getRackAwareLevel() pt.RackAwareLevel {
return c.cfg.RackAwareLevel
}
func (c *Cluster) GetExRacksByHosts(nodeType uint32, hosts []string) (exRacks []string) {
for _, host := range hosts {
if nodeType == TypeDataPartition {
node, err := c.dataNode(host)
if err != nil {
log.LogErrorf("getExRacksByHosts, node[%v] err[%v]", host, err.Error())
continue
}
exRacks = append(exRacks, node.Rack)
continue
}
node, err := c.metaNode(host)
if err != nil {
log.LogErrorf("getExRacksByHosts, node[%v] err[%v]", host, err.Error())
continue
}
exRacks = append(exRacks, node.Rack)
}
return
}
func (c *Cluster) tryToChangeLeaderByHost() error {
return c.partition.TryToLeader(1)
}
@ -2055,7 +2083,7 @@ func (c *Cluster) createDataPartition(volName string, mediaType uint32) (dp *Dat
} else {
zoneNum := c.decideZoneNum(vol, mediaType) // zoneNum scope [1,3]
if targetHosts, targetPeers, err = c.getHostFromNormalZone(TypeDataPartition, nil, nil, nil,
int(dpReplicaNum), zoneNum, zoneName, mediaType); err != nil {
int(dpReplicaNum), zoneNum, zoneName, mediaType, c.getRackAwareLevel()); err != nil {
goto errHandler
}
}
@ -2261,11 +2289,11 @@ func (c *Cluster) decideZoneNum(vol *Vol, mediaType uint32) (zoneNum int) {
return zoneNum
}
func (c *Cluster) chooseZone2Plus1(rsMgr *rsManager, zones []*Zone, excludeNodeSets []uint64, excludeHosts []string,
nodeType uint32, replicaNum int) (hosts []string, peers []proto.Peer, err error,
func (c *Cluster) chooseZone2Plus1(rsMgr *rsManager, zones []*Zone,
nodeType uint32, param *selectParam) (hosts []string, peers []proto.Peer, err error,
) {
if replicaNum < 2 || replicaNum > defaultReplicaNum {
return nil, nil, fmt.Errorf("action[chooseZone2Plus1] replicaNum [%v]", replicaNum)
if param.replicaNum < 2 || param.replicaNum > defaultReplicaNum {
return nil, nil, fmt.Errorf("action[chooseZone2Plus1] replicaNum [%v]", param.replicaNum)
}
zoneList := make([]*Zone, 2)
for i := range []int{1, 2} {
@ -2281,45 +2309,55 @@ func (c *Cluster) chooseZone2Plus1(rsMgr *rsManager, zones []*Zone, excludeNodeS
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))
paramCopy := param.copy()
paramCopy.excludeRacks = c.GetExRacksByHosts(nodeType, param.excludeHosts)
num := 1
for _, zone := range zoneList {
selectedHosts, selectedPeers, e := zone.getAvailNodeHosts(nodeType, excludeNodeSets, excludeHosts, num)
paramCopy.replicaNum = num
selectedHosts, selectedPeers, e := zone.getAvailNodeHosts(nodeType, paramCopy)
if e != nil {
log.LogErrorf("action[chooseZone2Plus1] getAvailNodeHosts error: [%v]", e)
log.LogErrorf("action[chooseZone2Plus1] getAvailNodeHosts param[%v] error: [%v]", paramCopy.String(), e)
return nil, nil, e
}
hosts = append(hosts, selectedHosts...)
peers = append(peers, selectedPeers...)
paramCopy.excludeHosts = append(paramCopy.excludeHosts, selectedHosts...)
paramCopy.excludeRacks = c.GetExRacksByHosts(nodeType, param.excludeHosts)
log.LogInfof("action[chooseZone2Plus1] zone [%v] left [%v] get hosts[%v]",
zone.name, zone.getSpaceLeft(nodeType), selectedHosts)
num = replicaNum - num
num = param.replicaNum - num
}
log.LogInfof("action[chooseZone2Plus1] finally get hosts[%v]", hosts)
return hosts, peers, nil
}
func (c *Cluster) chooseZoneNormal(zones []*Zone, excludeNodeSets []uint64, excludeHosts []string,
nodeType uint32, replicaNum int,
) (hosts []string, peers []proto.Peer, err error) {
log.LogInfof("action[chooseZoneNormal] zones[%s] nodeType[%d] replicaNum[%d]", printZonesName(zones), nodeType, replicaNum)
func (c *Cluster) chooseZoneNormal(zones []*Zone, nodeType uint32, param *selectParam) (hosts []string, peers []proto.Peer, err error) {
log.LogInfof("action[chooseZoneNormal] zones[%s] nodeType[%d] replicaNum[%d]", printZonesName(zones), nodeType, param.replicaNum)
c.zoneIdxMux.Lock()
defer c.zoneIdxMux.Unlock()
for i := 0; i < replicaNum; i++ {
paramCopy := param.copy()
paramCopy.excludeRacks = c.GetExRacksByHosts(nodeType, param.excludeHosts)
for i := 0; i < param.replicaNum; i++ {
// try all zone from zones list to get the available hosts
for j := 0; j < len(zones); j++ {
zone := zones[c.lastZoneIdxForNode]
c.lastZoneIdxForNode = (c.lastZoneIdxForNode + 1) % len(zones)
selectedHosts, selectedPeers, err := zone.getAvailNodeHosts(nodeType, excludeNodeSets, excludeHosts, 1)
paramCopy.replicaNum = 1
selectedHosts, selectedPeers, err := zone.getAvailNodeHosts(nodeType, paramCopy)
if err != nil {
// no zone available
if j == len(zones)-1 {
log.LogErrorf("action[chooseZoneNormal] error [%v]", err)
log.LogErrorf("action[chooseZoneNormal] param [%s] error [%v]", paramCopy.String(), err)
return nil, nil, err
}
continue
@ -2327,6 +2365,10 @@ func (c *Cluster) chooseZoneNormal(zones []*Zone, excludeNodeSets []uint64, excl
hosts = append(hosts, selectedHosts...)
peers = append(peers, selectedPeers...)
paramCopy.excludeHosts = append(paramCopy.excludeHosts, selectedHosts...)
paramCopy.excludeRacks = c.GetExRacksByHosts(nodeType, param.excludeHosts)
// if get the available hosts, choose zone for next replica
break
}
@ -2357,7 +2399,7 @@ func (c *Cluster) getSpecificZoneList(specifiedZone string) (zones []*Zone, err
func (c *Cluster) getHostFromNormalZone(nodeType uint32, excludeZones []string, excludeNodeSets []uint64,
excludeHosts []string, replicaNum int, zoneNumNeed int,
specifiedZoneName string, dataMediaType uint32) (hosts []string, peers []proto.Peer, err error,
specifiedZoneName string, dataMediaType uint32, rackLevel proto.RackAwareLevel) (hosts []string, peers []proto.Peer, err error,
) {
log.LogInfof("[getHostFromNormalZone] dataMediaType(%v) nodeType(%v) replicaNum(%v) zoneNumNeed(%v) specifiedZoneName(%v)",
proto.MediaTypeString(nodeType), nodeType, replicaNum, zoneNumNeed, specifiedZoneName)
@ -2387,9 +2429,17 @@ func (c *Cluster) getHostFromNormalZone(nodeType uint32, excludeZones []string,
}
}
param := &selectParam{
excludeNodeSets: excludeNodeSets,
replicaNum: replicaNum,
excludeHosts: excludeHosts,
rackLevel: rackLevel,
excludeRacks: c.GetExRacksByHosts(nodeType, excludeHosts),
}
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 {
if hosts, peers, err = zonesQualified[0].getAvailNodeHosts(nodeType, param); err != nil {
log.LogErrorf("action[getHostFromNormalZone] err[%v]", err)
return
}
@ -2399,13 +2449,14 @@ func (c *Cluster) getHostFromNormalZone(nodeType uint32, excludeZones []string,
if excludeHosts == nil {
excludeHosts = make([]string, 0)
}
// 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 || replicaNum == 1 {
if hosts, peers, err = c.chooseZoneNormal(zonesQualified, excludeNodeSets, excludeHosts, nodeType, replicaNum); err != nil {
if hosts, peers, err = c.chooseZoneNormal(zonesQualified, nodeType, param); err != nil {
return
}
} else {
if hosts, peers, err = c.chooseZone2Plus1(rsMgr, zonesQualified, excludeNodeSets, excludeHosts, nodeType, replicaNum); err != nil {
if hosts, peers, err = c.chooseZone2Plus1(rsMgr, zonesQualified, nodeType, param); err != nil {
return
}
}
@ -2988,6 +3039,14 @@ func (c *Cluster) migrateDataPartition(srcAddr, targetAddr string, dp *DataParti
replica, _ = dp.getReplica(srcAddr)
dp.RUnlock()
param := &selectParam{
excludeNodeSets: nil,
replicaNum: 1,
excludeHosts: dp.Hosts,
rackLevel: c.getRackAwareLevel(),
excludeRacks: c.GetExRacksByHosts(TypeDataPartition, dp.Hosts),
}
// delete if not normal data partition
if !proto.IsNormalDp(dp.PartitionType) {
c.vols[dp.VolName].deleteDataPartition(c, dp)
@ -3024,6 +3083,7 @@ func (c *Cluster) migrateDataPartition(srcAddr, targetAddr string, dp *DataParti
break
}
}
if err = c.checkMultipleReplicasOnSameMachine(finalHosts); err != nil {
goto errHandler
}
@ -3034,7 +3094,7 @@ func (c *Cluster) migrateDataPartition(srcAddr, targetAddr string, dp *DataParti
log.LogErrorf("[migrateDataPartition] check mediaType err: %v", err.Error())
goto errHandler
}
} else if targetHosts, _, err = ns.getAvailDataNodeHosts(dp.Hosts, 1); err != nil {
} else if targetHosts, _, err = ns.getAvailDataNodeHosts(param); err != nil {
if _, ok := c.vols[dp.VolName]; !ok {
log.LogWarnf("clusterID[%v] partitionID:%v on node:%v offline failed,PersistenceHosts:[%v]",
c.Name, dp.PartitionID, srcAddr, dp.Hosts)
@ -3047,10 +3107,17 @@ func (c *Cluster) migrateDataPartition(srcAddr, targetAddr string, dp *DataParti
}
// select data nodes from the other node set in same zone
excludeNodeSets = append(excludeNodeSets, ns.ID)
if targetHosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, excludeNodeSets, dp.Hosts, 1); err != nil {
param := &selectParam{
excludeNodeSets: excludeNodeSets,
replicaNum: 1,
excludeHosts: dp.Hosts,
rackLevel: c.getRackAwareLevel(),
excludeRacks: c.GetExRacksByHosts(TypeDataPartition, dp.Hosts),
}
if targetHosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, param); err != nil {
// select data nodes from the other zone
zones = dp.getLiveZones(srcAddr)
if targetHosts, _, err = c.getHostFromNormalZone(TypeDataPartition, zones, excludeNodeSets, dp.Hosts, 1, 1, "", dp.MediaType); err != nil {
if targetHosts, _, err = c.getHostFromNormalZone(TypeDataPartition, zones, excludeNodeSets, dp.Hosts, 1, 1, "", dp.MediaType, c.getRackAwareLevel()); err != nil {
goto errHandler
}
}

View File

@ -126,6 +126,14 @@ func (c *Cluster) migrateMetaPartition(srcAddr, targetAddr string, mp *MetaParti
}
mp.RUnlock()
param := &selectParam{
excludeNodeSets: nil,
replicaNum: 1,
excludeHosts: oldHosts,
rackLevel: c.getRackAwareLevel(),
excludeRacks: c.GetExRacksByHosts(TypeMetaPartition, oldHosts),
}
nodeType := TypeMetaPartition
if dstStoreMode == proto.StoreModeDef {
dstStoreMode, err = c.getMetaPartitionStoreMode(mp, srcAddr)
@ -157,7 +165,7 @@ func (c *Cluster) migrateMetaPartition(srcAddr, targetAddr string, mp *MetaParti
newPeers = []proto.Peer{{
Addr: targetAddr,
}}
} else if _, newPeers, err = ns.getAvailMetaNodeHosts(oldHosts, 1, dstStoreMode); err != nil {
} else if _, newPeers, err = ns.getAvailMetaNodeHosts(param, dstStoreMode); err != nil {
if _, ok := c.vols[mp.volName]; !ok {
log.LogWarnf("[migrateMetaPartition] clusterID[%v] partitionID:%v on node:[%v]",
c.Name, mp.PartitionID, mp.Hosts)
@ -170,7 +178,14 @@ func (c *Cluster) migrateMetaPartition(srcAddr, targetAddr string, mp *MetaParti
}
// choose a meta node in other node set in the same zone
excludeNodeSets = append(excludeNodeSets, ns.ID)
if _, newPeers, err = zone.getAvailNodeHosts(nodeType, excludeNodeSets, oldHosts, 1); err != nil {
param := &selectParam{
excludeNodeSets: excludeNodeSets,
replicaNum: 1,
excludeHosts: oldHosts,
rackLevel: c.getRackAwareLevel(),
excludeRacks: c.GetExRacksByHosts(nodeType, oldHosts),
}
if _, newPeers, err = zone.getAvailNodeHosts(nodeType, param); err != nil {
zones = mp.getLiveZones(srcAddr)
var excludeZone []string
if len(zones) == 0 {
@ -179,7 +194,7 @@ func (c *Cluster) migrateMetaPartition(srcAddr, targetAddr string, mp *MetaParti
excludeZone = append(excludeZone, zones[0])
}
// choose a meta node in other zone
if _, newPeers, err = c.getHostFromNormalZone(nodeType, excludeZone, excludeNodeSets, oldHosts, 1, 1, "", proto.MediaType_Unspecified); err != nil {
if _, newPeers, err = c.getHostFromNormalZone(nodeType, excludeZone, excludeNodeSets, oldHosts, 1, 1, "", proto.MediaType_Unspecified, c.getRackAwareLevel()); err != nil {
goto errHandler
}
}

View File

@ -2256,7 +2256,14 @@ func (partition *DataPartition) TryAcquireDecommissionToken(c *Cluster) bool {
log.LogDebugf("action[TryAcquireDecommissionToken]dp %v excludeHosts %v",
partition.PartitionID, excludeHosts)
// data nodes in a nodeset has the same mediaType
targetHosts, _, err = ns.getAvailDataNodeHosts(excludeHosts, 1)
param := &selectParam{
excludeNodeSets: nil,
replicaNum: 1,
excludeHosts: excludeHosts,
rackLevel: c.getRackAwareLevel(),
excludeRacks: c.GetExRacksByHosts(TypeDataPartition, excludeHosts),
}
targetHosts, _, err = ns.getAvailDataNodeHosts(param)
if err != nil {
if partition.DecommissionDstNodeSet != 0 {
log.LogWarnf("action[TryAcquireDecommissionToken] dp %v choose from given dst nodeset %v failed:%v",
@ -2278,13 +2285,20 @@ func (partition *DataPartition) TryAcquireDecommissionToken(c *Cluster) bool {
goto errHandler
}
excludeNodeSets = append(excludeNodeSets, ns.ID)
param := &selectParam{
excludeNodeSets: excludeNodeSets,
replicaNum: 1,
excludeHosts: excludeHosts,
rackLevel: c.getRackAwareLevel(),
excludeRacks: c.GetExRacksByHosts(TypeDataPartition, excludeHosts),
}
// data nodes in a zone has the same mediaType
if targetHosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, excludeNodeSets, excludeHosts, 1); err != nil {
if targetHosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, param); err != nil {
log.LogWarnf("action[TryAcquireDecommissionToken] dp %v choose from other nodeset failed:%v",
partition.PartitionID, err.Error())
// select data nodes from the other zone
zones = partition.getLiveZones(partition.DecommissionSrcAddr)
if targetHosts, _, err = c.getHostFromNormalZone(TypeDataPartition, zones, excludeNodeSets, excludeHosts, 1, 1, "", partition.MediaType); err != nil {
if targetHosts, _, err = c.getHostFromNormalZone(TypeDataPartition, zones, excludeNodeSets, excludeHosts, 1, 1, "", partition.MediaType, c.getRackAwareLevel()); err != nil {
log.LogWarnf("action[TryAcquireDecommissionToken] dp %v choose from other zone failed:%v",
partition.PartitionID, err.Error())
goto errHandler

View File

@ -120,6 +120,7 @@ type CarryWeightNodeSelector struct {
nodeType NodeType
carry map[uint64]float64
sync.RWMutex
}
func (s *CarryWeightNodeSelector) GetName() string {
@ -176,12 +177,17 @@ func (s *CarryWeightNodeSelector) getCarryNodes(nset *nodeSet, maxTotal uint64,
node.GetStorageInfo(), excludeHosts)
return true
}
s.RLock()
if s.carry[node.GetID()] >= 1.0 {
availCount++
}
s.RUnlock()
nt := new(weightedNode)
s.RLock()
nt.Carry = s.carry[node.GetID()]
s.RUnlock()
nt.Weight = float64(node.GetTotal()-node.GetUsed()) / float64(maxTotal)
nt.Ptr = node
nodeTabs = append(nodeTabs, nt)
@ -201,7 +207,9 @@ func (s *CarryWeightNodeSelector) setNodeCarry(nodes SortedWeightedNodes, availC
carry = 10.0
}
nt.Carry = carry
s.Lock()
s.carry[nt.Ptr.GetID()] = carry
s.Unlock()
if carry > 1.0 {
availCarryCount++
}
@ -212,7 +220,9 @@ func (s *CarryWeightNodeSelector) setNodeCarry(nodes SortedWeightedNodes, availC
func (s *CarryWeightNodeSelector) selectNodeForWrite(node Node) {
node.SelectNodeForWrite()
// decrease node weight
s.Lock()
s.carry[node.GetID()] -= 1.0
s.Unlock()
}
func (s *CarryWeightNodeSelector) Select(ns *nodeSet, excludeHosts []string, replicaNum int) (newHosts []string, peers []proto.Peer, err error) {
@ -612,34 +622,106 @@ func NewNodeSelector(name string, nodeType NodeType) NodeSelector {
}
}
func (ns *nodeSet) getAvailMetaNodeHosts(excludeHosts []string, replicaNum int, storeMode proto.StoreMode) (newHosts []string, peers []proto.Peer, err error) {
func (ns *nodeSet) getRackSets() nodeSetCollection {
ns.racksLock.RLock()
defer ns.racksLock.RUnlock()
rsets := make(nodeSetCollection, 0, len(ns.racks))
for _, rack := range ns.racks {
rsets = append(rsets, rack)
}
return rsets
}
func (ns *nodeSet) getAvailMetaNodeHosts(param *selectParam, storeMode proto.StoreMode) (newHosts []string, peers []proto.Peer, err error) {
ns.nodeSelectLock.Lock()
defer ns.nodeSelectLock.Unlock()
// we need a read lock to block the modification of node selector
ns.metaNodeSelectorLock.RLock()
defer ns.metaNodeSelectorLock.RUnlock()
if storeMode == proto.StoreModeRocksDb {
return ns.metaNodeRocksdbSelector.Select(ns, excludeHosts, replicaNum)
// If rack isolation is not enabled, use non-rack-aware selector directly
if param.rackLevel == proto.RackAwareNone {
return ns.getNodeSelector(MetaNodeType, storeMode).Select(ns, param.excludeHosts, param.replicaNum)
}
return ns.metaNodeMemorySelector.Select(ns, excludeHosts, replicaNum)
// Rack isolation enabled, prioritize strong constraint mode
return ns.selectNodesWithRack(param, MetaNodeType, storeMode)
}
func (ns *nodeSet) getAvailDataNodeHosts(excludeHosts []string, replicaNum int) (hosts []string, peers []proto.Peer, err error) {
// selectMetaNodesWithRack selects meta nodes with rack awareness
func (ns *nodeSet) selectNodesWithRack(param *selectParam, nodeType NodeType, storeMode proto.StoreMode) (newHosts []string, peers []proto.Peer, err error) {
rsets := ns.getRackSets()
paramCopy := param.copy()
// First attempt with strong rack awareness
paramCopy.rackLevel = proto.RackAwareStrong
paramCopy.replicaNum = 1
for {
rack, err := ns.getRackSelector(nodeType).Select(rsets, paramCopy)
if err != nil {
// if rack aware is not enabled or alreay weak aware, return error
if param.rackLevel == proto.RackAwareStrong || paramCopy.rackLevel == proto.RackAwareWeak {
return nil, nil, fmt.Errorf("strong rack aware selection failed for rack[%v], param[%v], err: %v",
rack, paramCopy.String(), err)
}
log.LogWarnf("action[getAvailMetaNodeHosts] weak rack aware selection failed for rack[%v], param %v, err: %v",
rack, paramCopy.String(), err.Error())
paramCopy.rackLevel = proto.RackAwareWeak
continue
}
// Select nodes
selector := rack.getNodeSelector(nodeType, storeMode)
rhosts, rpeers, err := selector.Select(rack, paramCopy.excludeHosts, 1)
if err != nil {
log.LogErrorf("action[getAvailMetaNodeHosts] node selection failed for rack[%v], param[%v], err: %v",
rack, paramCopy.String(), err.Error())
return nil, nil, fmt.Errorf("node selection failed for rack[%v], param[%v], err: %v",
rack.Rack, paramCopy.String(), err.Error())
}
// Update results
newHosts = append(newHosts, rhosts...)
peers = append(peers, rpeers...)
paramCopy.excludeHosts = append(paramCopy.excludeHosts, rhosts...)
paramCopy.excludeRacks = append(paramCopy.excludeRacks, rack.Rack)
// Check if replica number requirement is met
if len(newHosts) >= param.replicaNum {
return newHosts, peers, nil
}
}
}
func (ns *nodeSet) getAvailDataNodeHosts(param *selectParam) (hosts []string, peers []proto.Peer, err error) {
ns.nodeSelectLock.Lock()
defer ns.nodeSelectLock.Unlock()
// we need a read lock to block the modification of node selector
ns.dataNodeSelectorLock.Lock()
defer ns.dataNodeSelectorLock.Unlock()
return ns.dataNodeSelector.Select(ns, excludeHosts, replicaNum)
if param.rackLevel == proto.RackAwareNone {
return ns.dataNodeSelector.Select(ns, param.excludeHosts, param.replicaNum)
}
return ns.selectNodesWithRack(param, DataNodeType, proto.StoreModeDef)
}
func (s *CarryWeightNodeSelector) prepareCarryForDataNodes(nodes *sync.Map, total uint64) {
nodes.Range(func(key, value interface{}) bool {
node := value.(Node)
s.Lock()
if _, ok := s.carry[node.GetID()]; !ok {
// use available space to calculate initial weight
s.carry[node.GetID()] = float64(node.GetAvailableSpace()) / float64(total)
}
s.Unlock()
return true
})
}
@ -647,10 +729,12 @@ func (s *CarryWeightNodeSelector) prepareCarryForDataNodes(nodes *sync.Map, tota
func (s *CarryWeightNodeSelector) prepareCarryForMetaNodeMemory(nodes *sync.Map, total uint64) {
nodes.Range(func(key, value interface{}) bool {
metaNode := value.(*MetaNode)
s.Lock()
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)
}
s.Unlock()
return true
})
}
@ -658,10 +742,12 @@ func (s *CarryWeightNodeSelector) prepareCarryForMetaNodeMemory(nodes *sync.Map,
func (s *CarryWeightNodeSelector) prepareCarryForMetaNodeRocksdb(nodes *sync.Map, total uint64) {
nodes.Range(func(key, value interface{}) bool {
metaNode := value.(*MetaNode)
s.Lock()
if _, ok := s.carry[metaNode.ID]; !ok {
// use available space to calculate initial weight
s.carry[metaNode.ID] = float64(metaNode.GetRocksdbTotal()-metaNode.GetRocksdbUsed()) / float64(total)
}
s.Unlock()
return true
})
}

View File

@ -88,12 +88,12 @@ func (ns *nodeSet) getMetaNodeTotalAvailableSpace(nodeType NodeType) (space uint
return
}
func (ns *nodeSet) canWriteFor(nodeType NodeType, replica int) bool {
func (ns *nodeSet) canWriteFor(nodeType NodeType, param *selectParam) bool {
switch nodeType {
case DataNodeType:
return ns.canWriteForNode(ns.dataNodes, replica, nodeType)
return ns.canWriteForNode(ns.dataNodes, param, nodeType)
case MetaNodeType, RocksdbType:
return ns.canWriteForNode(ns.metaNodes, replica, nodeType)
return ns.canWriteForNode(ns.metaNodes, param, nodeType)
default:
panic("unknow node type")
}
@ -123,7 +123,7 @@ func (ns *nodeSet) getTotalAvailableSpaceOf(nodeType NodeType) uint64 {
type NodesetSelector interface {
GetName() string
Select(nsc nodeSetCollection, excludeNodeSets []uint64, replicaNum uint8) (ns *nodeSet, err error)
Select(nsc nodeSetCollection, param *selectParam) (ns *nodeSet, err error)
}
type RoundRobinNodesetSelector struct {
@ -132,7 +132,7 @@ type RoundRobinNodesetSelector struct {
nodeType NodeType
}
func (s *RoundRobinNodesetSelector) Select(nsc nodeSetCollection, excludeNodeSets []uint64, replicaNum uint8) (ns *nodeSet, err error) {
func (s *RoundRobinNodesetSelector) Select(nsc nodeSetCollection, param *selectParam) (ns *nodeSet, err error) {
// sort nodesets by id, so we can get a node list that is as stable as possible
sort.Slice(nsc, func(i, j int) bool {
return nsc[i].ID < nsc[j].ID
@ -145,10 +145,10 @@ func (s *RoundRobinNodesetSelector) Select(nsc nodeSetCollection, excludeNodeSet
ns = nsc[s.index]
s.index++
if containsID(excludeNodeSets, ns.ID) {
if containsID(param.excludeNodeSets, ns.ID) {
continue
}
if ns.canWriteFor(s.nodeType, int(replicaNum)) {
if ns.canWriteFor(s.nodeType, param) {
return
}
}
@ -205,11 +205,11 @@ func (s *CarryWeightNodesetSelector) prepareCarry(nsc nodeSetCollection, total u
}
}
func (s *CarryWeightNodesetSelector) getAvailNodesets(nsc nodeSetCollection, excludeNodeSets []uint64, replicaNum uint8) (newNsc nodeSetCollection) {
func (s *CarryWeightNodesetSelector) getAvailNodesets(nsc nodeSetCollection, param *selectParam) (newNsc nodeSetCollection) {
newNsc = make(nodeSetCollection, 0, nsc.Len())
for i := 0; i < nsc.Len(); i++ {
ns := nsc[i]
if ns.canWriteFor(s.nodeType, int(replicaNum)) && !containsID(excludeNodeSets, ns.ID) {
if ns.canWriteFor(s.nodeType, param) && !containsID(param.excludeNodeSets, ns.ID) {
newNsc = append(newNsc, ns)
}
}
@ -246,11 +246,11 @@ func (s *CarryWeightNodesetSelector) setNodesetCarry(nsc nodeSetCollection, tota
return count
}
func (s *CarryWeightNodesetSelector) Select(nsc nodeSetCollection, excludeNodeSets []uint64, replicaNum uint8) (ns *nodeSet, err error) {
func (s *CarryWeightNodesetSelector) Select(nsc nodeSetCollection, param *selectParam) (ns *nodeSet, err error) {
total := s.getMaxTotal(nsc)
// prepare weight of evert nodesets
s.prepareCarry(nsc, total)
nsc = s.getAvailNodesets(nsc, excludeNodeSets, replicaNum)
nsc = s.getAvailNodesets(nsc, param)
avaliCount := 0
if len(nsc) < 1 {
goto err
@ -263,17 +263,19 @@ func (s *CarryWeightNodesetSelector) Select(nsc nodeSetCollection, excludeNodeSe
// pick the first nodeset than has N writable node
for i := 0; i < avaliCount; i++ {
ns = nsc[i]
if ns.canWriteFor(s.nodeType, int(replicaNum)) && !containsID(excludeNodeSets, ns.ID) {
if ns.canWriteFor(s.nodeType, param) && !containsID(param.excludeNodeSets, ns.ID) {
break
}
}
if ns != nil {
if !ns.canWriteFor(s.nodeType, int(replicaNum)) || containsID(excludeNodeSets, ns.ID) {
if !ns.canWriteFor(s.nodeType, param) || containsID(param.excludeNodeSets, ns.ID) {
goto err
}
s.carrys[ns.ID] -= 1.0
}
return
err:
switch s.nodeType {
case DataNodeType:
@ -301,7 +303,7 @@ func (s *AvailableSpaceFirstNodesetSelector) GetName() string {
return AvailableSpaceFirstNodesetSelectorName
}
func (s *AvailableSpaceFirstNodesetSelector) Select(nsc nodeSetCollection, excludeNodeSets []uint64, replicaNum uint8) (ns *nodeSet, err error) {
func (s *AvailableSpaceFirstNodesetSelector) Select(nsc nodeSetCollection, param *selectParam) (ns *nodeSet, err error) {
// sort nodesets by available space
sort.Slice(nsc, func(i, j int) bool {
return nsc[i].getTotalAvailableSpaceOf(s.nodeType) > nsc[j].getTotalAvailableSpaceOf(s.nodeType)
@ -309,7 +311,7 @@ func (s *AvailableSpaceFirstNodesetSelector) Select(nsc nodeSetCollection, exclu
// pick the first nodeset that has N writable nodes
for i := 0; i < nsc.Len(); i++ {
ns = nsc[i]
if ns.canWriteFor(s.nodeType, int(replicaNum)) && !containsID(excludeNodeSets, ns.ID) {
if ns.canWriteFor(s.nodeType, param) && !containsID(param.excludeNodeSets, ns.ID) {
return
}
}
@ -348,10 +350,10 @@ func (s *StrawNodesetSelector) getWeight(ns *nodeSet) float64 {
return float64(ns.getTotalAvailableSpaceOf(s.nodeType) / util.GB)
}
func (s *StrawNodesetSelector) Select(nsc nodeSetCollection, excludeNodeSets []uint64, replicaNum uint8) (ns *nodeSet, err error) {
func (s *StrawNodesetSelector) Select(nsc nodeSetCollection, param *selectParam) (ns *nodeSet, err error) {
tmp := make(nodeSetCollection, 0)
for _, nodeset := range nsc {
if nodeset.canWriteFor(s.nodeType, int(replicaNum)) && !containsID(excludeNodeSets, nodeset.ID) {
if nodeset.canWriteFor(s.nodeType, param) && !containsID(param.excludeNodeSets, nodeset.ID) {
tmp = append(tmp, nodeset)
}
}

View File

@ -80,7 +80,14 @@ func NodesetSelectorTest(t *testing.T, selector NodesetSelector) {
return true
})
}
ns, err := selector.Select(nsc, nil, 1)
param := &selectParam{
excludeNodeSets: nil,
replicaNum: 1,
excludeHosts: nil,
rackLevel: proto.RackAwareNone,
excludeRacks: nil,
}
ns, err := selector.Select(nsc, param)
if err != nil {
t.Errorf("%v failed to select nodeset %v", selector.GetName(), err)
return
@ -154,7 +161,14 @@ func prepareMetaNodesetForBench(count int, initTotal uint64, grow uint64) (nsc n
func nodesetSelectorBench(selector NodesetSelector, nsc nodeSetCollection, onSelect func(id uint64)) (map[uint64]int, error) {
times := make(map[uint64]int)
for i := 0; i < loopNodeSelectorTestCount; i++ {
ns, err := selector.Select(nsc, nil, 1)
param := &selectParam{
excludeNodeSets: nil,
replicaNum: 1,
excludeHosts: nil,
rackLevel: proto.RackAwareNone,
excludeRacks: nil,
}
ns, err := selector.Select(nsc, param)
if err != nil {
return nil, err
}

View File

@ -29,6 +29,42 @@ import (
"github.com/cubefs/cubefs/util/log"
)
type selectParam struct {
excludeHosts []string
replicaNum int
rackLevel proto.RackAwareLevel
excludeRacks []string
excludeNodeSets []uint64
}
func (sp *selectParam) copy() *selectParam {
// deep copy slice fields
excludeHosts := make([]string, len(sp.excludeHosts))
copy(excludeHosts, sp.excludeHosts)
excludeRacks := make([]string, len(sp.excludeRacks))
copy(excludeRacks, sp.excludeRacks)
excludeNodeSets := make([]uint64, len(sp.excludeNodeSets))
copy(excludeNodeSets, sp.excludeNodeSets)
return &selectParam{
excludeHosts: excludeHosts,
replicaNum: sp.replicaNum,
rackLevel: sp.rackLevel,
excludeRacks: excludeRacks,
excludeNodeSets: excludeNodeSets,
}
}
func (sp *selectParam) String() string {
if sp == nil {
return "nil"
}
return fmt.Sprintf("excludeHosts: %v, replicaNum: %v, rackLevel: %v, excludeRacks: %v, excludeNodeSets: %v",
sp.excludeHosts, sp.replicaNum, sp.rackLevel, sp.excludeRacks, sp.excludeNodeSets)
}
type rsManager struct {
nodeType NodeType
nodes *sync.Map
@ -581,7 +617,15 @@ func (nsgm *DomainManager) getHostFromNodeSetGrpSpecific(domainGrpManager *Domai
}
if createType == TypeDataPartition {
if host, peer, err = ns.getAvailDataNodeHosts(nil, needNum); err != nil {
param := &selectParam{
excludeNodeSets: nil,
replicaNum: needNum,
excludeHosts: nil,
rackLevel: proto.RackAwareNone,
excludeRacks: nil,
}
if host, peer, err = ns.getAvailDataNodeHosts(param); err != nil {
log.LogErrorf("action[getHostFromNodeSetGrpSpecific] ns[%v] zone[%v] TypeDataPartition err[%v]", ns.ID, ns.zoneName, err)
// nsg.status = dataNodesUnAvailable
continue
@ -591,7 +635,14 @@ func (nsgm *DomainManager) getHostFromNodeSetGrpSpecific(domainGrpManager *Domai
if createType == TypeRocksdbPartition {
storeMode = proto.StoreModeRocksDb
}
if host, peer, err = ns.getAvailMetaNodeHosts(nil, needNum, storeMode); err != nil {
param := &selectParam{
excludeNodeSets: nil,
replicaNum: needNum,
excludeHosts: nil,
rackLevel: proto.RackAwareNone,
excludeRacks: nil,
}
if host, peer, err = ns.getAvailMetaNodeHosts(param, storeMode); err != nil {
log.LogErrorf("action[getHostFromNodeSetGrpSpecific] ns[%v] zone[%v] type(%d) err[%v]", ns.ID, ns.zoneName, createType, err)
// nsg.status = metaNodesUnAvailable
continue
@ -691,7 +742,14 @@ func (nsgm *DomainManager) getHostFromNodeSetGrp(domainId uint64, replicaNum uin
log.LogWarnf("action[getHostFromNodeSetGrp] ns[%v] zone[%v] dataNodesUnAvailable", ns.ID, ns.zoneName)
continue
}
if host, peer, err = ns.getAvailDataNodeHosts(hosts, 1); err != nil {
param := &selectParam{
excludeNodeSets: nil,
replicaNum: 1,
excludeHosts: hosts,
rackLevel: proto.RackAwareNone,
excludeRacks: nil,
}
if host, peer, err = ns.getAvailDataNodeHosts(param); err != nil {
log.LogWarnf("action[getHostFromNodeSetGrp] ns[%v] zone[%v] TypeDataPartition err[%v]", ns.ID, ns.zoneName, err)
// nsg.status = dataNodesUnAvailable
continue
@ -705,7 +763,14 @@ func (nsgm *DomainManager) getHostFromNodeSetGrp(domainId uint64, replicaNum uin
if createType == TypeRocksdbPartition {
storeMode = proto.StoreModeRocksDb
}
if host, peer, err = ns.getAvailMetaNodeHosts(hosts, 1, storeMode); err != nil {
param := &selectParam{
excludeNodeSets: nil,
replicaNum: 1,
excludeHosts: hosts,
rackLevel: proto.RackAwareNone,
excludeRacks: nil,
}
if host, peer, err = ns.getAvailMetaNodeHosts(param, storeMode); err != nil {
log.LogWarnf("action[getHostFromNodeSetGrp] ns[%v] zone[%v] ModeRocksDb err[%v]", ns.ID, ns.zoneName, err)
// nsg.status = metaNodesUnAvailable
continue
@ -988,8 +1053,10 @@ type nodeSet struct {
metaNodes *sync.Map
dataNodes *sync.Map
racks map[string]*nodeSet
racksLock sync.RWMutex
racks map[string]*nodeSet
racksLock sync.RWMutex
dataRackSelector NodesetSelector
metaRackSelector NodesetSelector
nodeSelectLock sync.Mutex
dataNodeSelectorLock sync.RWMutex
@ -1025,13 +1092,17 @@ type nodeSetDecommissionParallelStatus struct {
func newNodeSet(c *Cluster, id uint64, cap int, zoneName string, rack string) *nodeSet {
log.LogInfof("action[newNodeSet] id[%v]", id)
ns := &nodeSet{
ID: id,
Rack: rack,
Capacity: cap,
zoneName: zoneName,
metaNodes: new(sync.Map),
dataNodes: new(sync.Map),
racks: make(map[string]*nodeSet),
ID: id,
Rack: rack,
Capacity: cap,
zoneName: zoneName,
metaNodes: new(sync.Map),
dataNodes: new(sync.Map),
racks: make(map[string]*nodeSet),
dataRackSelector: NewNodesetSelector(StrawNodesetSelectorName, DataNodeType),
metaRackSelector: NewNodesetSelector(StrawNodesetSelectorName, MetaNodeType),
decommissionDataPartitionList: NewDecommissionDataPartitionList(c, id),
manualDecommissionDiskList: NewDecommissionDiskList(),
autoDecommissionDiskList: NewDecommissionDiskList(),
@ -1049,6 +1120,31 @@ func newNodeSet(c *Cluster, id uint64, cap int, zoneName string, rack string) *n
return ns
}
func (ns *nodeSet) getRackSelector(nodeType NodeType) NodesetSelector {
if nodeType == DataNodeType {
return ns.dataRackSelector
}
return ns.metaRackSelector
}
func (ns *nodeSet) getNodeSelector(nodeType NodeType, storeMode proto.StoreMode) NodeSelector {
if nodeType == DataNodeType {
return ns.dataNodeSelector
}
if storeMode == proto.StoreModeRocksDb {
return ns.metaNodeRocksdbSelector
}
return ns.metaNodeMemorySelector
}
func (ns *nodeSet) IsRackSet() bool {
return ns.Rack != ""
}
func (ns *nodeSet) IsNodeSet() bool {
return ns.Rack == ""
}
func (ns *nodeSet) GetDataNodeSelector() string {
ns.dataNodeSelectorLock.RLock()
defer ns.dataNodeSelectorLock.RUnlock()
@ -1137,28 +1233,82 @@ func (ns *nodeSet) deleteMetaNode(metaNode *MetaNode) {
}
func (ns *nodeSet) canWriteForNode(nodes *sync.Map, replicaNum int, nodeType NodeType) bool {
// Extract node check logic
func (ns *nodeSet) checkNodeWriteable(node interface{}, nodeType NodeType, param *selectParam) bool {
if nodeType == RocksdbType {
if metaNode, ok := node.(*MetaNode); ok {
return metaNode.IsRocksdbWriteAble() && metaNode.PartitionCntLimited() && !contains(param.excludeHosts, metaNode.Addr)
}
return false
}
if n, ok := node.(Node); ok {
return n.IsWriteAble() && n.PartitionCntLimited() && !contains(param.excludeHosts, n.GetAddr())
}
return false
}
// Optimized main function
func (ns *nodeSet) canWriteForNode(nodes *sync.Map, param *selectParam, nodeType NodeType) bool {
// Rack-aware strong consistency mode, select available rack set count >= replicaNum
if param.rackLevel == proto.RackAwareStrong && ns.IsNodeSet() {
return ns.checkRackAwareWriteable(param, nodeType)
}
// Rack-aware strong consistency mode, exclude rack
if ns.IsRackSet() && param.rackLevel == proto.RackAwareStrong && contains(param.excludeRacks, ns.Rack) {
return false
}
// Normal mode
return ns.checkNormalWriteable(nodes, param, nodeType)
}
// Rack-aware consistency check, select available rack set count >= replicaNum
func (ns *nodeSet) checkRackAwareWriteable(param *selectParam, nodeType NodeType) bool {
var count int
for _, rack := range ns.racks {
if contains(param.excludeRacks, rack.Rack) {
continue
}
rnodes := rack.metaNodes
if nodeType == DataNodeType {
rnodes = rack.dataNodes
}
rnodes.Range(func(key, value interface{}) bool {
if ns.checkNodeWriteable(value, nodeType, param) {
count++
return false
}
return true
})
if count >= param.replicaNum {
return true
}
}
return false
}
// Normal mode consistency check
func (ns *nodeSet) checkNormalWriteable(nodes *sync.Map, param *selectParam, nodeType NodeType) bool {
var count int
nodes.Range(func(key, value interface{}) bool {
if nodeType == RocksdbType {
node := value.(*MetaNode)
if node.IsRocksdbWriteAble() && node.PartitionCntLimited() {
count++
if ns.checkNodeWriteable(value, nodeType, param) {
count++
if count >= param.replicaNum {
return false // Found enough nodes, exit early
}
} else {
node := value.(Node)
if node.IsWriteAble() && node.PartitionCntLimited() {
count++
}
}
if count >= replicaNum {
return false
}
return true
})
log.LogInfof("canWriteForMetaNode zone[%v], ns[%v],count[%v] replicaNum[%v]",
ns.zoneName, ns.ID, count, replicaNum)
return count >= replicaNum
ns.zoneName, ns.ID, count, param.replicaNum)
return count >= param.replicaNum
}
func (ns *nodeSet) calcNodesForAlloc(nodes *sync.Map) (cnt int) {
@ -1846,7 +1996,7 @@ func (zone *Zone) deleteMetaNode(metaNode *MetaNode) (err error) {
return
}
func (zone *Zone) allocNodeSetForDataNode(excludeNodeSets []uint64, replicaNum uint8) (ns *nodeSet, err error) {
func (zone *Zone) allocNodeSetForDataNode(param *selectParam) (ns *nodeSet, err error) {
nset := zone.getAllNodeSet()
if nset == nil {
return nil, errors.NewError(proto.ErrNoNodeSetToCreateDataPartition)
@ -1858,16 +2008,38 @@ func (zone *Zone) allocNodeSetForDataNode(excludeNodeSets []uint64, replicaNum u
zone.dataNodesetSelectorLock.RLock()
defer zone.dataNodesetSelectorLock.RUnlock()
ns, err = zone.dataNodesetSelector.Select(nset, excludeNodeSets, replicaNum)
if err != nil {
log.LogErrorf("action[allocNodeSetForDataNode],nset len[%v],excludeNodeSets[%v],rNum[%v] err:%v",
nset.Len(), excludeNodeSets, replicaNum, proto.ErrNoNodeSetToCreateDataPartition)
paramCopy := param.copy()
if param.rackLevel == proto.RackAwareWeak {
paramCopy.rackLevel = proto.RackAwareStrong
}
// First attempt with strong rack awareness
ns, err = zone.dataNodesetSelector.Select(nset, paramCopy)
if err == nil {
return ns, nil
}
if param.rackLevel != proto.RackAwareWeak {
log.LogErrorf("action[allocNodeSetForDataNode], no available node set found, param[%v] err:%v",
param.String(), err.Error())
return nil, errors.NewError(proto.ErrNoNodeSetToCreateDataPartition)
}
log.LogWarnf("action[allocNodeSetForDataNode], no available node set found, try awake, param[%v] err:%v",
param.String(), err.Error())
paramCopy.rackLevel = proto.RackAwareWeak
ns, err = zone.dataNodesetSelector.Select(nset, paramCopy)
if err != nil {
log.LogErrorf("action[allocNodeSetForDataNode], still failed after change rack level,param[%v], err: %v",
param.String(), proto.ErrNoNodeSetToCreateDataPartition)
return nil, errors.NewError(proto.ErrNoNodeSetToCreateDataPartition)
}
return ns, nil
}
func (zone *Zone) allocNodeSetForMetaNode(excludeNodeSets []uint64, replicaNum uint8, storeMode proto.StoreMode) (ns *nodeSet, err error) {
func (zone *Zone) allocNodeSetForMetaNode(param *selectParam, storeMode proto.StoreMode) (ns *nodeSet, err error) {
nset := zone.getAllNodeSet()
if nset == nil {
return nil, proto.ErrNoNodeSetToCreateMetaPartition
@ -1878,16 +2050,42 @@ func (zone *Zone) allocNodeSetForMetaNode(excludeNodeSets []uint64, replicaNum u
// we need a read lock to block the modify of nodeset selector
zone.metaNodesetSelectorLock.RLock()
defer zone.metaNodesetSelectorLock.RUnlock()
// Select appropriate selector based on store mode
selector := zone.metaMemoryNodesetSelector
if storeMode == proto.StoreModeRocksDb {
selector = zone.metaRocksdbNodesetSelector
}
ns, err = selector.Select(nset, excludeNodeSets, replicaNum)
paramCopy := param.copy()
if param.rackLevel == proto.RackAwareWeak {
paramCopy.rackLevel = proto.RackAwareStrong
}
// First attempt with strong rack awareness
ns, err = selector.Select(nset, paramCopy)
if err == nil {
return ns, nil
}
if param.rackLevel != proto.RackAwareWeak {
log.LogErrorf("action[allocNodeSetForMetaNode], no available node set found, zone %s, param[%v] err:%v",
zone.name, param.String(), err.Error())
return nil, errors.NewError(proto.ErrNoNodeSetToCreateMetaPartition)
}
// Retry with weak rack awareness
log.LogWarnf("action[allocNodeSetForMetaNode], no available node set found, zone %s, param[%v] err:%v",
zone.name, param.String(), err.Error())
paramCopy.rackLevel = proto.RackAwareWeak
ns, err = selector.Select(nset, paramCopy)
if err != nil {
log.LogError(fmt.Sprintf("action[allocNodeSetForMetaNode],zone[%v],excludeNodeSets[%v],rNum[%v],err:%v",
zone.name, excludeNodeSets, replicaNum, proto.ErrNoNodeSetToCreateMetaPartition))
log.LogErrorf("action[allocNodeSetForMetaNode], still failed after change rack level, zone[%v], param[%v], err:%v",
zone.name, param.String(), proto.ErrNoNodeSetToCreateMetaPartition)
return nil, proto.ErrNoNodeSetToCreateMetaPartition
}
return ns, nil
}
@ -1983,19 +2181,19 @@ func (zone *Zone) getSpaceLeft(dataType uint32) (spaceLeft uint64) {
}
}
func (zone *Zone) getAvailNodeHosts(nodeType uint32, excludeNodeSets []uint64, excludeHosts []string, replicaNum int) (newHosts []string, peers []proto.Peer, err error) {
if replicaNum == 0 {
func (zone *Zone) getAvailNodeHosts(nodeType uint32, param *selectParam) (newHosts []string, peers []proto.Peer, err error) {
if param.replicaNum == 0 {
return
}
log.LogDebugf("[getAvailNodeHosts] get node host, zone(%s), nodeType(%d)", zone.name, nodeType)
if nodeType == TypeDataPartition {
ns, err := zone.allocNodeSetForDataNode(excludeNodeSets, uint8(replicaNum))
ns, err := zone.allocNodeSetForDataNode(param)
if err != nil {
return nil, nil, errors.Trace(err, "zone[%v] alloc node set,replicaNum[%v]", zone.name, replicaNum)
return nil, nil, errors.Trace(err, "zone[%v] alloc node set, param[%v], err %s", zone.name, param.String(), err.Error())
}
return ns.getAvailDataNodeHosts(excludeHosts, replicaNum)
return ns.getAvailDataNodeHosts(param)
}
storeMode := proto.StoreModeMem
@ -2003,12 +2201,12 @@ func (zone *Zone) getAvailNodeHosts(nodeType uint32, excludeNodeSets []uint64, e
storeMode = proto.StoreModeRocksDb
}
ns, err := zone.allocNodeSetForMetaNode(excludeNodeSets, uint8(replicaNum), storeMode)
ns, err := zone.allocNodeSetForMetaNode(param, storeMode)
if err != nil {
return nil, nil, errors.NewErrorf("zone[%v],err[%v]", zone.name, err)
return nil, nil, errors.NewErrorf("zone[%v], param[%v], err[%v]", zone.name, param.String(), err.Error())
}
return ns.getAvailMetaNodeHosts(excludeHosts, replicaNum, storeMode)
return ns.getAvailMetaNodeHosts(param, storeMode)
}
func (zone *Zone) updateNodesetSelector(cluster *Cluster, dataNodesetSelector string, metaNodesetSelector string) error {

1761
master/topology_rack_test.go Normal file

File diff suppressed because it is too large Load Diff

View File

@ -54,7 +54,15 @@ func TestSingleZone(t *testing.T) {
// single zone normal
zones, err = topo.allocZonesForNode(&topo.dataTopology, replicaNum, replicaNum, nil, []*Zone{}, proto.MediaType_Unspecified)
require.NoError(t, err)
newHosts, _, err := zones[0].getAvailNodeHosts(TypeDataPartition, nil, nil, replicaNum)
param := &selectParam{
excludeNodeSets: nil,
replicaNum: replicaNum,
excludeHosts: nil,
rackLevel: c.getRackAwareLevel(),
excludeRacks: nil,
}
newHosts, _, err := zones[0].getAvailNodeHosts(TypeDataPartition, param)
require.NoError(t, err)
t.Log(newHosts)
topo.deleteDataNode(createDataNodeForTopo(mds1Addr, zoneName, nodeSet))
@ -134,17 +142,17 @@ func TestAllocZones(t *testing.T) {
cluster.cfg = newClusterConfig()
// don't cross zone
hosts, _, err := cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, replicaNum, 1, "", proto.MediaType_Unspecified)
hosts, _, err := cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, replicaNum, 1, "", proto.MediaType_Unspecified, proto.RackAwareNone)
require.NoError(t, err)
t.Logf("ChooseTargetDataHosts in single zone,hosts[%v]", hosts)
// cross zone
_, _, err = cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, replicaNum, 2, "", proto.MediaType_Unspecified)
_, _, err = cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, replicaNum, 2, "", proto.MediaType_Unspecified, proto.RackAwareNone)
require.NoError(t, err)
// specific zone
hosts, _, err = cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, 3, 2, zoneName1+","+zoneName2, proto.MediaType_Unspecified)
hosts, _, err = cluster.getHostFromNormalZone(TypeDataPartition, nil, nil, nil, 3, 2, zoneName1+","+zoneName2, proto.MediaType_Unspecified, proto.RackAwareNone)
require.NoError(t, err)
require.EqualValues(t, getZoneCntFunc(hosts), 2)

View File

@ -1714,7 +1714,7 @@ func (vol *Vol) doCreateMetaPartition(c *Cluster, start, end uint64) (mp *MetaPa
var excludeZone []string
zoneNum := c.decideZoneNum(vol, proto.StorageClass_Unspecified)
if hosts, peers, err = c.getHostFromNormalZone(nodeType, excludeZone, nil, nil,
int(vol.mpReplicaNum), zoneNum, vol.zoneName, proto.StorageClass_Unspecified); err != nil {
int(vol.mpReplicaNum), zoneNum, vol.zoneName, proto.StorageClass_Unspecified, c.getRackAwareLevel()); err != nil {
log.LogErrorf("action[doCreateMetaPartition] getHostFromNormalZone err[%v]", err)
return nil, errors.NewError(err)
}