From 14fdd61aea5b013fe179657b2410b17eac95544a Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Thu, 8 Oct 2026 22:03:31 +0800 Subject: [PATCH] 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> --- weed/server/filer_grpc_server_sub_meta.go | 197 ++++++++++++++-------- weed/server/filer_subscribe_spin_test.go | 79 +++++++++ 2 files changed, 208 insertions(+), 68 deletions(-) create mode 100644 weed/server/filer_subscribe_spin_test.go diff --git a/weed/server/filer_grpc_server_sub_meta.go b/weed/server/filer_grpc_server_sub_meta.go index 858f64f4f..a40380633 100644 --- a/weed/server/filer_grpc_server_sub_meta.go +++ b/weed/server/filer_grpc_server_sub_meta.go @@ -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) { diff --git a/weed/server/filer_subscribe_spin_test.go b/weed/server/filer_subscribe_spin_test.go new file mode 100644 index 000000000..4cbcd91ce --- /dev/null +++ b/weed/server/filer_subscribe_spin_test.go @@ -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) + } +}