From 8c1be63c92867543f722c44d1d8b58aa1df47d59 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Fri, 25 Sep 2026 02:14:23 +0800 Subject: [PATCH] ecbalancer: let a non-overflow parity shard leave a data-bearing rack (#11438) The parity pass only queued shards past the per-type cap, so a single parity shard sharing a rack with data was never a move candidate even when an empty data-free rack existed (2+1 over 3 DCs settled 2/1/0). Non-overflow candidates now move too, but only to a rack without data; overflow shards keep the existing data-rack fallback. --- .../erasure_coding/ecbalancer/balancer.go | 25 ++++-- .../ecbalancer/small_ratio_regression_test.go | 87 +++++++++++++++++++ 2 files changed, 104 insertions(+), 8 deletions(-) create mode 100644 weed/storage/erasure_coding/ecbalancer/small_ratio_regression_test.go diff --git a/weed/storage/erasure_coding/ecbalancer/balancer.go b/weed/storage/erasure_coding/ecbalancer/balancer.go index 74cf233b2..59cdc51ac 100644 --- a/weed/storage/erasure_coding/ecbalancer/balancer.go +++ b/weed/storage/erasure_coding/ecbalancer/balancer.go @@ -385,19 +385,25 @@ func balanceShardTypeAcrossRacks(vk volKey, nodes map[string]*Node, racks map[st rackKeys := sortedKeys(racks) type pending struct { - shardID int - src *Node + shardID int + src *Node + avoidDataRack bool } var toMove []pending for _, rackID := range rackKeys { shards := append([]int(nil), shardsPerRack[rackID]...) - if len(shards) <= maxPerRack { - continue - } sort.Ints(shards) - for i := 0; i < len(shards)-maxPerRack; i++ { + overflow := max(0, len(shards)-maxPerRack) + for i := 0; i < len(shards); i++ { + // A parity shard can fit the per-type cap yet share a rack with + // data while a data-free rack is empty; such candidates may only + // move to a rack without data. + avoidDataRack := i >= overflow + if avoidDataRack && !antiAffinity[rackID] { + continue + } if src := nodeInRackHoldingShard(nodes, rackID, vk, shards[i]); src != nil { - toMove = append(toMove, pending{shards[i], src}) + toMove = append(toMove, pending{shards[i], src, avoidDataRack}) } } } @@ -405,7 +411,10 @@ func balanceShardTypeAcrossRacks(vk volKey, nodes map[string]*Node, racks map[st var moves []*move for _, pm := range toMove { destRack, ok := pickTarget(rackKeys, shardsPerRack, maxPerRack, antiAffinity, - func(r string) bool { return racks[r].freeSlots > 0 }, + func(r string) bool { + return r != pm.src.rack && racks[r].freeSlots > 0 && + (!pm.avoidDataRack || !antiAffinity[r]) + }, func(r string) bool { if rp == nil { return true diff --git a/weed/storage/erasure_coding/ecbalancer/small_ratio_regression_test.go b/weed/storage/erasure_coding/ecbalancer/small_ratio_regression_test.go new file mode 100644 index 000000000..948632756 --- /dev/null +++ b/weed/storage/erasure_coding/ecbalancer/small_ratio_regression_test.go @@ -0,0 +1,87 @@ +package ecbalancer + +import ( + "fmt" + "testing" + + "github.com/seaweedfs/seaweedfs/weed/storage/super_block" +) + +// Test every initial assignment of the three unique shards to three healthy DCs. +// This includes all shards on the encoder and the already-balanced-but-unsafe +// state with data+parity in one DC, data in another, and the third empty. +func TestSmallRatioAcrossThreeDCs(t *testing.T) { + for _, placement := range []string{"000", "200", "211", "011"} { + for a := 0; a < 3; a++ { + for b := 0; b < 3; b++ { + for p := 0; p < 3; p++ { + t.Run(fmt.Sprintf("%s/%d%d%d", placement, a, b, p), func(t *testing.T) { + topo := NewTopology() + for i := 0; i < 3; i++ { + id := fmt.Sprintf("node%d", i) + dc := fmt.Sprintf("dc%d", i) + n := topo.AddNode(id, dc, dc+":rack1", 100) + n.AddDisk(0, "ssd", 100, 0) + } + for shard, node := range []int{a, b, p} { + topo.nodes[fmt.Sprintf("node%d", node)].AddShards(7, "test", 0, bits(shard)) + } + for _, node := range topo.nodes { + occupied := volumeShardCount(node, volKey{collection: "test", vid: 7}) + node.freeSlots -= occupied + node.disks[0].freeSlots -= occupied + node.disks[0].shardCount = occupied + } + rp, err := super_block.NewReplicaPlacementFromString(placement) + if err != nil { + t.Fatal(err) + } + // Use the volume-specific ratio reported by Enterprise, even + // when the collection fallback has the OSS default of 10+4. + opts := Options{DiskType: "ssd", ReplicaPlacement: rp, Ratio: ratio(10, 4), + VolumeRatio: func(string, uint32) (int, int) { return 2, 1 }} + moves := Plan(topo, opts) + vk := volKey{collection: "test", vid: 7} + for id, node := range topo.nodes { + if got := volumeShardCount(node, vk); got != 1 { + t.Errorf("%s holds %d shards, want 1; moves=%+v", id, got, moves) + } + if node.freeSlots != 99 || node.disks[0].freeSlots != 99 || node.disks[0].shardCount != 1 { + t.Errorf("%s: inconsistent capacity after placement: node=%d disk=%+v", id, node.freeSlots, node.disks[0]) + } + } + for shard := 0; shard < 3; shard++ { + copies := 0 + for _, node := range topo.nodes { + if node.shards[vk] != nil && node.shards[vk].shardBits&bits(shard) != 0 { + copies++ + } + } + if copies != 1 { + t.Errorf("shard %d: got %d copies, want 1", shard, copies) + } + } + if next := Plan(topo, opts); len(next) != 0 { + t.Errorf("placement did not converge: %+v", next) + } + }) + } + } + } + } +} + +func TestSmallRatioDoesNotMoveParityToFullRack(t *testing.T) { + topo := NewTopology() + a := topo.AddNode("a", "a", "a:r", 100) + a.AddDisk(0, "ssd", 100, 2) + a.AddShards(1, "test", 0, bits(0, 2)) + b := topo.AddNode("b", "b", "b:r", 100) + b.AddDisk(0, "ssd", 100, 1) + b.AddShards(1, "test", 0, bits(1)) + c := topo.AddNode("c", "c", "c:r", 0) + c.AddDisk(0, "ssd", 0, 0) + if moves := Plan(topo, Options{DiskType: "ssd", Ratio: ratio(2, 1)}); len(moves) != 0 { + t.Fatalf("no eligible third rack: must not churn or move into a full destination: %+v", moves) + } +}