From 70d9dd5afe4b1c4b25331e75e0506649e7efab45 Mon Sep 17 00:00:00 2001 From: msementsov <47177265+m-sementsov@users.noreply.github.com> Date: Tue, 23 Jun 2026 20:48:33 +0300 Subject: [PATCH] volume.balance: add -volumesPerExec to cap moves per run Limit the number of volume moves performed in one command execution; re-run to continue. 0 = unlimited. --- weed/shell/command_volume_balance.go | 45 +++++++++++++----- weed/shell/command_volume_balance_test.go | 56 ++++++++++++++++++++++- 2 files changed, 89 insertions(+), 12 deletions(-) diff --git a/weed/shell/command_volume_balance.go b/weed/shell/command_volume_balance.go index d3c88df9b..3e2b7aff6 100644 --- a/weed/shell/command_volume_balance.go +++ b/weed/shell/command_volume_balance.go @@ -36,6 +36,8 @@ type commandVolumeBalance struct { commandEnv *CommandEnv volumeByActive *bool applyBalancing bool + volumesPerExec int + movedCount int } func (c *commandVolumeBalance) Name() string { @@ -45,7 +47,7 @@ func (c *commandVolumeBalance) Name() string { func (c *commandVolumeBalance) Help() string { return `balance all volumes among volume servers - volume.balance [-collection ALL_COLLECTIONS|EACH_COLLECTION|] [-apply] [-dataCenter=] [-racks=rack_name_one,rack_name_two] [-nodes=192.168.0.1:8080,192.168.0.2:8080] + volume.balance [-collection ALL_COLLECTIONS|EACH_COLLECTION|] [-apply] [-dataCenter=] [-racks=rack_name_one,rack_name_two] [-nodes=192.168.0.1:8080,192.168.0.2:8080] [-volumesPerExec=5] The -collection parameter supports: - ALL_COLLECTIONS: balance across all collections @@ -55,6 +57,10 @@ func (c *commandVolumeBalance) Help() string { * Match multiple buckets: volume.balance -collection="bucket.*" * Match all user collections: volume.balance -collection="user-.*" + The -volumesPerExec parameter limits the maximum number of volume moves in one command execution. + If unset - the command will try to balance all volumes at once. + It might be beneficial to set, if your cluster has lots of volumes growing and topology changes faster than balancing can occur. + Algorithm: For each type of volume server (different max volume count limit){ @@ -106,6 +112,8 @@ func (c *commandVolumeBalance) Do(args []string, commandEnv *CommandEnv, writer applyBalancing := balanceCommand.Bool("apply", false, "apply the balancing plan.") // TODO: remove this alias applyBalancingAlias := balanceCommand.Bool("force", false, "apply the balancing plan (alias for -apply)") + volumesPerExec := balanceCommand.Int("volumesPerExec", 0, "how many volumes to move in one run (default is 0 for unlimited)") + balanceCommand.Func("volumeBy", "only apply the balancing for ALL volumes and ACTIVE or FULL", func(flagValue string) error { if flagValue == "" { return nil @@ -123,6 +131,11 @@ func (c *commandVolumeBalance) Do(args []string, commandEnv *CommandEnv, writer } handleDeprecatedForceFlag(writer, balanceCommand, applyBalancingAlias, applyBalancing) c.applyBalancing = *applyBalancing + if *volumesPerExec < 0 { + return fmt.Errorf("volumesPerExec must be >= 0") + } + c.volumesPerExec = *volumesPerExec + c.movedCount = 0 infoAboutSimulationMode(writer, c.applyBalancing, "-apply") @@ -153,6 +166,9 @@ func (c *commandVolumeBalance) Do(args []string, commandEnv *CommandEnv, writer return err } for _, col := range collections { + if c.volumesPerExec > 0 && c.movedCount >= c.volumesPerExec { + break + } // Use direct string comparison for exact match (more efficient than regex) if err = c.balanceVolumeServers(diskTypes, volumeReplicas, volumeServers, nil, col); err != nil { return err @@ -179,6 +195,9 @@ func (c *commandVolumeBalance) Do(args []string, commandEnv *CommandEnv, writer func (c *commandVolumeBalance) balanceVolumeServers(diskTypes []types.DiskType, volumeReplicas map[uint32][]*VolumeReplica, nodes []*Node, collectionPattern *regexp.Regexp, collectionName string) error { for _, diskType := range diskTypes { + if c.volumesPerExec > 0 && c.movedCount >= c.volumesPerExec { + break + } if err := c.balanceVolumeServersByDiskType(diskType, volumeReplicas, nodes, collectionPattern, collectionName); err != nil { return err } @@ -208,7 +227,7 @@ func (c *commandVolumeBalance) balanceVolumeServersByDiskType(diskType types.Dis return selectVolumesByActive(v.Size, c.volumeByActive, c.volumeSizeLimitMb) }) } - if err := balanceSelectedVolume(c.commandEnv, diskType, volumeReplicas, nodes, sortWritableVolumes, c.volumeSizeLimitMb, c.applyBalancing); err != nil { + if err := c.balanceSelectedVolume(diskType, volumeReplicas, nodes, sortWritableVolumes); err != nil { return err } @@ -392,9 +411,10 @@ func selectVolumesByActive(volumeSize uint64, volumeByActive *bool, volumeSizeLi } } -func balanceSelectedVolume(commandEnv *CommandEnv, diskType types.DiskType, volumeReplicas map[uint32][]*VolumeReplica, nodes []*Node, sortCandidatesFn func(volumes []*master_pb.VolumeInformationMessage), volumeSizeLimitMb uint64, applyBalancing bool) (err error) { +func (c *commandVolumeBalance) balanceSelectedVolume(diskType types.DiskType, volumeReplicas map[uint32][]*VolumeReplica, nodes []*Node, sortCandidatesFn func(volumes []*master_pb.VolumeInformationMessage)) (err error) { selectedVolumeCount, volumeCapacities := uint64(0), float64(0) var nodesWithCapacity []*Node + volumeSizeLimitMb := c.volumeSizeLimitMb if volumeSizeLimitMb == 0 { volumeSizeLimitMb = util.VolumeSizeLimitGB * util.KiByte } @@ -414,16 +434,19 @@ func balanceSelectedVolume(commandEnv *CommandEnv, diskType types.DiskType, volu hasMoved := true - if commandEnv != nil && commandEnv.verbose { + if c.commandEnv != nil && c.commandEnv.verbose { fmt.Fprintf(os.Stdout, "selected nodes %d, volumes:%d, cap:%d, idealVolumeRatio %f\n", len(nodesWithCapacity), selectedVolumeCount, int64(volumeCapacities), idealVolumeRatio*100) } for hasMoved { hasMoved = false + if c.volumesPerExec > 0 && c.movedCount >= c.volumesPerExec { + break + } slices.SortFunc(nodesWithCapacity, func(a, b *Node) int { return cmp.Compare(a.localVolumeDensityRatio(capacityFunc), b.localVolumeDensityRatio(capacityFunc)) }) if len(nodesWithCapacity) == 0 { - if commandEnv != nil && commandEnv.verbose { + if c.commandEnv != nil && c.commandEnv.verbose { fmt.Fprintf(os.Stdout, "no volume server found with capacity for %s", diskType.ReadableString()) } return nil @@ -445,7 +468,7 @@ func balanceSelectedVolume(commandEnv *CommandEnv, diskType types.DiskType, volu candidateVolumes = append(candidateVolumes, v) } if fullNodeIndex == -1 { - if commandEnv != nil && commandEnv.verbose { + if c.commandEnv != nil && c.commandEnv.verbose { fmt.Fprintf(os.Stdout, "no nodes with capacity found for %s, nodes %d", diskType.ReadableString(), len(nodesWithCapacity)) } return nil @@ -453,20 +476,20 @@ func balanceSelectedVolume(commandEnv *CommandEnv, diskType types.DiskType, volu sortCandidatesFn(candidateVolumes) for _, emptyNode := range nodesWithCapacity[:fullNodeIndex] { if !(fullNode.localVolumeDensityNextRatio(capacityFunc) > idealVolumeRatio && emptyNode.localVolumeDensityNextRatio(capacityFunc) <= idealVolumeRatio) { - if commandEnv != nil && commandEnv.verbose { + if c.commandEnv != nil && c.commandEnv.verbose { fmt.Printf("no more volume servers with empty slots %s, idealVolumeRatio %f\n", emptyNode.info.Id, idealVolumeRatio) } break } fmt.Fprintf(os.Stdout, "%s %.2f %.2f:%.2f\t", diskType.ReadableString(), idealVolumeRatio, fullNode.localVolumeDensityRatio(capacityFunc), emptyNode.localVolumeDensityNextRatio(capacityFunc)) - if commandEnv != nil && commandEnv.verbose { + if c.commandEnv != nil && c.commandEnv.verbose { fmt.Fprintf(os.Stdout, "%s %.1f %.1f:%.1f\t", diskType.ReadableString(), idealVolumeRatio*100, fullNode.localVolumeDensityRatio(capacityFunc)*100, emptyNode.localVolumeDensityNextRatio(capacityFunc)*100) } - hasMoved, err = attemptToMoveOneVolume(commandEnv, volumeReplicas, fullNode, candidateVolumes, emptyNode, applyBalancing) + hasMoved, err = attemptToMoveOneVolume(c.commandEnv, volumeReplicas, fullNode, candidateVolumes, emptyNode, c.applyBalancing) if err != nil { - if commandEnv != nil && commandEnv.verbose { + if c.commandEnv != nil && c.commandEnv.verbose { fmt.Fprintf(os.Stdout, "attempt to move one volume error %+v\n", err) } if strings.Contains(err.Error(), util.ErrVolumeNoSpaceLeft) { @@ -475,7 +498,7 @@ func balanceSelectedVolume(commandEnv *CommandEnv, diskType types.DiskType, volu return } if hasMoved { - // moved one volume + c.movedCount++ break } } diff --git a/weed/shell/command_volume_balance_test.go b/weed/shell/command_volume_balance_test.go index 05f72511e..41e7c5fa1 100644 --- a/weed/shell/command_volume_balance_test.go +++ b/weed/shell/command_volume_balance_test.go @@ -341,7 +341,8 @@ func TestBalanceDoesNotDrainOntoOneNode(t *testing.T) { n.selectVolumes(func(v *master_pb.VolumeInformationMessage) bool { return true }) } - if err := balanceSelectedVolume(nil, types.HardDriveType, volumeReplicas, nodes, sortWritableVolumes, volumeSizeLimitMb, false); err != nil { + c := &commandVolumeBalance{volumeSizeLimitMb: volumeSizeLimitMb} + if err := c.balanceSelectedVolume(types.HardDriveType, volumeReplicas, nodes, sortWritableVolumes); err != nil { t.Fatalf("balanceSelectedVolume: %v", err) } @@ -355,6 +356,59 @@ func TestBalanceDoesNotDrainOntoOneNode(t *testing.T) { } } +// volumesPerExec caps the number of moves performed in a single execution. +func TestBalanceVolumesPerExec(t *testing.T) { + const mb = 1024 * 1024 + volumeSizeLimitMb := uint64(100) + + makeNode := func(id string, volumes []*master_pb.VolumeInformationMessage) *Node { + return &Node{ + info: &master_pb.DataNodeInfo{ + Id: id, + DiskInfos: map[string]*master_pb.DiskInfo{ + "": { + MaxVolumeCount: 10, + VolumeCount: int64(len(volumes)), + VolumeInfos: volumes, + }, + }, + }, + dc: "dc1", + rack: "rack1", + } + } + + var fullVolumes []*master_pb.VolumeInformationMessage + for id := uint32(1); id <= 6; id++ { + fullVolumes = append(fullVolumes, &master_pb.VolumeInformationMessage{Id: id, Size: 95 * mb}) + } + fullNode := makeNode("full", fullVolumes) + emptyNode := makeNode("empty", nil) + nodes := []*Node{fullNode, emptyNode} + + volumeReplicas := map[uint32][]*VolumeReplica{} + for _, v := range fullVolumes { + loc := newLocation("dc1", "rack1", fullNode.info) + volumeReplicas[v.Id] = []*VolumeReplica{{location: &loc, info: v}} + } + + for _, n := range nodes { + n.selectVolumes(func(v *master_pb.VolumeInformationMessage) bool { return true }) + } + + c := &commandVolumeBalance{volumeSizeLimitMb: volumeSizeLimitMb, volumesPerExec: 1} + if err := c.balanceSelectedVolume(types.HardDriveType, volumeReplicas, nodes, sortWritableVolumes); err != nil { + t.Fatalf("balanceSelectedVolume: %v", err) + } + + if c.movedCount != 1 { + t.Fatalf("expected exactly 1 move with volumesPerExec=1, got %d", c.movedCount) + } + if got := len(emptyNode.info.DiskInfos[""].VolumeInfos); got != 1 { + t.Fatalf("expected empty node to receive exactly 1 volume, got %d", got) + } +} + func TestVolumeSelection(t *testing.T) { topologyInfo := parseOutput(topoData)