From 16984b85693c471ea6e895dbd692cf91ef1996b8 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 5 Sep 2026 10:08:07 +0800 Subject: [PATCH 1/5] [common] Forward two-phase output stream creation to the resolved FileIO MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ResolvingFileIO overrode tryToWriteAtomic so atomic small writes reach the resolved FileIO, but not newTwoPhaseOutputStream: with resolving-file-io.enabled the interface default wrapped the resolver itself in a rename-based committer, bypassing the resolved OSS/S3 FileIO's native multipart two-phase commit — the exact bypass the class guards against elsewhere. Forward the call like the sibling methods. Adds a regression test asserting the delegate's stream is returned and no rename fallback runs (verified red on the original code). Assisted-by: GLM-5.3 --- .../org/apache/paimon/fs/ResolvingFileIO.java | 9 +++++++++ .../apache/paimon/fs/ResolvingFileIOTest.java | 18 ++++++++++++++++++ 2 files changed, 27 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 5568ba896cb3..ae06febd3b49 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 @@ -116,6 +116,15 @@ public boolean tryToWriteAtomic(Path path, String content) throws IOException { return wrap(() -> fileIO(path).tryToWriteAtomic(path, content)); } + @Override + public TwoPhaseOutputStream newTwoPhaseOutputStream(Path path, boolean overwrite) + throws IOException { + // Forward to the resolved FileIO so implementations with native multipart + // commits (object storage) keep them; the interface default would wrap this + // resolver in a rename-based committer instead. + return wrap(() -> fileIO(path).newTwoPhaseOutputStream(path, overwrite)); + } + @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 067c7da649aa..e5203833aa27 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 @@ -184,4 +184,22 @@ public void testTryToWriteAtomicReachesResolvedOverride() throws IOException { // the interface default would have written a temp file and renamed it instead verify(delegate, never()).rename(any(), any()); } + + @Test + public void testNewTwoPhaseOutputStreamReachesResolvedOverride() 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/data.parquet"); + TwoPhaseOutputStream mockStream = mock(TwoPhaseOutputStream.class); + when(delegate.newTwoPhaseOutputStream(target, false)).thenReturn(mockStream); + + assertEquals(mockStream, resolvingFileIO.newTwoPhaseOutputStream(target, false)); + verify(delegate).newTwoPhaseOutputStream(target, false); + // the interface default would have renamed a temp file on the resolver instead + verify(delegate, never()).rename(any(), any()); + } } From 2e23646b1d8211e2ba425cb60509185bd89973a5 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 5 Sep 2026 20:57:03 +0800 Subject: [PATCH 2/5] [common] Resolve a ResolvingFileIO before building the multipart upload store --- .../fs/BaseMultiPartUploadCommitter.java | 5 + .../org/apache/paimon/fs/ResolvingFileIO.java | 2 - .../fs/BaseMultiPartUploadCommitterTest.java | 97 +++++++++++++++++++ 3 files changed, 102 insertions(+), 2 deletions(-) create mode 100644 paimon-common/src/test/java/org/apache/paimon/fs/BaseMultiPartUploadCommitterTest.java diff --git a/paimon-common/src/main/java/org/apache/paimon/fs/BaseMultiPartUploadCommitter.java b/paimon-common/src/main/java/org/apache/paimon/fs/BaseMultiPartUploadCommitter.java index 5245dcc161d4..2734b4324868 100644 --- a/paimon-common/src/main/java/org/apache/paimon/fs/BaseMultiPartUploadCommitter.java +++ b/paimon-common/src/main/java/org/apache/paimon/fs/BaseMultiPartUploadCommitter.java @@ -110,6 +110,11 @@ private MultiPartUploadStore multiPartUploadStore(FileIO fileIO) throws IO RESTTokenFileIO restTokenFileIO = (RESTTokenFileIO) fileIO; fileIO = restTokenFileIO.fileIO(); } + if (fileIO instanceof ResolvingFileIO) { + // The upload was started on the FileIO this resolver resolved to, and + // multiPartUploadStore casts to that concrete type, so resolve again here. + fileIO = ((ResolvingFileIO) fileIO).fileIO(targetPath()); + } return multiPartUploadStore(fileIO, targetPath()); } } 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 ae06febd3b49..cc3aba497e65 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 @@ -18,7 +18,6 @@ package org.apache.paimon.fs; -import org.apache.paimon.annotation.VisibleForTesting; import org.apache.paimon.catalog.CatalogContext; import org.apache.paimon.data.BlobDescriptor; import org.apache.paimon.options.CatalogOptions; @@ -134,7 +133,6 @@ public String createBlobPresignedUrl( .createBlobPresignedUrl(tableRoot, descriptor, validity)); } - @VisibleForTesting public FileIO fileIO(Path path) throws IOException { CacheKey cacheKey = new CacheKey(path.toUri().getScheme(), path.toUri().getAuthority()); return fileIOMap.computeIfAbsent( diff --git a/paimon-common/src/test/java/org/apache/paimon/fs/BaseMultiPartUploadCommitterTest.java b/paimon-common/src/test/java/org/apache/paimon/fs/BaseMultiPartUploadCommitterTest.java new file mode 100644 index 000000000000..c23fa5896f23 --- /dev/null +++ b/paimon-common/src/test/java/org/apache/paimon/fs/BaseMultiPartUploadCommitterTest.java @@ -0,0 +1,97 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.fs; + +import org.apache.paimon.catalog.CatalogContext; +import org.apache.paimon.options.Options; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.io.IOException; +import java.util.Collections; + +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.when; + +/** Tests for {@link BaseMultiPartUploadCommitter}. */ +public class BaseMultiPartUploadCommitterTest { + + private static final Path TARGET = new Path("oss://bucket/table/data-0.parquet"); + + private FileIO resolved; + private ResolvingFileIO resolvingFileIO; + + @BeforeEach + public void setUp() throws IOException { + resolved = mock(FileIO.class); + FileIOLoader loader = mock(FileIOLoader.class); + when(loader.getScheme()).thenReturn("oss"); + when(loader.load(any())).thenReturn(resolved); + resolvingFileIO = new ResolvingFileIO(); + resolvingFileIO.configure(CatalogContext.create(new Options(), loader, null)); + } + + @Test + public void testCommitResolvesResolvingFileIO() throws IOException { + RecordingCommitter committer = new RecordingCommitter(); + committer.commit(resolvingFileIO); + // the subclasses cast this to their own concrete FileIO, so the resolver itself + // reaching them would be a ClassCastException at commit time + assertThat(committer.received).isSameAs(resolved); + } + + @Test + public void testDiscardStagingResolvesResolvingFileIO() throws IOException { + RecordingCommitter committer = new RecordingCommitter(); + committer.discardStaging(resolvingFileIO); + assertThat(committer.received).isSameAs(resolved); + } + + @Test + public void testConcreteFileIOIsPassedThroughUnchanged() throws IOException { + RecordingCommitter committer = new RecordingCommitter(); + committer.commit(resolved); + assertThat(committer.received).isSameAs(resolved); + } + + private static class RecordingCommitter extends BaseMultiPartUploadCommitter { + + private FileIO received; + + private RecordingCommitter() { + super( + "upload-id", + Collections.singletonList("part-1"), + "table/data-0.parquet", + 1L, + TARGET); + } + + @Override + @SuppressWarnings("unchecked") + protected MultiPartUploadStore multiPartUploadStore( + FileIO fileIO, Path targetPath) { + this.received = fileIO; + return mock(MultiPartUploadStore.class); + } + } +} From a6d5dfe4552d901812f24af23ae7fa609656f523 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Thu, 10 Sep 2026 15:55:10 +0800 Subject: [PATCH 3/5] test: add native S3 multipart two-phase commit IT through ResolvingFileIO --- .../s3/S3ResolvingTwoPhaseCommitITCase.java | 131 ++++++++++++++++++ 1 file changed, 131 insertions(+) create mode 100644 paimon-filesystems/paimon-s3-impl/src/test/java/org/apache/paimon/s3/S3ResolvingTwoPhaseCommitITCase.java diff --git a/paimon-filesystems/paimon-s3-impl/src/test/java/org/apache/paimon/s3/S3ResolvingTwoPhaseCommitITCase.java b/paimon-filesystems/paimon-s3-impl/src/test/java/org/apache/paimon/s3/S3ResolvingTwoPhaseCommitITCase.java new file mode 100644 index 000000000000..9cb921bc1732 --- /dev/null +++ b/paimon-filesystems/paimon-s3-impl/src/test/java/org/apache/paimon/s3/S3ResolvingTwoPhaseCommitITCase.java @@ -0,0 +1,131 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.s3; + +import org.apache.paimon.catalog.CatalogContext; +import org.apache.paimon.fs.FileIO; +import org.apache.paimon.fs.FileIOLoader; +import org.apache.paimon.fs.Path; +import org.apache.paimon.fs.RenamingTwoPhaseOutputStream; +import org.apache.paimon.fs.ResolvingFileIO; +import org.apache.paimon.fs.TwoPhaseOutputStream; +import org.apache.paimon.options.Options; +import org.apache.paimon.utils.InstantiationUtil; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.RegisterExtension; + +import java.nio.charset.StandardCharsets; +import java.util.UUID; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Integration test that a two-phase write routed through {@link ResolvingFileIO} uses S3's native + * multipart-upload commit end to end against a MinIO backend: the stream is the resolved override + * (not the rename fallback), its committer survives serialization, and + * commit/discard/discardStaging work when handed a fresh resolver that has to resolve the scheme + * before casting to {@link S3FileIO}. + */ +class S3ResolvingTwoPhaseCommitITCase { + + @RegisterExtension private static final MinioTestContainer MINIO = new MinioTestContainer(); + + // preferIO resolves the s3 scheme to a native S3FileIO without any ServiceLoader registration. + private static final FileIOLoader S3_LOADER = + new FileIOLoader() { + @Override + public String getScheme() { + return "s3"; + } + + @Override + public FileIO load(Path path) { + return new S3FileIO(); + } + }; + + private ResolvingFileIO newResolver() { + ResolvingFileIO resolver = new ResolvingFileIO(); + resolver.configure( + CatalogContext.create( + Options.fromMap(MINIO.getS3ConfigOptions()), S3_LOADER, null)); + return resolver; + } + + private Path target(String name) { + return new Path(MINIO.getS3UriForDefaultBucket() + "/two-phase/" + name); + } + + @Test + void nativeMultipartCommitThroughFreshResolver() throws Exception { + Path path = target(UUID.randomUUID() + ".data"); + String payload = "native-multipart-payload"; + + TwoPhaseOutputStream out = newResolver().newTwoPhaseOutputStream(path, true); + // The resolver must forward to S3's native stream, not fall back to a copy-and-rename one. + assertThat(out).isNotInstanceOf(RenamingTwoPhaseOutputStream.class); + out.write(payload.getBytes(StandardCharsets.UTF_8)); + TwoPhaseOutputStream.Committer committer = out.closeForCommit(); + + // The committer is handed across the commit boundary, so it has to serialize. + byte[] bytes = InstantiationUtil.serializeObject(committer); + TwoPhaseOutputStream.Committer restored = + InstantiationUtil.deserializeObject(bytes, getClass().getClassLoader()); + + // Not visible before commit; committing through a fresh resolver forces the resolve that + // precedes the (S3FileIO) cast in BaseMultiPartUploadCommitter. + assertThat(newResolver().exists(path)).isFalse(); + restored.commit(newResolver()); + + FileIO reader = newResolver(); + assertThat(reader.exists(path)).isTrue(); + assertThat(reader.readFileUtf8(path)).isEqualTo(payload); + reader.delete(path, false); + } + + @Test + void abortBeforeCompletionLeavesNoObject() throws Exception { + Path path = target(UUID.randomUUID() + ".data"); + + TwoPhaseOutputStream out = newResolver().newTwoPhaseOutputStream(path, true); + out.write("to-be-aborted".getBytes(StandardCharsets.UTF_8)); + TwoPhaseOutputStream.Committer committer = out.closeForCommit(); + + committer.discard(newResolver()); + assertThat(newResolver().exists(path)).isFalse(); + } + + @Test + void discardStagingAfterCommitPreservesObject() throws Exception { + Path path = target(UUID.randomUUID() + ".data"); + + TwoPhaseOutputStream out = newResolver().newTwoPhaseOutputStream(path, true); + out.write("committed".getBytes(StandardCharsets.UTF_8)); + TwoPhaseOutputStream.Committer committer = out.closeForCommit(); + + committer.commit(newResolver()); + assertThat(newResolver().exists(path)).isTrue(); + + // Aborting staged resources after a successful commit must never delete the object. + committer.discardStaging(newResolver()); + assertThat(newResolver().exists(path)).isTrue(); + newResolver().delete(path, false); + } +} From 45e290b5bfe0d34e840993cbd99db54943daaba8 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Fri, 11 Sep 2026 16:30:40 +0800 Subject: [PATCH 4/5] fix: supply a Hadoop configuration to the S3 IT resolver so ResolvingFileIO.configure succeeds --- .../apache/paimon/s3/S3ResolvingTwoPhaseCommitITCase.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/paimon-filesystems/paimon-s3-impl/src/test/java/org/apache/paimon/s3/S3ResolvingTwoPhaseCommitITCase.java b/paimon-filesystems/paimon-s3-impl/src/test/java/org/apache/paimon/s3/S3ResolvingTwoPhaseCommitITCase.java index 9cb921bc1732..0f35a6c0542f 100644 --- a/paimon-filesystems/paimon-s3-impl/src/test/java/org/apache/paimon/s3/S3ResolvingTwoPhaseCommitITCase.java +++ b/paimon-filesystems/paimon-s3-impl/src/test/java/org/apache/paimon/s3/S3ResolvingTwoPhaseCommitITCase.java @@ -28,6 +28,7 @@ import org.apache.paimon.options.Options; import org.apache.paimon.utils.InstantiationUtil; +import org.apache.hadoop.conf.Configuration; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.RegisterExtension; @@ -65,7 +66,10 @@ private ResolvingFileIO newResolver() { ResolvingFileIO resolver = new ResolvingFileIO(); resolver.configure( CatalogContext.create( - Options.fromMap(MINIO.getS3ConfigOptions()), S3_LOADER, null)); + Options.fromMap(MINIO.getS3ConfigOptions()), + new Configuration(), + S3_LOADER, + null)); return resolver; } From 6fed6b8712407543956384a34cdfdd768feed294 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Fri, 11 Sep 2026 21:18:33 +0800 Subject: [PATCH 5/5] fix: make native S3 two-phase multipart commit work end to end The new S3 IT is the first test to exercise the real S3 multipart network path, and exposed two bugs it depends on: - S3MultiPartUpload built its WriteOperationHelper by hand with a null callbacks arg (since #7187), so uploadPart and completeMultipartUpload dereferenced null against a real backend. Get the helper from S3AFileSystem.getWriteOperationHelper() so the callbacks are wired. - discardStaging rethrew when aborting an already-completed upload, but it runs after a commit that may have taken effect and must preserve the target; tolerate a not-found (FileNotFoundException / NoSuchUpload) upload. Co-Authored-By: Claude Code --- .../fs/BaseMultiPartUploadCommitter.java | 7 +++- .../apache/paimon/s3/S3MultiPartUpload.java | 32 ++++--------------- 2 files changed, 13 insertions(+), 26 deletions(-) diff --git a/paimon-common/src/main/java/org/apache/paimon/fs/BaseMultiPartUploadCommitter.java b/paimon-common/src/main/java/org/apache/paimon/fs/BaseMultiPartUploadCommitter.java index 2734b4324868..9d16620dd2d3 100644 --- a/paimon-common/src/main/java/org/apache/paimon/fs/BaseMultiPartUploadCommitter.java +++ b/paimon-common/src/main/java/org/apache/paimon/fs/BaseMultiPartUploadCommitter.java @@ -24,6 +24,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.io.FileNotFoundException; import java.io.IOException; import java.util.ArrayList; import java.util.List; @@ -80,8 +81,12 @@ public void discard(FileIO fileIO) throws IOException { @Override public void discardStaging(FileIO fileIO) throws IOException { try { - // Aborting an upload never deletes a possibly completed object. + // A completed or already-aborted upload has no staging left to release, so a + // not-found upload (e.g. S3 NoSuchUpload) is nothing to discard rather than a + // failure. Aborting never deletes a completed object, so the target stays intact. abortMultipartUpload(fileIO); + } catch (FileNotFoundException e) { + LOG.debug("Multipart upload {} already gone; nothing to discard.", uploadId); } catch (Exception e) { throw new IOException("Failed to discard multipart upload with ID: " + uploadId, e); } diff --git a/paimon-filesystems/paimon-s3-impl/src/main/java/org/apache/paimon/s3/S3MultiPartUpload.java b/paimon-filesystems/paimon-s3-impl/src/main/java/org/apache/paimon/s3/S3MultiPartUpload.java index c995dd088141..3398b54aa27d 100644 --- a/paimon-filesystems/paimon-s3-impl/src/main/java/org/apache/paimon/s3/S3MultiPartUpload.java +++ b/paimon-filesystems/paimon-s3-impl/src/main/java/org/apache/paimon/s3/S3MultiPartUpload.java @@ -26,9 +26,6 @@ import org.apache.hadoop.fs.s3a.S3AFileSystem; import org.apache.hadoop.fs.s3a.WriteOperationHelper; import org.apache.hadoop.fs.s3a.impl.PutObjectOptions; -import org.apache.hadoop.fs.s3a.statistics.S3AStatisticsContext; -import org.apache.hadoop.fs.store.audit.AuditSpan; -import org.apache.hadoop.fs.store.audit.AuditSpanSource; import software.amazon.awssdk.core.sync.RequestBody; import software.amazon.awssdk.services.s3.model.CompleteMultipartUploadResponse; import software.amazon.awssdk.services.s3.model.CompletedPart; @@ -48,18 +45,15 @@ public class S3MultiPartUpload implements MultiPartUploadStore { private final S3AFileSystem s3a; - private final InternalWriteOperationHelper s3accessHelper; + private final WriteOperationHelper s3accessHelper; public S3MultiPartUpload(S3AFileSystem s3a, Configuration conf) { - checkNotNull(s3a); - this.s3accessHelper = - new InternalWriteOperationHelper( - s3a, - checkNotNull(conf), - s3a.createStoreContext().getInstrumentation(), - s3a.getAuditSpanSource(), - s3a.getActiveAuditSpan()); - this.s3a = s3a; + this.s3a = checkNotNull(s3a); + // Take the helper from the file system instead of building it by hand: the hand-built + // one left WriteOperationHelperCallbacks null, so uploadPart and completeMultipartUpload + // dereferenced null against a real backend. getWriteOperationHelper wires the same audit + // span and statistics plus the callbacks the AWS SDK v2 path needs. + this.s3accessHelper = s3a.getWriteOperationHelper(); } @Override @@ -117,16 +111,4 @@ UploadPartRequest newUploadPartRequest( public void abortMultipartUpload(String destKey, String uploadId) throws IOException { s3accessHelper.abortMultipartUpload(destKey, uploadId, false, null); } - - private static final class InternalWriteOperationHelper extends WriteOperationHelper { - - InternalWriteOperationHelper( - S3AFileSystem owner, - Configuration conf, - S3AStatisticsContext statisticsContext, - AuditSpanSource auditSpanSource, - AuditSpan auditSpan) { - super(owner, conf, statisticsContext, auditSpanSource, auditSpan, null); - } - } }