From 2dc59c9b5111ba0a471b197585c8e3acfee5b1d7 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Sat, 8 Aug 2026 20:55:55 -0700 Subject: [PATCH] topology: track volume size only where writes can land (#10653) * topology: track volume size only where writes can land Size tracking decays pending assignment estimates so the master does not overfill a volume before heartbeats catch up. Nothing is ever assigned to a read-only volume, so an entry for one can never be consulted -- and in a tiered cluster that is most of them, which made this the volume layout's largest cost. A volume held out of the writable list for capacity is not read-only and keeps its entry: that entry is what enforces the recovery delay. 800k volumes, 90% read-only: sizeTracking 79.1MB -> 8.4MB, and the crowded set falls out with it because a read-only volume no longer reaches the threshold check at all. * topology: decide size tracking per volume, not per reporting replica A volume is unwritable if any replica is read-only, so asking the replica whose heartbeat happened to arrive made the answer depend on arrival order: a writable replica reporting after a read-only one put the tracking back. Ask the volume instead, which also drops the caller-supplied flag and the churn it caused. The crowded entry goes with the tracking, since leaving it behind would only move the memory this releases. Costs a map lookup per replica on a full-list heartbeat, about 19ms per 100k volumes and no allocations, on a path that is now rare. --- weed/topology/volume_layout.go | 23 ++- .../volume_layout_size_tracking_test.go | 163 ++++++++++++++++++ 2 files changed, 182 insertions(+), 4 deletions(-) create mode 100644 weed/topology/volume_layout_size_tracking_test.go diff --git a/weed/topology/volume_layout.go b/weed/topology/volume_layout.go index 32b0ceebf..91d7aa7c0 100644 --- a/weed/topology/volume_layout.go +++ b/weed/topology/volume_layout.go @@ -197,7 +197,9 @@ func (vl *VolumeLayout) RegisterVolume(v *storage.VolumeInfo, dn *DataNode) { defer vl.rememberOversizedVolume(v, dn) moveLookupOwnership(v.Id, vl.getOrCreateLocationList(v.Id).Set(dn), dn) - vl.initSizeTracking(v.Id, v.Size, v.CompactRevision) + if !v.ReadOnly { + vl.initSizeTracking(v.Id, v.Size, v.CompactRevision) + } // glog.V(4).Infof("volume %d added to %s len %d copy %d", v.Id, dn.Id(), vl.vid2location[v.Id].Length(), v.ReplicaPlacement.GetCopyCount()) for _, dn := range vl.vid2location[v.Id].list { if vInfo, err := dn.GetVolumesById(v.Id); err == nil { @@ -243,6 +245,19 @@ func (vl *VolumeLayout) UpdateVolumeSize(vid needle.VolumeId, reportedSize uint6 vl.accessLock.Lock() defer vl.accessLock.Unlock() + // Tracking exists to place writes, and nothing is written to a volume any + // replica reports read-only. Most volumes in a tiered cluster are read-only, + // so an entry each is the layout's largest cost. Asked of the volume rather + // than of the replica reporting, so which replica arrives first cannot + // decide the answer. A volume held out of the writable list for capacity is + // still all-writable and keeps its entry, which is what enforces the + // recovery delay. + if !vl.isAllWritable(vid) { + delete(vl.sizeTracking, vid) + vl.removeFromCrowded(vid) + return false + } + now := time.Now() st := vl.sizeTracking[vid] if st == nil { @@ -919,14 +934,14 @@ func (vl *VolumeLayout) SetVolumeAvailable(dn *DataNode, vid needle.VolumeId, is } // A disconnect during a long vacuum can drop the entry while the volume is - // still on the node; re-create it (and seed size tracking) instead of - // dereferencing a nil location, so the commit also repairs the split. + // still on the node; re-create it instead of dereferencing a nil location, + // so the commit also repairs the split. moveLookupOwnership(vid, vl.getOrCreateLocationList(vid).Set(dn), dn) - vl.initSizeTracking(vid, vInfo.Size, vInfo.CompactRevision) if vInfo.ReadOnly || isReadOnly || isFullCapacity { return false } + vl.initSizeTracking(vid, vInfo.Size, vInfo.CompactRevision) if vl.enoughCopies(vid) { becameWritable = vl.setVolumeWritable(vid) diff --git a/weed/topology/volume_layout_size_tracking_test.go b/weed/topology/volume_layout_size_tracking_test.go new file mode 100644 index 000000000..e5d375546 --- /dev/null +++ b/weed/topology/volume_layout_size_tracking_test.go @@ -0,0 +1,163 @@ +package topology + +import ( + "testing" + + "github.com/seaweedfs/seaweedfs/weed/pb/master_pb" + "github.com/seaweedfs/seaweedfs/weed/storage/needle" + "github.com/seaweedfs/seaweedfs/weed/storage/super_block" + "github.com/seaweedfs/seaweedfs/weed/storage/types" +) + +func sizeTrackingLayout(t *testing.T) (*Topology, *DataNode, *VolumeLayout) { + t.Helper() + topo := NewTopology("st", nil, 32*1024, 5, false) + dn := topo.GetOrCreateDataCenter("dc1").GetOrCreateRack("rack1"). + GetOrCreateDataNode("127.0.0.1", 8080, 18080, "", "", map[string]uint32{"": 100}) + rp, _ := super_block.NewReplicaPlacementFromString("000") + return topo, dn, topo.GetVolumeLayout("c", rp, needle.EMPTY_TTL, types.HardDriveType) +} + +func sizeTrackingVolume(id uint32, size uint64, readOnly bool) *master_pb.VolumeInformationMessage { + return &master_pb.VolumeInformationMessage{ + Id: id, Size: size, Collection: "c", Version: 3, ReadOnly: readOnly, + } +} + +func TestSizeTrackingSkipsReadOnlyVolumes(t *testing.T) { + topo, dn, vl := sizeTrackingLayout(t) + topo.SyncDataNodeRegistration([]*master_pb.VolumeInformationMessage{ + sizeTrackingVolume(1, 1000, false), sizeTrackingVolume(2, 1000, true), + }, dn) + + vl.accessLock.RLock() + _, writableTracked := vl.sizeTracking[needle.VolumeId(1)] + _, readOnlyTracked := vl.sizeTracking[needle.VolumeId(2)] + vl.accessLock.RUnlock() + + if !writableTracked { + t.Error("a writable volume lost its size tracking, so assignments stop being accounted for") + } + if readOnlyTracked { + t.Error("a read-only volume is tracked for writes it can never take") + } +} + +// A volume that goes read-only after the master already tracked it has to give +// the entry back, or a cluster that tiers its volumes keeps paying for them. +func TestSizeTrackingReleasedWhenAVolumeGoesReadOnly(t *testing.T) { + topo, dn, vl := sizeTrackingLayout(t) + topo.SyncDataNodeRegistration([]*master_pb.VolumeInformationMessage{sizeTrackingVolume(1, 1000, false)}, dn) + + topo.SyncDataNodeRegistration([]*master_pb.VolumeInformationMessage{sizeTrackingVolume(1, 1000, true)}, dn) + + vl.accessLock.RLock() + _, tracked := vl.sizeTracking[needle.VolumeId(1)] + vl.accessLock.RUnlock() + if tracked { + t.Error("a volume that became read-only kept its size tracking") + } +} + +func TestSizeTrackingReturnsWhenAVolumeBecomesWritable(t *testing.T) { + topo, dn, vl := sizeTrackingLayout(t) + topo.SyncDataNodeRegistration([]*master_pb.VolumeInformationMessage{sizeTrackingVolume(1, 1000, true)}, dn) + + topo.SyncDataNodeRegistration([]*master_pb.VolumeInformationMessage{sizeTrackingVolume(1, 2000, false)}, dn) + + vl.accessLock.RLock() + st, tracked := vl.sizeTracking[needle.VolumeId(1)] + vl.accessLock.RUnlock() + if !tracked { + t.Fatal("a volume that became writable was not tracked again") + } + if st.reportedSize != 2000 { + t.Errorf("tracking resumed at size %d, want the reported 2000", st.reportedSize) + } +} + +// A volume held out of the writable list for capacity is not read-only, and its +// entry is what enforces the recovery delay. +func TestSizeTrackingSurvivesAFullVolume(t *testing.T) { + topo, dn, vl := sizeTrackingLayout(t) + topo.SyncDataNodeRegistration([]*master_pb.VolumeInformationMessage{sizeTrackingVolume(1, 1000, false)}, dn) + + vl.RecordAssign(needle.VolumeId(1), int64(32*1024)) + + vl.accessLock.RLock() + st, tracked := vl.sizeTracking[needle.VolumeId(1)] + vl.accessLock.RUnlock() + if !tracked { + t.Fatal("a full volume lost its tracking, so nothing holds it out of the writable list") + } + if st.fullSince.IsZero() { + t.Error("a full volume did not record when it filled, so the recovery delay cannot apply") + } +} + +// A volume is unwritable if any replica is read-only, so the answer must not +// depend on which replica's heartbeat arrives first. +func TestSizeTrackingIgnoresReplicaReportOrder(t *testing.T) { + for _, readOnlyFirst := range []bool{true, false} { + name := "WritableFirst" + if readOnlyFirst { + name = "ReadOnlyFirst" + } + t.Run(name, func(t *testing.T) { + topo := NewTopology("st", nil, 32*1024, 5, false) + rack := topo.GetOrCreateDataCenter("dc1").GetOrCreateRack("rack1") + writableNode := rack.GetOrCreateDataNode("10.0.0.1", 8080, 18080, "", "a", map[string]uint32{"": 100}) + readOnlyNode := rack.GetOrCreateDataNode("10.0.0.2", 8080, 18080, "", "b", map[string]uint32{"": 100}) + + report := func(dn *DataNode, readOnly bool) { + topo.SyncDataNodeRegistration([]*master_pb.VolumeInformationMessage{ + {Id: 1, Size: 1000, Collection: "c", Version: 3, ReplicaPlacement: 1, ReadOnly: readOnly}, + }, dn) + } + if readOnlyFirst { + report(readOnlyNode, true) + report(writableNode, false) + } else { + report(writableNode, false) + report(readOnlyNode, true) + } + + rp, _ := super_block.NewReplicaPlacementFromString("001") + vl := topo.GetVolumeLayout("c", rp, needle.EMPTY_TTL, types.HardDriveType) + vl.accessLock.RLock() + _, tracked := vl.sizeTracking[needle.VolumeId(1)] + _, crowded := vl.crowded[needle.VolumeId(1)] + vl.accessLock.RUnlock() + + if tracked { + t.Error("a volume with a read-only replica is tracked for writes it cannot take") + } + if crowded { + t.Error("a volume with a read-only replica was left in the crowded set") + } + }) + } +} + +// The crowded entry has to go with the tracking, or the memory this releases +// is only moved. +func TestCrowdedEntryReleasedWithSizeTracking(t *testing.T) { + topo, dn, vl := sizeTrackingLayout(t) + topo.SyncDataNodeRegistration([]*master_pb.VolumeInformationMessage{sizeTrackingVolume(1, 31000, false)}, dn) + + vl.accessLock.RLock() + _, crowdedWhileWritable := vl.crowded[needle.VolumeId(1)] + vl.accessLock.RUnlock() + if !crowdedWhileWritable { + t.Fatal("expected a nearly full writable volume to be crowded") + } + + topo.SyncDataNodeRegistration([]*master_pb.VolumeInformationMessage{sizeTrackingVolume(1, 31000, true)}, dn) + + vl.accessLock.RLock() + _, stillCrowded := vl.crowded[needle.VolumeId(1)] + vl.accessLock.RUnlock() + if stillCrowded { + t.Error("a volume that became read-only kept its crowded entry") + } +}