mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-30 11:45:42 +00:00
filer: keep shared chunk buffers pinned while another stream reads them (#11502)
The ReaderCache is shared by all streams of a process (every S3 GET, for instance), but a ChunkReadAt released chunks as if it owned them: - moving on to the next chunk called UnCache on the previous one, destroying the buffer even when other streams were still inside it; - since #11384 a buffer is dropped once any reader has consumed it to the end and no read call is in flight. Streams copy out in slices (256 KiB in the S3 gateway), so between two calls a slower stream is not attached and loses the buffer to a faster one. Either way the slower stream refetches the whole chunk from the volume servers. With many clients downloading the same popular object at once, each chunk is fetched over and over; in production we saw the S3 gateway pull ~10 Gbit/s from volume servers while serving ~1 Gbit/s to clients. A ChunkReadAt now pins the chunk it is positioned in. The pin is taken and released only under the ReaderCache lock, since concurrent ReadAt calls on one ChunkReadAt (as in mount) share it. It is released when the stream reads the chunk to its end, moves to another chunk (including one served from the chunk cache), or falls back to random reads. A buffer is dropped once no stream pins it and no read is in progress, if it was consumed or its last stream left it; a read still in flight when the stream leaves drops it on detach, as UnCache did via destroy. Eviction by slot limit and memory budget is unchanged. lastChunkFid is now guarded as well: concurrent ReadAt calls raced on it. Tests: two ChunkReadAt instances streaming one object in interleaved slices fetch each chunk exactly once (2-3 times before); leaving a chunk for a chunk-cache hit or while another read is in flight releases it; concurrent ReadAt calls on one ChunkReadAt leave no pins behind under -race.
This commit is contained in:
+15
-12
@@ -32,7 +32,9 @@ type ChunkReadAt struct {
|
||||
fileSize int64
|
||||
readerCache *ReaderCache
|
||||
readerPattern *ReaderPattern
|
||||
lastChunkMu sync.Mutex // guards lastChunkFid; mount issues concurrent ReadAt calls
|
||||
lastChunkFid string
|
||||
stream chunkStream // chunk this reader is positioned in, pinned in the shared readerCache
|
||||
prefetchCount int // Number of chunks to prefetch ahead during sequential reads
|
||||
ctx context.Context // Context used for cancellation during chunk read operations
|
||||
}
|
||||
@@ -341,6 +343,7 @@ func (c *ChunkReadAt) doReadAt(ctx context.Context, p []byte, offset int64) (n i
|
||||
func (c *ChunkReadAt) readChunkSliceAt(ctx context.Context, buffer []byte, chunkView *ChunkView, nextChunkViews *Interval[*ChunkView], offset uint64) (n int, err error) {
|
||||
|
||||
if c.readerPattern.IsRandomMode() {
|
||||
c.readerCache.releaseStream(&c.stream)
|
||||
n, err := c.readerCache.chunkCache.ReadChunkAt(buffer, chunkView.FileId, offset)
|
||||
if n > 0 {
|
||||
return n, err
|
||||
@@ -350,20 +353,20 @@ func (c *ChunkReadAt) readChunkSliceAt(ctx context.Context, buffer []byte, chunk
|
||||
}
|
||||
|
||||
shouldCache := (uint64(chunkView.ViewOffset) + chunkView.ChunkSize) <= c.readerCache.chunkCache.GetMaxFilePartSizeInCache()
|
||||
n, err = c.readerCache.ReadChunkAt(ctx, buffer, chunkView.FileId, chunkView.CipherKey, chunkView.IsGzipped, int64(offset), int(chunkView.ChunkSize), shouldCache)
|
||||
if c.lastChunkFid != chunkView.FileId {
|
||||
if chunkView.OffsetInChunk == 0 { // start of a new chunk
|
||||
if c.lastChunkFid != "" {
|
||||
c.readerCache.UnCache(c.lastChunkFid)
|
||||
}
|
||||
if nextChunkViews != nil && c.prefetchCount > 0 {
|
||||
// Prefetch multiple chunks ahead for better sequential read throughput
|
||||
// This keeps the network pipeline full with parallel chunk fetches
|
||||
c.readerCache.MaybeCache(nextChunkViews, c.prefetchCount)
|
||||
}
|
||||
// The previous chunk is released through the stream pin rather than
|
||||
// UnCache: the buffer is shared, and other streams may still be reading it.
|
||||
n, err = c.readerCache.readChunkAt(ctx, &c.stream, buffer, chunkView.FileId, chunkView.CipherKey, chunkView.IsGzipped, int64(offset), int(chunkView.ChunkSize), shouldCache)
|
||||
c.lastChunkMu.Lock()
|
||||
enteredChunk := c.lastChunkFid != chunkView.FileId
|
||||
c.lastChunkFid = chunkView.FileId
|
||||
c.lastChunkMu.Unlock()
|
||||
if enteredChunk && chunkView.OffsetInChunk == 0 { // start of a new chunk
|
||||
if nextChunkViews != nil && c.prefetchCount > 0 {
|
||||
// Prefetch multiple chunks ahead for better sequential read throughput
|
||||
// This keeps the network pipeline full with parallel chunk fetches
|
||||
c.readerCache.MaybeCache(nextChunkViews, c.prefetchCount)
|
||||
}
|
||||
}
|
||||
c.lastChunkFid = chunkView.FileId
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,209 @@
|
||||
package filer
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/util/chunk_cache"
|
||||
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
|
||||
)
|
||||
|
||||
// Two streams (e.g. two S3 GETs) reading the same object through one shared
|
||||
// ReaderCache, interleaved in small slices, must share each chunk download:
|
||||
// one stream finishing a chunk, or moving on to the next one, must not drop a
|
||||
// buffer the other stream is still reading.
|
||||
func TestChunkReadAtConcurrentStreamsShareChunks(t *testing.T) {
|
||||
const chunkSize = 64 << 10
|
||||
const chunkCount = 3
|
||||
const sliceSize = 16 << 10
|
||||
|
||||
var mu sync.Mutex
|
||||
fetches := map[string]int{}
|
||||
rc := NewReaderCache(64, (*chunk_cache.TieredChunkCache)(nil), 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) {
|
||||
mu.Lock()
|
||||
fetches[fileId]++
|
||||
mu.Unlock()
|
||||
for i := range buffer {
|
||||
buffer[i] = fileId[len(fileId)-1]
|
||||
}
|
||||
return len(buffer), nil
|
||||
}
|
||||
|
||||
newStream := func() *ChunkReadAt {
|
||||
views := NewIntervalList[*ChunkView]()
|
||||
for i := 0; i < chunkCount; i++ {
|
||||
views.AppendInterval(&Interval[*ChunkView]{
|
||||
StartOffset: int64(i * chunkSize),
|
||||
StopOffset: int64((i + 1) * chunkSize),
|
||||
Value: &ChunkView{
|
||||
FileId: fmt.Sprintf("chunk%d", i),
|
||||
ViewSize: chunkSize,
|
||||
ViewOffset: int64(i * chunkSize),
|
||||
ChunkSize: chunkSize,
|
||||
},
|
||||
})
|
||||
}
|
||||
return NewChunkReaderAtFromClient(context.Background(), rc, views, chunkSize*chunkCount, 0)
|
||||
}
|
||||
streams := []*ChunkReadAt{newStream(), newStream()}
|
||||
|
||||
// The streams alternate slice by slice, the second one lagging one slice
|
||||
// behind, so the leader always finishes a chunk while the follower is
|
||||
// still inside it.
|
||||
offsets := make([]int64, len(streams))
|
||||
offsets[1] = -sliceSize
|
||||
for offsets[1] < chunkSize*chunkCount {
|
||||
for i, stream := range streams {
|
||||
if offsets[i] < 0 || offsets[i] >= chunkSize*chunkCount {
|
||||
offsets[i] += sliceSize
|
||||
continue
|
||||
}
|
||||
buf := make([]byte, sliceSize)
|
||||
n, err := stream.ReadAt(buf, offsets[i])
|
||||
if (err != nil && err != io.EOF) || n != sliceSize {
|
||||
t.Fatalf("stream %d at %d: n=%d err=%v", i, offsets[i], n, err)
|
||||
}
|
||||
if want := byte('0' + offsets[i]/chunkSize); buf[0] != want || buf[n-1] != want {
|
||||
t.Fatalf("stream %d at %d: got %q, want %q", i, offsets[i], buf[0], want)
|
||||
}
|
||||
offsets[i] += sliceSize
|
||||
}
|
||||
}
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
for i := 0; i < chunkCount; i++ {
|
||||
fileId := fmt.Sprintf("chunk%d", i)
|
||||
if fetches[fileId] != 1 {
|
||||
t.Errorf("%s fetched %d times for two concurrent streams, want 1", fileId, fetches[fileId])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func newStreamTestReaderCache(chunkCache chunk_cache.ChunkCache) *ReaderCache {
|
||||
rc := NewReaderCache(64, chunkCache, func(context.Context, string) ([]string, error) {
|
||||
return []string{"unused"}, nil
|
||||
}, nil)
|
||||
rc.fetchChunkDataFn = func(_ context.Context, buffer []byte, _ []string, _ []byte, _ bool, _ bool, _ int64, _ string, _ util_http.RefreshUrlsFunc) (int, error) {
|
||||
return len(buffer), nil
|
||||
}
|
||||
return rc
|
||||
}
|
||||
|
||||
func isRetained(rc *ReaderCache, fileId string) bool {
|
||||
rc.Lock()
|
||||
defer rc.Unlock()
|
||||
_, found := rc.downloaders[fileId]
|
||||
return found
|
||||
}
|
||||
|
||||
// A stream moving on to a chunk served from the chunk cache must still
|
||||
// release the chunk it was positioned in before.
|
||||
func TestChunkStreamReleasesPreviousChunkOnCacheHit(t *testing.T) {
|
||||
cache := newMockChunkCacheForReaderCache()
|
||||
cache.SetChunk("chunk1", make([]byte, 4<<10))
|
||||
rc := newStreamTestReaderCache(cache)
|
||||
defer rc.destroy()
|
||||
|
||||
stream := &chunkStream{}
|
||||
if _, err := rc.readChunkAt(context.Background(), stream, make([]byte, 1<<10), "chunk0", nil, false, 0, 4<<10, false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if n, err := rc.readChunkAt(context.Background(), stream, make([]byte, 1<<10), "chunk1", nil, false, 0, 4<<10, true); err != nil || n == 0 {
|
||||
t.Fatalf("cache hit read: n=%d err=%v", n, err)
|
||||
}
|
||||
if isRetained(rc, "chunk0") {
|
||||
t.Fatal("chunk0 still retained after the stream moved on to a cached chunk")
|
||||
}
|
||||
}
|
||||
|
||||
// A chunk the stream leaves while another read is in flight must be dropped
|
||||
// once that read ends, even if it did not read the chunk to the end.
|
||||
func TestChunkStreamDropsLeftChunkAfterInFlightRead(t *testing.T) {
|
||||
rc := newStreamTestReaderCache(newMockChunkCacheForReaderCache())
|
||||
defer rc.destroy()
|
||||
|
||||
stream := &chunkStream{}
|
||||
if _, err := rc.readChunkAt(context.Background(), stream, make([]byte, 1<<10), "chunk0", nil, false, 0, 4<<10, false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Another reader is in the middle of a partial read of chunk0.
|
||||
rc.Lock()
|
||||
other := rc.downloaders["chunk0"]
|
||||
other.wg.Add(1)
|
||||
atomic.AddInt32(&other.readers, 1)
|
||||
rc.Unlock()
|
||||
|
||||
if _, err := rc.readChunkAt(context.Background(), stream, make([]byte, 1<<10), "chunk1", nil, false, 0, 4<<10, false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !isRetained(rc, "chunk0") {
|
||||
t.Fatal("chunk0 dropped while a read was still in flight")
|
||||
}
|
||||
|
||||
// The in-flight read ends without reaching the end of the chunk.
|
||||
other.wg.Done()
|
||||
atomic.AddInt32(&other.readers, -1)
|
||||
rc.removeConsumed(other)
|
||||
if isRetained(rc, "chunk0") {
|
||||
t.Fatal("chunk0 retained after the stream left it and the last read ended")
|
||||
}
|
||||
}
|
||||
|
||||
// Concurrent ReadAt calls on one ChunkReadAt (as mount does) share its stream
|
||||
// pin; they must neither race on it nor leak or double-release pins.
|
||||
func TestChunkStreamConcurrentReadsOnOneReader(t *testing.T) {
|
||||
const chunkSize = 16 << 10
|
||||
const chunkCount = 4
|
||||
const sliceSize = 4 << 10
|
||||
rc := newStreamTestReaderCache((*chunk_cache.TieredChunkCache)(nil))
|
||||
defer rc.destroy()
|
||||
|
||||
views := NewIntervalList[*ChunkView]()
|
||||
for i := 0; i < chunkCount; i++ {
|
||||
views.AppendInterval(&Interval[*ChunkView]{
|
||||
StartOffset: int64(i * chunkSize),
|
||||
StopOffset: int64((i + 1) * chunkSize),
|
||||
Value: &ChunkView{
|
||||
FileId: fmt.Sprintf("chunk%d", i),
|
||||
ViewSize: chunkSize,
|
||||
ViewOffset: int64(i * chunkSize),
|
||||
ChunkSize: chunkSize,
|
||||
},
|
||||
})
|
||||
}
|
||||
reader := NewChunkReaderAtFromClient(context.Background(), rc, views, chunkSize*chunkCount, 0)
|
||||
|
||||
var wg sync.WaitGroup
|
||||
for g := 0; g < 8; g++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
for offset := int64(0); offset < chunkSize*chunkCount; offset += sliceSize {
|
||||
if _, err := reader.ReadAt(make([]byte, sliceSize), offset); err != nil && err != io.EOF {
|
||||
t.Error(err)
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
// Finishing the last chunk releases the stream's final pin.
|
||||
if _, err := reader.ReadAt(make([]byte, sliceSize), chunkSize*chunkCount-sliceSize); err != nil && err != io.EOF {
|
||||
t.Fatal(err)
|
||||
}
|
||||
rc.Lock()
|
||||
defer rc.Unlock()
|
||||
for fileId, cacher := range rc.downloaders {
|
||||
t.Errorf("%s retained after all reads finished: pins=%d readers=%d", fileId, atomic.LoadInt32(&cacher.pins), atomic.LoadInt32(&cacher.readers))
|
||||
}
|
||||
}
|
||||
@@ -35,6 +35,8 @@ 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
|
||||
@@ -112,7 +114,19 @@ func (rc *ReaderCache) MaybeCache(chunkViews *Interval[*ChunkView], count int) {
|
||||
return
|
||||
}
|
||||
|
||||
// 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()
|
||||
|
||||
@@ -130,8 +144,11 @@ retry:
|
||||
// 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
|
||||
}
|
||||
@@ -144,7 +161,10 @@ retry:
|
||||
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
|
||||
}
|
||||
}
|
||||
@@ -175,9 +195,77 @@ retry:
|
||||
rc.downloaders[fileId] = cacher
|
||||
cacher.wg.Add(1)
|
||||
atomic.AddInt32(&cacher.readers, 1)
|
||||
previous := stream.pin(cacher)
|
||||
rc.Unlock()
|
||||
rc.unpin(previous)
|
||||
|
||||
return cacher.readChunkAt(ctx, buffer, offset)
|
||||
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) {
|
||||
@@ -202,15 +290,17 @@ func (rc *ReaderCache) remove(downloader *SingleChunkCacher) {
|
||||
}
|
||||
}
|
||||
|
||||
// removeConsumed drops a cacher once its buffer was fully read and no
|
||||
// readers remain attached. 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.
|
||||
// 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.consumed) != 0
|
||||
atomic.LoadInt32(&downloader.pins) == 0 &&
|
||||
(atomic.LoadInt32(&downloader.consumed) != 0 || atomic.LoadInt32(&downloader.left) != 0)
|
||||
if removed {
|
||||
delete(rc.downloaders, downloader.chunkFileId)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user