filer: cut ReaderCache mutex contention on small reads (#11677)

* filer: skip cache lock when a cacher cannot be removed yet

SingleChunkCacher.readChunkAt defers removeConsumed on every read, and
every unpin retries removal too. Each call acquired the ReaderCache
mutex even though the removal conditions are plain atomics, so a busy
cache paid a lock acquisition per read just to discover readers>0.

Check the atomics before locking: when a cacher is not yet consumable
the removal is impossible and the lock round trip is pure contention.
Removals still run under the lock via removeConsumedLocked, so the
attach-vs-remove race keeps its existing serialization.

Ref #11676

* filer: guard stream position with a per-stream mutex, not the cache lock

chunkStream.cacher is only ever shared by concurrent ReadAt calls on one
ChunkReadAt, yet pin, unpin, releaseIfFinished, and releaseStream all
mutated it under the ReaderCache mutex that serializes every reader in
the process. Each read that switched chunks took that global lock to pin
the new chunk, released it, then reacquired it in unpin() just to drop
the old chunk's pin and retry removal; releaseIfFinished and
releaseStream did the same lock-detach-unlock-relock dance. At high
small-GET concurrency that turned per-request stream bookkeeping into a
global convoy (#11676).

Give chunkStream its own mutex and drop the cache lock from the pin
lifecycle entirely: pin/detach serialize on stream.mu, the pin counter
and consumable checks are already atomics, and removeConsumed only
takes the cache lock when a cacher is actually removable. The map lock
now guards only map membership and read registration.

* filer: look up cached chunks under the read lock

readChunkAt serialized every small read on the write lock even though
the common paths are read-only: an existing downloader just needs its
read registered, and a chunkCache hit needs no map access at all. Each
GET also paid the lock a second time to reach the chunkCache check.

Use the read lock for the downloader lookup and registration; the
registered read keeps the buffer alive against a concurrent destroy,
and only error eviction, insertion, and removal need the write lock.
The chunkCache probe runs lock-free, and the insert path re-checks the
map under the write lock to cover a downloader registered in between.

* filer: start the chunk download outside the map lock

The insert path held the ReaderCache write lock across goroutine spawn
and the cacheStartedCh handshake, so every downloader miss serialized
against the startup of a fetch goroutine. Register the cacher in the
map under the lock, then start the download after releasing it; a fetch
that fails early still lands in the map and is evicted by the next
reader's completed-error check.

* filer: test that stream pin lifecycle stays off the cache lock

Regression coverage for the contention fix: pin and releaseStream on a
non-removable cacher must complete while the ReaderCache lock is held
by another goroutine.

* filer: pin the stream under the map lock

Between read registration and stream.pin the cacher showed zero pins,
so a budget eviction in that gap could pick a chunk the stream was just
attaching to and the stream's next slice refetched it. The pin counter
is an atomic and stream.mu is never held while acquiring the cache
lock, so pinning inside the map hold is deadlock-free and closes the
window.
This commit is contained in:
Chris Lu authored and GitHub committed 2026-10-10 06:52:37 +08:00
1 parent 838c554e33
commit e337443176
2 files changed
+142 -72

No files matched your search

+111 -72
View File
@@ -25,7 +25,7 @@ type ReaderCache struct {
lookupFileIdFn wdclient.LookupFileIdFunctionType
cacheInvalidator CacheInvalidator
fetchChunkDataFn fetchChunkDataFnType
sync.Mutex
sync.RWMutex
downloaders map[string]*SingleChunkCacher
limit int
budget *ReaderCacheBudget
@@ -78,6 +78,13 @@ func (rc *ReaderCache) MaybeCache(chunkViews *Interval[*ChunkView], count int) {
count = 1
}
var started []*SingleChunkCacher
defer func() {
for _, cacher := range started {
<-cacher.cacheStartedCh
}
}()
rc.Lock()
defer rc.Unlock()
@@ -113,9 +120,10 @@ func (rc *ReaderCache) MaybeCache(chunkViews *Interval[*ChunkView], count int) {
// 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
cacher.wg.Add(1) // the fetch, so destroy() waits even before the goroutine runs
rc.downloaders[chunkView.FileId] = cacher
go cacher.startCaching()
started = append(started, cacher)
cached++
}
@@ -141,6 +149,7 @@ func (rc *ReaderCache) fetchChunkRange(ctx context.Context, buffer []byte, chunk
// moves elsewhere, so another stream finishing or leaving the same chunk does
// not drop the buffer from under it.
type chunkStream struct {
mu sync.Mutex // guards cacher; concurrent ReadAt calls on one ChunkReadAt share the stream
cacher *SingleChunkCacher
}
@@ -150,47 +159,56 @@ func (rc *ReaderCache) ReadChunkAt(ctx context.Context, buffer []byte, fileId st
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()
rc.RLock()
if cacher, found := rc.downloaders[fileId]; found {
if cacher.hasCompletedError() {
rc.RUnlock()
rc.remove(cacher)
goto retry
}
break
// 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.RUnlock()
n, err := rc.readFromCacher(ctx, stream, cacher, previous, buffer, offset, 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.
} else {
rc.RUnlock()
}
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)
rc.unpin(stream.detach(nil))
return n, err
}
}
rc.Lock()
// A downloader may have registered between the read lock and this one.
if cacher, found := rc.downloaders[fileId]; found {
if cacher.hasCompletedError() {
delete(rc.downloaders, fileId)
rc.Unlock()
cacher.destroy()
goto retry
}
cacher.wg.Add(1)
atomic.AddInt32(&cacher.readers, 1)
previous := stream.pin(cacher)
rc.Unlock()
n, err := rc.readFromCacher(ctx, stream, cacher, previous, buffer, offset, chunkSize)
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 {
@@ -224,40 +242,57 @@ 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)
cacher.wg.Add(2) // the fetch plus this read, so destroy() waits even before the goroutine runs
atomic.AddInt32(&cacher.readers, 1)
rc.downloaders[fileId] = cacher
previous := stream.pin(cacher)
rc.Unlock()
rc.unpin(previous)
n, err := cacher.readChunkAt(ctx, buffer, offset)
rc.releaseIfFinished(stream, cacher, offset, n, err, chunkSize)
go cacher.startCaching()
<-cacher.cacheStartedCh
n, err := rc.readFromCacher(ctx, stream, cacher, previous, buffer, offset, chunkSize)
return n, err
}
// readFromCacher serves the read and releases the pin if the stream finished
// the chunk. The caller must have registered the read (wg + readers) and
// pinned the stream under the map lock; previous is the chunk the stream left
// when it pinned.
func (rc *ReaderCache) readFromCacher(ctx context.Context, stream *chunkStream, cacher *SingleChunkCacher, previous *SingleChunkCacher, buffer []byte, offset int64, chunkSize int) (n int, err error) {
rc.unpin(previous)
n, err = cacher.readChunkAt(ctx, buffer, offset)
rc.releaseIfFinished(stream, cacher, offset, n, err, chunkSize)
return
}
// 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.
// pinned to before, which the caller unpins.
func (stream *chunkStream) pin(cacher *SingleChunkCacher) (previous *SingleChunkCacher) {
if stream == nil || stream.cacher == cacher {
if stream == nil {
return nil
}
stream.mu.Lock()
defer stream.mu.Unlock()
if stream.cacher == cacher {
return nil
}
atomic.AddInt32(&cacher.pins, 1)
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) {
// detach unpins the stream from its chunk and returns that chunk for the
// caller to unpin. A non-nil expected cacher makes the detach conditional on
// the stream still being positioned in it.
func (stream *chunkStream) detach(expected *SingleChunkCacher) (previous *SingleChunkCacher) {
if stream == nil {
return nil
}
stream.mu.Lock()
defer stream.mu.Unlock()
if expected != nil && stream.cacher != expected {
return nil
}
previous = stream.cacher
stream.cacher = nil
return previous
@@ -269,24 +304,12 @@ func (rc *ReaderCache) releaseIfFinished(stream *chunkStream, cacher *SingleChun
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)
rc.unpin(stream.detach(cacher))
}
// 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)
rc.unpin(stream.detach(nil))
}
// unpin drops one stream's pin. Once no stream is positioned in the chunk it
@@ -341,26 +364,41 @@ func (rc *ReaderCache) removeUnpinned(downloader *SingleChunkCacher) (removed bo
return
}
// consumable reports whether the cacher could be dropped: its buffer was
// fully read or every stream positioned in it has left, and no reader is
// attached or pinned. All fields are atomics, so callers can skip the
// ReaderCache lock whenever removal is impossible.
func (s *SingleChunkCacher) consumable() bool {
return atomic.LoadInt32(&s.readers) == 0 &&
atomic.LoadInt32(&s.pins) == 0 &&
(atomic.LoadInt32(&s.consumed) != 0 || atomic.LoadInt32(&s.left) != 0)
}
// 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)
if !downloader.consumable() {
return
}
rc.Lock()
removed := rc.removeConsumedLocked(downloader)
rc.Unlock()
if removed {
downloader.destroy()
}
}
func (rc *ReaderCache) removeConsumedLocked(downloader *SingleChunkCacher) (removed bool) {
if rc.downloaders[downloader.chunkFileId] == downloader && downloader.consumable() {
delete(rc.downloaders, downloader.chunkFileId)
return true
}
return false
}
func (rc *ReaderCache) destroy() {
rc.Lock()
downloaders := rc.downloaders
@@ -385,8 +423,9 @@ func newSingleChunkCacher(parent *ReaderCache, fileId string, cipherKey []byte,
}
// startCaching downloads a chunk shared by concurrent readers.
// The caller must s.wg.Add(1) before publishing the cacher in the
// downloaders map, so a removal can never observe the fetch as absent.
func (s *SingleChunkCacher) startCaching() {
s.wg.Add(1)
defer func() {
close(s.done)
s.wg.Done()
+31
View File
@@ -724,3 +724,34 @@ func TestSingleChunkCacherOneReaderCancelsOthersContinue(t *testing.T) {
t.Error("Other reader did not complete")
}
}
// TestReaderCacheBookkeepingOffCacheLock guards the contention fix: pin
// lifecycle calls that cannot remove a cacher must not block on the
// ReaderCache lock.
func TestReaderCacheBookkeepingOffCacheLock(t *testing.T) {
cache := newMockChunkCacheForReaderCache()
rc := NewReaderCache(10, cache, nil, nil)
defer rc.destroy()
cacher := &SingleChunkCacher{parent: rc, chunkFileId: "pinned-chunk"}
atomic.StoreInt32(&cacher.readers, 1) // a read in flight: not consumable
rc.Lock()
started := make(chan struct{})
done := make(chan struct{})
go func() {
close(started)
defer close(done)
var stream chunkStream
stream.pin(cacher)
rc.releaseStream(&stream)
}()
<-started
select {
case <-done:
case <-time.After(time.Second):
rc.Unlock()
t.Fatal("stream pin lifecycle blocked on the cache lock")
}
rc.Unlock()
}