Check the part limit in AioS3FileSystem transaction writes - #999
Conversation
#988 added _pipe_file_in_transaction() and _put_file_in_transaction() to AioS3FileSystem while #968 changed S3FileSystem.pipe_file() and put_file(), and the two were merged in parallel. Inside a transaction, pipe() and put() therefore did not reject data that needs more than 10,000 blocks before uploading, and put() sent block_size to the S3 API instead of open(). Apply the same check and the same block_size handling as the sync methods. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
| """ | ||
| 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) |
There was a problem hiding this comment.
Self-review round one (behavior and implementation) — base 92c9e3e67180eb3d52ef7b85579a24540a30fd64, head bab08235f0050aebdc4cf2d85a8fcba08b6dc59b. Result: CLEAN.
Covered: _pipe_file_in_transaction() and _put_file_in_transaction() against S3FileSystem.pipe_file()/put_file() (same check, same order: before the mode="create" existence check, before callback.set_size(), before open()); the block_size that pipe still passes to open() through kwargs; the error propagating through asyncio.to_thread and the generated sync wrappers; no file opened, so nothing joins transaction.files; non-transaction paths unchanged (they delegate to the sync methods). Tests: 3 new cases fail on the base and pass here; offline tests/pyathena/filesystem/ 262 passed vs 259 on the base, same 110 credential-less failures. Out of scope: the non-contiguous memoryview failure in the same with self.open(...) pattern (#997); max_workers routing (#969).
|
|
||
| callback.set_size(os.path.getsize(lpath)) | ||
| size = os.path.getsize(lpath) | ||
| block_size = kwargs.pop("block_size", None) or self._sync_fs.default_block_size |
There was a problem hiding this comment.
Self-review round two (claims, callers, operations) — base 92c9e3e67180eb3d52ef7b85579a24540a30fd64, head bab08235f0050aebdc4cf2d85a8fcba08b6dc59b. Result: CLEAN.
Claims checked: #988's branch head does not contain #968 (git merge-base --is-ancestor ae2d2a04 <#988 head> is false), so its transaction paths mirror the pre-#968 sync methods, as the PR body says; inside a transaction the data previously went through open() and failed at part 10,001 (S3File._upload_chunk guard), and put(..., block_size=...) placed block_size in s3_additional_kwargs. docs/filesystem.md ("put and pipe check the size before uploading anything") now holds for AioS3FileSystem transactions too, so no docs change. Callers: put(..., block_size=...) in an async transaction previously failed at the S3 API, so accepting it breaks nobody; no new API requests.
| """ | ||
| 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) |
There was a problem hiding this comment.
Independent review (relayed) — reviewer: Codex CLI 0.160.0, model gpt-6-astra, reasoning effort max, sandbox read-only, session 01a10115-dd19-7f23-90dc-25c82e2a475f. Base 92c9e3e67180eb3d52ef7b85579a24540a30fd64, head bab08235f0050aebdc4cf2d85a8fcba08b6dc59b, detached snapshot; given the diff, the intended behavior, and the fsspec source only. Static review only. Result: CLEAN.
Covered (reviewer): equivalence with the sync pipe_file()/put_file() (memoryview(value).nbytes, the default block size, the check before the mode="create" existence check and open(), block_size removed from the S3 parameters and passed to open(); directory/bucket early returns, MIME inference, progress, cache invalidation); errors through asyncio.to_thread, fsspec's wrapper generation, _runner, and sync, reaching async and generated sync callers; rejection before either file opens (fsspec registers transaction files only after _open() returns); docstrings and comments; the new tests fail against the base and cannot pass merely because open() is mocked.
Coverage limits noted by the reviewer (not defects): the forwarding test checks argument routing only (its 8-byte block works only with the mocked open()), and the tests omit multi-byte memoryviews, oversized data with mode="create", and an explicit block_size different from the default; the implementation handles these correctly. These paths share _check_multipart_upload_size(), which tests/pyathena/filesystem/test_s3.py covers for the sync methods, including the multi-byte memoryview.
WHAT
Apply #968's up-front part limit check and
put()block_sizehandling to the transaction paths ofAioS3FileSystem._pipe_file_in_transaction()callsS3FileSystem._check_multipart_upload_size()with the size in bytes (memoryview(value).nbytes) and the block size in use (block_sizeordefault_block_size), before themode="create"existence check and before opening the file, asS3FileSystem.pipe_file()does._put_file_in_transaction()popsblock_size, runs the same check on the local file size, and passesblock_sizetoopen(), asS3FileSystem.put_file()does. Before,block_sizewent intos3_additional_kwargsand reached the S3 API.Outside a transaction,
_pipe_file()/_put_file()delegate to the sync methods and already had this behavior.WHY
Follow-up to #953/#968, requested by the maintainer.
#988 added the two transaction paths (mirroring
S3FileSystem.pipe_file()/put_file()as they were before #968), and #968 changed the sync methods; the two PRs were merged in parallel, so inside anAioS3FileSystemtransactionpipe()/put()did not reject data needing more than 10,000 blocks before uploading (it failed at part 10,001 instead), andput(..., block_size=...)sentblock_sizeto S3.TEST
Tested commit: bab0823.
just formatandjust lint: pass.--noconftest):TestAioS3FileSystem::test_transaction_pipe_put_file_exceeding_max_parts(explicitblock_sizeanddefault_block_size): insidefs.transaction,pipe_file()andput_file()raiseValueErrorwithout callingopen()or any API.TestAioS3FileSystem::test_transaction_put_file_block_size:block_sizereachesopen(), nots3_additional_kwargs.tests/pyathena/filesystem/offline: 262 passed (259 on the base). The 110 failures are the same set as on the base: AWS integration tests stopped byNoCredentialsError, so no request was sent.🤖 Generated with Claude Code