mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-30 11:45:42 +00:00
filer: pin-aware reader cache eviction and stream release (#11503)
* filer: synchronize stream pins and release them on transitions Guard chunkStream.cacher with the ReaderCache lock everywhere: mount sections share one ChunkReadAt across concurrent reads, and unsynchronized release could double-unpin. Reads served from the chunk cache now detach the stream's pin instead of retaining the previous chunk. Eviction prefers unpinned downloaders so a pinned buffer is not dropped mid-stream. A new ReleaseStream lets callers drop their pin without destroying the shared cache; S3 and WebDAV readers use it. lastChunkFid becomes atomic since concurrent mount reads can update it. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: keep eviction bounded when every downloader is pinned Both eviction paths still fall back to a pinned victim when no unpinned one exists, so abandoned stream pins cannot bypass the downloader limit or stall the memory budget. Budget eviction also rechecks the pin under the ReaderCache lock at removal time: a stream that pinned the selected victim in between keeps it mapped and the selection retries. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * filer: restore budget bookkeeping when a victim gets pinned mid-eviction removeUnpinned losing the pin race left the victim out of the idle list while still holding its reservation, making it unevictable even as the pinned fallback. Push it back when the reservation is still live. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
co-authored by
Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
parent
3abdef3202
commit
9b3b12c607
@@ -169,10 +169,17 @@ func (c *ChunkReadAt) Size() int64 {
|
||||
}
|
||||
|
||||
func (c *ChunkReadAt) Close() error {
|
||||
c.ReleaseStream()
|
||||
c.readerCache.destroy()
|
||||
return nil
|
||||
}
|
||||
|
||||
// ReleaseStream drops this reader's hold on the chunk it is positioned in.
|
||||
// Unlike Close it leaves the (possibly shared) ReaderCache intact.
|
||||
func (c *ChunkReadAt) ReleaseStream() {
|
||||
c.readerCache.releaseStream(&c.stream)
|
||||
}
|
||||
|
||||
func (c *ChunkReadAt) ReadAt(p []byte, offset int64) (n int, err error) {
|
||||
|
||||
c.readerPattern.MonitorReadAt(offset, len(p))
|
||||
|
||||
@@ -207,3 +207,77 @@ func TestChunkStreamConcurrentReadsOnOneReader(t *testing.T) {
|
||||
t.Errorf("%s retained after all reads finished: pins=%d readers=%d", fileId, atomic.LoadInt32(&cacher.pins), atomic.LoadInt32(&cacher.readers))
|
||||
}
|
||||
}
|
||||
|
||||
// A chunk a stream is positioned in must outlast downloader-limit eviction:
|
||||
// otherwise a busy cache drops the buffer mid-stream and forces a refetch.
|
||||
func TestChunkReadAtPinnedChunkSurvivesEviction(t *testing.T) {
|
||||
const chunkSize = 64 << 10
|
||||
|
||||
var fetches int32
|
||||
rc := NewReaderCache(2, newMockChunkCacheForReaderCache(), func(context.Context, string) ([]string, error) {
|
||||
return []string{"unused"}, nil
|
||||
}, nil)
|
||||
defer rc.destroy()
|
||||
rc.fetchChunkDataFn = func(_ context.Context, buffer []byte, _ []string, _ []byte, _ bool, _ bool, _ int64, fileId string, _ util_http.RefreshUrlsFunc) (int, error) {
|
||||
if fileId == "chunk0" {
|
||||
atomic.AddInt32(&fetches, 1)
|
||||
}
|
||||
return len(buffer), nil
|
||||
}
|
||||
|
||||
stream := &chunkStream{}
|
||||
buf := make([]byte, 16<<10)
|
||||
// One slice in: the stream is positioned in chunk0 but has not left it.
|
||||
if _, err := rc.readChunkAt(context.Background(), stream, buf, "chunk0", nil, false, 0, chunkSize, false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Fill the downloader map past its limit with unpinned chunks.
|
||||
for _, fileId := range []string{"chunk1", "chunk2", "chunk3"} {
|
||||
if _, err := rc.ReadChunkAt(context.Background(), buf, fileId, nil, false, 0, chunkSize, true); err != nil {
|
||||
t.Fatalf("read %s: %v", fileId, err)
|
||||
}
|
||||
}
|
||||
|
||||
// The stream's next slice must come from the still-pinned buffer.
|
||||
if _, err := rc.readChunkAt(context.Background(), stream, buf, "chunk0", nil, false, 16<<10, chunkSize, false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := atomic.LoadInt32(&fetches); got != 1 {
|
||||
t.Errorf("chunk0 fetched %d times, want 1", got)
|
||||
}
|
||||
}
|
||||
|
||||
// When every downloader is pinned the limit still applies: the oldest pinned
|
||||
// buffer is evicted so abandoned streams cannot grow memory past the limit.
|
||||
func TestChunkReadAtPinnedEvictionFallsBackWhenAllPinned(t *testing.T) {
|
||||
const chunkSize = 64 << 10
|
||||
rc := NewReaderCache(2, newMockChunkCacheForReaderCache(), func(context.Context, string) ([]string, error) {
|
||||
return []string{"unused"}, nil
|
||||
}, nil)
|
||||
defer rc.destroy()
|
||||
rc.fetchChunkDataFn = func(_ context.Context, buffer []byte, _ []string, _ []byte, _ bool, _ bool, _ int64, _ string, _ util_http.RefreshUrlsFunc) (int, error) {
|
||||
return len(buffer), nil
|
||||
}
|
||||
|
||||
buf := make([]byte, 16<<10)
|
||||
streamA := &chunkStream{}
|
||||
streamB := &chunkStream{}
|
||||
if _, err := rc.readChunkAt(context.Background(), streamA, buf, "chunk0", nil, false, 0, chunkSize, false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := rc.readChunkAt(context.Background(), streamB, buf, "chunk1", nil, false, 0, chunkSize, false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Every downloader is now pinned; the next chunk must still get in.
|
||||
if _, err := rc.readChunkAt(context.Background(), &chunkStream{}, buf, "chunk2", nil, false, 0, chunkSize, false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if isRetained(rc, "chunk0") {
|
||||
t.Fatal("oldest pinned downloader was not evicted past the limit")
|
||||
}
|
||||
if !isRetained(rc, "chunk2") {
|
||||
t.Fatal("new downloader missing after pinned fallback eviction")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -169,14 +169,26 @@ retry:
|
||||
}
|
||||
}
|
||||
|
||||
// clean up old downloaders
|
||||
// clean up old downloaders; prefer one no stream is positioned in, but
|
||||
// fall back to a pinned one so abandoned pins cannot bypass the limit
|
||||
if len(rc.downloaders) >= rc.limit {
|
||||
oldestFid, oldestTime := "", time.Now().UnixNano()
|
||||
pinnedFid, pinnedTime := "", int64(0)
|
||||
for fid, downloader := range rc.downloaders {
|
||||
completedTime := atomic.LoadInt64(&downloader.completedTimeNew)
|
||||
if completedTime > 0 && completedTime < oldestTime {
|
||||
oldestFid, oldestTime = fid, completedTime
|
||||
if completedTime <= 0 {
|
||||
continue
|
||||
}
|
||||
if atomic.LoadInt32(&downloader.pins) == 0 {
|
||||
if completedTime < oldestTime {
|
||||
oldestFid, oldestTime = fid, completedTime
|
||||
}
|
||||
} else if pinnedFid == "" || completedTime < pinnedTime {
|
||||
pinnedFid, pinnedTime = fid, completedTime
|
||||
}
|
||||
}
|
||||
if oldestFid == "" {
|
||||
oldestFid = pinnedFid
|
||||
}
|
||||
if oldestFid != "" {
|
||||
oldDownloader := rc.downloaders[oldestFid]
|
||||
@@ -290,6 +302,23 @@ func (rc *ReaderCache) remove(downloader *SingleChunkCacher) {
|
||||
}
|
||||
}
|
||||
|
||||
// removeUnpinned drops a cacher only while no stream is positioned in it.
|
||||
// Budget eviction picks its victim under the budget lock, so the pin check
|
||||
// and the map removal must happen together under the ReaderCache lock.
|
||||
func (rc *ReaderCache) removeUnpinned(downloader *SingleChunkCacher) (removed bool) {
|
||||
rc.Lock()
|
||||
removed = rc.downloaders[downloader.chunkFileId] == downloader &&
|
||||
atomic.LoadInt32(&downloader.pins) == 0
|
||||
if removed {
|
||||
delete(rc.downloaders, downloader.chunkFileId)
|
||||
}
|
||||
rc.Unlock()
|
||||
if removed {
|
||||
downloader.destroy()
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
// removeConsumed drops a cacher once its buffer was fully read, or the
|
||||
// streams positioned in it have left, and no readers remain attached or
|
||||
// pinned. The checks run under the ReaderCache lock so a reader attaching
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"container/list"
|
||||
"fmt"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/util/mem"
|
||||
)
|
||||
@@ -51,12 +52,39 @@ func (b *ReaderCacheBudget) reserve(s *SingleChunkCacher) error {
|
||||
b.Unlock()
|
||||
return nil
|
||||
}
|
||||
if entry := b.idle.Front(); entry != nil {
|
||||
victim := entry.Value.(*SingleChunkCacher)
|
||||
// Prefer evicting an idle chunk no stream is positioned in; fall back
|
||||
// to the oldest pinned one so abandoned pins cannot block the budget.
|
||||
var victim *SingleChunkCacher
|
||||
var entry *list.Element
|
||||
pinnedVictim := false
|
||||
for e := b.idle.Front(); e != nil; e = e.Next() {
|
||||
c := e.Value.(*SingleChunkCacher)
|
||||
if atomic.LoadInt32(&c.pins) == 0 {
|
||||
victim, entry = c, e
|
||||
pinnedVictim = false
|
||||
break
|
||||
}
|
||||
if victim == nil {
|
||||
victim, entry = c, e
|
||||
pinnedVictim = true
|
||||
}
|
||||
}
|
||||
if entry != nil {
|
||||
b.idle.Remove(entry)
|
||||
delete(b.idleEntries, victim)
|
||||
b.Unlock()
|
||||
victim.parent.remove(victim)
|
||||
if pinnedVictim {
|
||||
victim.parent.remove(victim)
|
||||
} else if !victim.parent.removeUnpinned(victim) {
|
||||
// The victim was pinned between selection and removal: keep it
|
||||
// evictable so a pin abandoned in that gap cannot wedge the
|
||||
// budget, then retry the selection.
|
||||
b.Lock()
|
||||
if _, ok := b.reservations[victim]; ok && b.idleEntries[victim] == nil {
|
||||
b.idleEntries[victim] = b.idle.PushBack(victim)
|
||||
}
|
||||
b.Unlock()
|
||||
}
|
||||
continue
|
||||
}
|
||||
changed := b.changed
|
||||
|
||||
@@ -1103,6 +1103,7 @@ func (s3a *S3ApiServer) streamFromVolumeServers(w http.ResponseWriter, r *http.R
|
||||
tStreamPrep := time.Now()
|
||||
chunkViews := filer.ViewFromVisibleIntervals(visibleIntervals, offset, size)
|
||||
reader := filer.NewChunkReaderAtFromClient(ctx, s3a.readerCache, chunkViews, totalSize, filer.DefaultPrefetchCount)
|
||||
defer reader.ReleaseStream()
|
||||
streamPrepTime = time.Since(tStreamPrep)
|
||||
|
||||
// A cached chunk whose volume server is down, or whose needle was evicted and
|
||||
|
||||
@@ -130,7 +130,7 @@ type WebDavFile struct {
|
||||
off int64
|
||||
entry *filer_pb.Entry
|
||||
visibleIntervals *filer.IntervalList[*filer.VisibleInterval]
|
||||
reader io.ReaderAt
|
||||
reader *filer.ChunkReadAt
|
||||
bufWriter *buffered_writer.BufferedWriteCloser
|
||||
ctx context.Context
|
||||
}
|
||||
@@ -537,6 +537,9 @@ func (f *WebDavFile) Write(buf []byte) (int, error) {
|
||||
func (f *WebDavFile) Close() error {
|
||||
|
||||
glog.V(2).Infof("WebDavFileSystem.Close %v", f.name)
|
||||
if f.reader != nil {
|
||||
f.reader.ReleaseStream()
|
||||
}
|
||||
if f.bufWriter == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user