feat(master): hybrid cloud check destination datanode's media type must be the same

with the source node's when migration and decommission

Signed-off-by: true1064 <tangjingyu@oppo.com>
This commit is contained in:
true1064 2024-06-17 11:46:06 +08:00 committed by AmazingChi
parent 9fd51b5119
commit bf9e425de6
5 changed files with 77 additions and 18 deletions

View File

@ -2577,7 +2577,7 @@ func (c *Cluster) autoAddDataReplica(dp *DataPartition) (success bool, err error
if ns, err = zone.getNodeSet(nodeSets[0]); err != nil {
goto errHandler
}
if targetHosts, _, err = ns.getAvailDataNodeHosts(dp.Hosts, 1, zone.dataMediaType); err != nil {
if targetHosts, _, err = ns.getAvailDataNodeHosts(dp.Hosts, 1); err != nil {
goto errHandler
}
}
@ -2678,7 +2678,11 @@ func (c *Cluster) migrateDataPartition(srcAddr, targetAddr string, dp *DataParti
if targetAddr != "" {
targetHosts = []string{targetAddr}
} else if targetHosts, _, err = ns.getAvailDataNodeHosts(dp.Hosts, 1, dp.MediaType); err != nil {
if err = c.checkDataNodesMediaTypeForMigrate(dataNode, targetAddr); err != nil {
log.LogErrorf("[migrateDataPartition] check mediaType err: %v", err.Error())
goto errHandler
}
} else if targetHosts, _, err = ns.getAvailDataNodeHosts(dp.Hosts, 1); 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)
@ -2827,25 +2831,25 @@ func (c *Cluster) addDataReplica(dp *DataPartition, addr string, ignoreDecommiss
dp.addReplicaMutex.Lock()
defer dp.addReplicaMutex.Unlock()
dataNode, err := c.dataNode(addr)
targetDataNode, err := c.dataNode(addr)
if err != nil {
return
}
if dataNode.MediaType != dp.MediaType {
if targetDataNode.MediaType != dp.MediaType {
err = fmt.Errorf("target datanode mediaType(%v) not match datapartition mediaType(%v)",
proto.MediaTypeString(dataNode.MediaType), proto.MediaTypeString(dp.MediaType))
proto.MediaTypeString(targetDataNode.MediaType), proto.MediaTypeString(dp.MediaType))
log.LogErrorf("[addDataReplica] dpId(%v), err: %v", dp.PartitionID, err.Error())
return
}
addPeer := proto.Peer{ID: dataNode.ID, Addr: addr}
addPeer := proto.Peer{ID: targetDataNode.ID, Addr: addr}
if !proto.IsNormalDp(dp.PartitionType) {
return fmt.Errorf("action[addDataReplica] [%d] is not normal dp, not support add or delete replica", dp.PartitionID)
}
log.LogInfof("action[addDataReplica] dp %v dst addr %v try add raft member, node id %v", dp.PartitionID, addr, dataNode.ID)
log.LogInfof("action[addDataReplica] dp %v dst addr %v try add raft member, node id %v", dp.PartitionID, addr, targetDataNode.ID)
if err = c.addDataPartitionRaftMember(dp, addPeer); err != nil {
log.LogWarnf("action[addDataReplica] dp %v addr %v try add raft member err [%v]", dp.PartitionID, addr, err)
return
@ -3637,6 +3641,7 @@ func (c *Cluster) initDataPartitionsForCreateVol(vol *Vol, dpCount int, mediaTyp
}
for retryCount := 0; readWriteDataPartitions < defaultInitDataPartitionCnt && retryCount < 3; retryCount++ {
//TODO:tangjingyu: when retry, to create dp count should be dpCount - alreadyCreatedCount
err = vol.initDataPartitions(c, dpCount, mediaType)
if err != nil {
log.LogError("action[initDataPartitionsForCreateVol] init dataPartition error:",
@ -4677,11 +4682,55 @@ func (c *Cluster) TryDecommissionDataNode(dataNode *DataNode) {
auditlog.LogMasterOp("DataNodeDecommission", msg, nil)
}
func (c *Cluster) migrateDisk(dataNode *DataNode, diskPath, dstPath string, raftForce bool, limit int, diskDisable bool, migrateType uint32) (err error) {
func (c *Cluster) checkDataNodesMediaTypeForMigrate(srcNode *DataNode, dstAddr string) (err error) {
var dstNode *DataNode
if srcNode == nil {
log.LogErrorf("[checkDataNodesMediaTypeForMigrate] srcNode is nil")
return
}
if dstAddr != "" {
dstNode, err = c.dataNode(dstAddr)
if err != nil {
log.LogErrorf("[CheckDataNodesMediaTypeForMigrate] get dstNode(%v) failed: %v", dstAddr, err.Error())
return
}
if dstNode.MediaType != srcNode.MediaType {
err = fmt.Errorf("dstNode mediaType(%v) not match srcNode mediaType(%v)",
proto.MediaTypeString(dstNode.MediaType), proto.MediaTypeString(srcNode.MediaType))
log.LogErrorf("[CheckDataNodesMediaTypeForMigrate] %v", err.Error())
return
}
}
return nil
}
func (c *Cluster) checkDataNodeAddrMediaTypeForMigrate(srcAddr, dstAddr string) (err error) {
var srcNode *DataNode
srcNode, err = c.dataNode(srcAddr)
if err != nil {
log.LogErrorf("[CheckDataNodesMediaTypeForMigrate] get srcNode(%v) failed: %v", srcAddr, err.Error())
return
}
return c.checkDataNodesMediaTypeForMigrate(srcNode, dstAddr)
}
func (c *Cluster) migrateDisk(dataNode *DataNode, diskPath, dstAddr string, raftForce bool, limit int, diskDisable bool, migrateType uint32) (err error) {
var disk *DecommissionDisk
nodeAddr := dataNode.Addr
key := fmt.Sprintf("%s_%s", nodeAddr, diskPath)
if dstAddr != "" {
if err = c.checkDataNodeAddrMediaTypeForMigrate(nodeAddr, dstAddr); err != nil {
log.LogErrorf("[migrateDisk] check mediaType err: %v", err.Error())
return
}
}
key := fmt.Sprintf("%s_%s", nodeAddr, diskPath)
if value, ok := c.DecommissionDisks.Load(key); ok {
disk = value.(*DecommissionDisk)
status := disk.GetDecommissionStatus()
@ -4706,7 +4755,7 @@ func (c *Cluster) migrateDisk(dataNode *DataNode, diskPath, dstPath string, raft
disk.ResidualDecommissionDps = make([]proto.IgnoreDecommissionDP, 0)
disk.IgnoreDecommissionDps = make([]proto.IgnoreDecommissionDP, 0)
// disk should be decommission all the dp
disk.markDecommission(dstPath, raftForce, limit)
disk.markDecommission(dstAddr, raftForce, limit)
if err = c.syncAddDecommissionDisk(disk); err != nil {
err = fmt.Errorf("action[addDecommissionDisk],clusterID[%v] dataNodeAddr:%v diskPath:%v err:%v ",
c.Name, nodeAddr, diskPath, err.Error())

View File

@ -1129,6 +1129,13 @@ func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk
}
}()
if dstAddr != "" {
if err = c.checkDataNodeAddrMediaTypeForMigrate(srcAddr, dstAddr); err != nil {
log.LogErrorf("[MarkDecommissionStatus] check mediaType err: %v", err.Error())
return
}
}
var status uint32
// if mark discard, decommission it directly to delete replica
if partition.IsDiscard {
@ -1831,7 +1838,8 @@ func (partition *DataPartition) TryAcquireDecommissionToken(c *Cluster) bool {
}
log.LogDebugf("action[TryAcquireDecommissionToken]dp %v excludeHosts %v",
partition.PartitionID, excludeHosts)
targetHosts, _, err = ns.getAvailDataNodeHosts(excludeHosts, 1, partition.MediaType)
// data nodes in a nodeset has the same mediaType
targetHosts, _, err = ns.getAvailDataNodeHosts(excludeHosts, 1)
if err != nil {
log.LogWarnf("action[TryAcquireDecommissionToken] dp %v choose from src nodeset failed:%v",
partition.PartitionID, err.Error())
@ -1847,6 +1855,7 @@ func (partition *DataPartition) TryAcquireDecommissionToken(c *Cluster) bool {
goto errHandler
}
excludeNodeSets = append(excludeNodeSets, ns.ID)
// data nodes in a zone has the same mediaType
if targetHosts, _, err = zone.getAvailNodeHosts(TypeDataPartition, excludeNodeSets, excludeHosts, 1); err != nil {
log.LogWarnf("action[TryAcquireDecommissionToken] dp %v choose from other nodeset failed:%v",
partition.PartitionID, err.Error())

View File

@ -478,13 +478,13 @@ func (dd *DecommissionDisk) GetDecommissionFailedDP(c *Cluster) (error, []uint64
return nil, failedDps
}
func (dd *DecommissionDisk) markDecommission(dstPath string, raftForce bool, limit int) {
func (dd *DecommissionDisk) markDecommission(dstAddr string, raftForce bool, limit int) {
// if transfer from pause,do not change these attrs
if dd.GetDecommissionStatus() != DecommissionPause {
dd.DecommissionDpTotal = InvalidDecommissionDpCnt
dd.DecommissionDpCount = limit
dd.DecommissionRaftForce = raftForce
dd.DstAddr = dstPath
dd.DstAddr = dstAddr
dd.DecommissionTimes = 0
}
dd.DecommissionTerm = uint64(time.Now().Unix())

View File

@ -525,8 +525,7 @@ func (ns *nodeSet) getAvailMetaNodeHosts(excludeHosts []string, replicaNum int)
return ns.metaNodeSelector.Select(ns, excludeHosts, replicaNum)
}
//TODO:tangjingyu mediaType not used
func (ns *nodeSet) getAvailDataNodeHosts(excludeHosts []string, replicaNum int, mediaType uint32) (hosts []string, peers []proto.Peer, err error) {
func (ns *nodeSet) getAvailDataNodeHosts(excludeHosts []string, replicaNum int) (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

View File

@ -533,6 +533,7 @@ func (nsgm *DomainManager) buildNodeSetGrp(domainGrpManager *DomainNodeSetGrpMan
return nil
}
//TODO:tangjingyu param mediaType not used, need keep it?
func (nsgm *DomainManager) getHostFromNodeSetGrpSpecific(domainGrpManager *DomainNodeSetGrpManager, replicaNum uint8,
createType uint32, mediaType uint32) (
hosts []string,
@ -582,7 +583,7 @@ func (nsgm *DomainManager) getHostFromNodeSetGrpSpecific(domainGrpManager *Domai
}
if createType == TypeDataPartition {
if host, peer, err = ns.getAvailDataNodeHosts(nil, needNum, mediaType); err != nil {
if host, peer, err = ns.getAvailDataNodeHosts(nil, needNum); err != nil {
log.LogErrorf("action[getHostFromNodeSetGrpSpecific] ns[%v] zone[%v] TypeDataPartition err[%v]", ns.ID, ns.zoneName, err)
// nsg.status = dataNodesUnAvailable
continue
@ -611,6 +612,7 @@ func (nsgm *DomainManager) getHostFromNodeSetGrpSpecific(domainGrpManager *Domai
return nil, nil, fmt.Errorf("action[getHostFromNodeSetGrpSpecific] cann't alloc host")
}
//TODO:tangjingyu param mediaType not used, need keep it?
func (nsgm *DomainManager) getHostFromNodeSetGrp(domainId uint64, replicaNum uint8, createType uint32, mediaType uint32) (
hosts []string,
peers []proto.Peer,
@ -687,7 +689,7 @@ 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, mediaType); err != nil {
if host, peer, err = ns.getAvailDataNodeHosts(hosts, 1); err != nil {
log.LogWarnf("action[getHostFromNodeSetGrp] ns[%v] zone[%v] TypeDataPartition err[%v]", ns.ID, ns.zoneName, err)
// nsg.status = dataNodesUnAvailable
continue
@ -1865,7 +1867,7 @@ func (zone *Zone) getAvailNodeHosts(nodeType uint32, excludeNodeSets []uint64, e
if err != nil {
return nil, nil, errors.Trace(err, "zone[%v] alloc node set,replicaNum[%v]", zone.name, replicaNum)
}
return ns.getAvailDataNodeHosts(excludeHosts, replicaNum, zone.dataMediaType)
return ns.getAvailDataNodeHosts(excludeHosts, replicaNum)
}
ns, err := zone.allocNodeSetForMetaNode(excludeNodeSets, uint8(replicaNum))