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 5568ba896cb3..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; @@ -116,6 +115,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 { @@ -125,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); + } + } +} 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()); + } } 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..0f35a6c0542f --- /dev/null +++ b/paimon-filesystems/paimon-s3-impl/src/test/java/org/apache/paimon/s3/S3ResolvingTwoPhaseCommitITCase.java @@ -0,0 +1,135 @@ +/* + * 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.apache.hadoop.conf.Configuration; +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()), + new Configuration(), + 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); + } +}