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)) {

Reply via email to