Skip to content

Check the part limit in AioS3FileSystem transaction writes - #999

Merged
laughingman7743 merged 1 commit into
masterfrom
fix/aio-transaction-part-limit
Oct 3, 2026
Merged

laughingman7743 merged 1 commit into
masterfrom
fix/aio-transaction-part-limit

Conversation

@laughingman7743

Copy link
Copy Markdown
Member

WHAT

Apply #968's up-front part limit check and put() block_size handling to the transaction paths of AioS3FileSystem.

  • _pipe_file_in_transaction() calls S3FileSystem._check_multipart_upload_size() with the size in bytes (memoryview(value).nbytes) and the block size in use (block_size or default_block_size), before the mode="create" existence check and before opening the file, as S3FileSystem.pipe_file() does.
  • _put_file_in_transaction() pops block_size, runs the same check on the local file size, and passes block_size to open(), as S3FileSystem.put_file() does. Before, block_size went into s3_additional_kwargs and 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 an AioS3FileSystem transaction pipe()/put() did not reject data needing more than 10,000 blocks before uploading (it failed at part 10,001 instead), and put(..., block_size=...) sent block_size to S3.

TEST

Tested commit: bab0823.

  • just format and just lint: pass.
  • Offline unit tests, run with a dummy region and no credentials (--noconftest):
    • TestAioS3FileSystem::test_transaction_pipe_put_file_exceeding_max_parts (explicit block_size and default_block_size): inside fs.transaction, pipe_file() and put_file() raise ValueError without calling open() or any API.
    • TestAioS3FileSystem::test_transaction_put_file_block_size: block_size reaches open(), not s3_additional_kwargs.
    • All 3 cases fail on 92c9e3e and pass with this change.
  • 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 by NoCredentialsError, so no request was sent.
  • Not run locally: live AWS tests; the AWS CI runs once the PR is Ready.

🤖 Generated with Claude Code

#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)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@laughingman7743
laughingman7743 marked this pull request as ready for review October 3, 2026 09:31
@laughingman7743
laughingman7743 merged commit 4950996 into master Oct 3, 2026
12 checks passed
@laughingman7743
laughingman7743 deleted the fix/aio-transaction-part-limit branch October 3, 2026 09:51
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant