Files
seaweedfs/weed/filer/reader_cache.go
T
Chris LuandDevin 6c07a5fdd0 s3: keep small ranged GETs on range reads, no whole-chunk downloads (#11577)
* filer: keep a ranged read in random mode through its contiguous tail

A far ReadAt on a fresh ReaderPattern left the sequential counter at -1,
so the next buffer of the same ranged request landed on the frontier and
flipped the verdict straight back to sequential — readChunkSliceAt then
paid a whole-chunk fetch for the remainder of the range. Drop the
counter to -ModeChangeLimit when random mode is entered so the verdict
needs sustained sequential evidence to undo, matching the hysteresis an
established sequential stream already gets.

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

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

* s3: pin small ranged GETs to range reads

A ranged GET whose first read lands within SeqTolerance of offset 0 is
judged sequential immediately, and even a far-starting range could flip
back mid-request; either way readChunkSliceAt downloads each covered
chunk in full, multiplying disk reads for small ranged reads (measured
~7x). Pin random mode for ranged requests no larger than SeqTolerance so
all of the request's buffer reads stay range fetches. Larger ranges keep
the dynamic pattern, where whole-chunk fetches amortize.

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

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

* filer: fetch only the part of a chunk the view covers

Replaces the PinRandomMode size heuristic with a per-chunk coverage rule.
ViewFromVisibleIntervals already clips chunk views to the request window,
so a view that is not IsFullChunk() is one the request only partially
needs; fetch it as a range regardless of the detected read pattern.

This closes the holes a request-size pin left open: ranges larger than
SeqTolerance no longer revert to whole-chunk downloads once their buffers
look sequential, and ranges that fully cover a chunk keep the shared
whole-chunk path instead of fetching 256KiB slices piecemeal. Prefetch
(MaybeCache) skips clipped views so it cannot amplify a range read either.

PinRandomMode is dropped: no caller needs it once coverage drives the
fetch choice. Range fetches route through fetchChunkDataFn so tests
observe them the same way as whole-chunk downloads.

* filer: keep ciphered chunks on the whole-chunk path

A range fetch cannot save bytes for a ciphered chunk: readEncryptedUrl
always downloads and decrypts the whole blob before slicing. Sending
partial views of ciphered chunks through fetchChunkRange would repeat the
full download per buffer, so they keep the shared whole-chunk path where
one download serves every buffer. Prefetch stays enabled for them for
the same reason.

* filer: keep compressed chunks on the whole-chunk path

Like ciphered chunks, a range request on a compressed chunk makes the
volume server read and decompress the whole needle, so range-per-buffer
would repeat the full backend read for each 256KiB window. Route them
through the shared whole-chunk path via ChunkView.CanRangeFetch.

* filer: fall back to range fetch when a chunk exceeds the reader budget

A ciphered or compressed chunk larger than readerCacheSizeMB can never
be read through the whole-chunk path — the budget rejects the buffer —
so its partial views must still range-fetch or the GET fails outright.

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-10-03 15:01:04 +08:00

540 lines
17 KiB
Go

package filer
import (
"context"
"fmt"
"sync"
"sync/atomic"
"time"
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/util/chunk_cache"
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
"github.com/seaweedfs/seaweedfs/weed/util/mem"
"github.com/seaweedfs/seaweedfs/weed/wdclient"
)
type CacheInvalidator interface {
InvalidateCache(fileId string)
}
type fetchChunkDataFnType func(ctx context.Context, buffer []byte, urlStrings []string, cipherKey []byte, isGzipped bool, isFullChunk bool, offset int64, fileId string, refreshUrls util_http.RefreshUrlsFunc) (n int, err error)
type ReaderCache struct {
chunkCache chunk_cache.ChunkCache
lookupFileIdFn wdclient.LookupFileIdFunctionType
cacheInvalidator CacheInvalidator
fetchChunkDataFn fetchChunkDataFnType
sync.Mutex
downloaders map[string]*SingleChunkCacher
limit int
budget *ReaderCacheBudget
}
type SingleChunkCacher struct {
completedTimeNew int64
readers int32
consumed int32
pins int32 // streams currently positioned inside this chunk
left int32 // set once the last pinning stream moved on
sync.Mutex
parent *ReaderCache
chunkFileId string
data []byte
err error
cipherKey []byte
isGzipped bool
chunkSize int
shouldCache bool
wg sync.WaitGroup
cacheStartedCh chan struct{}
done chan struct{} // signals when download is complete
}
func NewReaderCache(limit int, chunkCache chunk_cache.ChunkCache, lookupFileIdFn wdclient.LookupFileIdFunctionType, cacheInvalidator CacheInvalidator, budgets ...*ReaderCacheBudget) *ReaderCache {
var budget *ReaderCacheBudget
if len(budgets) > 0 {
budget = budgets[0]
}
return &ReaderCache{
limit: limit,
budget: budget,
chunkCache: chunkCache,
lookupFileIdFn: lookupFileIdFn,
cacheInvalidator: cacheInvalidator,
fetchChunkDataFn: util_http.RetriedFetchChunkData,
downloaders: make(map[string]*SingleChunkCacher),
}
}
// MaybeCache prefetches up to 'count' chunks ahead in parallel.
// This improves read throughput for sequential reads by keeping the
// network pipeline full with parallel chunk fetches.
func (rc *ReaderCache) MaybeCache(chunkViews *Interval[*ChunkView], count int) {
if rc.lookupFileIdFn == nil {
return
}
if count <= 0 {
count = 1
}
rc.Lock()
defer rc.Unlock()
if len(rc.downloaders) >= rc.limit {
return
}
cached := 0
for x := chunkViews; x != nil && cached < count; x = x.Next {
chunkView := x.Value
if _, found := rc.downloaders[chunkView.FileId]; found {
continue
}
if rc.chunkCache.IsInCache(chunkView.FileId, true) {
glog.V(4).Infof("%s is in cache", chunkView.FileId)
continue
}
if len(rc.downloaders) >= rc.limit {
// abort when slots are filled
return
}
if (chunkView.CanRangeFetch() || !rc.budget.canFit(int(chunkView.ChunkSize))) && !chunkView.IsFullChunk() {
// the view is clipped to part of the chunk and will be
// range-fetched, so prefetching it whole would download bytes
// nobody needs; a ciphered or compressed partial view needs
// the whole blob anyway and is worth prefetching, but not when
// it cannot fit the budget at all
continue
}
// glog.V(4).Infof("prefetch %s offset %d", chunkView.FileId, chunkView.ViewOffset)
// cache this chunk if not yet
shouldCache := (uint64(chunkView.ViewOffset) + chunkView.ChunkSize) <= rc.chunkCache.GetMaxFilePartSizeInCache()
cacher := newSingleChunkCacher(rc, chunkView.FileId, chunkView.CipherKey, chunkView.IsGzipped, int(chunkView.ChunkSize), shouldCache)
go cacher.startCaching()
<-cacher.cacheStartedCh
rc.downloaders[chunkView.FileId] = cacher
cached++
}
return
}
// fetchChunkRange downloads only [offset, offset+len(buffer)) of a chunk,
// for views clipped to part of their chunk and for random-mode reads. It
// goes through fetchChunkDataFn so tests observe range fetches the same way
// they observe whole-chunk downloads.
func (rc *ReaderCache) fetchChunkRange(ctx context.Context, buffer []byte, chunkView *ChunkView, offset int64) (int, error) {
urlStrings, err := rc.lookupFileIdFn(ctx, chunkView.FileId)
if err != nil {
glog.ErrorfCtx(ctx, "operation LookupFileId %s failed, err: %v", chunkView.FileId, err)
return 0, err
}
return rc.fetchChunkDataFn(ctx, buffer, urlStrings, chunkView.CipherKey, chunkView.IsGzipped, false, offset, chunkView.FileId,
refreshUrls(ctx, rc.cacheInvalidator, rc.lookupFileIdFn, chunkView.FileId))
}
// chunkStream is one sequential reader's position in a shared ReaderCache.
// The chunk it is reading stays pinned until the stream reads it to the end or
// moves elsewhere, so another stream finishing or leaving the same chunk does
// not drop the buffer from under it.
type chunkStream struct {
cacher *SingleChunkCacher
}
func (rc *ReaderCache) ReadChunkAt(ctx context.Context, buffer []byte, fileId string, cipherKey []byte, isGzipped bool, offset int64, chunkSize int, shouldCache bool) (int, error) {
return rc.readChunkAt(ctx, nil, buffer, fileId, cipherKey, isGzipped, offset, chunkSize, shouldCache)
}
func (rc *ReaderCache) readChunkAt(ctx context.Context, stream *chunkStream, buffer []byte, fileId string, cipherKey []byte, isGzipped bool, offset int64, chunkSize int, shouldCache bool) (int, error) {
retry:
rc.Lock()
for {
if cacher, found := rc.downloaders[fileId]; found {
if cacher.hasCompletedError() {
delete(rc.downloaders, fileId)
rc.Unlock()
cacher.destroy()
rc.Lock()
continue
}
// Count this read on the cacher before releasing the map lock, so a
// concurrent destroy() (error eviction here, LRU, or UnCache) cannot
// start wg.Wait() on a zero counter while this read is about to register.
cacher.wg.Add(1)
atomic.AddInt32(&cacher.readers, 1)
previous := stream.pin(cacher)
rc.Unlock()
rc.unpin(previous)
n, err := cacher.readChunkAt(ctx, buffer, offset)
rc.releaseIfFinished(stream, cacher, offset, n, err, chunkSize)
if n > 0 || err != nil {
return n, err
}
// If n=0 and err=nil, the cacher couldn't provide data for this offset.
// Fall through to try chunkCache.
rc.Lock()
}
break
}
if shouldCache || rc.lookupFileIdFn == nil {
n, err := rc.chunkCache.ReadChunkAt(buffer, fileId, uint64(offset))
if n > 0 {
// Served from the chunk cache: the stream has left its pinned chunk.
previous := stream.unpinLocked()
rc.Unlock()
rc.unpin(previous)
return n, err
}
}
// 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 {
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]
delete(rc.downloaders, oldestFid)
rc.Unlock()
oldDownloader.destroy()
goto retry
}
}
// glog.V(4).Infof("cache1 %s", fileId)
cacher := newSingleChunkCacher(rc, fileId, cipherKey, isGzipped, chunkSize, shouldCache)
go cacher.startCaching()
<-cacher.cacheStartedCh
rc.downloaders[fileId] = cacher
cacher.wg.Add(1)
atomic.AddInt32(&cacher.readers, 1)
previous := stream.pin(cacher)
rc.Unlock()
rc.unpin(previous)
n, err := cacher.readChunkAt(ctx, buffer, offset)
rc.releaseIfFinished(stream, cacher, offset, n, err, chunkSize)
return n, err
}
// pin makes cacher the stream's current chunk and returns the chunk it was
// pinned to before, which the caller unpins once the ReaderCache lock is
// released. The stream is only touched under the ReaderCache lock, since
// concurrent ReadAt calls on one ChunkReadAt share it.
func (stream *chunkStream) pin(cacher *SingleChunkCacher) (previous *SingleChunkCacher) {
if stream == nil || stream.cacher == cacher {
return nil
}
previous = stream.cacher
stream.cacher = cacher
atomic.AddInt32(&cacher.pins, 1)
return previous
}
// unpinLocked detaches the stream from its chunk and returns that chunk for
// the caller to unpin once the ReaderCache lock is released.
func (stream *chunkStream) unpinLocked() (previous *SingleChunkCacher) {
if stream == nil {
return nil
}
previous = stream.cacher
stream.cacher = nil
return previous
}
// releaseIfFinished unpins the stream's chunk once the stream has read it to
// the end, since the stream will not come back to it.
func (rc *ReaderCache) releaseIfFinished(stream *chunkStream, cacher *SingleChunkCacher, offset int64, n int, err error, chunkSize int) {
if stream == nil || err != nil || offset+int64(n) < int64(chunkSize) {
return
}
var previous *SingleChunkCacher
rc.Lock()
if stream.cacher == cacher {
previous = stream.unpinLocked()
}
rc.Unlock()
rc.unpin(previous)
}
// releaseStream unpins whatever chunk the stream is positioned in.
func (rc *ReaderCache) releaseStream(stream *chunkStream) {
if stream == nil {
return
}
rc.Lock()
previous := stream.unpinLocked()
rc.Unlock()
rc.unpin(previous)
}
// unpin drops one stream's pin. Once no stream is positioned in the chunk it
// is dropped like UnCache would, as soon as no read is in progress either:
// here if none is, otherwise by the last read's removeConsumed.
func (rc *ReaderCache) unpin(downloader *SingleChunkCacher) {
if downloader == nil {
return
}
if atomic.AddInt32(&downloader.pins, -1) == 0 {
atomic.StoreInt32(&downloader.left, 1)
}
rc.removeConsumed(downloader)
}
func (rc *ReaderCache) UnCache(fileId string) {
rc.Lock()
downloader := rc.downloaders[fileId]
delete(rc.downloaders, fileId)
rc.Unlock()
if downloader != nil {
downloader.destroy()
}
}
func (rc *ReaderCache) remove(downloader *SingleChunkCacher) {
rc.Lock()
removed := rc.downloaders[downloader.chunkFileId] == downloader
if removed {
delete(rc.downloaders, downloader.chunkFileId)
}
rc.Unlock()
if removed {
downloader.destroy()
}
}
// 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
// at the same time either wins (the cacher stays and that reader's detach
// retries the removal) or misses the map and refetches.
func (rc *ReaderCache) removeConsumed(downloader *SingleChunkCacher) {
rc.Lock()
removed := rc.downloaders[downloader.chunkFileId] == downloader &&
atomic.LoadInt32(&downloader.readers) == 0 &&
atomic.LoadInt32(&downloader.pins) == 0 &&
(atomic.LoadInt32(&downloader.consumed) != 0 || atomic.LoadInt32(&downloader.left) != 0)
if removed {
delete(rc.downloaders, downloader.chunkFileId)
}
rc.Unlock()
if removed {
downloader.destroy()
}
}
func (rc *ReaderCache) destroy() {
rc.Lock()
downloaders := rc.downloaders
rc.downloaders = make(map[string]*SingleChunkCacher)
rc.Unlock()
for _, downloader := range downloaders {
downloader.destroy()
}
}
func newSingleChunkCacher(parent *ReaderCache, fileId string, cipherKey []byte, isGzipped bool, chunkSize int, shouldCache bool) *SingleChunkCacher {
return &SingleChunkCacher{
parent: parent,
chunkFileId: fileId,
cipherKey: cipherKey,
isGzipped: isGzipped,
chunkSize: chunkSize,
shouldCache: shouldCache,
cacheStartedCh: make(chan struct{}),
done: make(chan struct{}),
}
}
// startCaching downloads a chunk shared by concurrent readers.
func (s *SingleChunkCacher) startCaching() {
s.wg.Add(1)
defer func() {
close(s.done)
s.wg.Done()
if s.hasCompletedError() {
s.parent.remove(s)
} else {
s.parent.budget.complete(s)
}
}()
s.cacheStartedCh <- struct{}{}
if err := s.parent.budget.reserve(s); err != nil {
s.setError(err)
return
}
// Intentionally use context.Background(), not a request-specific context.
// The downloaded chunk is a shared resource: multiple concurrent readers may
// wait on this same download via s.done. A request-scoped context that got
// cancelled would abort the download and error every other waiting reader.
// The download always runs to completion once started; readers that cancel
// individually drop out via readChunkAt's select on ctx.Done().
urlStrings, err := s.parent.lookupFileIdFn(context.Background(), s.chunkFileId)
if err != nil {
s.setError(fmt.Errorf("operation LookupFileId %s failed, err: %v", s.chunkFileId, err))
return
}
if len(urlStrings) == 0 {
s.setError(fmt.Errorf("operation LookupFileId %s failed, err: urls not found", s.chunkFileId))
return
}
data, fetchErr := s.fetchChunkData(context.Background(), urlStrings)
if fetchErr != nil {
data, fetchErr = s.retryFetchAfterCacheInvalidation(context.Background(), urlStrings, fetchErr)
}
// Now acquire lock to update state
s.Lock()
atomic.StoreInt64(&s.completedTimeNew, time.Now().UnixNano())
if fetchErr != nil {
s.err = fetchErr
} else {
s.data = data
if s.shouldCache {
s.parent.chunkCache.SetChunk(s.chunkFileId, s.data)
}
}
s.Unlock()
}
func (s *SingleChunkCacher) setError(err error) {
s.Lock()
defer s.Unlock()
s.err = err
atomic.StoreInt64(&s.completedTimeNew, time.Now().UnixNano())
}
func (s *SingleChunkCacher) hasCompletedError() bool {
if atomic.LoadInt64(&s.completedTimeNew) == 0 {
return false
}
s.Lock()
defer s.Unlock()
return s.err != nil
}
func (s *SingleChunkCacher) fetchChunkData(ctx context.Context, urlStrings []string) ([]byte, error) {
// Allocate buffer and download without holding the lock.
// This allows multiple downloads to proceed in parallel.
data := mem.Allocate(s.chunkSize)
_, fetchErr := s.parent.fetchChunkDataFn(ctx, data, urlStrings, s.cipherKey, s.isGzipped, true, 0, s.chunkFileId, refreshUrls(ctx, s.parent.cacheInvalidator, s.parent.lookupFileIdFn, s.chunkFileId))
if fetchErr != nil {
mem.Free(data)
return nil, fetchErr
}
return data, nil
}
func (s *SingleChunkCacher) retryFetchAfterCacheInvalidation(ctx context.Context, oldUrlStrings []string, originalErr error) ([]byte, error) {
var data []byte
err := retryFetchWithFreshLocations(ctx, s.parent.cacheInvalidator, s.parent.lookupFileIdFn, s.chunkFileId, oldUrlStrings, originalErr, func(newUrls []string) error {
var fetchErr error
data, fetchErr = s.fetchChunkData(ctx, newUrls)
return fetchErr
})
if err != nil {
return nil, err
}
return data, nil
}
func (s *SingleChunkCacher) destroy() {
// wait for all reads to finish before destroying the data
s.wg.Wait()
s.Lock()
if s.data != nil {
mem.Free(s.data)
s.data = nil
}
s.Unlock()
s.parent.budget.release(s)
}
// readChunkAt reads data from the cached chunk.
// It waits for the download to complete if it's still in progress.
// The ctx parameter allows the reader to cancel its wait (but the download continues
// for other readers - see comment in startCaching about shared resource semantics).
// The caller must s.wg.Add(1) under the ReaderCache lock before calling; this only releases it.
func (s *SingleChunkCacher) readChunkAt(ctx context.Context, buf []byte, offset int64) (n int, err error) {
defer func() {
s.wg.Done()
atomic.AddInt32(&s.readers, -1)
s.parent.removeConsumed(s)
}()
// Wait for download to complete, but allow reader cancellation.
// Prioritize checking done first - if data is already available,
// return it even if context is also cancelled.
select {
case <-s.done:
// Download already completed, proceed immediately
default:
// Download not complete, wait for it or context cancellation
select {
case <-s.done:
// Download completed
case <-ctx.Done():
// Reader cancelled while waiting - download continues for other readers
return 0, ctx.Err()
}
}
s.Lock()
defer s.Unlock()
if s.err != nil {
return 0, s.err
}
if len(s.data) <= int(offset) {
return 0, nil
}
n = copy(buf, s.data[offset:])
if offset+int64(n) == int64(len(s.data)) {
atomic.StoreInt32(&s.consumed, 1)
}
return n, nil
}