aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorKonstantin Lebedev <9497591+kmlebedev@users.noreply.github.com>2022-07-12 11:33:08 +0500
committerKonstantin Lebedev <9497591+kmlebedev@users.noreply.github.com>2022-07-12 11:33:08 +0500
commit087fa1347f6348b829b95a03877367667cea5de8 (patch)
tree43a70ad58738ce524dfa1ea5ba29571d407ce439
parent4236c3659977dc38e42e64fc1fe74999b2e8d693 (diff)
downloadseaweedfs-087fa1347f6348b829b95a03877367667cea5de8.tar.xz
seaweedfs-087fa1347f6348b829b95a03877367667cea5de8.zip
volume server evacuate from rack
-rw-r--r--weed/shell/command_volume_server_evacuate.go103
1 files changed, 56 insertions, 47 deletions
diff --git a/weed/shell/command_volume_server_evacuate.go b/weed/shell/command_volume_server_evacuate.go
index f2c24a8b4..37fb29b14 100644
--- a/weed/shell/command_volume_server_evacuate.go
+++ b/weed/shell/command_volume_server_evacuate.go
@@ -19,6 +19,7 @@ func init() {
type commandVolumeServerEvacuate struct {
targetServer string
+ volumeRack string
}
func (c *commandVolumeServerEvacuate) Name() string {
@@ -47,6 +48,7 @@ func (c *commandVolumeServerEvacuate) Do(args []string, commandEnv *CommandEnv,
vsEvacuateCommand := flag.NewFlagSet(c.Name(), flag.ContinueOnError)
volumeServer := vsEvacuateCommand.String("node", "", "<host>:<port> of the volume server")
+ volumeRack := vsEvacuateCommand.String("rack", "", "rack for then volume servers")
targetServer := vsEvacuateCommand.String("target", "", "<host>:<port> of target volume")
skipNonMoveable := vsEvacuateCommand.Bool("skipNonMoveable", false, "skip volumes that can not be moved")
applyChange := vsEvacuateCommand.Bool("force", false, "actually apply the changes")
@@ -66,6 +68,9 @@ func (c *commandVolumeServerEvacuate) Do(args []string, commandEnv *CommandEnv,
if *targetServer != "" {
c.targetServer = *targetServer
}
+ if *volumeRack != "" {
+ c.volumeRack = *volumeRack
+ }
for i := 0; i < *retryCount+1; i++ {
if err = c.volumeServerEvacuate(commandEnv, *volumeServer, *skipNonMoveable, *applyChange, writer); err == nil {
return nil
@@ -102,41 +107,31 @@ func (c *commandVolumeServerEvacuate) volumeServerEvacuate(commandEnv *CommandEn
func (c *commandVolumeServerEvacuate) evacuateNormalVolumes(commandEnv *CommandEnv, topologyInfo *master_pb.TopologyInfo, volumeServer string, skipNonMoveable, applyChange bool, writer io.Writer) error {
// find this volume server
volumeServers := collectVolumeServersByDc(topologyInfo, "")
- thisNode, otherNodes := nodesOtherThan(volumeServers, volumeServer)
- if thisNode == nil {
+ thisNodes, otherNodes := c.nodesOtherThan(volumeServers, volumeServer)
+ if len(thisNodes) == 0 {
return fmt.Errorf("%s is not found in this cluster", volumeServer)
}
- if c.targetServer != "" {
- targetServerFound := false
- for _, otherNode := range otherNodes {
- if otherNode.info.Id == c.targetServer {
- otherNodes = []*Node{otherNode}
- targetServerFound = true
- break
- }
- }
- if !targetServerFound {
- return fmt.Errorf("target %s is not found in this cluster", c.targetServer)
- }
- }
// move away normal volumes
- volumeReplicas, _ := collectVolumeReplicaLocations(topologyInfo)
- for _, diskInfo := range thisNode.info.DiskInfos {
- for _, vol := range diskInfo.VolumeInfos {
- hasMoved, err := c.moveAwayOneNormalVolume(commandEnv, volumeReplicas, vol, thisNode, otherNodes, applyChange)
- if err != nil {
- fmt.Fprintf(writer, "move away volume %d from %s: %v", vol.Id, volumeServer, err)
- }
- if !hasMoved {
- if skipNonMoveable {
- replicaPlacement, _ := super_block.NewReplicaPlacementFromByte(byte(vol.ReplicaPlacement))
- fmt.Fprintf(writer, "skipping non moveable volume %d replication:%s\n", vol.Id, replicaPlacement.String())
- } else {
- return fmt.Errorf("failed to move volume %d from %s", vol.Id, volumeServer)
+ for _, thisNode := range thisNodes {
+ for _, diskInfo := range thisNode.info.DiskInfos {
+ volumeReplicas, _ := collectVolumeReplicaLocations(topologyInfo)
+ for _, vol := range diskInfo.VolumeInfos {
+ hasMoved, err := c.moveAwayOneNormalVolume(commandEnv, volumeReplicas, vol, thisNode, otherNodes, applyChange)
+ if err != nil {
+ fmt.Fprintf(writer, "move away volume %d from %s: %v", vol.Id, volumeServer, err)
+ }
+ if !hasMoved {
+ if skipNonMoveable {
+ replicaPlacement, _ := super_block.NewReplicaPlacementFromByte(byte(vol.ReplicaPlacement))
+ fmt.Fprintf(writer, "skipping non moveable volume %d replication:%s\n", vol.Id, replicaPlacement.String())
+ } else {
+ return fmt.Errorf("failed to move volume %d from %s", vol.Id, volumeServer)
+ }
}
}
}
+
}
return nil
}
@@ -144,23 +139,25 @@ func (c *commandVolumeServerEvacuate) evacuateNormalVolumes(commandEnv *CommandE
func (c *commandVolumeServerEvacuate) evacuateEcVolumes(commandEnv *CommandEnv, topologyInfo *master_pb.TopologyInfo, volumeServer string, skipNonMoveable, applyChange bool, writer io.Writer) error {
// find this ec volume server
ecNodes, _ := collectEcVolumeServersByDc(topologyInfo, "")
- thisNode, otherNodes := ecNodesOtherThan(ecNodes, volumeServer)
- if thisNode == nil {
+ thisNodes, otherNodes := c.ecNodesOtherThan(ecNodes, volumeServer)
+ if len(thisNodes) == 0 {
return fmt.Errorf("%s is not found in this cluster\n", volumeServer)
}
// move away ec volumes
- for _, diskInfo := range thisNode.info.DiskInfos {
- for _, ecShardInfo := range diskInfo.EcShardInfos {
- hasMoved, err := c.moveAwayOneEcVolume(commandEnv, ecShardInfo, thisNode, otherNodes, applyChange)
- if err != nil {
- fmt.Fprintf(writer, "move away volume %d from %s: %v", ecShardInfo.Id, volumeServer, err)
- }
- if !hasMoved {
- if skipNonMoveable {
- fmt.Fprintf(writer, "failed to move away ec volume %d from %s\n", ecShardInfo.Id, volumeServer)
- } else {
- return fmt.Errorf("failed to move away ec volume %d from %s", ecShardInfo.Id, volumeServer)
+ for _, thisNode := range thisNodes {
+ for _, diskInfo := range thisNode.info.DiskInfos {
+ for _, ecShardInfo := range diskInfo.EcShardInfos {
+ hasMoved, err := c.moveAwayOneEcVolume(commandEnv, ecShardInfo, thisNode, otherNodes, applyChange)
+ if err != nil {
+ fmt.Fprintf(writer, "move away volume %d from %s: %v", ecShardInfo.Id, volumeServer, err)
+ }
+ if !hasMoved {
+ if skipNonMoveable {
+ fmt.Fprintf(writer, "failed to move away ec volume %d from %s\n", ecShardInfo.Id, volumeServer)
+ } else {
+ return fmt.Errorf("failed to move away ec volume %d from %s", ecShardInfo.Id, volumeServer)
+ }
}
}
}
@@ -220,10 +217,16 @@ func (c *commandVolumeServerEvacuate) moveAwayOneNormalVolume(commandEnv *Comman
return
}
-func nodesOtherThan(volumeServers []*Node, thisServer string) (thisNode *Node, otherNodes []*Node) {
+func (c *commandVolumeServerEvacuate) nodesOtherThan(volumeServers []*Node, thisServer string) (thisNodes []*Node, otherNodes []*Node) {
for _, node := range volumeServers {
- if node.info.Id == thisServer {
- thisNode = node
+ if node.info.Id == thisServer || (c.volumeRack != "" && node.rack == c.volumeRack) {
+ thisNodes = append(thisNodes, node)
+ continue
+ }
+ if c.volumeRack != "" && c.volumeRack == node.rack {
+ continue
+ }
+ if c.targetServer != "" && c.targetServer != node.info.Id {
continue
}
otherNodes = append(otherNodes, node)
@@ -231,10 +234,16 @@ func nodesOtherThan(volumeServers []*Node, thisServer string) (thisNode *Node, o
return
}
-func ecNodesOtherThan(volumeServers []*EcNode, thisServer string) (thisNode *EcNode, otherNodes []*EcNode) {
+func (c *commandVolumeServerEvacuate) ecNodesOtherThan(volumeServers []*EcNode, thisServer string) (thisNodes []*EcNode, otherNodes []*EcNode) {
for _, node := range volumeServers {
- if node.info.Id == thisServer {
- thisNode = node
+ if node.info.Id == thisServer || (c.volumeRack != "" && string(node.rack) == c.volumeRack) {
+ thisNodes = append(thisNodes, node)
+ continue
+ }
+ if c.volumeRack != "" && c.volumeRack == string(node.rack) {
+ continue
+ }
+ if c.targetServer != "" && c.targetServer != node.info.Id {
continue
}
otherNodes = append(otherNodes, node)