feat(master): the dp decommission supports specifying the target nodeSet.

close:#1000211002

Signed-off-by: shuqiang-zheng <zhengshuqiang@oppo.com>
This commit is contained in:
shuqiang-zheng 2025-07-02 17:50:53 +08:00 committed by zhumingze1108
parent 9b3f91b475
commit 175ee90983
10 changed files with 82 additions and 21 deletions

View File

@ -154,6 +154,7 @@ const (
CliFlagEnablePersistAccessTime = "enablePersistAccessTime"
CliFlagDecommissionRaftForce = "raftForceDel"
CliFLagDecommissionWeight = "decommissionWeight"
CliFlagDecommissionDstNodeSet = "decommissionDstNodeSet"
CliFLagRecommissionType = "recommissionType"
CliFlagAllowedStorageClass = "allowedStorageClass"
CliFlagVolStorageClass = "volStorageClass"

View File

@ -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

View File

@ -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))

View File

@ -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 == "" {

View File

@ -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()

View File

@ -88,6 +88,7 @@ const (
forceKey = "force"
raftForceDelKey = "raftForceDel"
weightKey = "weight"
dstNodeSetKey = "dstNodeSet"
enablePosixAclKey = "enablePosixAcl"
enableTxMaskKey = "enableTxMask"
txTimeoutKey = "txTimeout"

View File

@ -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),

View File

@ -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(),

View File

@ -600,6 +600,7 @@ type DecommissionDataPartitionInfo struct {
SrcAddress string
SrcDiskPath string
DstAddress string
DstNodeSet uint64
Term uint64
Weight int
Replicas []string

View File

@ -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)