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 651120c443 [core] Fall back to renaming when multipart upload is
unavailable (#9610)
651120c443 is described below
commit 651120c44314829c1de381a2576b343e9cb27852
Author: Zouxxyy <[email protected]>
AuthorDate: Wed Sep 9 22:20:38 2026 +0800
[core] Fall back to renaming when multipart upload is unavailable (#9610)
---
.../java/org/apache/paimon/jindo/JindoFileIO.java | 10 ++++-
.../apache/paimon/jindo/JindoMultiPartUpload.java | 8 +++-
.../org/apache/paimon/jindo/JindoFileIOTest.java | 46 ++++++++++++++++++++++
3 files changed, 61 insertions(+), 3 deletions(-)
diff --git
a/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/JindoFileIO.java
b/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/JindoFileIO.java
index 58c9d2aac1..1e4646d4b9 100644
---
a/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/JindoFileIO.java
+++
b/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/JindoFileIO.java
@@ -35,6 +35,7 @@ import com.aliyun.jindodata.common.JindoHadoopSystem;
import com.aliyun.jindodata.dls.JindoDlsFileSystem;
import com.aliyun.jindodata.oss.JindoOssFileSystem;
import com.aliyun.jindodata.oss.auth.SimpleCredentialsProvider;
+import com.aliyun.jindodata.store.JindoMpuStore;
import com.aliyun.oss.OSSClient;
import com.aliyun.oss.OSSClientBuilder;
import org.apache.hadoop.conf.Configuration;
@@ -214,8 +215,15 @@ public class JindoFileIO extends HadoopCompliantFileIO
implements HadoopOptionsP
org.apache.hadoop.fs.Path hadoopPath = path(path);
Pair<JindoHadoopSystem, String> pair = getFileSystemPair(hadoopPath,
false);
JindoHadoopSystem fs = pair.getKey();
+ JindoMpuStore mpuStore = fs.getMpuStore(hadoopPath);
+ if (mpuStore == null) {
+ LOG.debug(
+ "Jindo multipart upload is unavailable for {}, falling
back to rename commit.",
+ path);
+ return super.newTwoPhaseOutputStream(path, overwrite);
+ }
return new JindoTwoPhaseOutputStream(
- new JindoMultiPartUpload(fs, hadoopPath), hadoopPath, path);
+ new JindoMultiPartUpload(mpuStore, fs.getWorkingDirectory()),
hadoopPath, path);
}
@Override
diff --git
a/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/JindoMultiPartUpload.java
b/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/JindoMultiPartUpload.java
index ff6ba49f75..00e88e531f 100644
---
a/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/JindoMultiPartUpload.java
+++
b/paimon-filesystems/paimon-jindo/src/main/java/org/apache/paimon/jindo/JindoMultiPartUpload.java
@@ -42,8 +42,12 @@ public class JindoMultiPartUpload
private final Path workingDirectory;
public JindoMultiPartUpload(JindoHadoopSystem fs, Path filePath) {
- this.workingDirectory = fs.getWorkingDirectory();
- this.mpuStore = fs.getMpuStore(filePath);
+ this(fs.getMpuStore(filePath), fs.getWorkingDirectory());
+ }
+
+ JindoMultiPartUpload(JindoMpuStore mpuStore, Path workingDirectory) {
+ this.workingDirectory = workingDirectory;
+ this.mpuStore = mpuStore;
}
@Override
diff --git
a/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/JindoFileIOTest.java
b/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/JindoFileIOTest.java
index 7978dfa2fb..91406badb4 100644
---
a/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/JindoFileIOTest.java
+++
b/paimon-filesystems/paimon-jindo/src/test/java/org/apache/paimon/jindo/JindoFileIOTest.java
@@ -21,11 +21,16 @@ package org.apache.paimon.jindo;
import org.apache.paimon.catalog.CatalogContext;
import org.apache.paimon.data.BlobDescriptor;
import org.apache.paimon.fs.Path;
+import org.apache.paimon.fs.RenamingTwoPhaseOutputStream;
+import org.apache.paimon.fs.TwoPhaseOutputStream;
import org.apache.paimon.options.Options;
+import org.apache.paimon.utils.Pair;
+import com.aliyun.jindodata.common.JindoHadoopSystem;
import com.aliyun.oss.HttpMethod;
import com.aliyun.oss.OSSClient;
import com.aliyun.oss.model.ObjectMetadata;
+import org.apache.hadoop.fs.FSDataOutputStream;
import org.junit.jupiter.api.Test;
import java.net.URI;
@@ -122,6 +127,47 @@ public class JindoFileIOTest {
verify(client).shutdown();
}
+ @Test
+ public void testFallbackToRenamingWhenMultipartUploadUnsupported() throws
Exception {
+ JindoHadoopSystem fs = mock(JindoHadoopSystem.class);
+ org.apache.hadoop.fs.Path hadoopPath =
mock(org.apache.hadoop.fs.Path.class);
+ when(fs.exists(any())).thenReturn(false);
+ when(fs.getMpuStore(any())).thenReturn(null);
+ when(fs.create(any(),
eq(false))).thenReturn(mock(FSDataOutputStream.class));
+
+ JindoFileIO fileIO = new TestingJindoFileIO(fs, hadoopPath);
+ TwoPhaseOutputStream stream =
+ fileIO.newTwoPhaseOutputStream(
+ new
Path("oss://bucket.cn-hangzhou.oss-dls.aliyuncs.com/table/file"),
+ false);
+
+ assertThat(stream).isInstanceOf(RenamingTwoPhaseOutputStream.class);
+ stream.close();
+ verify(fs).getMpuStore(hadoopPath);
+ }
+
+ private static class TestingJindoFileIO extends JindoFileIO {
+
+ private final JindoHadoopSystem fs;
+ private final org.apache.hadoop.fs.Path hadoopPath;
+
+ private TestingJindoFileIO(JindoHadoopSystem fs,
org.apache.hadoop.fs.Path hadoopPath) {
+ this.fs = fs;
+ this.hadoopPath = hadoopPath;
+ }
+
+ @Override
+ protected org.apache.hadoop.fs.Path path(Path path) {
+ return hadoopPath;
+ }
+
+ @Override
+ protected Pair<JindoHadoopSystem, String> getFileSystemPair(
+ org.apache.hadoop.fs.Path path, boolean enableCache) {
+ return Pair.of(fs, "dls");
+ }
+ }
+
private static String sha256Hex(byte[] bytes) throws Exception {
StringBuilder result = new StringBuilder();
for (byte value : MessageDigest.getInstance("SHA-256").digest(bytes)) {