From baf1f8b01ed1ee5bd4ff856203d68c62134b138b Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 01:40:07 +0900 Subject: [PATCH 1/5] Keep the multipart upload in commit() when its abort fails When completing a multipart upload failed, _finish_multipart_upload() aborted it and logged an abort failure, and S3File.commit() then dropped the upload ID in every case. A later discard(), as a transaction calls after a failed commit, had nothing to abort, so the incomplete upload was left behind. commit() now asks the helper not to abort and aborts with discard() itself. discard() clears the upload only after a successful abort, so a failed or interrupted abort keeps it for a later discard() to retry. The abort failure is logged so that the original error still propagates. Fixes #945. Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 34 +++++++++++++------ tests/pyathena/filesystem/test_s3.py | 50 +++++++++++++++++++++++----- 2 files changed, 64 insertions(+), 20 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index 0ecc0d0d..e1c97026 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -2298,14 +2298,15 @@ def _finish_multipart_upload( upload_id: str, futures: list[Future[S3MultipartUploadPart]], request_kwargs: Mapping[str, Any] | None = None, + abort: bool = True, ) -> S3CompleteMultipartUpload: """Collect the uploaded parts and complete the multipart upload. When any part or the completion fails, or the wait for them is interrupted, the parts that have not started are cancelled, the running ones are waited for, and the multipart upload is aborted so - that no incomplete upload or part is left behind. The original error - is then re-raised. + that no incomplete upload or part is left behind, unless ``abort`` is + false. The original error is then re-raised. Args: bucket: S3 bucket name. @@ -2315,6 +2316,8 @@ def _finish_multipart_upload( request_kwargs: Parameters of the upload, such as ``RequestPayer`` or the SSE-C parameters; the completion and the abort receive those that they accept. + abort: Whether to abort the multipart upload on failure. A caller + that keeps the upload to abort it itself passes false. Returns: S3CompleteMultipartUpload of the completed upload. @@ -2332,6 +2335,8 @@ def _finish_multipart_upload( **self._get_operation_kwargs("complete_multipart_upload", request_kwargs), ) except BaseException: + if not abort: + raise # A part that is still uploading when the upload is aborted may # be stored after the abort, so wait for the parts that could not # be cancelled first. @@ -3869,8 +3874,9 @@ def commit(self) -> None: Creates an empty object if nothing was written, uploads the buffered data with PutObject if no multipart upload part was submitted, and otherwise completes the multipart upload, which is aborted if the - completion fails or is interrupted. Invalidates the cache of the path - afterwards. + completion fails or is interrupted. If the abort also fails, the + upload is kept so that :meth:`discard` can abort it. Invalidates the + cache of the path afterwards. Raises: FileExistsError: If an object was created at the path after the @@ -3904,14 +3910,20 @@ def commit(self) -> None: upload_id=cast(str, self.multipart_upload.upload_id), futures=self.multipart_upload_parts, request_kwargs=self.s3_additional_kwargs, + abort=False, ) - except Exception: - # The multipart upload has been aborted by the helper; - # prevent discard() from aborting it again. An interrupt may - # have stopped the helper before the abort, so the upload is - # kept for discard() then. - self.multipart_upload = None - self.multipart_upload_parts = [] + except BaseException: + # discard() keeps the upload if the abort fails or is + # interrupted, so that a later discard(), as a transaction + # calls after a failed commit(), retries the abort. An abort + # failure is logged so that it does not mask the original + # error. + try: + self.discard() + except Exception: + _logger.exception( + f"Failed to abort multipart upload to s3://{self.bucket}/{self.key}." + ) raise self.fs.invalidate_cache(self.path) diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 49ee05eb..2fe16f61 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -3106,6 +3106,26 @@ def test_finish_multipart_upload_aborts_on_failure(self, error): UploadId="uploadid", ) + def test_finish_multipart_upload_without_abort(self): + # A caller that aborts the upload itself, as S3File.commit() does, + # gets the original error with the parts and the upload left alone. + fs = self._make_fs() + fs._complete_multipart_upload = mock.MagicMock() + failed: Future[SimpleNamespace] = Future() + failed.set_exception(RuntimeError("upload failed")) + pending: Future[SimpleNamespace] = Future() + + with pytest.raises(RuntimeError, match="upload failed"): + fs._finish_multipart_upload( + bucket="bucket", + key="key", + upload_id="uploadid", + futures=[failed, pending], + abort=False, + ) + assert not pending.cancelled() + fs._call.assert_not_called() + def test_finish_multipart_upload_abort_failure_does_not_mask_the_original_error(self): fs = self._make_fs() fs._complete_multipart_upload = mock.MagicMock() @@ -5526,21 +5546,33 @@ def wait_parts(futures): waited.assert_called_once_with([running]) assert pending.cancelled() - @pytest.mark.parametrize(("error", "aborts"), [(RuntimeError, 0), (KeyboardInterrupt, 1)]) - def test_commit_failure_and_discard(self, error, aborts): - # GH-1014: an error from _finish_multipart_upload() follows its - # abort, so a later discard(), as a transaction calls after a failed - # commit(), does not abort the upload again. An interrupt may have - # stopped it before the abort, so the upload is kept for discard(). + @pytest.mark.parametrize("abort_fails", [False, True]) + @pytest.mark.parametrize("error", [RuntimeError, KeyboardInterrupt]) + def test_commit_failure_and_discard(self, error, abort_fails): + # GH-1014: a failed or interrupted completion is aborted by commit(), + # so a later discard(), as a transaction calls after a failed + # commit(), does not abort the upload again. GH-945: if the abort + # also fails, the upload is kept so that discard() aborts it. file = self._make_multipart_write_file(b"x" * 16, autocommit=False) file._upload_chunk(final=True) - file.fs._finish_multipart_upload.side_effect = error("failed") + file.fs._finish_multipart_upload.side_effect = functools.partial( + S3FileSystem._finish_multipart_upload, file.fs + ) + file.fs._complete_multipart_upload.side_effect = error("complete failed") + if abort_fails: + file.fs._call.side_effect = [PermissionError("abort failed"), None] - with pytest.raises(error): + # The abort failure is logged, and the original error propagates. + with pytest.raises(error, match="complete failed"): file.commit() + assert (file.multipart_upload is not None) is abort_fails file.discard() - assert file.fs._call.call_count == aborts + assert file.fs._call.call_args_list == [ + mock.call("abort_multipart_upload", Bucket="bucket", Key="key.txt", UploadId="uploadid") + ] * (2 if abort_fails else 1) + assert file.multipart_upload is None + assert file.multipart_upload_parts == [] def test_discard_on_event_loop_thread(self): # GH-976: the parts that have not started are cancelled and not From f39c21c57c20e88042dd3f1ace36de97a0755c59 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 01:42:24 +0900 Subject: [PATCH 2/5] Do not tie the retry comment to fsspec's transaction internals Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 7 +++---- tests/pyathena/filesystem/test_s3.py | 6 +++--- 2 files changed, 6 insertions(+), 7 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index e1c97026..ea44346c 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -3914,10 +3914,9 @@ def commit(self) -> None: ) except BaseException: # discard() keeps the upload if the abort fails or is - # interrupted, so that a later discard(), as a transaction - # calls after a failed commit(), retries the abort. An abort - # failure is logged so that it does not mask the original - # error. + # interrupted, so that a later discard(), such as a + # transaction rollback, retries the abort. An abort failure + # is logged so that it does not mask the original error. try: self.discard() except Exception: diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 2fe16f61..66881e7f 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -5550,9 +5550,9 @@ def wait_parts(futures): @pytest.mark.parametrize("error", [RuntimeError, KeyboardInterrupt]) def test_commit_failure_and_discard(self, error, abort_fails): # GH-1014: a failed or interrupted completion is aborted by commit(), - # so a later discard(), as a transaction calls after a failed - # commit(), does not abort the upload again. GH-945: if the abort - # also fails, the upload is kept so that discard() aborts it. + # so a later discard(), such as a transaction rollback, does not + # abort the upload again. GH-945: if the abort also fails, the upload + # is kept so that discard() retries the abort. file = self._make_multipart_write_file(b"x" * 16, autocommit=False) file._upload_chunk(final=True) file.fs._finish_multipart_upload.side_effect = functools.partial( From 0deeb8a4430e1bb4ef5d05d54fac6b171449d3bd Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 01:43:42 +0900 Subject: [PATCH 3/5] Log the upload ID when commit() fails to abort a multipart upload Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 6 ++++-- tests/pyathena/filesystem/test_s3.py | 5 ++++- 2 files changed, 8 insertions(+), 3 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index ea44346c..d3e181d6 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -3903,11 +3903,12 @@ def commit(self) -> None: if not self.multipart_upload: raise RuntimeError("Multipart upload is not initialized.") + upload_id = cast(str, self.multipart_upload.upload_id) try: self.fs._finish_multipart_upload( bucket=self.bucket, key=self.key, - upload_id=cast(str, self.multipart_upload.upload_id), + upload_id=upload_id, futures=self.multipart_upload_parts, request_kwargs=self.s3_additional_kwargs, abort=False, @@ -3921,7 +3922,8 @@ def commit(self) -> None: self.discard() except Exception: _logger.exception( - f"Failed to abort multipart upload to s3://{self.bucket}/{self.key}." + f"Failed to abort multipart upload {upload_id} " + f"to s3://{self.bucket}/{self.key}." ) raise diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 66881e7f..571f1e31 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -5548,7 +5548,7 @@ def wait_parts(futures): @pytest.mark.parametrize("abort_fails", [False, True]) @pytest.mark.parametrize("error", [RuntimeError, KeyboardInterrupt]) - def test_commit_failure_and_discard(self, error, abort_fails): + def test_commit_failure_and_discard(self, caplog, error, abort_fails): # GH-1014: a failed or interrupted completion is aborted by commit(), # so a later discard(), such as a transaction rollback, does not # abort the upload again. GH-945: if the abort also fails, the upload @@ -5566,6 +5566,9 @@ def test_commit_failure_and_discard(self, error, abort_fails): with pytest.raises(error, match="complete failed"): file.commit() assert (file.multipart_upload is not None) is abort_fails + assert ( + "Failed to abort multipart upload uploadid to s3://bucket/key.txt." in caplog.text + ) is abort_fails file.discard() assert file.fs._call.call_args_list == [ From ebf798f6dd06f394a84c1cf054a0a7561dd9ef41 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 01:54:32 +0900 Subject: [PATCH 4/5] Test the kept part futures and an interrupted abort in commit() Co-Authored-By: Claude Opus 5.5 --- tests/pyathena/filesystem/test_s3.py | 23 +++++++++++++++++++++++ 1 file changed, 23 insertions(+) diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 571f1e31..7ae6efcc 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -5566,6 +5566,7 @@ def test_commit_failure_and_discard(self, caplog, error, abort_fails): with pytest.raises(error, match="complete failed"): file.commit() assert (file.multipart_upload is not None) is abort_fails + assert bool(file.multipart_upload_parts) is abort_fails assert ( "Failed to abort multipart upload uploadid to s3://bucket/key.txt." in caplog.text ) is abort_fails @@ -5577,6 +5578,28 @@ def test_commit_failure_and_discard(self, caplog, error, abort_fails): assert file.multipart_upload is None assert file.multipart_upload_parts == [] + def test_commit_failure_and_interrupted_abort(self): + # GH-945: if the abort after a failed completion is interrupted, the + # interrupt propagates and the upload is kept so that discard() + # retries the abort. + file = self._make_multipart_write_file(b"x" * 16, autocommit=False) + file._upload_chunk(final=True) + file.fs._finish_multipart_upload.side_effect = functools.partial( + S3FileSystem._finish_multipart_upload, file.fs + ) + file.fs._complete_multipart_upload.side_effect = RuntimeError("complete failed") + file.fs._call.side_effect = [KeyboardInterrupt, None] + + with pytest.raises(KeyboardInterrupt): + file.commit() + assert file.multipart_upload is not None + assert file.multipart_upload_parts + file.discard() + + assert file.fs._call.call_count == 2 + assert file.multipart_upload is None + assert file.multipart_upload_parts == [] + def test_discard_on_event_loop_thread(self): # GH-976: the parts that have not started are cancelled and not # waited for, so a rollback on the thread of the event loop that From 56bfa0b4f42b71476008cc81fc13f80b8f7ae07d Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 09:25:01 +0900 Subject: [PATCH 5/5] Keep the multipart upload when a failed write cannot abort it _close_without_commit() cleared the upload ID after a failed or interrupted abort, so that a deferred commit() could not complete the upload of a failed write. A later discard(), such as a transaction rollback, then had nothing to abort, and the incomplete upload was left behind. It now keeps the upload, as commit() does, and logs the upload ID. commit() returns early for a file whose written data was dropped, so the kept upload is never completed. Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 49 ++++++++++++++++------------ tests/pyathena/filesystem/test_s3.py | 23 +++++++++++-- 2 files changed, 49 insertions(+), 23 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index d3e181d6..1cfd91c8 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -3691,20 +3691,23 @@ def _close_without_commit(self) -> None: Drops the buffered data, so that neither close() nor a deferred commit() uploads it, and aborts the multipart upload, if any. An abort failure is logged instead of raised, so it does not mask the - error that the caller is handling. Even if the abort fails or is - interrupted, commit() does not complete the upload afterwards. The - executor is shut down here, as fsspec does not close a closed file - again when it is garbage collected. + error that the caller is handling. If the abort fails or is + interrupted, the upload is kept so that :meth:`discard` can abort it, + and commit() does not complete it. The executor is shut down here, as + fsspec does not close a closed file again when it is garbage + collected. """ self.buffer = None self.closed = True try: self.discard() except Exception: - _logger.exception(f"Failed to abort multipart upload to s3://{self.bucket}/{self.key}.") + # discard() keeps the upload when the abort fails. + upload_id = cast(S3MultipartUpload, self.multipart_upload).upload_id + _logger.exception( + f"Failed to abort multipart upload {upload_id} to s3://{self.bucket}/{self.key}." + ) finally: - self.multipart_upload = None - self.multipart_upload_parts = [] self._executor.shutdown() def _write_and_close(self, value: bytes | bytearray | memoryview) -> None: @@ -3876,7 +3879,8 @@ def commit(self) -> None: otherwise completes the multipart upload, which is aborted if the completion fails or is interrupted. If the abort also fails, the upload is kept so that :meth:`discard` can abort it. Invalidates the - cache of the path afterwards. + cache of the path afterwards. Does nothing for a file whose failed + write dropped the written data. Raises: FileExistsError: If an object was created at the path after the @@ -3884,21 +3888,24 @@ def commit(self) -> None: RuntimeError: If parts were submitted but no multipart upload is initialized. """ + if self.buffer is None: + # _close_without_commit() dropped the written data. A multipart + # upload that it failed to abort is kept for discard(), not + # completed. + return if self.tell() == 0: - if self.buffer is not None: - self.discard() - self.fs.touch(self.path, **self._get_request_kwargs("put_object")) + self.discard() + self.fs.touch(self.path, **self._get_request_kwargs("put_object")) elif not self.multipart_upload_parts: - if self.buffer is not None: - # Upload files smaller than block size. - self.buffer.seek(0) - data = self.buffer.read() - self.fs._put_object( - bucket=self.bucket, - key=self.key, - body=data, - **self._get_request_kwargs("put_object"), - ) + # Upload files smaller than block size. + self.buffer.seek(0) + data = self.buffer.read() + self.fs._put_object( + bucket=self.bucket, + key=self.key, + body=data, + **self._get_request_kwargs("put_object"), + ) else: if not self.multipart_upload: raise RuntimeError("Multipart upload is not initialized.") diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 7ae6efcc..d1edefa0 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -5267,10 +5267,11 @@ def test_write_exceeding_max_parts(self, existing, mode, writes, autocommit): fs._put_object.assert_not_called() @pytest.mark.parametrize("autocommit", [True, False]) - def test_write_exceeding_max_parts_abort_failure(self, autocommit): + def test_write_exceeding_max_parts_abort_failure(self, caplog, autocommit): # An abort failure is logged; the part limit error propagates, and # neither closing the file nor committing a deferred write retries - # the upload or completes it. + # the upload or completes it. GH-945: the upload is kept, so that + # discard() retries the abort. fs = self._make_append_fs(b"") fs.MULTIPART_UPLOAD_MAX_PARTS = 3 fs._call.side_effect = PermissionError("abort failed") @@ -5295,10 +5296,20 @@ def test_write_exceeding_max_parts_abort_failure(self, autocommit): fs._call.assert_called_once() fs._finish_multipart_upload.assert_not_called() fs._put_object.assert_not_called() + assert "Failed to abort multipart upload uploadid to s3://bucket/key.txt." in caplog.text + + assert f.multipart_upload is not None + fs._call.side_effect = None + f.discard() + assert fs._call.call_count == 2 + assert fs._call.call_args_list[1] == fs._call.call_args_list[0] + assert f.multipart_upload is None + assert f.multipart_upload_parts == [] def test_write_exceeding_max_parts_abort_interrupted(self): # GH-997: an interrupted abort propagates, and a deferred commit # still does not complete the upload; the executor is shut down. + # GH-945: the upload is kept, so that discard() retries the abort. fs = self._make_append_fs(b"") fs.MULTIPART_UPLOAD_MAX_PARTS = 3 fs._call.side_effect = KeyboardInterrupt @@ -5322,6 +5333,14 @@ def test_write_exceeding_max_parts_abort_interrupted(self): fs._finish_multipart_upload.assert_not_called() fs._put_object.assert_not_called() + assert f.multipart_upload is not None + fs._call.side_effect = None + f.discard() + assert fs._call.call_count == 2 + assert fs._call.call_args_list[1] == fs._call.call_args_list[0] + assert f.multipart_upload is None + assert f.multipart_upload_parts == [] + def test_write_exceeding_max_parts_without_close(self): # The executor of the closed file is shut down, as fsspec does not # close it again when it is garbage collected.