-
Notifications
You must be signed in to change notification settings - Fork 116
Leave the existing object unchanged when put_file() fails #1017
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We鈥檒l occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -16,15 +16,15 @@ | |
| from io import BytesIO | ||
| from multiprocessing import cpu_count | ||
| from re import Pattern | ||
| from typing import Any, cast | ||
| from typing import Any, BinaryIO, cast | ||
| from urllib.parse import unquote_plus | ||
|
|
||
| import botocore.exceptions | ||
| from boto3 import Session | ||
| from botocore import UNSIGNED | ||
| from botocore.client import BaseClient, Config | ||
| from fsspec import AbstractFileSystem | ||
| from fsspec.callbacks import _DEFAULT_CALLBACK | ||
| from fsspec.callbacks import _DEFAULT_CALLBACK, Callback | ||
| from fsspec.implementations.local import trailing_sep | ||
| from fsspec.spec import AbstractBufferedFile | ||
| from fsspec.utils import isfilelike, other_paths, tokenize | ||
|
|
@@ -1716,6 +1716,28 @@ def _write_and_close(f: S3File, value: bytes | bytearray | memoryview) -> None: | |
| raise | ||
| f.close() | ||
|
|
||
| @staticmethod | ||
| def _write_file_and_close(f: S3File, local: BinaryIO, callback: Callback) -> None: | ||
| """Write the rest of a local file to a file opened for writing and close it. | ||
|
|
||
| Unlike a ``with`` block, a failed read, write, or progress update | ||
| closes the file without committing it, so the existing object is | ||
| left unchanged. | ||
|
|
||
| Args: | ||
| f: The file to write to. | ||
| local: The local file to read from. | ||
| callback: Progress callback, updated with the size of each block. | ||
| """ | ||
| try: | ||
| while data := local.read(f.blocksize): | ||
| f.write(data) | ||
| callback.relative_update(len(data)) | ||
| except BaseException: | ||
| f._close_without_commit() | ||
| raise | ||
| f.close() | ||
|
|
||
| def pipe_file( | ||
| self, path: str, value: bytes | bytearray | memoryview, mode: str = "overwrite", **kwargs | ||
| ) -> None: | ||
|
|
@@ -1795,10 +1817,11 @@ def _finish_multipart_upload( | |
| ) -> S3CompleteMultipartUpload: | ||
| """Collect the uploaded parts and complete the multipart upload. | ||
|
|
||
| When any part or the completion fails, 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. | ||
| 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. | ||
|
|
||
| Args: | ||
| bucket: S3 bucket name. | ||
|
|
@@ -1824,7 +1847,7 @@ def _finish_multipart_upload( | |
| parts=parts, | ||
| **self._get_operation_kwargs("complete_multipart_upload", request_kwargs), | ||
| ) | ||
| except Exception: | ||
| except BaseException: | ||
| # 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. | ||
|
|
@@ -1937,7 +1960,9 @@ def put_file( | |
|
|
||
| Uploads a file from the local filesystem to an S3 location. Supports | ||
| automatic content type detection based on file extension and provides | ||
| progress callback functionality. | ||
| progress callback functionality. An upload that fails before it is | ||
| completed, including one of a local file that cannot be read, leaves | ||
| the existing object unchanged. | ||
|
|
||
| Args: | ||
| lpath: Local file path to upload. | ||
|
|
@@ -1984,19 +2009,20 @@ def put_file( | |
| if content_type is not None: | ||
| s3_additional_kwargs["ContentType"] = content_type | ||
|
|
||
| with ( | ||
| self.open( | ||
| rpath, | ||
| "xb" if mode == "create" else "wb", | ||
| block_size=block_size, | ||
| max_workers=max_workers, | ||
| s3_additional_kwargs=s3_additional_kwargs, | ||
| ) as remote, | ||
| open(lpath, "rb") as local, | ||
| ): | ||
| while data := local.read(remote.blocksize): | ||
| remote.write(data) | ||
| callback.relative_update(len(data)) | ||
| # The local file is opened first, so that an unreadable one fails | ||
| # before the remote file is opened. | ||
| with open(lpath, "rb") as local: | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Self-review round 2 (claims, callers, operations): FINDINGS (PR text only, repaired) Base Claims checked:
Finding (PR text, repaired): the behavior-change note said that a failure before the first block was flushed replaced the object with an empty one. A small file whose callback fails after its write had its written data uploaded, not an empty object. The note now says "the data written so far, or an empty object when nothing was written". It also states that the Callers: the Evidence limits: live S3 run (494 passed) covers the success paths only; the failure and interrupt paths are offline tests with mocked requests. |
||
| self._write_file_and_close( | ||
| self.open( | ||
| rpath, | ||
| "xb" if mode == "create" else "wb", | ||
| block_size=block_size, | ||
| max_workers=max_workers, | ||
| s3_additional_kwargs=s3_additional_kwargs, | ||
| ), | ||
| local, | ||
| callback, | ||
| ) | ||
|
|
||
| self.invalidate_cache(rpath) | ||
|
|
||
|
|
@@ -3160,7 +3186,8 @@ 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. Invalidates the cache of the path afterwards. | ||
| completion fails or is interrupted. Invalidates the cache of the path | ||
| afterwards. | ||
|
|
||
| Raises: | ||
| FileExistsError: If an object was created at the path after the | ||
|
|
@@ -3197,7 +3224,9 @@ def commit(self) -> None: | |
| ) | ||
| except Exception: | ||
| # The multipart upload has been aborted by the helper; | ||
| # prevent discard() from aborting it again. | ||
| # 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 = [] | ||
| raise | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -264,19 +264,19 @@ def _put_file_in_transaction( | |
| if content_type is not None: | ||
| s3_additional_kwargs["ContentType"] = content_type | ||
|
|
||
| with ( | ||
| self.open( | ||
| rpath, | ||
| "xb" if mode == "create" else "wb", | ||
| block_size=block_size, | ||
| max_workers=max_workers, | ||
| s3_additional_kwargs=s3_additional_kwargs, | ||
| ) as remote, | ||
| open(lpath, "rb") as local, | ||
| ): | ||
| while data := local.read(remote.blocksize): | ||
| remote.write(data) | ||
| callback.relative_update(len(data)) | ||
| # See S3FileSystem.put_file. | ||
| with open(lpath, "rb") as local: | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent review (relayed): pre-existing limitation (reviewer and scope: see the comment on
Author: deferred as pre-existing and out of scope for #1014. Every
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. |
||
| self._sync_fs._write_file_and_close( | ||
| self.open( | ||
| rpath, | ||
| "xb" if mode == "create" else "wb", | ||
| block_size=block_size, | ||
| max_workers=max_workers, | ||
| s3_additional_kwargs=s3_additional_kwargs, | ||
| ), | ||
| local, | ||
| callback, | ||
| ) | ||
| self.invalidate_cache(rpath) | ||
|
|
||
| async def _get_file(self, rpath: str, lpath: str, callback=_DEFAULT_CALLBACK, **kwargs) -> None: | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Self-review round 1 (behavior and implementation): CLEAN
Base
16aef64d10dce92cbf8b8d957c178bef9bd494b6(merge-base with master), head5027eba79c4238e0541d7c48deef3de8df0dbf83.Covered:
_write_file_and_close(),put_file(),_finish_multipart_upload(),S3File.commit()(pyathena/filesystem/s3.py),AioS3FileSystem._put_file_in_transaction()(pyathena/filesystem/s3_async.py), and the changed tests.Checked:
_close_without_commit(); withautocommit=Falsethe file stays in the transaction, and its latercommit()seesbuffer=None, no parts, and no multipart upload, so it sends nothing. A part-limitValueErrorfrom_upload_chunk()reaches this handler after_upload_chunk()already called_close_without_commit(); the second call is a no-op apart from the executor shutdown, which is idempotent (same aspipe_file()'s_write_and_close()since Write pipe_file() data without committing a failed write聽#1003).mode="create":FileExistsErrorfromself.open()is raised insidewith open(lpath), which only closes the local file;test_put_file_create_existingand the aio transaction test still pass.f.blocksize, write, callback, thenclose()); the transaction still defers the commit._finish_multipart_upload()onBaseException: also covers the multipart copy caller (_copy_object_with_multipart_upload), where aborting on an interrupt is equally correct. After an interrupt it waits for the parts that could not be cancelled before the abort; a second interrupt escapes that wait, and interpreter exit already joins the executor threads.commit()widened toBaseExceptionso a laterdiscard()does not abort the already-aborted upload again.chmod(0)(skipped as root; CI runs on Linux).Limitations (pre-existing, not changed): an interrupt during
close()before_finish_multipart_upload()is entered (e.g. whileCreateMultipartUploador the final part submission runs) can still leave an upload behind; this window also exists forpipe_file()andopen()and is outside #1014.There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Rebased onto master
de8cc52a(#1006, #1012, #1013, #1015 merged since the review base) as49edeea3e292405f4c9aec69cfba2090d7c2bff0; pushed with--force-with-leasefromd59f880b.fsspecimport lines inpyathena/filesystem/s3.py; resolved to keep both (Callbackfrom this PR,trailing_sepfrom master).git range-diff 16aef64d..d59f880b de8cc52a..49edeea3: the first commit differs only in that import line; the second is identical.put_file(),_write_and_close(),_finish_multipart_upload(),S3File.commit()/discard()/_close_without_commit(), orAioS3FileSystem._put_file_in_transaction(). TheS3Filehunk changes only append mode, andcp_file()now returnsTrue.49edeea3:just lintpassed;tests/pyathena/filesystem/530 passed against live S3. AWS CI is rerunning on this head.