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.
This commit is contained in:
Chris Lu
2026-08-08 20:55:55 -07:00
committed by GitHub
parent 0dd33bd7ec
commit 2dc59c9b51
2 changed files with 182 additions and 4 deletions
+19 -4
View File
@@ -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)
@@ -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")
}
}