Skip to content
Open
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 @@ -110,6 +110,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());
}
}
Loading