diff --git a/weed/ec/ec_common.go b/weed/ec/ec_common.go index c8ab8a90c..e7cb27b17 100644 --- a/weed/ec/ec_common.go +++ b/weed/ec/ec_common.go @@ -883,6 +883,9 @@ func (ecb *ecBalancer) balance(collections []string) error { // nodes fill proportionally (matching the worker). This is identical to raw // shard count when capacities are uniform. GlobalUtilizationBased: true, + RackTotalCapRaised: func(collection string, vid uint32, rackCap, evenCap int) { + fmt.Printf("ec volume %d (collection %q): rack-total-cap raised to %d shards per rack, above the even %d: the racks lack room for an even spread\n", vid, collection, rackCap, evenCap) + }, }) if len(ecb.volumeIds) > 0 { var deletions int diff --git a/weed/storage/erasure_coding/ecbalancer/balancer.go b/weed/storage/erasure_coding/ecbalancer/balancer.go index 59cdc51ac..feac3756d 100644 --- a/weed/storage/erasure_coding/ecbalancer/balancer.go +++ b/weed/storage/erasure_coding/ecbalancer/balancer.go @@ -106,6 +106,12 @@ type Options struct { // heterogeneous-capacity racks; when false, by raw shard count. Both the worker // and the shell enable it; the two metrics agree when capacities are uniform. GlobalUtilizationBased bool + // RackTotalCapRaised, when set, is called for each volume whose + // total-shards-per-rack cap had to be raised above ceil(shards/racks) + // because the racks lack room for an even spread. Losing a rack then loses + // more shards than an even spread would, so callers log it as a degraded + // durability floor. It is called on every Plan while the condition holds. + RackTotalCapRaised func(collection string, vid uint32, rackCap, evenCap int) } // move is the internal form carrying node pointers; converted to Move at the end. @@ -249,7 +255,7 @@ func Plan(topo *Topology, opts Options) []Move { all = append(all, m...) } for _, vk := range byCollection[collection] { - m := detectCrossRackImbalance(vk, nodes, racks, opts.DiskType, opts.ImbalanceThreshold, dataShardsByVolume[vk], parityShardsByVolume[vk], opts.ReplicaPlacement) + m := detectCrossRackImbalance(vk, nodes, racks, opts.DiskType, opts.ImbalanceThreshold, dataShardsByVolume[vk], parityShardsByVolume[vk], opts.ReplicaPlacement, opts.RackTotalCapRaised) applyMovesToTopology(m, racks) all = append(all, m...) } @@ -343,42 +349,108 @@ func detectDuplicateShards(vk volKey, nodes map[string]*Node) []*move { } // detectCrossRackImbalance spreads a volume's shards across racks in two passes -// (data, then parity with anti-affinity to data-bearing racks). Returns nil if -// the overall distribution is below the imbalance threshold. -func detectCrossRackImbalance(vk volKey, nodes map[string]*Node, racks map[string]*rack, diskType string, threshold float64, dataShards, parityShards int, rp *super_block.ReplicaPlacement) []*move { +// (data, then parity with anti-affinity to data-bearing racks), keeping each +// rack's total within planRackTotalCap. Returns nil if the overall distribution +// is below the imbalance threshold and no rack is above the total cap. +func detectCrossRackImbalance(vk volKey, nodes map[string]*Node, racks map[string]*rack, diskType string, threshold float64, dataShards, parityShards int, rp *super_block.ReplicaPlacement, onCapRaised func(collection string, vid uint32, rackCap, evenCap int)) []*move { numRacks := len(racks) if numRacks <= 1 { return nil } + // The per-type caps below spread data and parity independently, so together + // they can still stack e.g. 2 data + 1 parity on one rack. Cap the TOTAL per + // rack too, at the lowest value the racks' capacity allows. + rackShardCount := countShardsByRack(vk, nodes) + totalShards := 0 + for _, n := range rackShardCount { + totalShards += n + } + totalCap := planRackTotalCap(vk, racks, rackShardCount, totalShards, rp) + if evenCap := ceilDivide(totalShards, numRacks); totalCap > evenCap && onCapRaised != nil { + onCapRaised(vk.collection, vk.vid, totalCap, evenCap) + } + aboveTotalCap := func() bool { + for _, n := range rackShardCount { + if n > totalCap { + return true + } + } + return false + } + // Gate on per-type spread: act when data OR parity shards are unevenly // distributed across racks, even if the per-rack totals happen to be even. + // A rack above the total cap is a durability bound, not skew, so it bypasses + // the threshold. gateData, gateParity := shardsByGroup(vk, nodes, dataShards, func(n *Node) string { return n.rack }) - if !typeImbalanced(gateData, numRacks, threshold) && !typeImbalanced(gateParity, numRacks, threshold) { + if !aboveTotalCap() && !typeImbalanced(gateData, numRacks, threshold) && !typeImbalanced(gateParity, numRacks, threshold) { return nil } - rackShardCount := countShardsByRack(vk, nodes) var moves []*move + // A shard moves at most once per plan: each move runs as its own task, + // possibly in parallel with the others, so a second hop could start before + // the first has landed. + moved := make(map[int]bool) + for { + before := len(moves) - dataPerRack, _ := shardsByGroup(vk, nodes, dataShards, func(n *Node) string { return n.rack }) - moves = append(moves, balanceShardTypeAcrossRacks(vk, nodes, racks, diskType, dataShards, - dataPerRack, rackShardCount, ceilDivide(dataShards, numRacks), nil, rp)...) + // The data pass leaves the total cap out (0): ceil(data/racks) <= totalCap + // already bounds data, and rackShardCount still includes parity that the + // parity pass is about to shed, so a total bound here would block data + // moves that fit once that parity has gone. + dataPerRack, _ := shardsByGroup(vk, nodes, dataShards, func(n *Node) string { return n.rack }) + moves = append(moves, balanceShardTypeAcrossRacks(vk, nodes, racks, diskType, dataShards, + dataPerRack, rackShardCount, ceilDivide(dataShards, numRacks), 0, nil, rp, moved)...) - dataPerRack, parityPerRack := shardsByGroup(vk, nodes, dataShards, func(n *Node) string { return n.rack }) - antiAffinity := make(map[string]bool) - for rackID, shards := range dataPerRack { - if len(shards) > 0 { - antiAffinity[rackID] = true + dataPerRack, parityPerRack := shardsByGroup(vk, nodes, dataShards, func(n *Node) string { return n.rack }) + antiAffinity := make(map[string]bool) + for rackID, shards := range dataPerRack { + if len(shards) > 0 { + antiAffinity[rackID] = true + } + } + moves = append(moves, balanceShardTypeAcrossRacks(vk, nodes, racks, diskType, dataShards, + parityPerRack, rackShardCount, ceilDivide(parityShards, numRacks), totalCap, antiAffinity, rp, moved)...) + + // The data pass can run out of destinations before the parity pass frees + // slots on other racks. While a rack is still above the total cap, run + // another round over the shards not moved yet. A round that continues + // moved at least one more shard, so this ends. + if len(moves) == before || !aboveTotalCap() { + return moves } } - moves = append(moves, balanceShardTypeAcrossRacks(vk, nodes, racks, diskType, dataShards, - parityPerRack, rackShardCount, ceilDivide(parityShards, numRacks), antiAffinity, rp)...) - - return moves } -func balanceShardTypeAcrossRacks(vk volKey, nodes map[string]*Node, racks map[string]*rack, diskType string, dataShards int, shardsPerRack map[string][]int, rackShardCount map[string]int, maxPerRack int, antiAffinity map[string]bool, rp *super_block.ReplicaPlacement) []*move { +// planRackTotalCap returns the lowest per-rack total that the racks can hold a +// volume's totalShards under (see rackTotalCap). Plan may move any shard, so a +// rack's shards of the volume count as room there rather than as held in place. +// Moves of the volume shift slots between racks without changing a rack's +// shards-plus-free-slots, so the cap holds for the whole cross-rack phase. +// Under SameRackCount a node's free slots count only up to the shards of the +// volume it may still take, the limit pickBestNodeForVolume enforces. +func planRackTotalCap(vk volKey, racks map[string]*rack, rackShardCount map[string]int, totalShards int, rp *super_block.ReplicaPlacement) int { + room := make(map[string]int, len(racks)) + for rk, r := range racks { + room[rk] = rackShardCount[rk] + if rp == nil || rp.SameRackCount <= 0 { + room[rk] += max(r.freeSlots, 0) + continue + } + for _, n := range r.nodes { + room[rk] += min(max(n.freeSlots, 0), max(rp.SameRackCount-volumeShardCount(n, vk), 0)) + } + } + return rackTotalCap(sortedKeys(racks), nil, room, totalShards) +} + +// balanceShardTypeAcrossRacks spreads one shard type across racks, at most +// maxPerRack of the type per rack. A totalCap > 0 also bounds each rack's TOTAL +// shards of the volume: racks above it shed this type, and no move lands on a +// rack at it. Shards in moved stay put; the ones it moves are added to it. +func balanceShardTypeAcrossRacks(vk volKey, nodes map[string]*Node, racks map[string]*rack, diskType string, dataShards int, shardsPerRack map[string][]int, rackShardCount map[string]int, maxPerRack, totalCap int, antiAffinity map[string]bool, rp *super_block.ReplicaPlacement, moved map[int]bool) []*move { if maxPerRack < 1 { maxPerRack = 1 } @@ -392,14 +464,25 @@ func balanceShardTypeAcrossRacks(vk volKey, nodes map[string]*Node, racks map[st var toMove []pending for _, rackID := range rackKeys { shards := append([]int(nil), shardsPerRack[rackID]...) - sort.Ints(shards) + // Shards this plan already moved go last, so the shed takes the others. + sort.Slice(shards, func(i, j int) bool { + if moved[shards[i]] != moved[shards[j]] { + return !moved[shards[i]] + } + return shards[i] < shards[j] + }) overflow := max(0, len(shards)-maxPerRack) + if totalCap > 0 { + // A rack above the total cap sheds this type until it fits, even + // when the type is within its own cap there. + overflow = max(overflow, min(len(shards), rackShardCount[rackID]-totalCap)) + } 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] { + if moved[shards[i]] || (avoidDataRack && !antiAffinity[rackID]) { continue } if src := nodeInRackHoldingShard(nodes, rackID, vk, shards[i]); src != nil { @@ -408,22 +491,32 @@ func balanceShardTypeAcrossRacks(vk volKey, nodes map[string]*Node, racks map[st } } + withinLimit := func(r string) bool { + if totalCap > 0 && rackShardCount[r] >= totalCap { + return false + } + if rp != nil && rp.DiffRackCount > 0 && rackShardCount[r] >= rp.DiffRackCount { + return false + } + return true + } + var moves []*move for _, pm := range toMove { - destRack, ok := pickTarget(rackKeys, shardsPerRack, maxPerRack, antiAffinity, - 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 - } - if rp.DiffRackCount > 0 && rackShardCount[r] >= rp.DiffRackCount { - return false - } - return true - }) + hasRoom := func(r string) bool { + return r != pm.src.rack && racks[r].freeSlots > 0 && + (!pm.avoidDataRack || !antiAffinity[r]) + } + destRack, ok := pickTarget(rackKeys, shardsPerRack, maxPerRack, antiAffinity, hasRoom, withinLimit) + if !ok && totalCap > 0 && rackShardCount[pm.src.rack] > totalCap { + // No rack is under the per-type cap. A shard must not stay on a rack + // above the total cap while another rack is under it, so relax the + // per-type cap up to the total cap; anti-affinity stays a preference + // inside pickTarget. Only a source above the total cap may relax: + // otherwise the destination just takes over the per-type overflow + // and the next Plan moves it back. + destRack, ok = pickTarget(rackKeys, shardsPerRack, totalCap, antiAffinity, hasRoom, withinLimit) + } if !ok { continue } @@ -455,6 +548,7 @@ func balanceShardTypeAcrossRacks(vk volKey, nodes map[string]*Node, racks map[st // limited-capacity destination across successive moves. destNode.freeSlots-- pm.src.freeSlots++ + moved[pm.shardID] = true } return moves } diff --git a/weed/storage/erasure_coding/ecbalancer/balancer_rack_cap_test.go b/weed/storage/erasure_coding/ecbalancer/balancer_rack_cap_test.go new file mode 100644 index 000000000..b83a616dd --- /dev/null +++ b/weed/storage/erasure_coding/ecbalancer/balancer_rack_cap_test.go @@ -0,0 +1,314 @@ +package ecbalancer + +import ( + "fmt" + "testing" + + "github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding" + "github.com/seaweedfs/seaweedfs/weed/storage/super_block" +) + +// Total-shards-per-rack cap for Plan. +// +// The per-type caps (ceil(data/racks) data, ceil(parity/racks) parity) spread +// each type evenly on their own, but together they allow e.g. 2 data + 1 parity +// on one rack. When a rack is one physical disk, losing two such racks loses 6 +// of a 10+4 volume's shards. Plan also caps the TOTAL per rack, sized like +// Place's rackTotalCap: ceil(shards/racks) unless the racks' capacity forces it +// higher. + +const capTestVid = uint32(1) + +var capTestVk = volKey{collection: "c1", vid: capTestVid} + +// rackCapThresholds are the imbalance thresholds of the shell (0) and the +// worker's default (0.2); both reach Plan. +var rackCapThresholds = []float64{0, 0.2} + +// addCapNode adds a node with one disk of `free` slots, before its shards. +func addCapNode(topo *Topology, id, rackKey string, free int) *Node { + n := topo.AddNode(id, "dc1", rackKey, free) + n.AddDisk(0, "", free, 0) + return n +} + +// putShards places the volume's shards on the node's disk 0 and takes their slots. +func putShards(n *Node, ids ...int) { + n.AddShards(capTestVid, "c1", 0, bits(ids...)) + n.freeSlots -= len(ids) + n.disks[0].freeSlots -= len(ids) + n.disks[0].shardCount += len(ids) +} + +func shardTotalsPerRack(topo *Topology) map[string]int { + totals := map[string]int{} + for _, n := range topo.nodes { + totals[n.rack] += volumeShardCount(n, capTestVk) + } + return totals +} + +// planAndCheckRackTotals runs Plan, then checks every shard is still placed, +// that no rack holds more than maxPerRack, and that a second Plan is a no-op. +func planAndCheckRackTotals(t *testing.T, topo *Topology, opts Options, maxPerRack int) { + t.Helper() + Plan(topo, opts) + totals := shardTotalsPerRack(topo) + t.Logf("total shards per rack: %v", totals) + sum := 0 + for rk, n := range totals { + sum += n + if n > maxPerRack { + t.Errorf("rack %s holds %d shards, want at most %d", rk, n, maxPerRack) + } + } + if sum != erasure_coding.TotalShardsCount { + t.Errorf("%d shards placed, want %d", sum, erasure_coding.TotalShardsCount) + } + if again := Plan(topo, opts); len(again) != 0 { + t.Errorf("second Plan is not a no-op: %+v", again) + } +} + +// raisedCaps records Options.RackTotalCapRaised reports. +type raisedCaps []string + +func (r *raisedCaps) record(collection string, vid uint32, rackCap, evenCap int) { + *r = append(*r, fmt.Sprintf("%s/%d cap=%d even=%d", collection, vid, rackCap, evenCap)) +} + +func forEachThreshold(t *testing.T, run func(t *testing.T, threshold float64)) { + for _, th := range rackCapThresholds { + t.Run(fmt.Sprintf("threshold=%v", th), func(t *testing.T) { run(t, th) }) + } +} + +// TestPlanTotalShardsPerRackCap: 10+4 encoded onto one of 8 racks spreads to at +// most ceil(14/8) = 2 shards per rack, not 2 data + 1 parity. +func TestPlanTotalShardsPerRackCap(t *testing.T) { + forEachThreshold(t, func(t *testing.T, th float64) { + topo := NewTopology() + putShards(addCapNode(topo, "n1", "dc1:rack1", 100), 0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13) + for r := 2; r <= 8; r++ { + addCapNode(topo, fmt.Sprintf("n%d", r), fmt.Sprintf("dc1:rack%d", r), 100) + } + planAndCheckRackTotals(t, topo, Options{ImbalanceThreshold: th, Ratio: ratio(10, 4)}, 2) + }) +} + +// TestPlanTotalShardsPerRackCapTwoNodesPerRack: the cap is per rack, not per +// node, so two volume servers sharing a rack (one physical disk) hold at most 2 +// between them. +func TestPlanTotalShardsPerRackCapTwoNodesPerRack(t *testing.T) { + forEachThreshold(t, func(t *testing.T, th float64) { + topo := NewTopology() + for r := 1; r <= 8; r++ { + for i := 0; i < 2; i++ { + n := addCapNode(topo, fmt.Sprintf("n%d-%d", r, i), fmt.Sprintf("dc1:rack%d", r), 100) + if r == 1 && i == 0 { + putShards(n, 0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13) + } + } + } + planAndCheckRackTotals(t, topo, Options{ImbalanceThreshold: th, Ratio: ratio(10, 4)}, 2) + }) +} + +// TestPlanTotalCapFromClumpedStart: data 2,2,2,2,2,0,0,0 with all parity on +// rack1. Racks 6-8 take one parity each, after which no rack is under the parity +// cap of 1; the last parity must still leave rack1 (2 data + 1 parity) for a +// data-bearing rack under the total cap rather than stay put. +func TestPlanTotalCapFromClumpedStart(t *testing.T) { + forEachThreshold(t, func(t *testing.T, th float64) { + topo := NewTopology() + layout := [][]int{{0, 1, 10, 11, 12, 13}, {2, 3}, {4, 5}, {6, 7}, {8, 9}, nil, nil, nil} + for r, ids := range layout { + n := addCapNode(topo, fmt.Sprintf("n%d", r+1), fmt.Sprintf("dc1:rack%d", r+1), 100) + if len(ids) > 0 { + putShards(n, ids...) + } + } + planAndCheckRackTotals(t, topo, Options{ImbalanceThreshold: th, Ratio: ratio(10, 4)}, 2) + }) +} + +// TestPlanTotalCapAllRacksBearData: every rack holds data, so parity has to +// share a rack with data. Anti-affinity stays a preference and the plan still +// reaches at most 2 per rack. +func TestPlanTotalCapAllRacksBearData(t *testing.T) { + forEachThreshold(t, func(t *testing.T, th float64) { + topo := NewTopology() + layout := [][]int{{0, 1, 10, 11, 12, 13}, {2, 3}, {4}, {5}, {6}, {7}, {8}, {9}} + for r, ids := range layout { + putShards(addCapNode(topo, fmt.Sprintf("n%d", r+1), fmt.Sprintf("dc1:rack%d", r+1), 100), ids...) + } + planAndCheckRackTotals(t, topo, Options{ImbalanceThreshold: th, Ratio: ratio(10, 4)}, 2) + }) +} + +// TestPlanTotalCapParityOnDataRacks: each type is within its own cap (2 data, +// 1 parity per rack) yet four racks hold 3. The non-overflow parity shards that +// may only move to a data-free rack find just two of those; the over-cap ones +// that remain must go to a data-bearing rack under the total cap. +func TestPlanTotalCapParityOnDataRacks(t *testing.T) { + forEachThreshold(t, func(t *testing.T, th float64) { + topo := NewTopology() + layout := [][]int{{0, 1, 10}, {2, 3, 11}, {4, 5, 12}, {6, 7, 13}, {8}, {9}, nil, nil} + for r, ids := range layout { + n := addCapNode(topo, fmt.Sprintf("n%d", r+1), fmt.Sprintf("dc1:rack%d", r+1), 100) + if len(ids) > 0 { + putShards(n, ids...) + } + } + planAndCheckRackTotals(t, topo, Options{ImbalanceThreshold: th, Ratio: ratio(10, 4)}, 2) + }) +} + +// TestPlanTotalCapIgnoresImbalanceThreshold: a rack above the total cap is a +// durability problem, so a threshold high enough to pass both per-type checks +// must not stop Plan from fixing it. +func TestPlanTotalCapIgnoresImbalanceThreshold(t *testing.T) { + topo := NewTopology() + layout := [][]int{{0, 1, 10}, {2, 3, 11}, {4, 5}, {6, 7}, {8}, {9}, {12}, {13}} + for r, ids := range layout { + putShards(addCapNode(topo, fmt.Sprintf("n%d", r+1), fmt.Sprintf("dc1:rack%d", r+1), 100), ids...) + } + planAndCheckRackTotals(t, topo, Options{ImbalanceThreshold: 100, Ratio: ratio(10, 4)}, 2) +} + +// TestPlanTotalCapWithNearlyFullRack: 8 racks of two volume servers each, plus +// a small, nearly full SSD rack already holding one shard. The SSD rack's lack +// of room must not raise the cap above ceil(14/9) = 2 for the HDD racks. The +// start is the layout Place produced before it capped the total per rack. +func TestPlanTotalCapWithNearlyFullRack(t *testing.T) { + forEachThreshold(t, func(t *testing.T, th float64) { + topo := NewTopology() + layout := [][]int{{0, 1, 10}, {2, 3, 11}, {4, 5}, {6, 7}, {8}, {9}, {12}, nil} + for r, ids := range layout { + for i := 0; i < 2; i++ { + n := addCapNode(topo, fmt.Sprintf("n%d-%d", r, i), fmt.Sprintf("dc1:hdd%d", r), 100) + if i == 0 && len(ids) > 0 { + putShards(n, ids...) + } + } + } + putShards(addCapNode(topo, "nvme", "dc1:ssd", 2), 13) + var raised raisedCaps + planAndCheckRackTotals(t, topo, Options{ImbalanceThreshold: th, Ratio: ratio(10, 4), RackTotalCapRaised: raised.record}, 2) + if len(raised) != 0 { + t.Errorf("cap reported as raised: %v", raised) + } + }) +} + +// TestPlanTotalCapRisesWhenRacksAreFull: two of 8 racks are full, so the other +// six can hold at best 3 each. A cap fixed at ceil(14/8) = 2 would block every +// destination and leave all four parity shards on rack1; the capacity-sized cap +// of 3 lets them spread as far as the room allows. +func TestPlanTotalCapRisesWhenRacksAreFull(t *testing.T) { + forEachThreshold(t, func(t *testing.T, th float64) { + topo := NewTopology() + layout := [][]int{{10, 11, 12, 13}, {0, 1}, {2, 3}, {4, 5}, {6, 7}, {8, 9}} + for r, ids := range layout { + putShards(addCapNode(topo, fmt.Sprintf("n%d", r+1), fmt.Sprintf("dc1:rack%d", r+1), 100), ids...) + } + addCapNode(topo, "full7", "dc1:rack7", 0) + addCapNode(topo, "full8", "dc1:rack8", 0) + var raised raisedCaps + planAndCheckRackTotals(t, topo, Options{ImbalanceThreshold: th, Ratio: ratio(10, 4), RackTotalCapRaised: raised.record}, 3) + // Reported on every Plan (two here) while the cap stays raised. + want := "c1/1 cap=3 even=2" + if len(raised) != 2 || raised[0] != want || raised[1] != want { + t.Errorf("raised-cap reports = %v, want [%s %s]", raised, want, want) + } + }) +} + +// TestPlanRackTotalCap: Plan may move every shard, so shards already on a rack +// count as room there, not as held in place. All 14 on one rack still gives 2. +func TestPlanRackTotalCap(t *testing.T) { + build := func(free []int, held map[int]int) (map[string]*rack, map[string]int) { + topo := NewTopology() + next := 0 + for r, f := range free { + n := addCapNode(topo, fmt.Sprintf("n%d", r), fmt.Sprintf("dc1:rack%d", r), f) + for i := 0; i < held[r]; i++ { + n.AddShards(capTestVid, "c1", 0, bits(next)) + next++ + } + } + return buildRacks(topo.nodes), countShardsByRack(capTestVk, topo.nodes) + } + cases := []struct { + name string + free []int + held map[int]int + want int + }{ + {name: "all on one rack", free: []int{0, 50, 50, 50, 50, 50, 50, 50}, held: map[int]int{0: 14}, want: 2}, + {name: "even over 5", free: []int{50, 50, 50, 50, 50}, held: map[int]int{0: 14}, want: 3}, + {name: "two full empty racks", free: []int{50, 50, 50, 50, 50, 50, 0, 0}, held: map[int]int{0: 14}, want: 3}, + {name: "small nearly full ninth rack", free: []int{50, 50, 50, 50, 50, 50, 50, 50, 1}, held: map[int]int{0: 13, 8: 1}, want: 2}, + {name: "only the holding rack has room", free: []int{0, 0}, held: map[int]int{0: 14}, want: 14}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + racks, held := build(tc.free, tc.held) + if got := planRackTotalCap(capTestVk, racks, held, erasure_coding.TotalShardsCount, nil); got != tc.want { + t.Errorf("planRackTotalCap = %d, want %d", got, tc.want) + } + }) + } +} + +// TestPlanTotalCapDataRackWaitsForParitySlot: the data pass runs out of +// destinations, and the only slots left for rack5's surplus data open up when +// the parity pass moves parity off rack3 and rack4. One Plan must still bring +// every rack within the cap (3), using those freed slots, without moving any +// shard twice. +func TestPlanTotalCapDataRackWaitsForParitySlot(t *testing.T) { + forEachThreshold(t, func(t *testing.T, th float64) { + topo := NewTopology() + layout := [][]int{{0, 1, 2}, {3, 4}, {10, 11}, {12}, {5, 6, 7, 8, 9, 13}} + free := []int{0, 4, 1, 1, 4} + for r, ids := range layout { + n := addCapNode(topo, fmt.Sprintf("n%d", r+1), fmt.Sprintf("dc1:rack%d", r+1), free[r]+len(ids)) + putShards(n, ids...) + } + moves := Plan(topo, Options{ImbalanceThreshold: th, Ratio: ratio(10, 4)}) + seen := map[int]bool{} + for _, m := range moves { + if seen[m.ShardID] { + t.Errorf("shard %d moved twice in one plan: %+v", m.ShardID, moves) + } + seen[m.ShardID] = true + } + totals := shardTotalsPerRack(topo) + t.Logf("total shards per rack: %v", totals) + for rk, n := range totals { + if n > 3 { + t.Errorf("rack %s holds %d shards, want at most 3", rk, n) + } + } + }) +} + +// TestPlanTotalCapCountsSameRackCount: with SameRackCount=1 a node takes at +// most one shard of the volume, so racks of one node that already hold a shard +// have no room however many free slots they report. Rack A's seven shards can't +// spread, the cap is 7, and Plan reports it as raised. +func TestPlanTotalCapCountsSameRackCount(t *testing.T) { + topo := NewTopology() + for i := 0; i < 7; i++ { + putShards(addCapNode(topo, fmt.Sprintf("a%d", i), "dc1:rackA", 100), i) + } + for i, rk := range []string{"B", "C", "D", "E", "F", "G", "H"} { + putShards(addCapNode(topo, "n"+rk, "dc1:rack"+rk, 100), 7+i) + } + var raised raisedCaps + rp := &super_block.ReplicaPlacement{SameRackCount: 1} + Plan(topo, Options{ReplicaPlacement: rp, Ratio: ratio(10, 4), RackTotalCapRaised: raised.record}) + if want := "c1/1 cap=7 even=2"; len(raised) != 1 || raised[0] != want { + t.Errorf("raised-cap reports = %v, want [%s]", raised, want) + } +} diff --git a/weed/worker/tasks/ec_balance/detection.go b/weed/worker/tasks/ec_balance/detection.go index abc6d0a58..9d147ae3c 100644 --- a/weed/worker/tasks/ec_balance/detection.go +++ b/weed/worker/tasks/ec_balance/detection.go @@ -87,6 +87,9 @@ func Detection( GlobalMaxMovesPerRack: 10, // Balance heterogeneous-capacity racks by fractional fullness. GlobalUtilizationBased: true, + RackTotalCapRaised: func(collection string, vid uint32, rackCap, evenCap int) { + glog.Warningf("EC volume %d (collection %q): rack-total-cap raised to %d shards per rack, above the even %d: the racks lack room for an even spread", vid, collection, rackCap, evenCap) + }, }) if len(moves) == 0 { return nil, false, nil