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());
+ }
}