weijietong commented on code in PR #9606:
URL: https://github.com/apache/paimon/pull/9606#discussion_r3987445949


##########
paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java:
##########
@@ -268,6 +283,322 @@ public boolean tryToWriteAtomic(Path path, String 
content) throws IOException {
         }
     }
 
+    /**
+     * Rename across buckets, mirroring {@link AliyunOSSFileSystem#rename}.
+     *
+     * <p>The underlying {@code AliyunOSSFileSystem.rename} only supports copy 
within the same
+     * bucket, because {@code AliyunOSSFileSystemStore} pins a single {@code 
bucketName} and calls
+     * {@code copyObject(bucketName, srcKey, bucketName, dstKey)}. The {@link 
OSSClient} held by
+     * that store, however, accepts distinct source/destination buckets as 
long as they are in the
+     * same region. So for cross-bucket renames we bypass the Hadoop layer and 
drive {@code
+     * OSSClient} directly, but keep exactly the same rename semantics as 
{@code
+     * AliyunOSSFileSystem.rename}: root/path-containment checks, destination 
existence/type
+     * handling, directory rewrite to {@code dst/srcName}, empty-directory 
marker creation,
+     * file/directory dispatch and delete-after-copy. Copy itself is a single 
{@code copyObject}
+     * first, falling back to multipart {@code uploadPartCopy} when the object 
is larger than 1GB or
+     * a shallow copy is not supported. Cross-region copy is not supported by 
the OSS server-side
+     * copy API and will surface as an {@link OSSException}.
+     *
+     * @param src source path.
+     * @param dst destination path.
+     * @return true if the rename succeeds.
+     */
+    @Override
+    public boolean rename(Path src, Path dst) throws IOException {
+        URI srcUri = src.toUri();
+        URI dstUri = dst.toUri();
+        String srcBucket = srcUri.getHost();
+        String dstBucket = dstUri.getHost();
+        // Same bucket (or authority missing): fall back to the Hadoop rename 
path, which keeps
+        // the original semantics including directory handling and 
delete-after-copy.
+        if (srcBucket == null || dstBucket == null || 
srcBucket.equals(dstBucket)) {
+            return super.rename(src, dst);
+        }
+
+        OSSClient ossClient;
+        try {
+            ossClient = ossClient(src);
+        } catch (IOException e) {
+            throw e;
+        } catch (Exception e) {
+            throw new IOException("Failed to access OSSClient for cross-bucket 
rename", e);
+        }
+        return renameCrossBucket(ossClient, srcBucket, dstBucket, src, dst);
+    }
+
+    /**
+     * Cross-bucket rename whose control flow mirrors {@link 
AliyunOSSFileSystem#rename}. Existence
+     * and type checks go through this {@link OSSFileIO} (which selects the 
right per-bucket {@link
+     * AliyunOSSFileSystem}), while the actual copy is driven on {@code 
ossClient}.
+     */
+    private boolean renameCrossBucket(
+            OSSClient ossClient, String srcBucket, String dstBucket, Path src, 
Path dst)
+            throws IOException {
+        // Cannot rename the root of a filesystem.
+        if (src.getParent() == null) {
+            LOG.debug("Cannot rename the root of a filesystem");
+            return false;
+        }
+        // Reject renaming a directory into a subdirectory of itself.
+        Path parent = dst.getParent();
+        while (parent != null && !src.equals(parent)) {
+            parent = parent.getParent();
+        }
+        if (parent != null) {
+            return false;
+        }
+
+        FileStatus srcStatus = getFileStatus(src);
+        FileStatus dstStatus;
+        try {
+            dstStatus = getFileStatus(dst);
+        } catch (FileNotFoundException fnde) {
+            dstStatus = null;
+        }
+
+        if (dstStatus == null) {
+            // If dst doesn't exist, its parent must exist and be a directory.
+            FileStatus dstParentStatus = getFileStatus(dst.getParent());
+            if (!dstParentStatus.isDir()) {
+                throw new IOException(
+                        String.format(
+                                "Failed to rename %s to %s, %s is a file",
+                                src, dst, dst.getParent()));
+            }
+        } else {
+            if (srcStatus.getPath().equals(dstStatus.getPath())) {
+                return !srcStatus.isDir();
+            } else if (dstStatus.isDir()) {
+                // If dst is a directory, rewrite to dst/srcName.
+                dst = new Path(dst, src.getName());
+                FileStatus[] statuses;
+                try {
+                    statuses = listStatus(dst);
+                } catch (FileNotFoundException fnde) {
+                    statuses = null;
+                }
+                if (statuses != null && statuses.length > 0) {
+                    // If dst exists and not a directory / not empty.
+                    throw new FileAlreadyExistsException(
+                            String.format(
+                                    "Failed to rename %s to %s, file already 
exists or not empty!",
+                                    src, dst));
+                }
+            } else {
+                // If dst is not a directory.
+                throw new FileAlreadyExistsException(
+                        String.format("Failed to rename %s to %s, file already 
exists!", src, dst));
+            }
+        }
+
+        boolean succeed;
+        if (srcStatus.isDir()) {
+            succeed = copyDirectoryCrossBucket(ossClient, srcBucket, 
dstBucket, src, dst);
+        } else {
+            succeed =
+                    copyFileCrossBucket(
+                            ossClient,
+                            srcBucket,
+                            pathToObjectKey(src.toUri()),
+                            dstBucket,
+                            pathToObjectKey(dst.toUri()),
+                            srcStatus.getLen());
+        }
+        return src.equals(dst) || (succeed && delete(src, true));
+    }
+
+    /** OSS object keys are the URI path with the leading '/' stripped. */
+    private static String pathToObjectKey(URI uri) {
+        String path = uri.getPath();
+        if (path == null || path.isEmpty()) {
+            return "";
+        }
+        return path.startsWith("/") ? path.substring(1) : path;
+    }
+
+    /** Ensures {@code key} carries a trailing slash, treating it as an OSS 
directory prefix. */
+    private static String maybeAddTrailingSlash(String key) {
+        if (key.isEmpty()) {
+            return key;
+        }
+        return key.endsWith("/") ? key : key + "/";
+    }
+
+    /**
+     * Copy a single file across buckets, mirroring {@code 
AliyunOSSFileSystemStore.copyFile}:
+     * single {@code copyObject} first, fall back to multipart {@code 
uploadPartCopy} on failure
+     * (object larger than 1GB or shallow copy not supported).
+     */
+    private boolean copyFileCrossBucket(
+            OSSClient ossClient,
+            String srcBucket,
+            String srcKey,
+            String dstBucket,
+            String dstKey,
+            long contentLength)
+            throws IOException {
+        SseConfig sse = configuredSse();
+        try {
+            CopyObjectRequest request = new CopyObjectRequest(srcBucket, 
srcKey, dstBucket, dstKey);
+            if (sse != null) {
+                applySse(request, sse);
+            }
+            ossClient.copyObject(request);
+            return true;
+        } catch (OSSException e) {
+            LOG.debug(
+                    "Single cross-bucket copy failed for {} -> {}, fallback to 
multipartCopy: {}",
+                    srcKey,
+                    dstKey,
+                    e.getMessage());
+            return multipartCopyCrossBucket(
+                    ossClient, srcBucket, srcKey, dstBucket, dstKey, 
contentLength, sse);
+        }
+    }
+
+    /**
+     * Copy a single object across buckets via multipart upload-copy. Required 
when the object is
+     * larger than 1GB or when the OSS server does not support a shallow 
single-copy between the
+     * source and destination (e.g. differing storage classes). Preserves 
content-type and user
+     * metadata, which multipart copy does not copy by default.
+     */
+    private boolean multipartCopyCrossBucket(
+            OSSClient ossClient,
+            String srcBucket,
+            String srcKey,
+            String dstBucket,
+            String dstKey,
+            long contentLength,
+            SseConfig sse)
+            throws IOException {
+        ObjectMetadata srcMeta = ossClient.getObjectMetadata(srcBucket, 
srcKey);
+        ObjectMetadata newMeta = new ObjectMetadata();
+        if (srcMeta.getContentType() != null) {
+            newMeta.setContentType(srcMeta.getContentType());
+        }
+        if (srcMeta.getUserMetadata() != null && 
!srcMeta.getUserMetadata().isEmpty()) {
+            newMeta.setUserMetadata(srcMeta.getUserMetadata());
+        }
+        if (sse != null) {
+            newMeta = applySse(newMeta, sse);
+        }
+
+        InitiateMultipartUploadRequest initiateRequest =
+                new InitiateMultipartUploadRequest(dstBucket, dstKey);
+        initiateRequest.setObjectMetadata(newMeta);
+        InitiateMultipartUploadResult initiateResult =
+                ossClient.initiateMultipartUpload(initiateRequest);
+        String uploadId = initiateResult.getUploadId();
+
+        List<PartETag> partETags = new ArrayList<>();
+        try {
+            long partSize = CROSS_BUCKET_COPY_PART_SIZE;
+            long remaining = contentLength;

Review Comment:
   done



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to