mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-08-16 20:26:45 +00:00
shell: add a reusable target picker for volume moves (#10579)
* shell: add a reusable target picker for volume moves Picks the emptiest node near the source: locality first (same rack, then same data center), emptiest within a tier. Free bytes decide where the cluster reports them, free slots break the tie. The pick is spent in the passed topology, so planning several moves from one snapshot spreads them instead of stacking every one on whichever node started emptiest. Claude-Session: https://claude.ai/code/session_01Ks16jnt4S7gdDk8cheQ3xu * shell: order targets on one metric, and reserve the real volume size Comparing some pairs on free bytes and others on free slots is intransitive, so the winner depended on sort order. One node too old to report filesystem bytes now puts every candidate on slots. Reserving the tier average let a batch of large volumes overcommit a destination; callers pass what the move actually consumes. Claude-Session: https://claude.ai/code/session_01Ks16jnt4S7gdDk8cheQ3xu
This commit is contained in:
@@ -0,0 +1,147 @@
|
||||
package shell
|
||||
|
||||
import (
|
||||
"sort"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
)
|
||||
|
||||
// PlacementPreference describes where a volume that has to leave its current
|
||||
// node may land.
|
||||
type PlacementPreference struct {
|
||||
// Source is the node the volume is leaving. It is never a candidate, and
|
||||
// its rack and data center are what PreferLocality measures against.
|
||||
Source string
|
||||
// DiskType is the medium the volume must land on.
|
||||
DiskType types.DiskType
|
||||
// AnchorDataCenter restricts candidates to one data center when set.
|
||||
AnchorDataCenter string
|
||||
// Exclude names nodes already spoken for by the same plan, so a caller
|
||||
// placing several copies does not stack them on one node.
|
||||
Exclude map[string]bool
|
||||
// VolumeBytes is what this move will actually consume on the destination.
|
||||
// Zero falls back to the tier's average volume size, which is all a caller
|
||||
// planning a not-yet-created volume can know.
|
||||
VolumeBytes uint64
|
||||
}
|
||||
|
||||
// PickTarget chooses where a volume that has to leave its node should land:
|
||||
// the emptiest node near it.
|
||||
//
|
||||
// Locality ranks first -- same rack as the source, then same data center, then
|
||||
// anywhere -- because moving a volume past a node that could have held it costs
|
||||
// cross-rack bandwidth for nothing. Within a locality tier the emptiest wins, so
|
||||
// a burst of moves still spreads instead of filling whichever node the topology
|
||||
// happens to list first. Free bytes decide when the cluster reports them and
|
||||
// free slots break the tie, which also leaves servers too old to report
|
||||
// filesystem bytes ordered sensibly.
|
||||
//
|
||||
// Locality is not optional. A caller that genuinely wants to scatter should
|
||||
// widen what it passes as Source, not ask placement to forget where the volume
|
||||
// came from.
|
||||
//
|
||||
// The chosen node's volume count is spent in topo before returning. Planning
|
||||
// several moves from one snapshot otherwise sends them all to the same node,
|
||||
// since each pick would see the capacity its predecessors already took. Callers
|
||||
// therefore pass a snapshot they own and may mutate.
|
||||
func PickTarget(topo *master_pb.TopologyInfo, pref PlacementPreference) *master_pb.DataNodeInfo {
|
||||
nodes := collectVolumeServersByDcRackNode(topo, pref.AnchorDataCenter, "", "")
|
||||
|
||||
var srcDc, srcRack string
|
||||
for _, n := range nodes {
|
||||
if n.info.Id == pref.Source {
|
||||
srcDc, srcRack = n.dc, n.rack
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
type candidate struct {
|
||||
node *Node
|
||||
locality int // 0 same rack, 1 same dc, 2 elsewhere
|
||||
bytes uint64
|
||||
hasBytes bool
|
||||
slots int64
|
||||
}
|
||||
|
||||
var candidates []candidate
|
||||
for _, n := range nodes {
|
||||
if n.info.Id == pref.Source || pref.Exclude[n.info.Id] {
|
||||
continue
|
||||
}
|
||||
if !n.hasFreeVolumeSlot(pref.DiskType) {
|
||||
continue
|
||||
}
|
||||
c := candidate{node: n}
|
||||
_, free, ok := n.diskBytes(pref.DiskType)
|
||||
c.bytes, c.hasBytes = free, ok
|
||||
if d, found := n.info.DiskInfos[string(pref.DiskType)]; found && d != nil {
|
||||
c.slots = d.MaxVolumeCount - d.VolumeCount
|
||||
}
|
||||
switch {
|
||||
case srcRack != "" && n.rack == srcRack && n.dc == srcDc:
|
||||
c.locality = 0
|
||||
case srcDc != "" && n.dc == srcDc:
|
||||
c.locality = 1
|
||||
default:
|
||||
c.locality = 2
|
||||
}
|
||||
candidates = append(candidates, c)
|
||||
}
|
||||
if len(candidates) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
// One metric decides the whole ordering. Comparing some pairs on bytes and
|
||||
// others on slots is intransitive -- with A,B reporting bytes and C not, the
|
||||
// order of A and C depends on which the sort happens to compare first -- so a
|
||||
// single node too old to report filesystem bytes puts every candidate on
|
||||
// slots rather than silently mixing the two.
|
||||
byBytes := true
|
||||
for _, c := range candidates {
|
||||
if !c.hasBytes {
|
||||
byBytes = false
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
sort.SliceStable(candidates, func(i, j int) bool {
|
||||
a, b := candidates[i], candidates[j]
|
||||
if a.locality != b.locality {
|
||||
return a.locality < b.locality
|
||||
}
|
||||
if byBytes && a.bytes != b.bytes {
|
||||
return a.bytes > b.bytes
|
||||
}
|
||||
return a.slots > b.slots
|
||||
})
|
||||
|
||||
chosen := candidates[0].node.info
|
||||
if d, ok := chosen.DiskInfos[string(pref.DiskType)]; ok && d != nil {
|
||||
d.VolumeCount++
|
||||
if d.DiskTotalBytes != 0 && d.DiskFreeBytes != 0 {
|
||||
// Spend the space this move actually takes, so the next pick from
|
||||
// this snapshot sees the node as fuller by the right amount. A big
|
||||
// volume charged the tier average would let a batch overcommit the
|
||||
// destination; a small one would divert later moves off a node that
|
||||
// still had room.
|
||||
spend := pref.VolumeBytes
|
||||
if spend == 0 {
|
||||
spend = d.DiskTotalBytes / uint64(max64(d.MaxVolumeCount, 1))
|
||||
}
|
||||
if d.DiskFreeBytes > spend {
|
||||
d.DiskFreeBytes -= spend
|
||||
} else {
|
||||
d.DiskFreeBytes = 0
|
||||
}
|
||||
}
|
||||
}
|
||||
return chosen
|
||||
}
|
||||
|
||||
func max64(a, b int64) int64 {
|
||||
if a > b {
|
||||
return a
|
||||
}
|
||||
return b
|
||||
}
|
||||
@@ -0,0 +1,226 @@
|
||||
package shell
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
)
|
||||
|
||||
func placementNode(id string, maxVolumes, used int64, totalBytes, freeBytes uint64) *master_pb.DataNodeInfo {
|
||||
return &master_pb.DataNodeInfo{Id: id, DiskInfos: map[string]*master_pb.DiskInfo{
|
||||
"ssd": {
|
||||
MaxVolumeCount: maxVolumes, VolumeCount: used,
|
||||
DiskTotalBytes: totalBytes, DiskFreeBytes: freeBytes,
|
||||
},
|
||||
}}
|
||||
}
|
||||
|
||||
// placementTopo builds a one-datacenter topology, taking each node's first
|
||||
// character as its rack so the tests read compactly.
|
||||
func placementTopo(nodes ...*master_pb.DataNodeInfo) *master_pb.TopologyInfo {
|
||||
byRack := map[string][]*master_pb.DataNodeInfo{}
|
||||
var order []string
|
||||
for _, n := range nodes {
|
||||
rack := n.Id[:1]
|
||||
if _, seen := byRack[rack]; !seen {
|
||||
order = append(order, rack)
|
||||
}
|
||||
byRack[rack] = append(byRack[rack], n)
|
||||
}
|
||||
var racks []*master_pb.RackInfo
|
||||
for _, r := range order {
|
||||
racks = append(racks, &master_pb.RackInfo{Id: r, DataNodeInfos: byRack[r]})
|
||||
}
|
||||
return &master_pb.TopologyInfo{DataCenterInfos: []*master_pb.DataCenterInfo{{Id: "dc1", RackInfos: racks}}}
|
||||
}
|
||||
|
||||
func TestPickTargetPrefersTheEmptiestByFreeBytes(t *testing.T) {
|
||||
topo := placementTopo(
|
||||
placementNode("a1", 10, 1, 1000, 200),
|
||||
placementNode("a2", 10, 1, 1000, 900),
|
||||
placementNode("a3", 10, 1, 1000, 500),
|
||||
)
|
||||
got := PickTarget(topo, PlacementPreference{Source: "z9", DiskType: types.SsdType})
|
||||
if got.GetId() != "a2" {
|
||||
t.Fatalf("got %q, want the node with the most free bytes (a2)", got.GetId())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPickTargetFallsBackToSlotsWithoutReportedBytes(t *testing.T) {
|
||||
// A volume server too old to report filesystem bytes must still be ordered
|
||||
// sensibly rather than read as having zero free.
|
||||
topo := placementTopo(
|
||||
placementNode("a1", 10, 9, 0, 0),
|
||||
placementNode("a2", 10, 2, 0, 0),
|
||||
)
|
||||
got := PickTarget(topo, PlacementPreference{Source: "z9", DiskType: types.SsdType})
|
||||
if got.GetId() != "a2" {
|
||||
t.Fatalf("got %q, want the node with the most free slots (a2)", got.GetId())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPickTargetSkipsSourceExcludedAndFull(t *testing.T) {
|
||||
topo := placementTopo(
|
||||
placementNode("a1", 10, 1, 1000, 900), // the source itself
|
||||
placementNode("a2", 10, 1, 1000, 800), // already spoken for by the plan
|
||||
placementNode("a3", 5, 5, 1000, 700), // no free slot
|
||||
placementNode("a4", 10, 1, 1000, 100),
|
||||
)
|
||||
got := PickTarget(topo, PlacementPreference{
|
||||
Source: "a1", DiskType: types.SsdType, Exclude: map[string]bool{"a2": true},
|
||||
})
|
||||
if got.GetId() != "a4" {
|
||||
t.Fatalf("got %q, want a4", got.GetId())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPickTargetReturnsNilWhenNothingFits(t *testing.T) {
|
||||
topo := placementTopo(placementNode("a1", 5, 5, 1000, 900))
|
||||
if got := PickTarget(topo, PlacementPreference{Source: "z9", DiskType: types.SsdType}); got != nil {
|
||||
t.Fatalf("got %q, want nil", got.GetId())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPickTargetHonoursAnchorDataCenter(t *testing.T) {
|
||||
topo := &master_pb.TopologyInfo{DataCenterInfos: []*master_pb.DataCenterInfo{
|
||||
{Id: "dc1", RackInfos: []*master_pb.RackInfo{{Id: "r1", DataNodeInfos: []*master_pb.DataNodeInfo{
|
||||
placementNode("a1", 10, 1, 1000, 100),
|
||||
}}}},
|
||||
{Id: "dc2", RackInfos: []*master_pb.RackInfo{{Id: "r2", DataNodeInfos: []*master_pb.DataNodeInfo{
|
||||
placementNode("b1", 10, 1, 1000, 900),
|
||||
}}}},
|
||||
}}
|
||||
got := PickTarget(topo, PlacementPreference{
|
||||
Source: "z9", DiskType: types.SsdType, AnchorDataCenter: "dc1",
|
||||
})
|
||||
if got.GetId() != "a1" {
|
||||
t.Fatalf("got %q, want a1 -- the anchor outranks an emptier node elsewhere", got.GetId())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPickTargetPrefersTheSourceRackOverAnEmptierStranger(t *testing.T) {
|
||||
topo := placementTopo(
|
||||
placementNode("a1", 10, 1, 1000, 500), // source, rack a
|
||||
placementNode("a2", 10, 1, 1000, 100), // same rack, tighter
|
||||
placementNode("b1", 10, 1, 1000, 900), // other rack, emptiest
|
||||
)
|
||||
got := PickTarget(topo, PlacementPreference{Source: "a1", DiskType: types.SsdType})
|
||||
if got.GetId() != "a2" {
|
||||
t.Fatalf("got %q, want a2 -- same rack outranks emptier", got.GetId())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPickTargetLeavesTheRackWhenItCannotHoldTheVolume(t *testing.T) {
|
||||
// Locality is a preference over candidates, never a reason to overfill.
|
||||
topo := placementTopo(
|
||||
placementNode("a1", 10, 1, 1000, 500), // source
|
||||
placementNode("a2", 5, 5, 1000, 900), // same rack, no free slot
|
||||
placementNode("b1", 10, 1, 1000, 100),
|
||||
)
|
||||
got := PickTarget(topo, PlacementPreference{Source: "a1", DiskType: types.SsdType})
|
||||
if got.GetId() != "b1" {
|
||||
t.Fatalf("got %q, want b1 -- a full same-rack node is not a candidate", got.GetId())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPickTargetPrefersSameDataCenterOverAnother(t *testing.T) {
|
||||
topo := &master_pb.TopologyInfo{DataCenterInfos: []*master_pb.DataCenterInfo{
|
||||
{Id: "dc1", RackInfos: []*master_pb.RackInfo{
|
||||
{Id: "r1", DataNodeInfos: []*master_pb.DataNodeInfo{placementNode("a1", 10, 1, 1000, 500)}},
|
||||
{Id: "r2", DataNodeInfos: []*master_pb.DataNodeInfo{placementNode("a2", 10, 1, 1000, 100)}},
|
||||
}},
|
||||
{Id: "dc2", RackInfos: []*master_pb.RackInfo{
|
||||
{Id: "r3", DataNodeInfos: []*master_pb.DataNodeInfo{placementNode("b1", 10, 1, 1000, 900)}},
|
||||
}},
|
||||
}}
|
||||
got := PickTarget(topo, PlacementPreference{Source: "a1", DiskType: types.SsdType})
|
||||
if got.GetId() != "a2" {
|
||||
t.Fatalf("got %q, want a2 -- another rack in the same dc beats another dc", got.GetId())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPickTargetFallsBackToCapacityForAnUnknownSource(t *testing.T) {
|
||||
// A source that is not in the topology has no rack to be near, so ordering
|
||||
// degrades to capacity rather than to topology order.
|
||||
topo := placementTopo(
|
||||
placementNode("a1", 10, 1, 1000, 100),
|
||||
placementNode("b1", 10, 1, 1000, 900),
|
||||
)
|
||||
got := PickTarget(topo, PlacementPreference{Source: "gone:8080", DiskType: types.SsdType})
|
||||
if got.GetId() != "b1" {
|
||||
t.Fatalf("got %q, want b1", got.GetId())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPickTargetSpendsCapacitySoABatchSpreads(t *testing.T) {
|
||||
// The reservation this exists for: planning several moves from one snapshot
|
||||
// must not stack them all on whichever node started emptiest.
|
||||
topo := placementTopo(
|
||||
placementNode("a1", 4, 0, 4000, 4000),
|
||||
placementNode("a2", 4, 0, 4000, 3000),
|
||||
)
|
||||
seen := map[string]int{}
|
||||
for i := 0; i < 4; i++ {
|
||||
got := PickTarget(topo, PlacementPreference{Source: "z9", DiskType: types.SsdType})
|
||||
if got == nil {
|
||||
t.Fatalf("pick %d returned nil", i)
|
||||
}
|
||||
seen[got.GetId()]++
|
||||
}
|
||||
if seen["a1"] == 4 || seen["a2"] == 4 {
|
||||
t.Fatalf("all four picks landed on one node: %v", seen)
|
||||
}
|
||||
if seen["a1"]+seen["a2"] != 4 {
|
||||
t.Fatalf("unexpected spread: %v", seen)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPickTargetStopsWhenTheSnapshotIsSpent(t *testing.T) {
|
||||
// Reservation has to make a node stop being a candidate, not merely rank it
|
||||
// lower, or a batch could overcommit a tier.
|
||||
topo := placementTopo(placementNode("a1", 2, 0, 2000, 2000))
|
||||
for i := 0; i < 2; i++ {
|
||||
if got := PickTarget(topo, PlacementPreference{Source: "z9", DiskType: types.SsdType}); got == nil {
|
||||
t.Fatalf("pick %d returned nil while slots remained", i)
|
||||
}
|
||||
}
|
||||
if got := PickTarget(topo, PlacementPreference{Source: "z9", DiskType: types.SsdType}); got != nil {
|
||||
t.Fatalf("got %q, want nil once the snapshot is spent", got.GetId())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPickTargetUsesOneMetricWhenReportingIsMixed(t *testing.T) {
|
||||
// b1 has the most free bytes but the fewest slots; a2 the reverse. With a3
|
||||
// reporting no bytes at all, every candidate must be ordered on slots, or
|
||||
// the comparator is intransitive and the winner depends on sort order.
|
||||
topo := placementTopo(
|
||||
placementNode("a1", 10, 1, 1000, 500), // source
|
||||
placementNode("a2", 10, 1, 1000, 100),
|
||||
placementNode("a3", 10, 2, 0, 0), // too old to report bytes
|
||||
placementNode("a4", 10, 8, 1000, 900),
|
||||
)
|
||||
got := PickTarget(topo, PlacementPreference{Source: "a1", DiskType: types.SsdType})
|
||||
if got.GetId() != "a2" {
|
||||
t.Fatalf("got %q, want a2 -- the most free slots once bytes are unusable", got.GetId())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPickTargetReservesTheVolumeSize(t *testing.T) {
|
||||
// A volume far bigger than the tier average must be charged at its real
|
||||
// size, or a batch of them overcommits the destination.
|
||||
topo := placementTopo(
|
||||
placementNode("a1", 10, 0, 1000, 1000),
|
||||
placementNode("a2", 10, 0, 1000, 900),
|
||||
)
|
||||
first := PickTarget(topo, PlacementPreference{Source: "z9", DiskType: types.SsdType, VolumeBytes: 800})
|
||||
if first.GetId() != "a1" {
|
||||
t.Fatalf("first pick %q, want a1", first.GetId())
|
||||
}
|
||||
// a1 now reports 200 free against a2's 900, so the next move must go there
|
||||
// rather than to the node that merely started emptiest.
|
||||
second := PickTarget(topo, PlacementPreference{Source: "z9", DiskType: types.SsdType, VolumeBytes: 800})
|
||||
if second.GetId() != "a2" {
|
||||
t.Fatalf("second pick %q, want a2 after 800 bytes were spent on a1", second.GetId())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user