mateczagany commented on code in PR #29132:
URL: https://github.com/apache/flink/pull/29132#discussion_r4103033255
##########
flink-filesystems/flink-s3-fs-native/src/main/java/org/apache/flink/fs/s3native/writer/NativeS3RecoverableFsDataOutputStream.java:
##########
@@ -233,11 +233,13 @@ public Committer closeForCommit() throws IOException {
new NativeS3Recoverable(
key, uploadId, new
ArrayList<>(completedParts), numBytesInParts);
} catch (IOException e) {
- // The commit failed after the multipart upload had been
created and parts may
- // already have been uploaded. Abort it so it does not leak as
an orphan upload.
+ // Only local resources are released. The upload is
deliberately left open: a
+ // previous persist() may have handed it out in a recoverable
that a completed
+ // checkpoint references, and aborting it would break recovery
from that
+ // checkpoint. See the class-level Javadoc.
closed = true;
try {
- tryAbortUploadAndReleaseResources();
+ releaseLocalResources();
Review Comment:
Thank you for the review, I've addressed your comments!
--
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]