filer: stop the aggregated metadata subscribe loop rescanning an exhausted persisted log (#11644)

* fix(filer): gate the aggregated metadata disk pass on real change

A subscriber whose start position is past the end of the local persisted
log re-ran the whole persisted-log pass - store listings, file opens,
readahead - on every loop iteration. Each iteration is paced only by the
shortest wake (the 20ms hold floor on a busy watermark), so one parked
subscriber kept a full CPU core busy for the life of the stream.

The aggregated loop now mirrors the local loop's gate: the disk pass
runs on the first pass and afterwards only when something it cannot
miss changed - a local flush landed, the peers' flush low-watermark
advanced (more content admitted, or new files in a shared store), the
cursor moved, or a disk hold is pending (the ring read that follows an
empty pass parks internally, so skipping there would strand a held
entry).

Regression test: a subscriber parked past the persisted-log tail holds
the listing rate near zero and still delivers once peers report
progress.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: re-arm the aggregated disk pass on unobserved change

Review found three staleness classes the gate could not see: the flush
low-watermark only catching rises (a joining peer lowers the minimum and
invalidates an earlier pass's proof), a peer past the minimum landing a
file without moving it, and a chunk subscriber's refs-stop bound
advancing with wall time. Re-read when the low-watermark moves in either
direction, when the chunk listing bound admits more files, and on a slow
re-probe cadence for files no watermark can signal.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: unwind the parked ring read so the disk re-probe runs, and re-read on cursor rewinds

A caught-up subscriber parks inside LoopProcessLogData's wait loop, so
the re-probe interval in the outer disk gate could never elapse there;
the callback now unwinds the read once the cadence is due so the gate
re-evaluates. The cursor trigger also needs to notice rewinds, not just
advances, since ResumeFromDiskError moves the cursor backward.

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
Chris LuandDevin authored and GitHub committed 2026-10-08 22:03:31 +08:00
1 parent 49c25890ad
commit 14fdd61aea
2 files changed
+208 -68

No files matched your search

+129 -68
View File
@@ -27,6 +27,11 @@ var (
// (possibly-unflushed) gap, in case the flush notification is missed.
unflushedGapRetryInterval = 2 * time.Second
// aggDiskReprobeInterval paces the aggregated persisted-log re-listing for
// files no watermark signals: a peer past the flush low-watermark can land
// a file without moving the minimum.
aggDiskReprobeInterval = 2 * time.Second
// gapStallWarnInterval paces the warning for a subscriber that stays parked.
gapStallWarnInterval = time.Minute
@@ -666,8 +671,9 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
var lastHeartbeatNs int64
baseEachLogEntryFn := eachLogEntryFn(req, sender, eachEventNotificationFn, &unsyncedEvents)
// heldAtTsNs remembers the entry a read was held at (for the log line);
// the rewind target is the last entry actually delivered.
var heldAtTsNs int64
// diskHeldAtTsNs is the same marker for the disk pass alone: a pending
// disk hold keeps the pass re-reading until the entry is served.
var heldAtTsNs, diskHeldAtTsNs int64
// Each read path holds at its own watermark: persisted logs are complete
// only up to every peer's flush watermark, the ring only up to every
// peer's delivery watermark.
@@ -698,7 +704,14 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
// What the last disk pass proved covered: flushed on every peer AND inside
// the pass's listing, so an empty pass proves (cursor, proven] empty.
var diskPassProvenTsNs int64
diskEachLogEntryFn := guardedEachLogEntryFn(func() int64 { return diskPassHoldTsNs })
diskBaseEachLogEntryFn := guardedEachLogEntryFn(func() int64 { return diskPassHoldTsNs })
diskEachLogEntryFn := func(logEntry *filer_pb.LogEntry) (bool, error) {
isDone, err := diskBaseEachLogEntryFn(logEntry)
if errors.Is(err, errHeldByPeerWatermark) {
diskHeldAtTsNs = logEntry.TsNs
}
return isDone, err
}
memEachLogEntryFn := guardedEachLogEntryFn(holdMemTsNs)
// waitHeld pauses a held read until a peer reports further progress, or
// the retry interval elapses (a peer dropped past its grace, or a log file
@@ -733,6 +746,11 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
var readPersistedLogErr error
var readInMemoryLogErr error
var isDone bool
var lastCheckedFlushTsNs int64 = -1 // Track the last local flush we read the disk under
var lastCheckedFlushLowTsNs int64 = -1 // Track the last peer flush low-watermark we read the disk under
var lastDiskReadTsNs int64 = -1 // Track the last read position we used for disk read
var lastDiskRefsStopTsNs int64 = -1 // Track the last chunk listing bound we read the disk under
var lastDiskPassAt time.Time // Paces re-probes for files no watermark signals
sentRefs := make(map[string]sentRefState)
aggBuffer := fs.filer.MetaAggregator.MetaLogBuffer
@@ -765,78 +783,110 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
// diskPassHoldTsNs above).
diskPassFlushLowTsNs = fs.filer.MetaAggregator.PeerLowFlushWatermarkTsNs()
diskPassHoldTsNs = resolveAggReadHoldTsNs(diskPassFlushLowTsNs, time.Now().UnixNano(), metadataGapSettledHorizon)
diskPassProvenTsNs = diskPassFlushLowTsNs
// Re-read the disk only when something changed it cannot miss: a local
// flush landed, the peers' flush low-watermark moved in either
// direction (a joining peer invalidates what an earlier pass proved),
// the cursor moved, or a disk hold is pending (the ring read after an
// empty pass parks internally, so skipping would strand the held
// entry). A peer past the low-watermark can still land a file without
// moving it, so the listing is also re-probed at a slow cadence;
// chunk listings re-arm as soon as their read bound admits more files.
currentFlushTsNs := fs.filer.LocalMetaLogBuffer.GetLastFlushTsNs()
currentReadTsNs := lastReadTime.Time.UnixNano()
var currentRefsStopTsNs int64
if req.ClientSupportsMetadataChunks {
refsStopTsNs := chunkRefsStopTsNs(diskPassHoldTsNs, req.UntilNs)
// Nothing above the listing bound is proven by this pass.
if refsStopTsNs < diskPassProvenTsNs {
diskPassProvenTsNs = refsStopTsNs
}
if refsStopTsNs > lastReadTime.Time.UnixNano() {
processedTsNs, isDone, readPersistedLogErr = fs.chunkDiskPass(ctx, sender, lastReadTime, refsStopTsNs, sentRefs, nil)
} else {
processedTsNs, isDone, readPersistedLogErr = 0, false, nil
}
} else {
processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, diskEachLogEntryFn)
currentRefsStopTsNs = chunkRefsStopTsNs(diskPassHoldTsNs, req.UntilNs)
}
if errors.Is(readPersistedLogErr, errHeldByPeerWatermark) {
// Stay at the last delivered entry; the held entry is re-read (and
// re-checked) by the next pass.
if processedTsNs > 0 {
shouldReadFromDisk := lastCheckedFlushTsNs == -1 ||
currentFlushTsNs > lastCheckedFlushTsNs ||
diskPassFlushLowTsNs != lastCheckedFlushLowTsNs ||
currentReadTsNs != lastDiskReadTsNs ||
currentRefsStopTsNs > lastDiskRefsStopTsNs ||
diskHeldAtTsNs != 0 ||
time.Since(lastDiskPassAt) >= aggDiskReprobeInterval
diskAdvanced := false
if shouldReadFromDisk {
lastCheckedFlushTsNs = currentFlushTsNs
lastCheckedFlushLowTsNs = diskPassFlushLowTsNs
lastDiskReadTsNs = currentReadTsNs
lastDiskRefsStopTsNs = currentRefsStopTsNs
lastDiskPassAt = time.Now()
diskHeldAtTsNs = 0
diskPassProvenTsNs = diskPassFlushLowTsNs
if req.ClientSupportsMetadataChunks {
refsStopTsNs := currentRefsStopTsNs
// Nothing above the listing bound is proven by this pass.
if refsStopTsNs < diskPassProvenTsNs {
diskPassProvenTsNs = refsStopTsNs
}
if refsStopTsNs > lastReadTime.Time.UnixNano() {
processedTsNs, isDone, readPersistedLogErr = fs.chunkDiskPass(ctx, sender, lastReadTime, refsStopTsNs, sentRefs, nil)
} else {
processedTsNs, isDone, readPersistedLogErr = 0, false, nil
}
} else {
processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, diskEachLogEntryFn)
}
if errors.Is(readPersistedLogErr, errHeldByPeerWatermark) {
// Stay at the last delivered entry; the held entry is re-read (and
// re-checked) by the next pass.
if processedTsNs > 0 {
lastReadTime = log_buffer.NewMessagePosition(processedTsNs, gapResumeCursorOffset)
if processedTsNs > diskAnchorTsNs {
diskAnchorTsNs = processedTsNs
}
}
// A hold is not a gap: clear any stale ResumeFromDiskError so the
// next pass's disk-miss handling cannot skip past the held entry.
readInMemoryLogErr = nil
if !waitHeld("disk", flushChan) {
return nil
}
continue
}
if readPersistedLogErr != nil {
return fmt.Errorf("reading from persisted logs: %w", readPersistedLogErr)
}
if isDone {
return nil
}
glog.V(4).Infof("processed to %v: %v", clientName, processedTsNs)
diskAdvanced = diskReadAdvanced(processedTsNs, lastReadTime)
// Read after the disk read (an eviction landing mid-read must count) and
// in received-ts space: the ring's bumped stopTimes exceed anything on
// any peer's disk, and gating disk cursors on them parks subscribers
// that drained every peer's log.
lastEvictedTsNs := fs.filer.MetaAggregator.MetaLogBuffer.GetLastEvictedOriginalTsNs()
if diskAdvanced {
gapStall.resumed()
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, processedTsNs, lastEvictedTsNs, diskPassFlushLowTsNs, clientName, req.PathPrefix)
lastReadTime = log_buffer.NewMessagePosition(processedTsNs, gapResumeCursorOffset)
if processedTsNs > diskAnchorTsNs {
diskAnchorTsNs = processedTsNs
}
}
// A hold is not a gap: clear any stale ResumeFromDiskError so the
// next pass's disk-miss handling cannot skip past the held entry.
readInMemoryLogErr = nil
if !waitHeld("disk", flushChan) {
return nil
}
continue
}
if readPersistedLogErr != nil {
return fmt.Errorf("reading from persisted logs: %w", readPersistedLogErr)
}
if isDone {
return nil
}
glog.V(4).Infof("processed to %v: %v", clientName, processedTsNs)
diskAdvanced := diskReadAdvanced(processedTsNs, lastReadTime)
// Read after the disk read (an eviction landing mid-read must count) and
// in received-ts space: the ring's bumped stopTimes exceed anything on
// any peer's disk, and gating disk cursors on them parks subscribers
// that drained every peer's log.
lastEvictedTsNs := fs.filer.MetaAggregator.MetaLogBuffer.GetLastEvictedOriginalTsNs()
if diskAdvanced {
gapStall.resumed()
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, processedTsNs, lastEvictedTsNs, diskPassFlushLowTsNs, clientName, req.PathPrefix)
lastReadTime = log_buffer.NewMessagePosition(processedTsNs, gapResumeCursorOffset)
if processedTsNs > diskAnchorTsNs {
diskAnchorTsNs = processedTsNs
}
} else if readInMemoryLogErr == nil {
// Nothing on disk and memory never spoke: scan forward for the next
// day that has logs.
nextDayTs := util.GetNextDayTsNano(lastReadTime.Time.UnixNano())
// The day jump delivers nothing; stay put until the hold point
// covers the skipped range.
if nextDayTs <= diskPassHoldTsNs {
position := log_buffer.NewMessagePosition(nextDayTs, gapResumeCursorOffset)
found, err := fs.filer.HasPersistedLogFiles(position)
if err != nil {
return fmt.Errorf("checking persisted log files: %w", err)
}
if found {
gapStall.resumed()
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, nextDayTs, lastEvictedTsNs, diskPassFlushLowTsNs, clientName, req.PathPrefix)
lastReadTime = position
if nextDayTs > diskAnchorTsNs {
diskAnchorTsNs = nextDayTs
} else if readInMemoryLogErr == nil {
// Nothing on disk and memory never spoke: scan forward for the next
// day that has logs.
nextDayTs := util.GetNextDayTsNano(lastReadTime.Time.UnixNano())
// The day jump delivers nothing; stay put until the hold point
// covers the skipped range.
if nextDayTs <= diskPassHoldTsNs {
position := log_buffer.NewMessagePosition(nextDayTs, gapResumeCursorOffset)
found, err := fs.filer.HasPersistedLogFiles(position)
if err != nil {
return fmt.Errorf("checking persisted log files: %w", err)
}
if found {
gapStall.resumed()
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, nextDayTs, lastEvictedTsNs, diskPassFlushLowTsNs, clientName, req.PathPrefix)
lastReadTime = position
if nextDayTs > diskAnchorTsNs {
diskAnchorTsNs = nextDayTs
}
}
}
}
@@ -867,6 +917,7 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
// every event whose original timestamp is at or below this.
preMemDeliveryLowTsNs := fs.filer.MetaAggregator.PeerLowWatermarkTsNs()
diskReprobeDue := false
lastReadTime, isDone, readInMemoryLogErr = fs.filer.MetaAggregator.MetaLogBuffer.LoopProcessLogData(aggReaderName, lastReadTime, req.UntilNs, func() bool {
select {
case <-ctx.Done():
@@ -876,6 +927,13 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
if !fs.hasClient(req.ClientId, req.ClientEpoch) {
return false
}
// Caught-up readers park in the inner wait loop; the outer disk
// gate never runs again unless this read returns, so unwind to
// re-probe the persisted logs on the slow cadence.
if time.Since(lastDiskPassAt) >= aggDiskReprobeInterval {
diskReprobeDue = true
return false
}
// Contiguous and caught up: advance the anchor to the delivery
// low-watermark so long live tails keep eviction rewinds short.
// Only once the run is connected to the ring - the empty-ring
@@ -919,6 +977,9 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
}
}
if isDone {
if diskReprobeDue {
continue
}
return nil
}
if !fs.hasClient(req.ClientId, req.ClientEpoch) {
+79
View File
@@ -0,0 +1,79 @@
package weed_server
// The persisted log may end before a subscriber's start position (a filer
// whose own journal is older than the position a backup client resumes from).
// With nothing on disk and every ring entry held by a peer watermark, the
// aggregated loop used to re-list and re-read the persisted log on every wake
// - about a full CPU core per such subscriber. The disk pass may only re-run
// when something it cannot miss has changed.
import (
"context"
"strings"
"sync/atomic"
"testing"
"time"
"github.com/seaweedfs/seaweedfs/weed/filer"
"github.com/seaweedfs/seaweedfs/weed/util"
)
type countingStore struct {
filer.FilerStore
logLists *atomic.Int64
}
func (s *countingStore) ListDirectoryPrefixedEntries(ctx context.Context, dirPath util.FullPath, startFileName string, includeStartFile bool, limit int64, prefix string, eachEntryFunc filer.ListEachEntryFunc) (lastFileName string, err error) {
if strings.HasPrefix(string(dirPath), filer.SystemLogDir) {
s.logLists.Add(1)
}
return s.FilerStore.ListDirectoryPrefixedEntries(ctx, dirPath, startFileName, includeStartFile, limit, prefix, eachEntryFunc)
}
func TestSubscribeLoop_AggregatedNoPersistedEntryAfterStart(t *testing.T) {
h := newSubscribeHarness(t)
lists := &atomic.Int64{}
h.f.SetStore(&countingStore{FilerStore: h.f.GetStore(), logLists: lists})
// Persisted log ends at T1: one flushed window, nothing after.
h.append(h.tsAt(0, 0))
h.append(h.tsAt(0, 1))
h.f.LocalMetaLogBuffer.ForceFlush()
waitForFlushedFiles(t, h, h.tsAt(0, 1))
// Client cursor sits just past the last persisted entry.
cursor := h.tsAt(0, 1) + int64(time.Millisecond)
ma := h.startAggregator()
// Aggregated ring holds only much newer entries (peer events).
recent := time.Now().UnixNano()
h.appendAggregated(recent)
h.appendAggregated(recent + int64(time.Millisecond))
// Peers' watermarks are stuck at the old log tail.
reportPeersAt(ma, h.tsAt(0, 1), h.tsAt(0, 1))
r := h.subscribeAggregated(cursor)
// Warm up: the first pass plus the cursor-move re-read after the gap
// machinery re-arms the cursor are legitimate.
time.Sleep(150 * time.Millisecond)
before := lists.Load()
time.Sleep(500 * time.Millisecond)
rate := float64(lists.Load()-before) / 0.5
if rate > 4 {
t.Fatalf("%.1f persisted-log listings per second while parked; the disk pass re-ran on every wake", rate)
}
// Peer progress through the held entries releases the read: the held
// events are delivered without another disk pass.
lists.Store(0)
reportPeersAt(ma, recent+int64(time.Millisecond), h.tsAt(0, 1))
waitForEvents(t, r, []int64{recent, recent + int64(time.Millisecond)}, 3*time.Second)
if got := lists.Load(); got > 4 {
t.Fatalf("%d listings while draining held entries; delivery should come from the ring", got)
}
}