From 78e428758b868a7360f1af1153547d7474610167 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Wed, 8 Jul 2026 20:03:01 -0700 Subject: [PATCH] volume.fix.replication: parallelize under-replicated volume copies (#10275) * volume.fix.replication: parallelize under-replicated volume copies Fan out fixOneUnderReplicatedVolume up to -maxParallelization at a time. Destination selection and free-slot accounting stay atomic behind a scheduler mutex, so concurrent fixes see each other's reservations; a failed copy returns its reserved slot. A per-server in-flight cap (-maxParallelizationPerServer, default 1) keeps many simultaneous copies from swamping a single destination server; when every eligible destination is at the cap the fix waits for a slot instead of failing. * volume.fix.replication: isolate the scheduler's location ordering Clone allLocations before the per-volume fan-out so the scheduler's re-sorting cannot alias the slice shared with the concurrent delete phases, and make the test's per-iteration location copy explicit. --- weed/shell/command_volume_fix_replication.go | 172 +++++++++++++----- ...nd_volume_fix_replication_parallel_test.go | 144 +++++++++++++++ 2 files changed, 270 insertions(+), 46 deletions(-) create mode 100644 weed/shell/command_volume_fix_replication_parallel_test.go diff --git a/weed/shell/command_volume_fix_replication.go b/weed/shell/command_volume_fix_replication.go index e8e54be37..f2903c3aa 100644 --- a/weed/shell/command_volume_fix_replication.go +++ b/weed/shell/command_volume_fix_replication.go @@ -6,6 +6,7 @@ import ( "io" "strconv" "strings" + "sync" "time" "slices" @@ -49,8 +50,10 @@ func (c *commandVolumeFixReplication) Help() string { Note: * each time this will only add back one replica for each volume id that is under replicated. If there are multiple replicas are missing, e.g. replica count is > 2, you may need to run this multiple times. - * do not run this too quickly within seconds, since the new volume replica may take a few seconds + * do not run this too quickly within seconds, since the new volume replica may take a few seconds to register itself to the master. + * under-replicated volumes are copied up to -maxParallelization at a time, with at most + -maxParallelizationPerServer concurrent copies onto any single destination server. ` } @@ -70,6 +73,7 @@ func (c *commandVolumeFixReplication) Do(args []string, commandEnv *CommandEnv, doDelete := volFixReplicationCommand.Bool("doDelete", true, "Also delete over-replicated volumes besides fixing under-replication") doCheck := volFixReplicationCommand.Bool("doCheck", true, "Also check synchronization before deleting") maxParallelization := volFixReplicationCommand.Int("maxParallelization", DefaultMaxParallelization, "run up to X tasks in parallel, whenever possible") + maxParallelizationPerServer := volFixReplicationCommand.Int("maxParallelizationPerServer", 1, "run up to X volume copies onto the same destination server in parallel") retryCount := volFixReplicationCommand.Int("retry", 5, "how many times to retry") volumesPerStep := volFixReplicationCommand.Int("volumesPerStep", 0, "how many volumes to fix in one cycle") @@ -154,7 +158,7 @@ func (c *commandVolumeFixReplication) Do(args []string, commandEnv *CommandEnv, ewg.Reset() ewg.Add(func() error { // find the most underpopulated data nodes - fixedVolumeReplicas, err = c.fixUnderReplicatedVolumes(commandEnv, writer, *applyChanges, underReplicatedVolumeIds, volumeReplicas, allLocations, *retryCount, *volumesPerStep) + fixedVolumeReplicas, err = c.fixUnderReplicatedVolumes(commandEnv, writer, *applyChanges, underReplicatedVolumeIds, volumeReplicas, allLocations, *retryCount, *volumesPerStep, *maxParallelization, *maxParallelizationPerServer) return err }) if *doDelete { @@ -336,7 +340,7 @@ func (c *commandVolumeFixReplication) deleteOneVolume(commandEnv *CommandEnv, wr return nil } -func (c *commandVolumeFixReplication) fixUnderReplicatedVolumes(commandEnv *CommandEnv, writer io.Writer, applyChanges bool, volumeIds []uint32, volumeReplicas map[uint32][]*VolumeReplica, allLocations []location, retryCount int, volumesPerStep int) (fixedVolumes map[string]int, err error) { +func (c *commandVolumeFixReplication) fixUnderReplicatedVolumes(commandEnv *CommandEnv, writer io.Writer, applyChanges bool, volumeIds []uint32, volumeReplicas map[uint32][]*VolumeReplica, allLocations []location, retryCount int, volumesPerStep int, maxParallelization int, maxParallelizationPerServer int) (fixedVolumes map[string]int, err error) { fixedVolumes = map[string]int{} if len(volumeIds) == 0 { @@ -346,61 +350,137 @@ func (c *commandVolumeFixReplication) fixUnderReplicatedVolumes(commandEnv *Comm if len(volumeIds) > volumesPerStep && volumesPerStep > 0 { volumeIds = volumeIds[0:volumesPerStep] } + + // own a private copy of the locations list: the scheduler re-sorts it on + // every reservation, and the caller's slice is shared with the concurrent + // delete phases + allLocations = slices.Clone(allLocations) + + scheduler := newVolumeCopyScheduler(maxParallelizationPerServer) + var fixedVolumesMu sync.Mutex + ewg := NewErrorWaitGroup(maxParallelization) for _, vid := range volumeIds { - for i := 0; i < retryCount+1; i++ { - var copied bool - if copied, err = c.fixOneUnderReplicatedVolume(commandEnv, writer, applyChanges, volumeReplicas, vid, allLocations); err == nil { - if applyChanges && copied { - fixedVolumes[strconv.FormatUint(uint64(vid), 10)] = len(volumeReplicas[vid]) + ewg.Add(func() error { + for i := 0; i < retryCount+1; i++ { + if copied, err := c.fixOneUnderReplicatedVolume(commandEnv, writer, applyChanges, volumeReplicas, vid, allLocations, scheduler); err == nil { + if applyChanges && copied { + fixedVolumesMu.Lock() + fixedVolumes[strconv.FormatUint(uint64(vid), 10)] = len(volumeReplicas[vid]) + fixedVolumesMu.Unlock() + } + break + } else { + fmt.Fprintf(writer, "fixing under replicated volume %d: %v\n", vid, err) } - break - } else { - fmt.Fprintf(writer, "fixing under replicated volume %d: %v\n", vid, err) } - } + return nil + }) } - return fixedVolumes, nil + return fixedVolumes, ewg.Wait() } -func (c *commandVolumeFixReplication) fixOneUnderReplicatedVolume(commandEnv *CommandEnv, writer io.Writer, applyChanges bool, volumeReplicas map[uint32][]*VolumeReplica, vid uint32, allLocations []location) (bool, error) { +// volumeCopyScheduler serializes destination selection for concurrent volume +// copies: selection and free-slot accounting are atomic so parallel fixes see +// each other's reservations, and the per-server cap keeps many simultaneous +// copies from swamping one destination's disks. +type volumeCopyScheduler struct { + mu sync.Mutex + cond *sync.Cond + inflight map[string]int // destination dataNode.Id -> copies in flight + maxPerServer int +} + +func newVolumeCopyScheduler(maxPerServer int) *volumeCopyScheduler { + if maxPerServer <= 0 { + maxPerServer = 1 + } + s := &volumeCopyScheduler{ + inflight: make(map[string]int), + maxPerServer: maxPerServer, + } + s.cond = sync.NewCond(&s.mu) + return s +} + +// reserveTarget picks the emptiest data node satisfying the replica placement +// and reserves a volume slot on it. When every eligible destination is at the +// per-server copy cap it waits for a copy to finish instead of failing. +// Returns nil only when no data node can accept the replica at all. With +// countInflight=false (simulation) the slot is reserved but no copy is +// counted in flight. +func (s *volumeCopyScheduler) reserveTarget(replicaPlacement *super_block.ReplicaPlacement, replicas []*VolumeReplica, allLocations []location, diskType string, countInflight bool) *location { + s.mu.Lock() + defer s.mu.Unlock() + fn := capacityByFreeVolumeCount(types.ToDiskType(diskType)) + for { + keepDataNodesSorted(allLocations, types.ToDiskType(diskType)) + eligibleButBusy := false + for _, dst := range allLocations { + // check whether data nodes satisfy the constraints + if fn(dst.dataNode) <= 0 || !satisfyReplicaPlacement(replicaPlacement, replicas, dst) { + continue + } + if countInflight && s.inflight[dst.dataNode.Id] >= s.maxPerServer { + eligibleButBusy = true + continue + } + addVolumeCount(dst.dataNode.DiskInfos[diskType], 1) + if countInflight { + s.inflight[dst.dataNode.Id]++ + } + return &dst + } + if !eligibleButBusy { + return nil + } + s.cond.Wait() + } +} + +// releaseTarget ends a copy counted by reserveTarget. A failed copy also +// returns the reserved volume slot, so retries do not drain the topology's +// free-slot accounting. +func (s *volumeCopyScheduler) releaseTarget(dst *location, diskType string, copied bool) { + s.mu.Lock() + defer s.mu.Unlock() + s.inflight[dst.dataNode.Id]-- + if s.inflight[dst.dataNode.Id] <= 0 { + delete(s.inflight, dst.dataNode.Id) + } + if !copied { + addVolumeCount(dst.dataNode.DiskInfos[diskType], -1) + } + s.cond.Broadcast() +} + +func (c *commandVolumeFixReplication) fixOneUnderReplicatedVolume(commandEnv *CommandEnv, writer io.Writer, applyChanges bool, volumeReplicas map[uint32][]*VolumeReplica, vid uint32, allLocations []location, scheduler *volumeCopyScheduler) (bool, error) { replicas := volumeReplicas[vid] replica := pickOneReplicaToCopyFrom(replicas) replicaPlacement, _ := super_block.NewReplicaPlacementFromByte(byte(replica.info.ReplicaPlacement)) - foundNewLocation := false - keepDataNodesSorted(allLocations, types.ToDiskType(replica.info.DiskType)) - fn := capacityByFreeVolumeCount(types.ToDiskType(replica.info.DiskType)) - for _, dst := range allLocations { - // check whether data nodes satisfy the constraints - if fn(dst.dataNode) > 0 && satisfyReplicaPlacement(replicaPlacement, replicas, dst) { - // ask the volume server to replicate the volume - foundNewLocation = true - fmt.Fprintf(writer, "replicating volume %d %s from %s to dataNode %s ...\n", replica.info.Id, replicaPlacement, replica.location.dataNode.Id, dst.dataNode.Id) - if !applyChanges { - // adjust volume count - addVolumeCount(dst.dataNode.DiskInfos[replica.info.DiskType], 1) - return true, nil - } - - err := replicateVolumeToServer(commandEnv.option.GrpcDialOption, writer, needle.VolumeId(replica.info.Id), - pb.NewServerAddressFromDataNode(replica.location.dataNode), - pb.NewServerAddressFromDataNode(dst.dataNode), - replica.info.DiskType) - - if err != nil { - return false, err - } - - // adjust volume count - addVolumeCount(dst.dataNode.DiskInfos[replica.info.DiskType], 1) - return true, nil - } - } - - if !foundNewLocation { + dst := scheduler.reserveTarget(replicaPlacement, replicas, allLocations, replica.info.DiskType, applyChanges) + if dst == nil { fmt.Fprintf(writer, "failed to place volume %d replica as %s, existing:%+v\n", replica.info.Id, replicaPlacement, len(replicas)) + return false, nil } - return false, nil + + // ask the volume server to replicate the volume + fmt.Fprintf(writer, "replicating volume %d %s from %s to dataNode %s ...\n", replica.info.Id, replicaPlacement, replica.location.dataNode.Id, dst.dataNode.Id) + + if !applyChanges { + return true, nil + } + + err := replicateVolumeToServer(commandEnv.option.GrpcDialOption, writer, needle.VolumeId(replica.info.Id), + pb.NewServerAddressFromDataNode(replica.location.dataNode), + pb.NewServerAddressFromDataNode(dst.dataNode), + replica.info.DiskType) + scheduler.releaseTarget(dst, replica.info.DiskType, err == nil) + if err != nil { + return false, err + } + + return true, nil } func addVolumeCount(info *master_pb.DiskInfo, count int) { diff --git a/weed/shell/command_volume_fix_replication_parallel_test.go b/weed/shell/command_volume_fix_replication_parallel_test.go new file mode 100644 index 000000000..4eeaaf7fc --- /dev/null +++ b/weed/shell/command_volume_fix_replication_parallel_test.go @@ -0,0 +1,144 @@ +package shell + +import ( + "io" + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/pb/master_pb" + "github.com/seaweedfs/seaweedfs/weed/storage/super_block" +) + +func testFixLocation(dc, rack, id string, maxVolumes int64) location { + return location{ + dc: dc, + rack: rack, + dataNode: &master_pb.DataNodeInfo{ + Id: id, + DiskInfos: map[string]*master_pb.DiskInfo{ + "": {MaxVolumeCount: maxVolumes, FreeVolumeCount: maxVolumes}, + }, + }, + } +} + +func TestReserveTargetSpreadsAcrossServers(t *testing.T) { + src := testFixLocation("dc1", "r1", "dn0", 0) + allLocations := []location{ + src, + testFixLocation("dc1", "r1", "dn1", 10), + testFixLocation("dc1", "r1", "dn2", 5), + testFixLocation("dc1", "r1", "dn3", 3), + } + // keepDataNodesSorted reorders allLocations in place, so keep stable + // per-node copies; the shared dataNode pointers carry the accounting + byId := make(map[string]*location) + for _, loc := range allLocations { + loc := loc + byId[loc.dataNode.Id] = &loc + } + replicas := []*VolumeReplica{ + {location: &src, info: &master_pb.VolumeInformationMessage{Id: 1}}, + } + rp, _ := super_block.NewReplicaPlacementFromString("001") + + // with all copies in flight, consecutive reservations must not converge on + // the emptiest server + s := newVolumeCopyScheduler(1) + var reserved []*location + got := make(map[string]bool) + for i := 0; i < 3; i++ { + dst := s.reserveTarget(rp, replicas, allLocations, "", true) + if dst == nil { + t.Fatalf("reservation %d found no destination", i) + } + reserved = append(reserved, dst) + got[dst.dataNode.Id] = true + } + for _, id := range []string{"dn1", "dn2", "dn3"} { + if !got[id] { + t.Errorf("expected a reservation on %s, got %v", id, got) + } + } + + // every eligible destination is at the copy cap: the next reservation + // waits for a free slot instead of failing + done := make(chan *location) + go func() { + done <- s.reserveTarget(rp, replicas, allLocations, "", true) + }() + select { + case dst := <-done: + t.Fatalf("reserveTarget should wait while all destinations are at the copy cap, got %s", dst.dataNode.Id) + case <-time.After(100 * time.Millisecond): + } + s.releaseTarget(byId["dn1"], "", true) + var waited *location + select { + case waited = <-done: + case <-time.After(5 * time.Second): + t.Fatal("reserveTarget did not wake up after a copy slot freed") + } + if waited == nil || waited.dataNode.Id != "dn1" { + t.Fatalf("expected the freed dn1 to take the waiting copy, got %+v", waited) + } + + // a successful copy keeps its volume slot, a failed one returns it + s.releaseTarget(waited, "", true) + if count := byId["dn1"].dataNode.DiskInfos[""].VolumeCount; count != 2 { + t.Errorf("dn1 should keep 2 reserved slots, got %d", count) + } + for _, dst := range reserved { + if dst.dataNode.Id != "dn1" { + s.releaseTarget(dst, "", false) + if count := dst.dataNode.DiskInfos[""].VolumeCount; count != 0 { + t.Errorf("%s should have its failed reservation returned, got volume count %d", dst.dataNode.Id, count) + } + } + } + if len(s.inflight) != 0 { + t.Errorf("all copies released, but %d still in flight", len(s.inflight)) + } +} + +func TestFixUnderReplicatedVolumesInParallel(t *testing.T) { + src := testFixLocation("dc1", "r1", "dn0", 0) + allLocations := []location{ + src, + testFixLocation("dc1", "r1", "dn1", 4), + testFixLocation("dc1", "r1", "dn2", 4), + testFixLocation("dc1", "r1", "dn3", 4), + } + rp, _ := super_block.NewReplicaPlacementFromString("001") + + volumeReplicas := make(map[uint32][]*VolumeReplica) + var volumeIds []uint32 + for vid := uint32(1); vid <= 12; vid++ { + volumeReplicas[vid] = []*VolumeReplica{ + {location: &src, info: &master_pb.VolumeInformationMessage{Id: vid, ReplicaPlacement: uint32(rp.Byte())}}, + } + volumeIds = append(volumeIds, vid) + } + + c := &commandVolumeFixReplication{collectionPattern: new(string)} + fixedVolumes, err := c.fixUnderReplicatedVolumes(nil, io.Discard, false, volumeIds, volumeReplicas, allLocations, 0, 0, 8, 1) + if err != nil { + t.Fatalf("fixUnderReplicatedVolumes: %v", err) + } + if len(fixedVolumes) != 0 { + t.Errorf("simulation should not record fixed volumes, got %d", len(fixedVolumes)) + } + + // 12 volumes must exactly fill the 3x4 free slots without over-reserving + // any single destination + for _, loc := range allLocations { + if loc.dataNode.Id == "dn0" { + continue + } + diskInfo := loc.dataNode.DiskInfos[""] + if diskInfo.VolumeCount != 4 || diskInfo.FreeVolumeCount != 0 { + t.Errorf("%s expected exactly 4 reserved slots, got volume count %d, free %d", + loc.dataNode.Id, diskInfo.VolumeCount, diskInfo.FreeVolumeCount) + } + } +}