diff --git a/cli/cmd/const.go b/cli/cmd/const.go index d6cf21314..2103695c9 100644 --- a/cli/cmd/const.go +++ b/cli/cmd/const.go @@ -154,6 +154,7 @@ const ( CliFlagEnablePersistAccessTime = "enablePersistAccessTime" CliFlagDecommissionRaftForce = "raftForceDel" CliFLagDecommissionWeight = "decommissionWeight" + CliFlagDecommissionDstNodeSet = "decommissionDstNodeSet" CliFLagRecommissionType = "recommissionType" CliFlagAllowedStorageClass = "allowedStorageClass" CliFlagVolStorageClass = "volStorageClass" diff --git a/cli/cmd/datapartition.go b/cli/cmd/datapartition.go index 70df83f34..55a940bbf 100644 --- a/cli/cmd/datapartition.go +++ b/cli/cmd/datapartition.go @@ -413,6 +413,7 @@ The "reset" command will be released in next version`, func newDataPartitionDecommissionCmd(client *master.MasterClient) *cobra.Command { var raftForceDel bool var weight int + var dstNodeSet uint64 var decommissionType string var clientIDKey string cmd := &cobra.Command{ @@ -432,7 +433,7 @@ func newDataPartitionDecommissionCmd(client *master.MasterClient) *cobra.Command if err != nil { return } - if err := client.AdminAPI().DecommissionDataPartition(partitionID, address, raftForceDel, weight, clientIDKey, decommissionType); err != nil { + if err := client.AdminAPI().DecommissionDataPartition(partitionID, address, dstNodeSet, raftForceDel, weight, clientIDKey, decommissionType); err != nil { stdout(fmt.Sprintf("failed:err(%v)\n", err.Error())) return } @@ -447,6 +448,7 @@ func newDataPartitionDecommissionCmd(client *master.MasterClient) *cobra.Command } cmd.Flags().BoolVarP(&raftForceDel, CliFlagDecommissionRaftForce, "r", false, "true for raftForceDel") cmd.Flags().IntVar(&weight, CliFLagDecommissionWeight, lowPriorityDecommissionWeight, "decommission weight") + cmd.Flags().Uint64Var(&dstNodeSet, CliFlagDecommissionDstNodeSet, 0, "decommission dst nodeSet") cmd.Flags().StringVar(&decommissionType, "decommissionType", "1", "decommission type") cmd.Flags().StringVar(&clientIDKey, CliFlagClientIDKey, client.ClientIDKey(), CliUsageClientIDKey) return cmd diff --git a/cli/cmd/fmt.go b/cli/cmd/fmt.go index c5d3c874c..a55bf7293 100644 --- a/cli/cmd/fmt.go +++ b/cli/cmd/fmt.go @@ -1386,6 +1386,7 @@ func formatDataPartitionDecommissionProgress(info *proto.DecommissionDataPartiti sb.WriteString(fmt.Sprintf("SrcAddress: %v\n", info.SrcAddress)) sb.WriteString(fmt.Sprintf("SrcDiskPath: %v\n", info.SrcDiskPath)) sb.WriteString(fmt.Sprintf("DstAddress: %v\n", info.DstAddress)) + sb.WriteString(fmt.Sprintf("DstNodeSet: %v\n", info.DstNodeSet)) sb.WriteString(fmt.Sprintf("Term: %v\n", info.Term)) sb.WriteString(fmt.Sprintf("Weight: %v\n", info.Weight)) sb.WriteString(fmt.Sprintf("Replicas: %v\n", info.Replicas)) diff --git a/master/api_service.go b/master/api_service.go index 1d9557fc2..092cc0ee1 100644 --- a/master/api_service.go +++ b/master/api_service.go @@ -2092,6 +2092,7 @@ func (m *Server) decommissionDataPartition(w http.ResponseWriter, r *http.Reques rstMsg string dp *DataPartition addr string + dstNodeSet uint64 partitionID uint64 raftForce bool weight int @@ -2112,6 +2113,12 @@ func (m *Server) decommissionDataPartition(w http.ResponseWriter, r *http.Reques sendErrReply(w, r, newErrHTTPReply(proto.ErrDataPartitionNotExists)) return } + + dstNodeSet, err = parseDstNodeSet(r) + if err != nil { + sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()}) + return + } raftForce, err = parseRaftForce(r) if err != nil { sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()}) @@ -2143,7 +2150,7 @@ func (m *Server) decommissionDataPartition(w http.ResponseWriter, r *http.Reques sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: rstMsg}) return } - err = m.cluster.markDecommissionDataPartition(dp, node, raftForce, uint32(decommissionType), weight) + err = m.cluster.markDecommissionDataPartition(dp, node, dstNodeSet, raftForce, uint32(decommissionType), weight) if err != nil { sendErrReply(w, r, newErrHTTPReply(err)) return @@ -2309,6 +2316,7 @@ func (m *Server) queryDataPartitionDecommissionStatus(w http.ResponseWriter, r * SrcAddress: dp.DecommissionSrcAddr, SrcDiskPath: dp.DecommissionSrcDiskPath, DstAddress: dp.DecommissionDstAddr, + DstNodeSet: dp.DecommissionDstNodeSet, Term: dp.DecommissionTerm, Weight: dp.DecommissionWeight, Replicas: replicas, @@ -5703,6 +5711,20 @@ func parseWeight(r *http.Request) (int, error) { return newVal, nil } +func parseDstNodeSet(r *http.Request) (uint64, error) { + val := r.FormValue(dstNodeSetKey) + if val == "" { + return 0, nil + } + + newVal, err := strconv.ParseUint(val, 10, 64) + if err != nil { + return 0, fmt.Errorf("parse %s uint64 vaal err, err %s", dstNodeSetKey, err.Error()) + } + + return newVal, nil +} + func extractPosixAcl(r *http.Request) (enablePosix bool, err error) { var value string if value = r.FormValue(enablePosixAclKey); value == "" { diff --git a/master/cluster.go b/master/cluster.go index 4f539ac5f..4a5367797 100644 --- a/master/cluster.go +++ b/master/cluster.go @@ -5258,7 +5258,7 @@ func (c *Cluster) handleDataNodeBadDisk(dataNode *DataNode) { log.LogInfof("[handleDataNodeBadDisk] data node(%v) not found in dp(%v) maybe decommissioned?", dataNode.Addr, dpId) continue } - err = c.markDecommissionDataPartition(dp, dataNode, false, AutoDecommission, highPriorityDecommissionWeight) + err = c.markDecommissionDataPartition(dp, dataNode, 0, false, AutoDecommission, highPriorityDecommissionWeight) if err != nil { log.LogErrorf("[handleDataNodeBadDisk] failed to decommssion dp(%v) on data node(%v) disk(%v), err(%v)", dataNode.Addr, disk.DiskPath, dp.PartitionID, err) continue @@ -5378,7 +5378,7 @@ func (c *Cluster) TryDecommissionDisk(disk *DecommissionDisk) { ignoreIDs = append(ignoreIDs, dp.PartitionID) continue } - if err = dp.MarkDecommissionStatus(node.Addr, disk.DstAddr, disk.DiskPath, disk.DecommissionRaftForce, + if err = dp.MarkDecommissionStatus(node.Addr, disk.DstAddr, disk.DiskPath, 0, disk.DecommissionRaftForce, disk.DecommissionTerm, disk.Type, disk.DecommissionWeight, c, ns); err != nil { if strings.Contains(err.Error(), proto.ErrDecommissionDiskErrDPFirst.Error()) { c.syncUpdateDataPartition(dp) @@ -6195,7 +6195,7 @@ func (c *Cluster) rangeAllParitions(f func(d *DataPartition) bool) { } } -func (c *Cluster) markDecommissionDataPartition(dp *DataPartition, src *DataNode, raftForce bool, migrateType uint32, weight int) (err error) { +func (c *Cluster) markDecommissionDataPartition(dp *DataPartition, src *DataNode, dstNodeSetID uint64, raftForce bool, migrateType uint32, weight int) (err error) { addr := src.Addr replica, err := dp.getReplica(addr) if err != nil { @@ -6213,7 +6213,7 @@ func (c *Cluster) markDecommissionDataPartition(dp *DataPartition, src *DataNode return } - if err = dp.MarkDecommissionStatus(addr, "", replica.DiskPath, raftForce, uint64(time.Now().Unix()), migrateType, weight, c, ns); err != nil { + if err = dp.MarkDecommissionStatus(addr, "", replica.DiskPath, dstNodeSetID, raftForce, uint64(time.Now().Unix()), migrateType, weight, c, ns); err != nil { if !strings.Contains(err.Error(), proto.ErrDecommissionDiskErrDPFirst.Error()) { dp.markRollbackFailed(false) dp.DecommissionErrorMessage = err.Error() diff --git a/master/const.go b/master/const.go index 71eb4f179..6c2f34575 100644 --- a/master/const.go +++ b/master/const.go @@ -88,6 +88,7 @@ const ( forceKey = "force" raftForceDelKey = "raftForceDel" weightKey = "weight" + dstNodeSetKey = "dstNodeSet" enablePosixAclKey = "enablePosixAcl" enableTxMaskKey = "enableTxMask" txTimeoutKey = "txTimeout" diff --git a/master/data_partition.go b/master/data_partition.go index e5ab3bb8b..cedf13256 100644 --- a/master/data_partition.go +++ b/master/data_partition.go @@ -69,6 +69,7 @@ type DataPartition struct { DecommissionSrcDiskPath string DecommissionTerm uint64 DecommissionDstAddrSpecify bool // if DecommissionDstAddrSpecify is true, donot rollback when add replica fail + DecommissionDstNodeSet uint64 DecommissionNeedRollback bool DecommissionNeedRollbackTimes uint32 DecommissionErrorMessage string @@ -116,6 +117,7 @@ func newDataPartition(ID uint64, replicaNum uint8, volName string, volID uint64, partition.DecommissionStatus = DecommissionInitial partition.SpecialReplicaDecommissionStep = SpecialDecommissionInitial partition.DecommissionDstAddrSpecify = false + partition.DecommissionDstNodeSet = 0 partition.LeaderReportTime = now partition.RepairBlockSize = util.DefaultDataPartitionSize partition.RestoreReplica = RestoreReplicaMetaStop @@ -1290,9 +1292,9 @@ errHandle: partition.markRollbackFailed(false) partition.DecommissionErrorMessage = err.Error() log.LogWarnf("action[AcquireDecommissionFirstHostToken] clusterID[%v] vol[%v] partitionID[%v]"+ - " retry [%v] status [%v] DecommissionDstAddrSpecify [%v] DecommissionDstAddr [%v] failed", + " retry [%v] status [%v] DecommissionDstAddrSpecify [%v] DecommissionDstAddr [%v] DecommissionDstNodeSet [%v] failed", c.Name, partition.VolName, partition.PartitionID, partition.DecommissionRetry, partition.GetDecommissionStatus(), - partition.DecommissionDstAddrSpecify, partition.DecommissionDstAddr) + partition.DecommissionDstAddrSpecify, partition.DecommissionDstAddr, partition.DecommissionDstNodeSet) return false } @@ -1305,7 +1307,7 @@ func isReplicasContainsHost(replicas []*DataReplica, host string) bool { return false } -func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk string, raftForce bool, term uint64, +func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk string, dstNodeSetID uint64, raftForce bool, term uint64, migrateType uint32, weight int, c *Cluster, ns *nodeSet, ) (err error) { defer func() { @@ -1322,6 +1324,13 @@ func (partition *DataPartition) MarkDecommissionStatus(srcAddr, dstAddr, srcDisk } } + if dstNodeSetID != 0 { + if _, err = c.t.getNodeSetByNodeSetId(dstNodeSetID); err != nil { + log.LogErrorf("[MarkDecommissionStatus] check dstNodeSetID(%v) err: %v", dstNodeSetID, err.Error()) + return + } + } + var status uint32 // if mark discard, decommission it directly to delete replica if partition.IsDiscard { @@ -1567,6 +1576,7 @@ directly: partition.DecommissionSrcAddr = srcAddr partition.DecommissionDstAddr = dstAddr partition.DecommissionSrcDiskPath = srcDisk + partition.DecommissionDstNodeSet = dstNodeSetID partition.DecommissionRaftForce = raftForce partition.DecommissionTerm = term partition.DecommissionWeight = weight @@ -1952,6 +1962,7 @@ func (partition *DataPartition) ResetDecommissionStatus() { partition.DecommissionTerm = 0 partition.DecommissionWeight = 0 partition.DecommissionDstAddrSpecify = false + partition.DecommissionDstNodeSet = 0 partition.DecommissionNeedRollback = false atomic.StoreUint32(&partition.DecommissionNeedRollbackTimes, 0) partition.SetDecommissionStatus(DecommissionInitial) @@ -2049,10 +2060,10 @@ func (partition *DataPartition) addToDecommissionList(c *Cluster) { log.LogWarnf("action[addToDecommissionList]dataNode[%v] nodeSet is nil:%v", dataNode.Addr, err.Error()) return } - log.LogInfof("action[addToDecommissionList]ready to add dp[%v] decommission src[%v] Disk[%v] dst[%v] status[%v] specialStep[%v],"+ + log.LogInfof("action[addToDecommissionList]ready to add dp[%v] decommission srcAddr[%v] Disk[%v] dstAddr[%v] dstNodeSet[%v] status[%v] specialStep[%v],"+ " RollbackTimes(%v) isRecover(%v) host[%v] to decommission list[%v]", partition.PartitionID, partition.DecommissionSrcAddr, partition.DecommissionSrcDiskPath, - partition.DecommissionDstAddr, partition.GetDecommissionStatus(), partition.GetSpecialReplicaDecommissionStep(), + partition.DecommissionDstAddr, partition.DecommissionDstNodeSet, partition.GetDecommissionStatus(), partition.GetSpecialReplicaDecommissionStep(), partition.DecommissionNeedRollbackTimes, partition.isRecover, partition.Hosts, ns.ID) ns.AddToDecommissionDataPartitionList(partition, c) } @@ -2195,13 +2206,23 @@ func (partition *DataPartition) TryAcquireDecommissionToken(c *Cluster) bool { // the first time for dst addr not specify if !partition.DecommissionDstAddrSpecify && partition.DecommissionDstAddr == "" { - // try to find available data node in src nodeset - ns, zone, err = getTargetNodeset(partition.DecommissionSrcAddr, c) - if err != nil { - log.LogWarnf("action[TryAcquireDecommissionToken] dp %v find src nodeset failed:%v", - partition.PartitionID, err.Error()) - goto errHandler + if partition.DecommissionDstNodeSet != 0 { + ns, err = c.t.getNodeSetByNodeSetId(partition.DecommissionDstNodeSet) + if err != nil { + log.LogWarnf("action[TryAcquireDecommissionToken]dp %v find given dst nodeset %v failed:%v", + partition.PartitionID, partition.DecommissionDstNodeSet, err.Error()) + goto errHandler + } + } else { + // try to find available data node in src nodeset + ns, zone, err = getTargetNodeset(partition.DecommissionSrcAddr, c) + if err != nil { + log.LogWarnf("action[TryAcquireDecommissionToken] dp %v find src nodeset failed:%v", + partition.PartitionID, err.Error()) + goto errHandler + } } + if partition.isSpecialReplicaCnt() && ns.HasDecommissionToken(partition.PartitionID) { log.LogDebugf("action[TryAcquireDecommissionToken]dp %v has token when reloading meta from nodeset %v", partition.PartitionID, ns.ID) @@ -2221,6 +2242,12 @@ func (partition *DataPartition) TryAcquireDecommissionToken(c *Cluster) bool { // data nodes in a nodeset has the same mediaType targetHosts, _, err = ns.getAvailDataNodeHosts(excludeHosts, 1) if err != nil { + if partition.DecommissionDstNodeSet != 0 { + log.LogWarnf("action[TryAcquireDecommissionToken] dp %v choose from given dst nodeset %v failed:%v", + partition.PartitionID, partition.DecommissionDstNodeSet, err.Error()) + goto errHandler + } + log.LogWarnf("action[TryAcquireDecommissionToken] dp %v choose from src nodeset failed:%v", partition.PartitionID, err.Error()) if _, ok := c.vols[partition.VolName]; !ok { @@ -2752,7 +2779,7 @@ func (partition *DataPartition) checkReplicaMeta(c *Cluster) (err error) { partition.PartitionID, addr) return nil } - err = c.markDecommissionDataPartition(partition, node, false, AutoAddReplica, highPriorityDecommissionWeight) + err = c.markDecommissionDataPartition(partition, node, 0, false, AutoAddReplica, highPriorityDecommissionWeight) auditMsg = fmt.Sprintf("dp(%v) ReplicaNum %v hostsNum %v auto add replica", partition.PartitionID, partition.ReplicaNum, len(partition.Hosts)) log.LogDebugf("action[checkReplicaMeta]%v: err %v", auditMsg, err) @@ -2803,9 +2830,9 @@ func (partition *DataPartition) decommissionInfo() string { replicas = append(replicas, replica.Addr) } - return fmt.Sprintf("vol(%v)_dp(%v)_replicaNum(%v)_src(%v)_dst(%v)_hosts(%v)_retry(%v)_isRecover(%v)_status(%v)_specialStatus(%v)"+ + return fmt.Sprintf("vol(%v)_dp(%v)_replicaNum(%v)_srcAddr(%v)_dstAddr(%v)_dstNodeSet(%v)_hosts(%v)_retry(%v)_isRecover(%v)_status(%v)_specialStatus(%v)"+ "_needRollback(%v)_rollbackTimes(%v)_force(%v)_type(%v)_RestoreReplica(%v)_errMsg(%v)_discard(%v)_term(%v)_weight(%v)_firstHostDiskTokenKey(%v)_replica(%v)_recoverStartTime(%v)_addr(%p)", - partition.VolName, partition.PartitionID, partition.ReplicaNum, partition.DecommissionSrcAddr, partition.DecommissionDstAddr, + partition.VolName, partition.PartitionID, partition.ReplicaNum, partition.DecommissionSrcAddr, partition.DecommissionDstAddr, partition.DecommissionDstNodeSet, partition.Hosts, partition.DecommissionRetry, partition.isRecover, GetDecommissionStatusMessage(partition.GetDecommissionStatus()), GetSpecialDecommissionStatusMessage(partition.GetSpecialReplicaDecommissionStep()), partition.DecommissionNeedRollback, partition.DecommissionNeedRollbackTimes, partition.DecommissionRaftForce, GetDecommissionTypeMessage(partition.DecommissionType), diff --git a/master/metadata_fsm_op.go b/master/metadata_fsm_op.go index 25abb57a0..f7e2ea7b9 100644 --- a/master/metadata_fsm_op.go +++ b/master/metadata_fsm_op.go @@ -195,6 +195,7 @@ type dataPartitionValue struct { DecommissionWeight int SpecialReplicaDecommissionStep uint32 DecommissionDstAddrSpecify bool + DecommissionDstNodeSet uint64 DecommissionNeedRollback bool RecoverStartTime int64 RecoverUpdateTime int64 @@ -233,6 +234,7 @@ func (dpv *dataPartitionValue) Restore(c *Cluster) (dp *DataPartition) { dp.DecommissionWeight = dpv.DecommissionWeight dp.SpecialReplicaDecommissionStep = dpv.SpecialReplicaDecommissionStep dp.DecommissionDstAddrSpecify = dpv.DecommissionDstAddrSpecify + dp.DecommissionDstNodeSet = dpv.DecommissionDstNodeSet dp.DecommissionNeedRollback = dpv.DecommissionNeedRollback dp.RecoverStartTime = time.Unix(dpv.RecoverStartTime, 0) dp.RecoverUpdateTime = time.Unix(dpv.RecoverUpdateTime, 0) @@ -292,6 +294,7 @@ func newDataPartitionValue(dp *DataPartition) (dpv *dataPartitionValue) { DecommissionWeight: dp.DecommissionWeight, SpecialReplicaDecommissionStep: dp.SpecialReplicaDecommissionStep, DecommissionDstAddrSpecify: dp.DecommissionDstAddrSpecify, + DecommissionDstNodeSet: dp.DecommissionDstNodeSet, DecommissionNeedRollback: dp.DecommissionNeedRollback, RecoverStartTime: dp.RecoverStartTime.Unix(), RecoverUpdateTime: dp.RecoverUpdateTime.Unix(), diff --git a/proto/model.go b/proto/model.go index 859fe6778..e3e7bd906 100644 --- a/proto/model.go +++ b/proto/model.go @@ -600,6 +600,7 @@ type DecommissionDataPartitionInfo struct { SrcAddress string SrcDiskPath string DstAddress string + DstNodeSet uint64 Term uint64 Weight int Replicas []string diff --git a/sdk/master/api_admin.go b/sdk/master/api_admin.go index d53ae23c5..1b7f09d16 100644 --- a/sdk/master/api_admin.go +++ b/sdk/master/api_admin.go @@ -186,10 +186,13 @@ func (api *AdminAPI) CreateDataPartition(volName string, count int, clientIDKey )) } -func (api *AdminAPI) DecommissionDataPartition(dataPartitionID uint64, nodeAddr string, raftForce bool, weight int, clientIDKey, decommissionType string) (err error) { +func (api *AdminAPI) DecommissionDataPartition(dataPartitionID uint64, nodeAddr string, dstNodeSet uint64, raftForce bool, weight int, clientIDKey, decommissionType string) (err error) { request := newRequest(get, proto.AdminDecommissionDataPartition).Header(api.h) request.addParam("id", strconv.FormatUint(dataPartitionID, 10)) request.addParam("addr", nodeAddr) + if dstNodeSet != 0 { + request.addParam("dstNodeSet", strconv.FormatUint(dstNodeSet, 10)) + } request.addParam("raftForceDel", strconv.FormatBool(raftForce)) request.addParam("weight", strconv.Itoa(weight)) request.addParam("clientIDKey", clientIDKey)