filer: stream persisted log files when serving metadata subscriptions (#9821)

* filer: stream persisted log files when serving metadata subscriptions

readFileEntries buffered every LogEntry of a whole log file into memory
before returning them one by one, making a subscription read O(entries in
one log file) instead of O(one entry). On a filer with large per-entry
metadata, many concurrent SubscribeMetadata streams each loading a full
log file exhausted memory.

Keep a current LogFileIterator and return one entry at a time, advancing
files as each is exhausted. The deleted-volume skip is preserved.

* filer: close the log file iterator on read errors too

A genuine read error returned early without closing the current
LogFileIterator, leaving its ChunkStreamReader alive until GC. Close on
every exit path and propagate only a real error.

* filer: close persisted-log iterators when a subscription stops early

The streaming iterator keeps a log file reader open across calls, so a
subscription that returns before EOF (early stop, cancellation) left the
reader alive until GC. Add idempotent Close on LogFileQueueIterator and
OrderedLogVisitor, and have ReadPersistedLogBuffer wait for the readahead
goroutine and close the visitor on the way out.
This commit is contained in:
Chris Lu
2026-06-04 13:27:25 -07:00
committed by GitHub
parent 6e8002f065
commit 8c2d9f466f
2 changed files with 62 additions and 50 deletions
+9 -2
View File
@@ -220,8 +220,10 @@ func (f *Filer) ReadPersistedLogBuffer(startPosition log_buffer.MessagePosition,
}
ch := make(chan entryOrErr, readaheadSize)
stopReadahead := make(chan struct{})
readaheadDone := make(chan struct{})
go func() {
defer close(ch)
defer close(readaheadDone)
for {
entry, readErr := visitor.GetNext()
if readErr != nil {
@@ -240,7 +242,13 @@ func (f *Filer) ReadPersistedLogBuffer(startPosition log_buffer.MessagePosition,
}
}
}()
defer close(stopReadahead)
// Stop the readahead goroutine, wait for it to exit, then release any log
// file readers it left open (e.g. on early return or cancellation).
defer func() {
close(stopReadahead)
<-readaheadDone
visitor.Close()
}()
for item := range ch {
if item.err != nil {
@@ -261,4 +269,3 @@ func (f *Filer) ReadPersistedLogBuffer(startPosition log_buffer.MessagePosition,
return
}
+53 -48
View File
@@ -197,6 +197,15 @@ func (o *OrderedLogVisitor) GetNext() (logEntry *filer_pb.LogEntry, err error) {
return item.Entry, nil
}
// Close releases any log file readers still open across the per-filer
// iterators, e.g. when a subscription stops before reaching the end. Safe to
// call more than once.
func (o *OrderedLogVisitor) Close() {
for _, it := range o.perFilerIteratorMap {
it.Close()
}
}
func getFilerId(name string) string {
idx := strings.LastIndex(name, ".")
if idx < 0 {
@@ -329,12 +338,11 @@ func (c *LogFileEntryCollector) collectMore(v *OrderedLogVisitor) (err error) {
// ----------
type LogFileQueueIterator struct {
q *util.Queue[*LogFileEntry]
masterClient *wdclient.MasterClient
startTsNs int64
stopTsNs int64
pendingEntries []*filer_pb.LogEntry
pendingIndex int
q *util.Queue[*LogFileEntry]
masterClient *wdclient.MasterClient
startTsNs int64
stopTsNs int64
currentFileIterator *LogFileIterator
}
func newLogFileQueueIterator(masterClient *wdclient.MasterClient, q *util.Queue[*LogFileEntry], startTsNs, stopTsNs int64) *LogFileQueueIterator {
@@ -346,20 +354,48 @@ func newLogFileQueueIterator(masterClient *wdclient.MasterClient, q *util.Queue[
}
}
// getNext will return io.EOF when done
// Close releases the current log file reader, if any. Safe to call more than once.
func (iter *LogFileQueueIterator) Close() {
if iter.currentFileIterator != nil {
if err := iter.currentFileIterator.Close(); err != nil {
glog.Warningf("close log file %s: %v", iter.currentFileIterator.filePath, err)
}
iter.currentFileIterator = nil
}
}
// getNext streams one log entry at a time from the current file, advancing to
// the next file as each is exhausted. It returns io.EOF when done. Entries are
// not buffered per file, so memory stays O(1) regardless of log file size.
func (iter *LogFileQueueIterator) getNext(v *OrderedLogVisitor) (logEntry *filer_pb.LogEntry, err error) {
for {
// return pending entries first
if iter.pendingIndex < len(iter.pendingEntries) {
logEntry = iter.pendingEntries[iter.pendingIndex]
iter.pendingIndex++
return logEntry, nil
if iter.currentFileIterator != nil {
logEntry, err = iter.currentFileIterator.getNext()
if err == nil {
return logEntry, nil
}
// The current file is done (io.EOF), its volume was deleted, or it is
// unreadable. Close it on every path so the reader is not left alive
// until GC; only a genuine read error is propagated.
readErr := err
switch {
case readErr == io.EOF:
readErr = nil
case isChunkNotFoundError(readErr):
// Volume or chunk was deleted, skip the rest of this log file
glog.Warningf("skipping rest of %s: %v", iter.currentFileIterator.filePath, readErr)
readErr = nil
}
if closeErr := iter.currentFileIterator.Close(); closeErr != nil {
glog.Warningf("close log file %s: %v", iter.currentFileIterator.filePath, closeErr)
}
iter.currentFileIterator = nil
if readErr != nil {
return nil, readErr
}
}
// reset for next file
iter.pendingEntries = nil
iter.pendingIndex = 0
// read entries from next file
// advance to the next file
if iter.q.Len() == 0 {
return nil, io.EOF
}
@@ -382,38 +418,7 @@ func (iter *LogFileQueueIterator) getNext(v *OrderedLogVisitor) (logEntry *filer
if next != nil && next.TsNs <= iter.startTsNs {
continue
}
// read all entries from this file
iter.pendingEntries, err = iter.readFileEntries(t.FileEntry)
if err != nil {
return nil, err
}
}
}
// readFileEntries reads all log entries from a single file
func (iter *LogFileQueueIterator) readFileEntries(fileEntry *Entry) (entries []*filer_pb.LogEntry, err error) {
fileIterator := newLogFileIterator(iter.masterClient, fileEntry, iter.startTsNs, iter.stopTsNs)
defer func() {
if closeErr := fileIterator.Close(); closeErr != nil && err == nil {
err = closeErr
}
}()
for {
logEntry, err := fileIterator.getNext()
if err == io.EOF {
return entries, nil
}
if err != nil {
if isChunkNotFoundError(err) {
// Volume or chunk was deleted, skip the rest of this log file
glog.Warningf("skipping rest of %s: %v", fileIterator.filePath, err)
return entries, nil
}
return nil, err
}
entries = append(entries, logEntry)
iter.currentFileIterator = newLogFileIterator(iter.masterClient, t.FileEntry, iter.startTsNs, iter.stopTsNs)
}
}