|
From 3e9633eaf0fbc52d93b40948e8c289b88d1b1a32 Mon Sep 17 00:00:00 2001 |
|
From: Mic92 <joerg@thalheim.io> |
|
Date: Fri, 25 Sep 2026 16:19:15 +0200 |
|
Subject: [PATCH] filer: skip the remote stat for names missing from a remote |
|
listing |
|
MIME-Version: 1.0 |
|
Content-Type: text/plain; charset=UTF-8 |
|
Content-Transfer-Encoding: 8bit |
|
|
|
A local miss under a remote mount makes the filer stat the object in |
|
the remote bucket. Most misses are for objects that do not exist, and |
|
each one costs a request that some providers rate limit. |
|
|
|
With filer.options.remote_key_index enabled, list every mount once in |
|
the background after the mounts load. Once a mount's listing is stored, |
|
a miss for a name outside it returns not-found without a remote |
|
request. Before the listing finishes, or if it fails, lookups behave as |
|
before. Objects written to the remote afterwards stay hidden until a |
|
restart or remote.meta.sync. |
|
|
|
filer.options.remote_key_index_dir additionally keeps each listing on |
|
disk. A restart then gates lookups from the start with the last |
|
snapshot, and the fresh listing replaces it once it completes. |
|
|
|
Signed-off-by: Jörg Thalheim <joerg@thalheim.io> |
|
--- |
|
weed/filer/filer.go | 22 ++-- |
|
weed/filer/filer_lazy_remote.go | 6 ++ |
|
weed/filer/filer_lazy_remote_test.go | 140 +++++++++++++++++++++++- |
|
weed/filer/filer_on_meta_event.go | 1 + |
|
weed/filer/filer_remote_key_index.go | 156 +++++++++++++++++++++++++++ |
|
weed/server/filer_server.go | 2 + |
|
6 files changed, 320 insertions(+), 7 deletions(-) |
|
create mode 100644 weed/filer/filer_remote_key_index.go |
|
|
|
diff --git a/weed/filer/filer.go b/weed/filer/filer.go |
|
index ef9b18f..45d8a1d 100644 |
|
--- a/weed/filer/filer.go |
|
+++ b/weed/filer/filer.go |
|
@@ -7,6 +7,7 @@ import ( |
|
"os" |
|
"sort" |
|
"strings" |
|
+ "sync" |
|
"time" |
|
|
|
"github.com/seaweedfs/seaweedfs/weed/remote_storage" |
|
@@ -67,12 +68,21 @@ type Filer struct { |
|
lazyListGroup singleflight.Group |
|
Dlm *lock_manager.DistributedLockManager |
|
MaxFilenameLength uint32 |
|
- deletionQuit chan struct{} |
|
- DeletionRetryQueue *DeletionRetryQueue |
|
- EmptyFolderCleaner *empty_folder_cleanup.EmptyFolderCleaner |
|
- EmptyFolderCleanupDelay time.Duration |
|
- persistedLogCache *persistedLogCache |
|
- metaLogInflight metaLogInflight |
|
+ // EnableRemoteKeyIndex lists each remote mount once and skips the remote |
|
+ // stat for names missing from that listing. Objects written to the remote |
|
+ // behind the filer stay hidden until a restart or remote.meta.sync. |
|
+ EnableRemoteKeyIndex bool |
|
+ // RemoteKeyIndexDir keeps each listing across restarts, so lookups are |
|
+ // gated from the start instead of after the first listing. |
|
+ RemoteKeyIndexDir string |
|
+ remoteKeys sync.Map // mount dir -> *remoteKeySet |
|
+ remoteKeysLoading sync.Map // mount dir -> struct{} |
|
+ deletionQuit chan struct{} |
|
+ DeletionRetryQueue *DeletionRetryQueue |
|
+ EmptyFolderCleaner *empty_folder_cleanup.EmptyFolderCleaner |
|
+ EmptyFolderCleanupDelay time.Duration |
|
+ persistedLogCache *persistedLogCache |
|
+ metaLogInflight metaLogInflight |
|
} |
|
|
|
func NewFiler(masters pb.ServerDiscovery, grpcDialOption grpc.DialOption, filerHost pb.ServerAddress, filerGroup string, collection string, replication string, dataCenter string, maxFilenameLength uint32, notifyFn func()) *Filer { |
|
diff --git a/weed/filer/filer_lazy_remote.go b/weed/filer/filer_lazy_remote.go |
|
index 1b6983b..7d332b5 100644 |
|
--- a/weed/filer/filer_lazy_remote.go |
|
+++ b/weed/filer/filer_lazy_remote.go |
|
@@ -48,6 +48,12 @@ func (f *Filer) maybeLazyFetchFromRemote(ctx context.Context, p util.FullPath) ( |
|
return nil, nil |
|
} |
|
|
|
+ if keys, ok := f.remoteKeys.Load(string(mountDir)); ok { |
|
+ if _, known := keys.(*remoteKeySet).keys[p]; !known { |
|
+ return nil, nil |
|
+ } |
|
+ } |
|
+ |
|
relPath := strings.TrimPrefix(string(p), string(mountDir)) |
|
if relPath != "" && !strings.HasPrefix(relPath, "/") { |
|
relPath = "/" + relPath |
|
diff --git a/weed/filer/filer_lazy_remote_test.go b/weed/filer/filer_lazy_remote_test.go |
|
index f9b3442..cc00037 100644 |
|
--- a/weed/filer/filer_lazy_remote_test.go |
|
+++ b/weed/filer/filer_lazy_remote_test.go |
|
@@ -201,6 +201,7 @@ type stubRemoteClient struct { |
|
deleteCalls []*remote_pb.RemoteStorageLocation |
|
removeCalls []*remote_pb.RemoteStorageLocation |
|
|
|
+ traverseFn func(loc *remote_pb.RemoteStorageLocation, visitFn remote_storage.VisitFunc) error |
|
listDirFn func(loc *remote_pb.RemoteStorageLocation, visitFn remote_storage.VisitFunc) error |
|
listDirCalls int |
|
} |
|
@@ -208,7 +209,10 @@ type stubRemoteClient struct { |
|
func (c *stubRemoteClient) StatFile(*remote_pb.RemoteStorageLocation) (*filer_pb.RemoteEntry, error) { |
|
return c.statResult, c.statErr |
|
} |
|
-func (c *stubRemoteClient) Traverse(*remote_pb.RemoteStorageLocation, remote_storage.VisitFunc) error { |
|
+func (c *stubRemoteClient) Traverse(loc *remote_pb.RemoteStorageLocation, visitFn remote_storage.VisitFunc) error { |
|
+ if c.traverseFn != nil { |
|
+ return c.traverseFn(loc, visitFn) |
|
+ } |
|
return nil |
|
} |
|
func (c *stubRemoteClient) ReadFile(*remote_pb.RemoteStorageLocation, int64, int64) ([]byte, error) { |
|
@@ -1300,3 +1304,137 @@ func TestMaybeLazyListFromRemote_ContextGuardPreventsRecursion(t *testing.T) { |
|
f.maybeLazyListFromRemote(fetchCtx, util.FullPath("/buckets/mybucket")) |
|
assert.Equal(t, 0, stub.listDirCalls) |
|
} |
|
+ |
|
+func newKeyIndexFiler(t *testing.T, storageType string, keys []string) (*Filer, *countingRemoteClient, func()) { |
|
+ t.Helper() |
|
+ stub := &countingRemoteClient{ |
|
+ stubRemoteClient: stubRemoteClient{ |
|
+ statResult: &filer_pb.RemoteEntry{RemoteMtime: 1, RemoteSize: 1}, |
|
+ traverseFn: func(_ *remote_pb.RemoteStorageLocation, visit remote_storage.VisitFunc) error { |
|
+ for _, k := range keys { |
|
+ dir, name := util.FullPath("/" + k).DirAndName() |
|
+ if err := visit(dir, name, false, &filer_pb.RemoteEntry{}); err != nil { |
|
+ return err |
|
+ } |
|
+ } |
|
+ return nil |
|
+ }, |
|
+ }, |
|
+ } |
|
+ restore := registerStubMaker(t, storageType, stub) |
|
+ conf := &remote_pb.RemoteConf{Name: storageType, Type: storageType} |
|
+ rs := NewFilerRemoteStorage() |
|
+ rs.storageNameToConf[conf.Name] = conf |
|
+ rs.mapDirectoryToRemoteStorage("/buckets/mybucket", &remote_pb.RemoteStorageLocation{ |
|
+ Name: storageType, |
|
+ Bucket: "mybucket", |
|
+ Path: "/", |
|
+ }) |
|
+ f := newTestFiler(t, newStubFilerStore(), rs) |
|
+ f.EnableRemoteKeyIndex = true |
|
+ return f, stub, restore |
|
+} |
|
+ |
|
+func TestRemoteKeyIndex_UnknownNameSkipsRemoteStat(t *testing.T) { |
|
+ f, stub, restore := newKeyIndexFiler(t, "stub_keyidx_unknown", []string{"known.narinfo", "log/a"}) |
|
+ defer restore() |
|
+ require.NoError(t, f.buildRemoteKeyIndexes()) |
|
+ |
|
+ entry, err := f.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/other.narinfo") |
|
+ require.NoError(t, err) |
|
+ assert.Nil(t, entry) |
|
+ assert.Equal(t, 0, stub.statCalls, "a name absent from the listing must not reach the remote") |
|
+} |
|
+ |
|
+func TestRemoteKeyIndex_KnownNameIsFetched(t *testing.T) { |
|
+ f, stub, restore := newKeyIndexFiler(t, "stub_keyidx_known", []string{"known.narinfo", "log/a"}) |
|
+ defer restore() |
|
+ require.NoError(t, f.buildRemoteKeyIndexes()) |
|
+ |
|
+ for _, p := range []util.FullPath{"/buckets/mybucket/known.narinfo", "/buckets/mybucket/log/a"} { |
|
+ entry, err := f.maybeLazyFetchFromRemote(context.Background(), p) |
|
+ require.NoError(t, err) |
|
+ assert.NotNil(t, entry, p) |
|
+ } |
|
+ assert.Equal(t, 2, stub.statCalls) |
|
+} |
|
+ |
|
+func TestRemoteKeyIndex_BeforeListingLooksUpRemote(t *testing.T) { |
|
+ f, stub, restore := newKeyIndexFiler(t, "stub_keyidx_early", nil) |
|
+ defer restore() |
|
+ |
|
+ entry, err := f.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/x.narinfo") |
|
+ require.NoError(t, err) |
|
+ assert.NotNil(t, entry) |
|
+ assert.Equal(t, 1, stub.statCalls) |
|
+} |
|
+ |
|
+func TestRemoteKeyIndex_FailedListingKeepsLookingUpRemote(t *testing.T) { |
|
+ f, stub, restore := newKeyIndexFiler(t, "stub_keyidx_fail", nil) |
|
+ defer restore() |
|
+ stub.traverseFn = func(*remote_pb.RemoteStorageLocation, remote_storage.VisitFunc) error { |
|
+ return errors.New("list failed") |
|
+ } |
|
+ require.Error(t, f.buildRemoteKeyIndexes()) |
|
+ |
|
+ entry, err := f.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/x.narinfo") |
|
+ require.NoError(t, err) |
|
+ assert.NotNil(t, entry) |
|
+} |
|
+ |
|
+func TestRemoteKeyIndex_DisabledByDefault(t *testing.T) { |
|
+ f, stub, restore := newKeyIndexFiler(t, "stub_keyidx_off", []string{"known.narinfo"}) |
|
+ defer restore() |
|
+ f.EnableRemoteKeyIndex = false |
|
+ require.NoError(t, f.buildRemoteKeyIndexes()) |
|
+ |
|
+ entry, err := f.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/other.narinfo") |
|
+ require.NoError(t, err) |
|
+ assert.NotNil(t, entry) |
|
+ assert.Equal(t, 1, stub.statCalls) |
|
+} |
|
+ |
|
+func TestRemoteKeyIndex_SnapshotGatesLookupsBeforeListing(t *testing.T) { |
|
+ dir := t.TempDir() |
|
+ f, _, restore := newKeyIndexFiler(t, "stub_keyidx_snap1", []string{"known.narinfo"}) |
|
+ defer restore() |
|
+ f.RemoteKeyIndexDir = dir |
|
+ require.NoError(t, f.buildRemoteKeyIndexes()) |
|
+ |
|
+ f2, stub2, restore2 := newKeyIndexFiler(t, "stub_keyidx_snap1b", nil) |
|
+ defer restore2() |
|
+ f2.RemoteKeyIndexDir = dir |
|
+ stub2.traverseFn = func(*remote_pb.RemoteStorageLocation, remote_storage.VisitFunc) error { |
|
+ return errors.New("b2 unreachable") |
|
+ } |
|
+ require.Error(t, f2.buildRemoteKeyIndexes()) |
|
+ |
|
+ entry, err := f2.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/other.narinfo") |
|
+ require.NoError(t, err) |
|
+ assert.Nil(t, entry) |
|
+ entry, err = f2.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/known.narinfo") |
|
+ require.NoError(t, err) |
|
+ assert.NotNil(t, entry) |
|
+ assert.Equal(t, 1, stub2.statCalls, "only the snapshotted name reaches the remote") |
|
+} |
|
+ |
|
+func TestRemoteKeyIndex_ListingReplacesSnapshot(t *testing.T) { |
|
+ dir := t.TempDir() |
|
+ f, _, restore := newKeyIndexFiler(t, "stub_keyidx_snap2", []string{"old.narinfo"}) |
|
+ defer restore() |
|
+ f.RemoteKeyIndexDir = dir |
|
+ require.NoError(t, f.buildRemoteKeyIndexes()) |
|
+ |
|
+ f2, stub2, restore2 := newKeyIndexFiler(t, "stub_keyidx_snap2b", []string{"new.narinfo"}) |
|
+ defer restore2() |
|
+ f2.RemoteKeyIndexDir = dir |
|
+ require.NoError(t, f2.buildRemoteKeyIndexes()) |
|
+ |
|
+ entry, err := f2.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/old.narinfo") |
|
+ require.NoError(t, err) |
|
+ assert.Nil(t, entry) |
|
+ entry, err = f2.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/new.narinfo") |
|
+ require.NoError(t, err) |
|
+ assert.NotNil(t, entry) |
|
+ assert.Equal(t, 1, stub2.statCalls) |
|
+} |
|
diff --git a/weed/filer/filer_on_meta_event.go b/weed/filer/filer_on_meta_event.go |
|
index d57810f..30cc0a0 100644 |
|
--- a/weed/filer/filer_on_meta_event.go |
|
+++ b/weed/filer/filer_on_meta_event.go |
|
@@ -134,6 +134,7 @@ func (f *Filer) LoadRemoteStorageConfAndMapping() { |
|
glog.Errorf("read remote conf and mapping: %v", err) |
|
return |
|
} |
|
+ f.startRemoteKeyIndexes() |
|
} |
|
func (f *Filer) maybeReloadRemoteStorageConfigurationAndMapping(event *filer_pb.SubscribeMetadataResponse) { |
|
if !filer_pb.MetadataEventTouchesDirectory(event, DirectoryEtcRemote) { |
|
diff --git a/weed/filer/filer_remote_key_index.go b/weed/filer/filer_remote_key_index.go |
|
new file mode 100644 |
|
index 0000000..4e87d2f |
|
--- /dev/null |
|
+++ b/weed/filer/filer_remote_key_index.go |
|
@@ -0,0 +1,156 @@ |
|
+package filer |
|
+ |
|
+import ( |
|
+ "bufio" |
|
+ "context" |
|
+ "fmt" |
|
+ "net/url" |
|
+ "os" |
|
+ "path/filepath" |
|
+ |
|
+ "github.com/seaweedfs/seaweedfs/weed/glog" |
|
+ "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" |
|
+ "github.com/seaweedfs/seaweedfs/weed/pb/remote_pb" |
|
+ "github.com/seaweedfs/seaweedfs/weed/util" |
|
+) |
|
+ |
|
+// remoteKeySet holds the object paths of one mount. A set loaded from a |
|
+// snapshot is not listed yet and is replaced by the first listing. |
|
+type remoteKeySet struct { |
|
+ keys map[util.FullPath]struct{} |
|
+ listed bool |
|
+} |
|
+ |
|
+func (rs *FilerRemoteStorage) mounts() map[util.FullPath]*remote_pb.RemoteStorageLocation { |
|
+ rs.mu.RLock() |
|
+ defer rs.mu.RUnlock() |
|
+ mounts := make(map[util.FullPath]*remote_pb.RemoteStorageLocation) |
|
+ rs.rules.Walk(func(key []byte, loc *remote_pb.RemoteStorageLocation) bool { |
|
+ mounts[util.FullPath(key[:len(key)-1])] = loc |
|
+ return true |
|
+ }) |
|
+ return mounts |
|
+} |
|
+ |
|
+// startRemoteKeyIndexes lists every mount that has no listing yet in the |
|
+// background. Until a listing or snapshot is stored, lookups behave as |
|
+// without an index. |
|
+func (f *Filer) startRemoteKeyIndexes() { |
|
+ if !f.EnableRemoteKeyIndex || f.RemoteStorage == nil { |
|
+ return |
|
+ } |
|
+ go func() { |
|
+ if err := f.buildRemoteKeyIndexes(); err != nil { |
|
+ glog.Warningf("remote key index: %v", err) |
|
+ } |
|
+ }() |
|
+} |
|
+ |
|
+func (f *Filer) buildRemoteKeyIndexes() error { |
|
+ if !f.EnableRemoteKeyIndex || f.RemoteStorage == nil { |
|
+ return nil |
|
+ } |
|
+ var firstErr error |
|
+ for mountDir, loc := range f.RemoteStorage.mounts() { |
|
+ if set, ok := f.remoteKeys.Load(string(mountDir)); ok && set.(*remoteKeySet).listed { |
|
+ continue |
|
+ } |
|
+ if _, busy := f.remoteKeysLoading.LoadOrStore(string(mountDir), struct{}{}); busy { |
|
+ continue |
|
+ } |
|
+ err := f.buildRemoteKeyIndex(mountDir, loc) |
|
+ f.remoteKeysLoading.Delete(string(mountDir)) |
|
+ if err != nil && firstErr == nil { |
|
+ firstErr = fmt.Errorf("list %s: %w", mountDir, err) |
|
+ } |
|
+ } |
|
+ return firstErr |
|
+} |
|
+ |
|
+func (f *Filer) buildRemoteKeyIndex(mountDir util.FullPath, loc *remote_pb.RemoteStorageLocation) error { |
|
+ if _, ok := f.remoteKeys.Load(string(mountDir)); !ok { |
|
+ if keys, err := f.loadRemoteKeySnapshot(mountDir); err != nil { |
|
+ glog.Warningf("remote key index: snapshot for %s: %v", mountDir, err) |
|
+ } else if keys != nil { |
|
+ f.remoteKeys.Store(string(mountDir), &remoteKeySet{keys: keys}) |
|
+ glog.V(0).Infof("remote key index: %d objects under %s from snapshot", len(keys), mountDir) |
|
+ } |
|
+ } |
|
+ |
|
+ conf, found := f.RemoteStorage.GetRemoteStorageConf(loc.Name) |
|
+ if !found { |
|
+ return fmt.Errorf("no remote storage %q", loc.Name) |
|
+ } |
|
+ client, err := f.buildRemoteStorageClient(context.Background(), conf) |
|
+ if err != nil { |
|
+ return err |
|
+ } |
|
+ keys := make(map[util.FullPath]struct{}) |
|
+ err = client.Traverse(loc, func(dir string, name string, _ bool, _ *filer_pb.RemoteEntry) error { |
|
+ remotePath := util.FullPath(dir).Child(name) |
|
+ keys[MapRemoteStorageLocationPathToFullPath(mountDir, loc, string(remotePath))] = struct{}{} |
|
+ return nil |
|
+ }) |
|
+ if err != nil { |
|
+ return err |
|
+ } |
|
+ f.remoteKeys.Store(string(mountDir), &remoteKeySet{keys: keys, listed: true}) |
|
+ glog.V(0).Infof("remote key index: %d objects under %s", len(keys), mountDir) |
|
+ if err := f.saveRemoteKeySnapshot(mountDir, keys); err != nil { |
|
+ glog.Warningf("remote key index: save snapshot for %s: %v", mountDir, err) |
|
+ } |
|
+ return nil |
|
+} |
|
+ |
|
+func (f *Filer) remoteKeySnapshotPath(mountDir util.FullPath) string { |
|
+ return filepath.Join(f.RemoteKeyIndexDir, url.PathEscape(string(mountDir))) |
|
+} |
|
+ |
|
+func (f *Filer) loadRemoteKeySnapshot(mountDir util.FullPath) (map[util.FullPath]struct{}, error) { |
|
+ if f.RemoteKeyIndexDir == "" { |
|
+ return nil, nil |
|
+ } |
|
+ file, err := os.Open(f.remoteKeySnapshotPath(mountDir)) |
|
+ if os.IsNotExist(err) { |
|
+ return nil, nil |
|
+ } |
|
+ if err != nil { |
|
+ return nil, err |
|
+ } |
|
+ defer file.Close() |
|
+ keys := make(map[util.FullPath]struct{}) |
|
+ scanner := bufio.NewScanner(file) |
|
+ for scanner.Scan() { |
|
+ keys[util.FullPath(scanner.Text())] = struct{}{} |
|
+ } |
|
+ return keys, scanner.Err() |
|
+} |
|
+ |
|
+// saveRemoteKeySnapshot renames a finished file into place so a crash never |
|
+// leaves a truncated snapshot that would hide existing objects. |
|
+func (f *Filer) saveRemoteKeySnapshot(mountDir util.FullPath, keys map[util.FullPath]struct{}) error { |
|
+ if f.RemoteKeyIndexDir == "" { |
|
+ return nil |
|
+ } |
|
+ if err := os.MkdirAll(f.RemoteKeyIndexDir, 0o755); err != nil { |
|
+ return err |
|
+ } |
|
+ tmp, err := os.CreateTemp(f.RemoteKeyIndexDir, ".snapshot.") |
|
+ if err != nil { |
|
+ return err |
|
+ } |
|
+ defer os.Remove(tmp.Name()) |
|
+ w := bufio.NewWriter(tmp) |
|
+ for k := range keys { |
|
+ w.WriteString(string(k)) |
|
+ w.WriteByte('\n') |
|
+ } |
|
+ if err := w.Flush(); err != nil { |
|
+ tmp.Close() |
|
+ return err |
|
+ } |
|
+ if err := tmp.Close(); err != nil { |
|
+ return err |
|
+ } |
|
+ return os.Rename(tmp.Name(), f.remoteKeySnapshotPath(mountDir)) |
|
+} |
|
diff --git a/weed/server/filer_server.go b/weed/server/filer_server.go |
|
index 7d998ad..2f2107e 100644 |
|
--- a/weed/server/filer_server.go |
|
+++ b/weed/server/filer_server.go |
|
@@ -228,6 +228,8 @@ func NewFilerServer(defaultMux, readonlyMux *http.ServeMux, option *FilerOption) |
|
glog.V(0).Infof("max_file_name_length %d", maxFilenameLength) |
|
fs.filer = filer.NewFiler(*option.Masters, fs.grpcDialOption, option.Host, option.FilerGroup, option.Collection, option.DefaultReplication, option.DataCenter, maxFilenameLength, nil) |
|
fs.filer.Cipher = option.Cipher |
|
+ fs.filer.EnableRemoteKeyIndex = v.GetBool("filer.options.remote_key_index") |
|
+ fs.filer.RemoteKeyIndexDir = v.GetString("filer.options.remote_key_index_dir") |
|
fs.filer.DefaultDiskType = option.DiskType |
|
fs.filer.BuildGuardedRemoteClient = BuildGuardedRemoteStorageClient |
|
fs.filer.AllowUntrustedRemoteEndpoints = option.AllowUntrustedRemoteEndpoints |
|
-- |
|
2.55.0 |
|
|