This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 60b973f032 [fs] Delegate tryToWriteAtomic through FileIO wrappers 
(#9034)
60b973f032 is described below

commit 60b973f032b2eff2657e92561eeb87083e335b43
Author: Jiajia Li <[email protected]>
AuthorDate: Wed Aug 5 23:15:27 2026 +0800

    [fs] Delegate tryToWriteAtomic through FileIO wrappers (#9034)
---
 .../java/org/apache/paimon/fs/ResolvingFileIO.java |  7 +++++
 .../org/apache/paimon/fs/cache/CachingFileIO.java  |  6 ++++
 .../org/apache/paimon/rest/RESTTokenFileIO.java    |  7 +++++
 .../org/apache/paimon/fs/ResolvingFileIOTest.java  | 19 ++++++++++++
 .../apache/paimon/fs/cache/CachingFileIOTest.java  | 20 +++++++++++++
 .../apache/paimon/rest/RESTTokenFileIOTest.java    | 34 ++++++++++++++++++++++
 6 files changed, 93 insertions(+)

diff --git 
a/paimon-common/src/main/java/org/apache/paimon/fs/ResolvingFileIO.java 
b/paimon-common/src/main/java/org/apache/paimon/fs/ResolvingFileIO.java
index 93bc4a4f6f..5568ba896c 100644
--- a/paimon-common/src/main/java/org/apache/paimon/fs/ResolvingFileIO.java
+++ b/paimon-common/src/main/java/org/apache/paimon/fs/ResolvingFileIO.java
@@ -109,6 +109,13 @@ public class ResolvingFileIO implements FileIO {
         return wrap(() -> fileIO(src).rename(src, dst));
     }
 
+    @Override
+    public boolean tryToWriteAtomic(Path path, String content) throws 
IOException {
+        // the interface default (temp file + rename) would bypass the 
resolved FileIO's atomic
+        // override
+        return wrap(() -> fileIO(path).tryToWriteAtomic(path, content));
+    }
+
     @Override
     public String createBlobPresignedUrl(
             Path tableRoot, BlobDescriptor descriptor, Duration validity) 
throws IOException {
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/fs/cache/CachingFileIO.java 
b/paimon-common/src/main/java/org/apache/paimon/fs/cache/CachingFileIO.java
index 28d5276db4..65eeaa3ebf 100644
--- a/paimon-common/src/main/java/org/apache/paimon/fs/cache/CachingFileIO.java
+++ b/paimon-common/src/main/java/org/apache/paimon/fs/cache/CachingFileIO.java
@@ -180,6 +180,12 @@ public class CachingFileIO implements FileIO {
         return delegate.rename(src, dst);
     }
 
+    @Override
+    public boolean tryToWriteAtomic(Path path, String content) throws 
IOException {
+        // the interface default (temp file + rename) would bypass the 
delegate's atomic override
+        return delegate.tryToWriteAtomic(path, content);
+    }
+
     @Override
     public String createBlobPresignedUrl(
             Path tableRoot, BlobDescriptor descriptor, Duration validity) 
throws IOException {
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java 
b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java
index 757916cdde..73b5541d3c 100644
--- a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java
+++ b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java
@@ -152,6 +152,13 @@ public class RESTTokenFileIO implements FileIO {
         return fileIO().rename(src, dst);
     }
 
+    @Override
+    public boolean tryToWriteAtomic(Path path, String content) throws 
IOException {
+        // the interface default (temp file + rename) would bypass the inner 
FileIO's atomic
+        // override
+        return fileIO().tryToWriteAtomic(path, content);
+    }
+
     @Override
     public String createBlobPresignedUrl(
             Path tableRoot, BlobDescriptor descriptor, Duration validity) 
throws IOException {
diff --git 
a/paimon-common/src/test/java/org/apache/paimon/fs/ResolvingFileIOTest.java 
b/paimon-common/src/test/java/org/apache/paimon/fs/ResolvingFileIOTest.java
index c550b84df1..067c7da649 100644
--- a/paimon-common/src/test/java/org/apache/paimon/fs/ResolvingFileIOTest.java
+++ b/paimon-common/src/test/java/org/apache/paimon/fs/ResolvingFileIOTest.java
@@ -36,8 +36,10 @@ import java.util.concurrent.Future;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertInstanceOf;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
@@ -165,4 +167,21 @@ public class ResolvingFileIOTest {
                 resolvingFileIO.createBlobPresignedUrl(tableRoot, descriptor, 
validity));
         verify(delegate).createBlobPresignedUrl(tableRoot, descriptor, 
validity);
     }
+
+    @Test
+    public void testTryToWriteAtomicReachesResolvedOverride() throws 
IOException {
+        FileIO delegate = mock(FileIO.class);
+        FileIOLoader loader = mock(FileIOLoader.class);
+        when(loader.load(any())).thenReturn(delegate);
+        when(loader.getScheme()).thenReturn("oss");
+        resolvingFileIO.configure(CatalogContext.create(new Options(), loader, 
null));
+
+        Path target = new Path("oss://bucket/table/snapshot/LATEST");
+        when(delegate.tryToWriteAtomic(target, "content")).thenReturn(true);
+
+        assertTrue(resolvingFileIO.tryToWriteAtomic(target, "content"));
+        verify(delegate).tryToWriteAtomic(target, "content");
+        // the interface default would have written a temp file and renamed it 
instead
+        verify(delegate, never()).rename(any(), any());
+    }
 }
diff --git 
a/paimon-common/src/test/java/org/apache/paimon/fs/cache/CachingFileIOTest.java 
b/paimon-common/src/test/java/org/apache/paimon/fs/cache/CachingFileIOTest.java
index dec0d7b7d4..ae0bec51f0 100644
--- 
a/paimon-common/src/test/java/org/apache/paimon/fs/cache/CachingFileIOTest.java
+++ 
b/paimon-common/src/test/java/org/apache/paimon/fs/cache/CachingFileIOTest.java
@@ -49,7 +49,9 @@ import static 
org.apache.paimon.options.CatalogOptions.LOCAL_CACHE_ENABLED;
 import static org.apache.paimon.options.CatalogOptions.LOCAL_CACHE_MAX_SIZE;
 import static org.apache.paimon.options.CatalogOptions.LOCAL_CACHE_WHITELIST;
 import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
@@ -87,6 +89,24 @@ class CachingFileIOTest {
         verify(delegate).createBlobPresignedUrl(tableRoot, descriptor, 
validity);
     }
 
+    @Test
+    void testTryToWriteAtomicReachesDelegateOverride() throws IOException {
+        FileIO delegate = mock(FileIO.class);
+        CachingFileIO cachingIO =
+                newCachingFileIO(
+                        delegate,
+                        new LocalMemoryCacheManager(1024, 64),
+                        EnumSet.of(FileType.DATA),
+                        64);
+        Path target = new Path("oss://bucket/table/snapshot/LATEST");
+        when(delegate.tryToWriteAtomic(target, "content")).thenReturn(true);
+
+        assertThat(cachingIO.tryToWriteAtomic(target, "content")).isTrue();
+        verify(delegate).tryToWriteAtomic(target, "content");
+        // the interface default would have written a temp file and renamed it 
instead
+        verify(delegate, never()).rename(any(), any());
+    }
+
     private CachingFileIO newCachingFileIO(
             FileIO delegate, LocalCacheManager cache, EnumSet<FileType> 
whitelist, int blockSize) {
         return new CachingFileIO(delegate, cache, whitelist);
diff --git 
a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java 
b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java
index 2234d66335..42e1746700 100644
--- 
a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java
+++ 
b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java
@@ -32,11 +32,13 @@ import org.junit.jupiter.api.Test;
 import java.io.IOException;
 import java.time.Duration;
 import java.util.Collections;
+import java.util.UUID;
 
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
@@ -78,4 +80,36 @@ class RESTTokenFileIOTest {
                 .isInstanceOf(IOException.class)
                 .hasMessageContaining("bound table root");
     }
+
+    @Test
+    void testTryToWriteAtomicReachesInnerOverride() throws IOException {
+        Path tableRoot = new Path("oss://bucket/table");
+        FileIO delegate = mock(FileIO.class);
+        FileIOLoader loader = mock(FileIOLoader.class);
+        when(loader.load(any())).thenReturn(delegate);
+        when(loader.getScheme()).thenReturn("oss");
+        RESTApi api = mock(RESTApi.class);
+        Identifier identifier = Identifier.create("db", "table");
+        // a unique token, so the static token-keyed FileIO cache cannot serve 
another test's
+        // delegate
+        when(api.loadTableToken(identifier))
+                .thenReturn(
+                        new GetTableTokenResponse(
+                                Collections.singletonMap("token", 
UUID.randomUUID().toString()),
+                                Long.MAX_VALUE));
+        RESTTokenFileIO fileIO =
+                new RESTTokenFileIO(
+                        CatalogContext.create(new Options(), loader, null),
+                        api,
+                        identifier,
+                        tableRoot);
+
+        Path target = new Path("oss://bucket/table/snapshot/LATEST");
+        when(delegate.tryToWriteAtomic(target, "content")).thenReturn(true);
+
+        assertThat(fileIO.tryToWriteAtomic(target, "content")).isTrue();
+        verify(delegate).tryToWriteAtomic(target, "content");
+        // the interface default would have written a temp file and renamed it 
instead
+        verify(delegate, never()).rename(any(), any());
+    }
 }

Reply via email to