-
Notifications
You must be signed in to change notification settings - Fork 116
Move the write-and-close helpers from S3FileSystem to S3File #1035
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’ll 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 |
|---|---|---|
|
|
@@ -1879,49 +1879,6 @@ def _check_multipart_upload_size(self, path: str, size: int, block_size: int) -> | |
| f"{min_block_size} bytes." | ||
| ) | ||
|
|
||
| @staticmethod | ||
| def _write_and_close(f: S3File, value: bytes | bytearray | memoryview) -> None: | ||
| """Write the whole value to a file opened for writing and close it. | ||
|
|
||
| Unlike a ``with`` block, a failed write closes the file without | ||
| committing it, so the existing object is left unchanged. | ||
|
|
||
| Args: | ||
| f: The file to write to. | ||
| value: The bytes to write. | ||
| """ | ||
| try: | ||
| if isinstance(value, memoryview) and not value.c_contiguous: | ||
| # The buffer of the file cannot write a non-contiguous memoryview. | ||
| value = value.tobytes() | ||
| f.write(value) | ||
| except BaseException: | ||
| f._close_without_commit() | ||
| 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: | ||
|
|
@@ -1975,9 +1932,7 @@ def pipe_file( | |
| # Defer to the buffered open() path, which keeps the | ||
| # deferred-commit semantics of fsspec transactions and uploads | ||
| # large data as a parallel multipart upload. | ||
| self._write_and_close( | ||
| self.open(path, "xb" if mode == "create" else "wb", **kwargs), value | ||
| ) | ||
| self.open(path, "xb" if mode == "create" else "wb", **kwargs)._write_and_close(value) | ||
| return | ||
| bucket, key, version_id = self.parse_path(path) | ||
| if version_id: | ||
|
|
@@ -2210,17 +2165,13 @@ def put_file( | |
| # The local file is opened first, so that an unreadable one fails | ||
| # before the remote file is opened. | ||
| with open(lpath, "rb") as local: | ||
| 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.open( | ||
| rpath, | ||
| "xb" if mode == "create" else "wb", | ||
| block_size=block_size, | ||
| max_workers=max_workers, | ||
| s3_additional_kwargs=s3_additional_kwargs, | ||
| )._write_file_and_close(local, callback) | ||
|
|
||
| self.invalidate_cache(rpath) | ||
|
|
||
|
|
@@ -3403,6 +3354,45 @@ def _close_without_commit(self) -> None: | |
| self.multipart_upload_parts = [] | ||
| self._executor.shutdown() | ||
|
|
||
| def _write_and_close(self, value: bytes | bytearray | memoryview) -> None: | ||
|
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 1 (behavior and implementation): CLEAN Base Covered: the moved
|
||
| """Write the whole value and close the file. | ||
|
|
||
| Unlike a ``with`` block, a failed write closes the file without | ||
| committing it, so the existing object is left unchanged. | ||
|
|
||
| Args: | ||
| value: The bytes to write. | ||
| """ | ||
| try: | ||
| if isinstance(value, memoryview) and not value.c_contiguous: | ||
| # The buffer of the file cannot write a non-contiguous memoryview. | ||
| value = value.tobytes() | ||
| self.write(value) | ||
| except BaseException: | ||
| self._close_without_commit() | ||
| raise | ||
| self.close() | ||
|
|
||
| def _write_file_and_close(self, local: BinaryIO, callback: Callback) -> None: | ||
| """Write the rest of a local file and close the file. | ||
|
|
||
| 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: | ||
| local: The local file to read from. | ||
| callback: Progress callback, updated with the size of each block. | ||
| """ | ||
| try: | ||
| while data := local.read(self.blocksize): | ||
| self.write(data) | ||
| callback.relative_update(len(data)) | ||
| except BaseException: | ||
| self._close_without_commit() | ||
| raise | ||
| self.close() | ||
|
|
||
| def _initiate_upload(self) -> None: | ||
| if not self.append_block and self.tell() < self.blocksize: | ||
| # Files smaller than block size in size cannot be multipart uploaded. | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -205,9 +205,7 @@ def _pipe_file_in_transaction( | |
| block_size = kwargs.get("block_size") or self._sync_fs.default_block_size | ||
| # The size in bytes; the length of a memoryview counts its items. | ||
| self._sync_fs._check_multipart_upload_size(path, memoryview(value).nbytes, block_size) | ||
| self._sync_fs._write_and_close( | ||
| self.open(path, "xb" if mode == "create" else "wb", **kwargs), value | ||
| ) | ||
| self.open(path, "xb" if mode == "create" else "wb", **kwargs)._write_and_close(value) | ||
|
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: "No behavior change" (the bodies match the removed static methods line for line apart from Finding (PR text, repaired): the TEST section said that all the listed regression tests patch |
||
|
|
||
| async def _put_file( | ||
| self, | ||
|
|
@@ -273,17 +271,13 @@ def _put_file_in_transaction( | |
|
|
||
| # See S3FileSystem.put_file. | ||
| with open(lpath, "rb") as local: | ||
| 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.open( | ||
| rpath, | ||
| "xb" if mode == "create" else "wb", | ||
| block_size=block_size, | ||
| max_workers=max_workers, | ||
| s3_additional_kwargs=s3_additional_kwargs, | ||
| )._write_file_and_close(local, callback) | ||
| self.invalidate_cache(rpath) | ||
|
|
||
| async def _get_file(self, rpath: str, lpath: str, callback=_DEFAULT_CALLBACK, **kwargs) -> None: | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.
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.
Independent review (relayed): FINDINGS
Reviewer: Codex CLI 0.160.0 (
codex exec --sandbox read-only, model reportedgpt-6-astra), session01a1022e-eefe-7122-9d90-c62474c23aaf. Static review only. Basefe21250a3e228e3d38e61898334c5baacc1285c2, head8011487869ccf0cde9a48ce9b129b59bfce2e74c, detached snapshot (unchanged afterwards). The prompt held the diff and conventions, without the PR number, PR text, or prior findings.Covered: sync and aio
pipe_file/put_filein and outside transactions incl.mode="create";open/_open, compression wrappers, deferred commit, discard, exception cleanup, executor shutdown;AbstractBufferedFileattributes andAioS3Fileinheritance (no name collisions or shadowing); changed tests; docstrings and conventions.Author verification (offline, mocked requests):
pipe_file(..., compression="gzip")on the buffered path (in a transaction, or larger than the block size) uploads a valid gzip body on masterfe21250aand before Write pipe_file() data without committing a failed write #1003 (5a60dd19); on80114878it raisesAttributeError: 'GzipFile' object has no attribute '_write_and_close'.put_file()does not passcompressiontoopen()and is unaffected.AttributeError: 'GzipFile' object has no attribute '_close_without_commit'(masking the original error) and still uploads a 10-byte object holding only the gzip header, so Write pipe_file() data without committing a failed write #1003's guarantee does not hold with compression. The single-request path (small values outside a transaction) passescompressionto PutObject, which botocore parameter validation rejects (ParamValidationError: Unknown parameter in input: "compression").The repair changes how
pipe_file()handlescompression, which is a design decision; this PR stays Draft until it is decided.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.
Resolved by #1038 (merged as
abc99b0984f9e19755580e5be494003b725e61cb, maintainer's choice "compress up front"):pipe_file()andAioS3FileSystem._pipe_file_in_transaction()now popcompression, compress the value withCompressedBuffer.compress(), and never passcompressiontoopen(), soopen()returns theS3File. #1038 also fixed the pre-existing failure path noted above.Rebased onto
abc99b09as880ba8a018353cec2fdcd025151ce06176056ff9(--force-with-leasefrom80114878); no conflicts, andgit range-diff fe21250a..80114878 abc99b09..880ba8a0shows the commit unchanged.Self-review after the rebase:
_write_and_close()/_write_file_and_close()now receives the file fromopen()in binary write mode withoutcompression(pipe_file()pops it;put_file()and_put_file_in_transaction()pass onlyblock_size,max_workers, ands3_additional_kwargs), so the object is always anS3File/AioS3File.880ba8a0:just lintpassed;tests/pyathena/filesystem/616 passed against live S3, including the 23 compression tests of Compress pipe_file() values up front, and route and key them consistently #1038 (compressed buffered writes in and outside transactions, multipart).An independent re-review of the rebased PR follows.
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.
Independent re-review of the rebased PR (relayed): CLEAN
Reviewer: Codex CLI 0.160.0 (
codex exec --sandbox read-only, model reportedgpt-6-astra), session01a102b6-5215-7f12-874c-824af00ceb4f. Static review only. Scopegit diff abc99b09..880ba8a0(full PR on the new base); snapshot at880ba8a018353cec2fdcd025151ce06176056ff9, unchanged afterwards. The prompt held the diff and conventions, without the PR number, PR text, or prior findings, and limited reading to the snapshot and the fsspec source.