fix(filer): bound concurrent persisted-log replays

Each server-side replay holds an open chunk reader per source filer plus a
readahead buffer, so a reconnect storm of clients that predate the
metadata-chunks offload multiplies into many GB. Gate replays with a
semaphore; abort the acquire when the subscriber's stream is gone so
cancelled clients do not pile up parked goroutines.
This commit is contained in:
Chris Lu
2026-06-09 09:52:41 -07:00
parent 8725c06420
commit 8840b6ca09
2 changed files with 23 additions and 3 deletions
+21 -1
View File
@@ -200,7 +200,27 @@ func isChunkNotFoundError(err error) bool {
httpNotFoundPattern.MatchString(errMsg)
}
func (f *Filer) ReadPersistedLogBuffer(startPosition log_buffer.MessagePosition, stopTsNs int64, eachLogEntryFn log_buffer.EachLogEntryFuncType) (lastTsNs int64, isDone bool, err error) {
// persistedLogReplayLimit caps concurrent legacy replays; each holds a chunk
// reader per source filer, so a reconnect storm of pre-offload clients would
// otherwise pin many GB. Metadata-chunks clients take sendLogFileRefs and never
// reach this path.
const persistedLogReplayLimit = 16
var persistedLogReplaySem = make(chan struct{}, persistedLogReplayLimit)
func (f *Filer) ReadPersistedLogBuffer(ctx context.Context, startPosition log_buffer.MessagePosition, stopTsNs int64, eachLogEntryFn log_buffer.EachLogEntryFuncType) (lastTsNs int64, isDone bool, err error) {
// Cap concurrent replays; bail if the stream is already gone so cancelled
// clients do not park on the semaphore.
if err := ctx.Err(); err != nil {
return 0, false, err
}
select {
case persistedLogReplaySem <- struct{}{}:
defer func() { <-persistedLogReplaySem }()
case <-ctx.Done():
return 0, false, ctx.Err()
}
visitor, visitErr := f.collectPersistedLogBuffer(startPosition, stopTsNs)
if visitErr != nil {
+2 -2
View File
@@ -215,7 +215,7 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
if req.ClientSupportsMetadataChunks {
processedTsNs, isDone, readPersistedLogErr = fs.sendLogFileRefs(ctx, stream, lastReadTime, req.UntilNs)
} else {
processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(lastReadTime, req.UntilNs, eachLogEntryFn)
processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, eachLogEntryFn)
}
if readPersistedLogErr != nil {
return fmt.Errorf("reading from persisted logs: %w", readPersistedLogErr)
@@ -367,7 +367,7 @@ func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq
if req.ClientSupportsMetadataChunks {
processedTsNs, isDone, readPersistedLogErr = fs.sendLogFileRefs(ctx, stream, lastReadTime, req.UntilNs)
} else {
processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(lastReadTime, req.UntilNs, eachLogEntryFn)
processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, eachLogEntryFn)
}
if readPersistedLogErr != nil {
glog.V(0).Infof("read on disk %v local subscribe %s from %+v: %v", clientName, req.PathPrefix, lastReadTime, readPersistedLogErr)