From 44e546a9338736e3123d239b707df48c94bd6b8f Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Wed, 5 Aug 2026 00:27:14 -0700 Subject: [PATCH] shell: pick tier.move replica targets with the shared placement picker (#10582) The command chose destinations by walking its location list in order, so it neither preferred a node near the source nor spread a burst of copies. Replica placement and "this node already holds the volume" move into the Accept predicate; the ranking and the per-pick reservation come from placement. Claude-Session: https://claude.ai/code/session_01Ks16jnt4S7gdDk8cheQ3xu --- weed/shell/command_volume_tier_move.go | 76 ++++++++++++++++++++------ 1 file changed, 58 insertions(+), 18 deletions(-) diff --git a/weed/shell/command_volume_tier_move.go b/weed/shell/command_volume_tier_move.go index fad43853c..3115c26d6 100644 --- a/weed/shell/command_volume_tier_move.go +++ b/weed/shell/command_volume_tier_move.go @@ -4,6 +4,7 @@ import ( "context" "flag" "fmt" + "github.com/seaweedfs/seaweedfs/weed/placement" "io" "path/filepath" "sync" @@ -489,26 +490,36 @@ func (c *commandVolumeTierMove) ensureReplicationFulfilled(commandEnv *CommandEn fmt.Fprintf(writer, "volume %d: creating %d additional replica(s) for replication %s\n", vid, additionalCopiesNeeded, replicaPlacement) - fn := capacityByFreeVolumeCount(toDiskType) + // One picker decides where copies go, so this command spreads and stays near + // the source the same way every other mover does. Constraints it cannot model + // -- replica placement, and a node already holding the volume -- stay here. + topo := topologyFromLocations(allLocations) + // The picker answers with a node; replica placement is judged on where that + // node sits, so keep the mapping back to its rack and data center. + placeOf := make(map[string]location, len(allLocations)) + for _, l := range allLocations { + placeOf[l.dataNode.Id] = l + } + taken := make(map[string]bool) copiesMade := 0 - for _, candidateDst := range allLocations { - if copiesMade >= additionalCopiesNeeded { + for copiesMade < additionalCopiesNeeded { + dst := placement.PickTarget(topo, placement.PlacementPreference{ + Source: sourceAddress.String(), + DiskType: toDiskType, + Exclude: taken, + Accept: func(dn *master_pb.DataNodeInfo, dc, rack string) bool { + if nodesWithVolume[dn.Id] { + return false + } + return satisfyReplicaPlacement(replicaPlacement, targetTierReplicas, newLocation(dc, rack, dn)) + }, + }) + if dst == nil { break } - if fn(candidateDst.dataNode) <= 0 { - continue - } - // Skip nodes that already host this volume on any disk type to avoid - // VolumeCopy conflicts (e.g., same volume on source tier and target tier). - if nodesWithVolume[candidateDst.dataNode.Id] { - continue - } - if !satisfyReplicaPlacement(replicaPlacement, targetTierReplicas, candidateDst) { - continue - } - - candidateAddress := pb.NewServerAddressFromDataNode(candidateDst.dataNode) - fmt.Fprintf(writer, "volume %d: replicating from %s to %s\n", vid, sourceAddress, candidateDst.dataNode.Id) + taken[dst.Id] = true + candidateDst := placeOf[dst.Id] + candidateAddress := pb.NewServerAddressFromDataNode(dst) if copyErr := replicateVolumeToServer(context.Background(), commandEnv.option.GrpcDialOption, writer, vid, sourceAddress, candidateAddress, toDiskType.ReadableString()); copyErr != nil { return nil, fmt.Errorf("replicate volume %d to %s: %v", vid, candidateDst.dataNode.Id, copyErr) @@ -527,7 +538,8 @@ func (c *commandVolumeTierMove) ensureReplicationFulfilled(commandEnv *CommandEn location: &candidateDst, info: targetTierReplicas[0].info, }) - addVolumeCount(candidateDst.dataNode.DiskInfos[string(toDiskType)], 1) + // PickTarget already spent the slot in topo, which shares these DataNodeInfo + // pointers with allLocations, so counting it again here would double it. copiesMade++ } @@ -591,3 +603,31 @@ func collectVolumeIdsForTierChange(topologyInfo *master_pb.TopologyInfo, volumeS return } + +// topologyFromLocations rebuilds a topology snapshot from the locations a +// command already collected, so placement sees exactly the candidate set the +// command would have iterated. +func topologyFromLocations(locations []location) *master_pb.TopologyInfo { + dcs := make(map[string]map[string][]*master_pb.DataNodeInfo) + var dcOrder []string + rackOrder := make(map[string][]string) + for _, l := range locations { + if _, ok := dcs[l.dc]; !ok { + dcs[l.dc] = make(map[string][]*master_pb.DataNodeInfo) + dcOrder = append(dcOrder, l.dc) + } + if _, ok := dcs[l.dc][l.rack]; !ok { + rackOrder[l.dc] = append(rackOrder[l.dc], l.rack) + } + dcs[l.dc][l.rack] = append(dcs[l.dc][l.rack], l.dataNode) + } + topo := &master_pb.TopologyInfo{} + for _, dc := range dcOrder { + dcInfo := &master_pb.DataCenterInfo{Id: dc} + for _, rack := range rackOrder[dc] { + dcInfo.RackInfos = append(dcInfo.RackInfos, &master_pb.RackInfo{Id: rack, DataNodeInfos: dcs[dc][rack]}) + } + topo.DataCenterInfos = append(topo.DataCenterInfos, dcInfo) + } + return topo +}