mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-05 22:12:04 +02:00
ecbalancer: honour total-shards-per-rack cap in Place / PlaceDurabilityFirst (#11553)
* ecbalancer: honour total-shards-per-rack cap in Place / PlaceDurabilityFirst
Worker auto-EC encode places via Topology.Place, which capped each shard
type independently (ceil(data/racks), ceil(parity/racks)). On an 8-rack
topology that permits 3 total shards on one rack, so losing two racks
strands 6/14 and a 10+4 volume becomes unreadable.
- tryPlace caps the total shards (data + parity) per rack in both modes,
whether or not ReplicaPlacement is set.
- rackTotalCap picks the smallest per-rack total the racks' real room
(free slots, bounded by the per-disk cap and node free slots, counting
shards already placed) can satisfy. On a uniform cluster it is
ceil(shards/racks); a nearly full rack raises it just enough that the
cap alone never fails an encode.
- PlaceDurabilityFirst gets a last rung that drops the rack cap
("rack-total-cap" in Relaxed), so it fails only when no disk has room.
PlaceStrict keeps the cap as a hard limit.
- chooseShardDest tries the next rack when the chosen one has no node
that fits, and room checks count the per-disk cap, so a rack whose
disks are all at the cap is no longer picked and then failed on
(pre-existing: 3-node rack + single-disk rack failed at shard 9).
- Docs no longer claim the cap guarantees surviving rack loss; the
placement error names the caps in effect; the encode warning no longer
says replica placement when other constraints were relaxed.
place_rack_cap_test.go covers 10+4 over 8 racks (max 2/rack, 3/rack on
master), a starved rack, nearly full racks, the preferred-tag tier, the
full-disk rack, and rackTotalCap directly.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
* ecbalancer: size the rack total cap from room left under SameRackCount
The rack total cap counted each rack's free disk room, but attempts that
enforce ReplicaPlacement also stop a node at SameRackCount shards. With
SameRackCount=1, four one-node racks and four three-node racks got cap 2,
which fits only 12 of 14 shards: strict placement failed and
durability-first relaxed replica placement although 1 per small rack and
up to 3 per large rack fits.
Attempts that enforce ReplicaPlacement now use a cap sized from each
node's remaining SameRackCount allowance; attempts that relax it keep the
disk-room cap.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
1 parent
164c3db606
commit
2a42d56437
3 files changed
+455
-65
No files matched your search
@@ -81,15 +81,17 @@ type PlaceResult struct {
|
||||
type PlacementMode int
|
||||
|
||||
const (
|
||||
// PlaceStrict: caps and ReplicaPlacement are hard. Place fails rather than
|
||||
// violate them, so the caller can defer (leave the volume as-is and retry).
|
||||
// PlaceStrict: caps, the total-shards-per-rack cap and ReplicaPlacement are
|
||||
// hard. Place fails rather than violate them, so the caller can defer (leave
|
||||
// the volume as-is and retry).
|
||||
PlaceStrict PlacementMode = iota
|
||||
// PlaceDurabilityFirst (used by both encode and repair): relax per-type caps ->
|
||||
// data/parity anti-affinity -> ReplicaPlacement, in that order, until each shard
|
||||
// lands, reporting what was relaxed in PlaceResult.Relaxed. The per-disk
|
||||
// durability cap (<= parityShards per disk) is never relaxed. Fails only if no
|
||||
// disk has free capacity. Encode places best-effort this way and rebalancing
|
||||
// tightens the spread afterward.
|
||||
// data/parity anti-affinity -> ReplicaPlacement -> the total-shards-per-rack
|
||||
// cap, in that order, until each shard lands, reporting what was relaxed in
|
||||
// PlaceResult.Relaxed. The per-disk durability cap (<= parityShards per disk)
|
||||
// is never relaxed. Fails only if no eligible disk has room for another shard
|
||||
// of the volume. Encode places best-effort this way and rebalancing tightens
|
||||
// the spread afterward.
|
||||
PlaceDurabilityFirst
|
||||
)
|
||||
|
||||
@@ -101,6 +103,7 @@ type relaxation struct {
|
||||
caps bool
|
||||
antiAffinity bool
|
||||
rp bool
|
||||
rackTotal bool
|
||||
}
|
||||
|
||||
func (r relaxation) relaxedNames() []string {
|
||||
@@ -114,16 +117,23 @@ func (r relaxation) relaxedNames() []string {
|
||||
if !r.rp {
|
||||
n = append(n, "replica-placement")
|
||||
}
|
||||
if !r.rackTotal {
|
||||
n = append(n, "rack-total-cap")
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
var strictAttempts = []relaxation{{caps: true, antiAffinity: true, rp: true}}
|
||||
var strictAttempts = []relaxation{{caps: true, antiAffinity: true, rp: true, rackTotal: true}}
|
||||
|
||||
// The total-shards-per-rack cap is sized from real capacity (see rackTotalCap),
|
||||
// so it is relaxed last and only as a safety net; with it dropped, the per-disk
|
||||
// cap is the only constraint left.
|
||||
var durabilityAttempts = []relaxation{
|
||||
{caps: true, antiAffinity: true, rp: true},
|
||||
{caps: false, antiAffinity: true, rp: true},
|
||||
{caps: false, antiAffinity: false, rp: true},
|
||||
{caps: false, antiAffinity: false, rp: false},
|
||||
{caps: true, antiAffinity: true, rp: true, rackTotal: true},
|
||||
{caps: false, antiAffinity: true, rp: true, rackTotal: true},
|
||||
{caps: false, antiAffinity: false, rp: true, rackTotal: true},
|
||||
{caps: false, antiAffinity: false, rp: false, rackTotal: true},
|
||||
{caps: false, antiAffinity: false, rp: false, rackTotal: false},
|
||||
}
|
||||
|
||||
type placedEntry struct {
|
||||
@@ -236,12 +246,17 @@ func (t *Topology) tryPlace(vk volKey, need []int, dataShards, parityShards int,
|
||||
}
|
||||
}
|
||||
|
||||
// Even per-rack caps divide by racks that actually have an eligible free disk,
|
||||
// not all racks (the snapshot keeps every disk type/tag), so a valid tiered
|
||||
// cluster — e.g. SSDs in only 2 of 4 racks — is not capped impossibly low.
|
||||
// Even per-rack caps divide by racks that can actually take another shard of
|
||||
// the volume, not all racks (the snapshot keeps every disk type/tag), so a
|
||||
// valid tiered cluster — e.g. SSDs in only 2 of 4 racks — is not capped
|
||||
// impossibly low.
|
||||
rackRoom := make(map[string]int, len(rackKeys))
|
||||
rackRoomRP := make(map[string]int, len(rackKeys))
|
||||
numEligibleRacks := 0
|
||||
for _, rk := range rackKeys {
|
||||
if rackHasFreeDisk(racks[rk], eligible) {
|
||||
rackRoom[rk] = rackShardRoom(racks[rk], vk, eligible, parityShards)
|
||||
rackRoomRP[rk] = rackShardRoomUnderRP(racks[rk], vk, eligible, parityShards, rp)
|
||||
if rackRoom[rk] > 0 {
|
||||
numEligibleRacks++
|
||||
}
|
||||
}
|
||||
@@ -249,6 +264,24 @@ func (t *Topology) tryPlace(vk volKey, need []int, dataShards, parityShards int,
|
||||
numEligibleRacks = 1
|
||||
}
|
||||
|
||||
// The per-type caps 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
|
||||
// as well, at the lowest value the racks' free capacity allows. Attempts that
|
||||
// enforce ReplicaPlacement size it from the room left under SameRackCount, so
|
||||
// a rack of few nodes does not count disk room its nodes may not use.
|
||||
totalShards := len(need)
|
||||
for _, n := range rackShardCount {
|
||||
totalShards += n
|
||||
}
|
||||
maxTotalPerRack := rackTotalCap(rackKeys, rackShardCount, rackRoom, totalShards)
|
||||
maxTotalPerRackRP := rackTotalCap(rackKeys, rackShardCount, rackRoomRP, totalShards)
|
||||
rackCap := func(rl relaxation) int {
|
||||
if rl.rp {
|
||||
return maxTotalPerRackRP
|
||||
}
|
||||
return maxTotalPerRack
|
||||
}
|
||||
|
||||
attempts := strictAttempts
|
||||
if mode == PlaceDurabilityFirst {
|
||||
attempts = durabilityAttempts
|
||||
@@ -264,7 +297,7 @@ func (t *Topology) tryPlace(vk volKey, need []int, dataShards, parityShards int,
|
||||
typeTotal = parityShards
|
||||
}
|
||||
for _, rl := range attempts {
|
||||
node, diskID, spilled, ok := chooseShardDest(vk, sid, isData, dataShards, typeTotal, numEligibleRacks, parityShards, racks, rackKeys, rp, eligible, prefer, shardsPerRack[isData], rackShardCount, bearing, rl)
|
||||
node, diskID, spilled, ok := chooseShardDest(vk, sid, isData, dataShards, typeTotal, numEligibleRacks, parityShards, rackCap(rl), racks, rackKeys, rp, eligible, prefer, shardsPerRack[isData], rackShardCount, bearing, rl)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
@@ -303,7 +336,7 @@ func (t *Topology) tryPlace(vk volKey, need []int, dataShards, parityShards int,
|
||||
e.node.freeSlots++
|
||||
racks[e.rackKey].freeSlots++
|
||||
}
|
||||
return nil, fmt.Errorf("cannot place EC shard %d of volume %d (collection %q)", sid, vk.vid, vk.collection)
|
||||
return nil, fmt.Errorf("cannot place EC shard %d of volume %d (collection %q) (rack total cap %d, per-disk cap %d)", sid, vk.vid, vk.collection, rackCap(attempts[len(attempts)-1]), parityShards)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -316,11 +349,13 @@ func (t *Topology) tryPlace(vk volKey, need []int, dataShards, parityShards int,
|
||||
}
|
||||
|
||||
// chooseShardDest selects a (node, disk) for one shard at the given relaxation
|
||||
// level: pick a rack (even per-type cap + ReplicaPlacement caps + two-pass
|
||||
// anti-affinity to the opposite type), then the least-loaded eligible node, then
|
||||
// the best eligible disk. The third return reports whether the disk spilled off
|
||||
// the soft-preferred type. ok=false when no rack/node/disk fits.
|
||||
func chooseShardDest(vk volKey, sid int, isData bool, dataShards, typeTotal, numEligibleRacks, maxPerDisk int, racks map[string]*rack, rackKeys []string, rp *super_block.ReplicaPlacement, eligible func(*disk) bool, prefer func(*disk) bool, shardsPerRackType map[string][]int, rackShardCount map[string]int, bearing map[bool]map[string]bool, rl relaxation) (*Node, uint32, bool, bool) {
|
||||
// level: pick a rack (total-shards-per-rack cap + even per-type cap +
|
||||
// ReplicaPlacement caps + two-pass anti-affinity to the opposite type), then the
|
||||
// least-loaded eligible node, then the best eligible disk. If no node in the
|
||||
// chosen rack fits (e.g. every node is at SameRackCount), the next-best rack is
|
||||
// tried. The third return reports whether the disk spilled off the
|
||||
// soft-preferred type. ok=false when no rack/node/disk fits.
|
||||
func chooseShardDest(vk volKey, sid int, isData bool, dataShards, typeTotal, numEligibleRacks, maxPerDisk, maxTotalPerRack int, racks map[string]*rack, rackKeys []string, rp *super_block.ReplicaPlacement, eligible func(*disk) bool, prefer func(*disk) bool, shardsPerRackType map[string][]int, rackShardCount map[string]int, bearing map[bool]map[string]bool, rl relaxation) (*Node, uint32, bool, bool) {
|
||||
maxPerRack := numEligibleRacks*typeTotal + 1 // effectively unlimited when caps are relaxed
|
||||
if rl.caps {
|
||||
if maxPerRack = ceilDivide(typeTotal, numEligibleRacks); maxPerRack < 1 {
|
||||
@@ -336,74 +371,127 @@ func chooseShardDest(vk volKey, sid int, isData bool, dataShards, typeTotal, num
|
||||
if !rl.rp {
|
||||
rp = nil
|
||||
}
|
||||
// A rack is eligible only if it is under the per-rack shard cap (DiffRackCount),
|
||||
// enforced only when set (and relaxed with rp).
|
||||
// A rack is eligible only if it is under the total-shards-per-rack cap and,
|
||||
// when set, the per-rack shard cap (DiffRackCount). The total cap does not
|
||||
// depend on ReplicaPlacement, so it also holds when rp is nil.
|
||||
withinLimit := func(r string) bool {
|
||||
if rp == nil {
|
||||
return true
|
||||
if rl.rackTotal && rackShardCount[r] >= maxTotalPerRack {
|
||||
return false
|
||||
}
|
||||
if rp.DiffRackCount > 0 && rackShardCount[r] >= rp.DiffRackCount {
|
||||
if rp != nil && rp.DiffRackCount > 0 && rackShardCount[r] >= rp.DiffRackCount {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
destRack, ok := pickTarget(rackKeys, shardsPerRackType, maxPerRack, anti,
|
||||
func(r string) bool { return racks[r].freeSlots > 0 && rackHasFreeDisk(racks[r], eligible) },
|
||||
withinLimit)
|
||||
if !ok {
|
||||
return nil, 0, false, false
|
||||
tried := map[string]bool{}
|
||||
hasRoom := func(r string) bool {
|
||||
return !tried[r] && racks[r].freeSlots > 0 && rackShardRoom(racks[r], vk, eligible, maxPerDisk) > 0
|
||||
}
|
||||
node := pickNodeInRackEligible(racks[destRack], vk, rp, eligible)
|
||||
if node == nil {
|
||||
return nil, 0, false, false
|
||||
for {
|
||||
destRack, ok := pickTarget(rackKeys, shardsPerRackType, maxPerRack, anti, hasRoom, withinLimit)
|
||||
if !ok {
|
||||
return nil, 0, false, false
|
||||
}
|
||||
if node := pickNodeInRackEligible(racks[destRack], vk, rp, eligible, maxPerDisk); node != nil {
|
||||
if diskID, ok, spilled := pickBestDiskEligible(node, vk, eligible, prefer, sid, dataShards, maxPerDisk); ok {
|
||||
return node, diskID, spilled, true
|
||||
}
|
||||
}
|
||||
tried[destRack] = true
|
||||
}
|
||||
diskID, ok, spilled := pickBestDiskEligible(node, vk, eligible, prefer, sid, dataShards, maxPerDisk)
|
||||
if !ok {
|
||||
return nil, 0, false, false
|
||||
}
|
||||
return node, diskID, spilled, true
|
||||
}
|
||||
|
||||
// nodeHasFreeDisk reports whether the node has a free disk satisfying eligible.
|
||||
func nodeHasFreeDisk(n *Node, eligible func(*disk) bool) bool {
|
||||
// rackTotalCap returns the smallest per-rack total c such that the racks can
|
||||
// hold totalShards shards of the volume with none above c, given each rack's
|
||||
// shards already placed (held) and its room for more. On a uniform cluster this
|
||||
// is ceil(totalShards/racks); a nearly full rack raises it just enough for the
|
||||
// other racks to absorb its share, so the cap alone never makes a feasible
|
||||
// placement fail.
|
||||
//
|
||||
// It is the most even spread the free capacity allows, not a durability
|
||||
// guarantee: losing k racks loses up to k*c shards, which the volume survives
|
||||
// only while k*c <= parityShards. With few racks no placement can achieve that.
|
||||
func rackTotalCap(rackKeys []string, held, room map[string]int, totalShards int) int {
|
||||
for c := 1; c < totalShards; c++ {
|
||||
fits := 0
|
||||
for _, rk := range rackKeys {
|
||||
fits += max(held[rk], min(c, held[rk]+room[rk]))
|
||||
}
|
||||
if fits >= totalShards {
|
||||
return c
|
||||
}
|
||||
}
|
||||
return totalShards
|
||||
}
|
||||
|
||||
// diskShardRoom returns how many more shards of the volume the disk can take:
|
||||
// its free slots, bounded by the per-disk durability cap (maxPerDisk shards of
|
||||
// one volume per disk; <= 0 disables it).
|
||||
func diskShardRoom(n *Node, d *disk, vk volKey, maxPerDisk int) int {
|
||||
room := d.freeSlots
|
||||
if maxPerDisk > 0 {
|
||||
held := 0
|
||||
if info := n.shards[vk]; info != nil {
|
||||
held = info.diskShardBits[d.diskID].Count()
|
||||
}
|
||||
room = min(room, maxPerDisk-held)
|
||||
}
|
||||
return max(room, 0)
|
||||
}
|
||||
|
||||
// nodeShardRoom returns how many more shards of the volume the node's eligible
|
||||
// disks can take, bounded by the node's free slots.
|
||||
func nodeShardRoom(n *Node, vk volKey, eligible func(*disk) bool, maxPerDisk int) int {
|
||||
room := 0
|
||||
for _, d := range n.disks {
|
||||
if d.freeSlots > 0 && eligible(d) {
|
||||
return true
|
||||
if eligible(d) {
|
||||
room += diskShardRoom(n, d, vk, maxPerDisk)
|
||||
}
|
||||
}
|
||||
return false
|
||||
return max(min(room, n.freeSlots), 0)
|
||||
}
|
||||
|
||||
// rackHasFreeDisk reports whether any node in the rack has a free eligible disk.
|
||||
func rackHasFreeDisk(r *rack, eligible func(*disk) bool) bool {
|
||||
// rackShardRoom returns how many more shards of the volume the rack can take.
|
||||
func rackShardRoom(r *rack, vk volKey, eligible func(*disk) bool, maxPerDisk int) int {
|
||||
room := 0
|
||||
for _, n := range r.nodes {
|
||||
if n.freeSlots > 0 && nodeHasFreeDisk(n, eligible) {
|
||||
return true
|
||||
}
|
||||
room += nodeShardRoom(n, vk, eligible, maxPerDisk)
|
||||
}
|
||||
return false
|
||||
return room
|
||||
}
|
||||
|
||||
// pickNodeInRackEligible is pickNodeInRack restricted to nodes that have a free
|
||||
// eligible disk. FromActiveTopology keeps all disk types/tags in the snapshot, so
|
||||
// without this a node with free volume slots but no eligible disk could be chosen.
|
||||
// rackShardRoomUnderRP is rackShardRoom with each node further bounded by the
|
||||
// shards it may still take under ReplicaPlacement's SameRackCount (max shards of
|
||||
// the volume per node), which pickNodeInRackEligible enforces.
|
||||
func rackShardRoomUnderRP(r *rack, vk volKey, eligible func(*disk) bool, maxPerDisk int, rp *super_block.ReplicaPlacement) int {
|
||||
if rp == nil || rp.SameRackCount <= 0 {
|
||||
return rackShardRoom(r, vk, eligible, maxPerDisk)
|
||||
}
|
||||
room := 0
|
||||
for _, n := range r.nodes {
|
||||
room += min(nodeShardRoom(n, vk, eligible, maxPerDisk), max(rp.SameRackCount-volumeShardCount(n, vk), 0))
|
||||
}
|
||||
return room
|
||||
}
|
||||
|
||||
// pickNodeInRackEligible is pickNodeInRack restricted to nodes with an eligible
|
||||
// disk that can take another shard of the volume (free slot, under maxPerDisk).
|
||||
// FromActiveTopology keeps all disk types/tags in the snapshot, so without this a
|
||||
// node with free volume slots but no eligible disk could be chosen.
|
||||
//
|
||||
// Among eligible nodes it ranks by fewest shards of the volume per machine, then per
|
||||
// node, with free capacity breaking ties. The free-capacity tie-break (not sorted id)
|
||||
// keeps the lowest-id machine from winning every volume's first shard against the
|
||||
// shared encode snapshot and piling up load.
|
||||
func pickNodeInRackEligible(r *rack, vk volKey, rp *super_block.ReplicaPlacement, eligible func(*disk) bool) *Node {
|
||||
func pickNodeInRackEligible(r *rack, vk volKey, rp *super_block.ReplicaPlacement, eligible func(*disk) bool, maxPerDisk int) *Node {
|
||||
machineShards := countShardsByHost(vk, r.nodes)
|
||||
machineFree := freeSlotsByHost(r.nodes)
|
||||
var best *Node
|
||||
var bestMCount, bestMFree, bestNCount, bestNFree int
|
||||
for _, id := range sortedNodeKeys(r.nodes) {
|
||||
node := r.nodes[id]
|
||||
if node.freeSlots <= 0 {
|
||||
continue
|
||||
}
|
||||
if !nodeHasFreeDisk(node, eligible) {
|
||||
if nodeShardRoom(node, vk, eligible, maxPerDisk) <= 0 {
|
||||
continue
|
||||
}
|
||||
count := volumeShardCount(node, vk)
|
||||
|
||||
@@ -0,0 +1,301 @@
|
||||
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 Place / PlaceDurabilityFirst.
|
||||
//
|
||||
// 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. Place also caps the TOTAL per rack at the lowest value the racks'
|
||||
// free capacity allows (rackTotalCap), under both modes and with a nil
|
||||
// ReplicaPlacement.
|
||||
|
||||
// placedTotalsPerRack checks every shard was placed and returns the number of
|
||||
// shards per rack.
|
||||
func placedTotalsPerRack(t *testing.T, res *PlaceResult) map[string]int {
|
||||
t.Helper()
|
||||
if len(res.Destinations) != erasure_coding.TotalShardsCount {
|
||||
t.Fatalf("placed %d shards, want %d", len(res.Destinations), erasure_coding.TotalShardsCount)
|
||||
}
|
||||
totals := map[string]int{}
|
||||
for _, d := range res.Destinations {
|
||||
totals[d.Rack]++
|
||||
}
|
||||
t.Logf("total shards per rack: %v", totals)
|
||||
return totals
|
||||
}
|
||||
|
||||
func assertMaxPerRack(t *testing.T, totals map[string]int, maxAllowed int) {
|
||||
t.Helper()
|
||||
for rk, n := range totals {
|
||||
if n > maxAllowed {
|
||||
t.Errorf("rack %s holds %d shards, want at most %d", rk, n, maxAllowed)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func assertNotRelaxed(t *testing.T, res *PlaceResult, name string) {
|
||||
t.Helper()
|
||||
for _, r := range res.Relaxed {
|
||||
if r == name {
|
||||
t.Errorf("unexpected %q relaxation, got %v", name, res.Relaxed)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// setRackFreeSlots leaves the rack with `slots` free slots, all on the first
|
||||
// node's disk 0 (buildPlaceTopo gives each node one disk). Node and disk free
|
||||
// slots are kept consistent, since rack capacity is derived from them.
|
||||
func setRackFreeSlots(topo *Topology, rackKey string, slots int) {
|
||||
first := true
|
||||
for _, id := range sortedNodeKeys(topo.nodes) {
|
||||
n := topo.nodes[id]
|
||||
if n.rack != rackKey {
|
||||
continue
|
||||
}
|
||||
free := 0
|
||||
if first {
|
||||
free, first = slots, false
|
||||
}
|
||||
for _, d := range n.disks {
|
||||
d.freeSlots = free
|
||||
}
|
||||
n.freeSlots = free
|
||||
}
|
||||
}
|
||||
|
||||
// TestPlaceTotalShardsPerRackCap: 10+4 over 8 racks of two nodes each places at
|
||||
// most ceil(14/8) = 2 shards per rack, in both modes. Before the cap this was
|
||||
// 3,3,2,2,1,1,1,1: two racks with 3 each.
|
||||
func TestPlaceTotalShardsPerRackCap(t *testing.T) {
|
||||
modes := []struct {
|
||||
name string
|
||||
mode PlacementMode
|
||||
}{{"strict", PlaceStrict}, {"durability-first", PlaceDurabilityFirst}}
|
||||
for _, m := range modes {
|
||||
t.Run(m.name, func(t *testing.T) {
|
||||
topo := buildPlaceTopo(8, 2, 50)
|
||||
res, err := topo.Place(1, "c1", allShards(), Constraints{}, m.mode)
|
||||
if err != nil {
|
||||
t.Fatalf("Place: %v", err)
|
||||
}
|
||||
assertMaxPerRack(t, placedTotalsPerRack(t, res), 2)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestPlaceTotalCapWithStarvedRack: a rack with a single free slot still takes
|
||||
// one shard and the other seven racks keep to 2 each.
|
||||
func TestPlaceTotalCapWithStarvedRack(t *testing.T) {
|
||||
topo := buildPlaceTopo(8, 2, 50)
|
||||
setRackFreeSlots(topo, "dc1:rack0", 1)
|
||||
|
||||
res, err := topo.Place(1, "c1", allShards(), Constraints{}, PlaceDurabilityFirst)
|
||||
if err != nil {
|
||||
t.Fatalf("Place: %v", err)
|
||||
}
|
||||
totals := placedTotalsPerRack(t, res)
|
||||
assertMaxPerRack(t, totals, 2)
|
||||
assertNotRelaxed(t, res, "rack-total-cap")
|
||||
}
|
||||
|
||||
// TestPlaceTotalCapRisesForNearlyFullRacks: when nearly full racks cannot take
|
||||
// an even share, the cap rises so the roomy racks absorb it instead of failing.
|
||||
// A plain ceil(14/eligibleRacks) cap failed both cases.
|
||||
func TestPlaceTotalCapRisesForNearlyFullRacks(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
racks int
|
||||
starved map[string]int // rack -> free slots
|
||||
wantCap int
|
||||
wantHeld map[string]int // exact shard count expected on starved racks
|
||||
}{
|
||||
{
|
||||
// ceil(14/3) = 5 would fit only 5+5+1 = 11.
|
||||
name: "3 racks, one with 1 slot",
|
||||
racks: 3,
|
||||
starved: map[string]int{"dc1:rack2": 1},
|
||||
wantCap: 7,
|
||||
wantHeld: map[string]int{"rack2": 1},
|
||||
},
|
||||
{
|
||||
// ceil(14/8) = 2 would fit only 2+7 = 9.
|
||||
name: "8 racks, seven with 1 slot",
|
||||
racks: 8,
|
||||
starved: map[string]int{
|
||||
"dc1:rack1": 1, "dc1:rack2": 1, "dc1:rack3": 1, "dc1:rack4": 1,
|
||||
"dc1:rack5": 1, "dc1:rack6": 1, "dc1:rack7": 1,
|
||||
},
|
||||
wantCap: 7,
|
||||
wantHeld: map[string]int{
|
||||
"rack1": 1, "rack2": 1, "rack3": 1, "rack4": 1,
|
||||
"rack5": 1, "rack6": 1, "rack7": 1,
|
||||
},
|
||||
},
|
||||
}
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
topo := buildPlaceTopo(tc.racks, 2, 50)
|
||||
for rk, slots := range tc.starved {
|
||||
setRackFreeSlots(topo, rk, slots)
|
||||
}
|
||||
res, err := topo.Place(1, "c1", allShards(), Constraints{}, PlaceDurabilityFirst)
|
||||
if err != nil {
|
||||
t.Fatalf("Place: %v", err)
|
||||
}
|
||||
totals := placedTotalsPerRack(t, res)
|
||||
assertMaxPerRack(t, totals, tc.wantCap)
|
||||
for rk, want := range tc.wantHeld {
|
||||
if totals[rk] != want {
|
||||
t.Errorf("rack %s holds %d shards, want %d", rk, totals[rk], want)
|
||||
}
|
||||
}
|
||||
assertNotRelaxed(t, res, "rack-total-cap")
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestPlaceTotalCapKeepsPreferredTagTier: a preferred-tag tier whose racks have
|
||||
// uneven room must still place the whole volume on the tagged disks rather than
|
||||
// spill to untagged ones.
|
||||
func TestPlaceTotalCapKeepsPreferredTagTier(t *testing.T) {
|
||||
topo := NewTopology()
|
||||
for r := 0; r < 3; r++ {
|
||||
for n := 0; n < 2; n++ {
|
||||
taggedFree := 50
|
||||
if r == 2 {
|
||||
taggedFree = 0
|
||||
if n == 0 {
|
||||
taggedFree = 2
|
||||
}
|
||||
}
|
||||
node := topo.AddNode(fmt.Sprintf("10.0.%d.%d:8080", r, n), "dc1", fmt.Sprintf("dc1:rack%d", r), taggedFree+50)
|
||||
node.AddDisk(0, "", taggedFree, 0)
|
||||
node.AddDiskTags(0, []string{"ec-pool"})
|
||||
node.AddDisk(1, "", 50, 0)
|
||||
}
|
||||
}
|
||||
|
||||
res, err := topo.Place(1, "c1", allShards(), Constraints{PreferredTags: []string{"ec-pool"}}, PlaceDurabilityFirst)
|
||||
if err != nil {
|
||||
t.Fatalf("Place: %v", err)
|
||||
}
|
||||
if res.SpilledOutsidePreferredTags {
|
||||
t.Fatal("placement spilled outside the preferred tags although the tagged disks can hold every shard")
|
||||
}
|
||||
for sid, d := range res.Destinations {
|
||||
if d.DiskID != 0 {
|
||||
t.Errorf("shard %d landed on untagged disk %d of %s", sid, d.DiskID, d.Node)
|
||||
}
|
||||
}
|
||||
placedTotalsPerRack(t, res)
|
||||
}
|
||||
|
||||
// TestPlaceSkipsRackWithFullDisks: once a rack's disks hit the per-disk cap
|
||||
// (parityShards shards of the volume), Place moves on to a rack that still has
|
||||
// room instead of failing the shard.
|
||||
func TestPlaceSkipsRackWithFullDisks(t *testing.T) {
|
||||
topo := buildPlaceTopo(1, 3, 50) // rack0: 3 nodes
|
||||
small := topo.AddNode("10.0.1.0:8080", "dc1", "dc1:rack1", 50)
|
||||
small.AddDisk(0, "", 50, 0) // rack1: a single disk, so at most 4 shards
|
||||
|
||||
res, err := topo.Place(1, "c1", allShards(), Constraints{}, PlaceDurabilityFirst)
|
||||
if err != nil {
|
||||
t.Fatalf("Place: %v", err)
|
||||
}
|
||||
totals := placedTotalsPerRack(t, res)
|
||||
if totals["rack1"] > erasure_coding.ParityShardsCount {
|
||||
t.Errorf("rack1 holds %d shards on one disk, want at most %d", totals["rack1"], erasure_coding.ParityShardsCount)
|
||||
}
|
||||
}
|
||||
|
||||
// TestPlaceTotalCapCountsSameRackCount: with SameRackCount=1, a one-node rack
|
||||
// takes a single shard however much disk room it has, so the cap must be sized
|
||||
// from the room left under that limit. Four one-node racks and four three-node
|
||||
// racks fit 10+4 as 1 per small rack and up to 3 per large rack; a cap sized from
|
||||
// disk room alone is 2, which fits only 4 + 4*2 = 12, so strict placement failed
|
||||
// and durability-first relaxed replica placement.
|
||||
func TestPlaceTotalCapCountsSameRackCount(t *testing.T) {
|
||||
rp := &super_block.ReplicaPlacement{SameRackCount: 1} // <=1 shard per node
|
||||
modes := []struct {
|
||||
name string
|
||||
mode PlacementMode
|
||||
}{{"strict", PlaceStrict}, {"durability-first", PlaceDurabilityFirst}}
|
||||
for _, m := range modes {
|
||||
t.Run(m.name, func(t *testing.T) {
|
||||
topo := NewTopology()
|
||||
for r := 0; r < 8; r++ {
|
||||
nodes := 1
|
||||
if r >= 4 {
|
||||
nodes = 3
|
||||
}
|
||||
for n := 0; n < nodes; n++ {
|
||||
node := topo.AddNode(fmt.Sprintf("10.0.%d.%d:8080", r, n), "dc1", fmt.Sprintf("dc1:rack%d", r), 50)
|
||||
node.AddDisk(0, "", 50, 0)
|
||||
}
|
||||
}
|
||||
res, err := topo.Place(1, "c1", allShards(), Constraints{ReplicaPlacement: rp}, m.mode)
|
||||
if err != nil {
|
||||
t.Fatalf("Place: %v", err)
|
||||
}
|
||||
assertMaxPerRack(t, placedTotalsPerRack(t, res), 3)
|
||||
assertNotRelaxed(t, res, "replica-placement")
|
||||
assertNotRelaxed(t, res, "rack-total-cap")
|
||||
perNode := map[string]int{}
|
||||
for _, d := range res.Destinations {
|
||||
perNode[d.Node]++
|
||||
}
|
||||
for node, n := range perNode {
|
||||
if n > rp.SameRackCount {
|
||||
t.Errorf("node %s holds %d shards, want at most %d", node, n, rp.SameRackCount)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestRackTotalCap(t *testing.T) {
|
||||
keys := func(n int) []string {
|
||||
out := make([]string, n)
|
||||
for i := range out {
|
||||
out[i] = fmt.Sprintf("r%d", i)
|
||||
}
|
||||
return out
|
||||
}
|
||||
uniform := func(n, room int) map[string]int {
|
||||
out := map[string]int{}
|
||||
for _, k := range keys(n) {
|
||||
out[k] = room
|
||||
}
|
||||
return out
|
||||
}
|
||||
cases := []struct {
|
||||
name string
|
||||
racks int
|
||||
held map[string]int
|
||||
room map[string]int
|
||||
total int
|
||||
want int
|
||||
}{
|
||||
{name: "even over 8", racks: 8, room: uniform(8, 50), total: 14, want: 2},
|
||||
{name: "even over 5", racks: 5, room: uniform(5, 50), total: 14, want: 3},
|
||||
{name: "one rack with 1 slot", racks: 3, room: map[string]int{"r0": 50, "r1": 50, "r2": 1}, total: 14, want: 7},
|
||||
{name: "one rack full", racks: 3, room: map[string]int{"r0": 50, "r1": 50}, total: 14, want: 7},
|
||||
// Repair: r0 already holds 5 survivors; the other racks share the rest.
|
||||
{name: "held above even share", racks: 3, held: map[string]int{"r0": 5}, room: uniform(3, 50), total: 14, want: 5},
|
||||
{name: "not enough room", racks: 2, room: uniform(2, 3), total: 14, want: 14},
|
||||
}
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
if got := rackTotalCap(keys(tc.racks), tc.held, tc.room, tc.total); got != tc.want {
|
||||
t.Errorf("rackTotalCap = %d, want %d", got, tc.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -485,9 +485,10 @@ func buildNodeAddressMap(at *topology.ActiveTopology) map[string]string {
|
||||
// in the same detection cycle see the reduced capacity. Rebuilding it per volume
|
||||
// is O(volumes × topology) and times out on large clusters.
|
||||
//
|
||||
// Encode is lenient (PlaceDurabilityFirst): it relaxes caps/anti-affinity/RP as
|
||||
// needed rather than fail, and prefers the source disk type but spills if that
|
||||
// type can't hold every shard. rp is the resolved replica placement (may be nil).
|
||||
// Encode is lenient (PlaceDurabilityFirst): it relaxes caps/anti-affinity/RP and,
|
||||
// last, the total-shards-per-rack cap as needed, failing only when no eligible
|
||||
// disk has room. It prefers the source disk type but spills if that type can't
|
||||
// hold every shard. rp is the resolved replica placement (may be nil).
|
||||
func planECDestinations(snap *ecbalancer.Topology, nodeAddresses map[string]string, metric *types.VolumeHealthMetrics, ecConfig *Config, rp *super_block.ReplicaPlacement, dataShards, parityShards int) (*topology.MultiDestinationPlan, [][]uint32, error) {
|
||||
if snap == nil {
|
||||
return nil, nil, fmt.Errorf("EC placement snapshot not available")
|
||||
@@ -529,8 +530,8 @@ func planECDestinations(snap *ecbalancer.Topology, nodeAddresses map[string]stri
|
||||
if len(res.Relaxed) > 0 {
|
||||
// Encode is best-effort (PlaceDurabilityFirst): it relaxes these constraints
|
||||
// rather than defer when the cluster can't satisfy them. Surface it so a tight
|
||||
// replica placement isn't silently weakened; rebalancing tightens the spread.
|
||||
glog.Warningf("EC volume %d: placed with relaxed constraints %v; replica placement not fully satisfied (rebalancing will adjust)", metric.VolumeID, res.Relaxed)
|
||||
// replica placement or rack spread isn't silently weakened.
|
||||
glog.Warningf("EC volume %d: placed with relaxed placement constraints %v", metric.VolumeID, res.Relaxed)
|
||||
}
|
||||
|
||||
// Group the per-shard destinations into one plan per (node,disk), iterating
|
||||
|
||||
Reference in new issue
Block a user