Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -110,6 +115,11 @@ private MultiPartUploadStore<T, C> 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());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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 {
Expand All @@ -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(
Expand Down
Original file line number Diff line number Diff line change
@@ -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<String, String> {

private FileIO received;

private RecordingCommitter() {
super(
"upload-id",
Collections.singletonList("part-1"),
"table/data-0.parquet",
1L,
TARGET);
}

@Override
@SuppressWarnings("unchecked")
protected MultiPartUploadStore<String, String> multiPartUploadStore(
FileIO fileIO, Path targetPath) {
this.received = fileIO;
return mock(MultiPartUploadStore.class);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -48,18 +45,15 @@ public class S3MultiPartUpload
implements MultiPartUploadStore<CompletedPart, CompleteMultipartUploadResponse> {

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
Expand Down Expand Up @@ -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);
}
}
}
Original file line number Diff line number Diff line change
@@ -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);
}
}
Loading