Align AioS3FileSystem with S3FileSystem in transactions, open() and touch() - #988
Conversation
| return | ||
| await asyncio.to_thread(self._sync_fs.put_file, lpath, rpath, callback=callback, **kwargs) | ||
|
|
||
| def _put_file_in_transaction(self, lpath: str, rpath: str, callback, **kwargs) -> None: |
There was a problem hiding this comment.
Self-review round 1 (implementation behavior): CLEAN, no repairs.
Scope: git diff aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe..a060c55078a5064752c57582be4886942ddc1093: pyathena/filesystem/s3_async.py, pyathena/filesystem/s3_executor.py, tests/pyathena/filesystem/test_s3_async.py, docs/aio.md.
Checked:
- Transaction path vs
S3FileSystem.pipe_file→AbstractFileSystem.pipe_file/open(): samemode="create"check, path stripping,autocommit=False, registration intransaction.files;Transaction.completecommits/discards theS3File(its thread pool is shut down onclose()after the parts finish, as in the sync path). put_filehelper vsS3FileSystem.put_file(s3.py:1476): same directory/bucket-only early returns, size callback,ContentTypeguess,s3_additional_kwargs, cache invalidation.- Starvation: an
AioS3Filehere would wait in a default-executor thread on parts queued to the same executor; the internalS3Fileuses its own pool, so concurrentput()(fsspec batches up to 128) cannot deadlock. - fsspec sets
_loop = Noneexactly whenasynchronous=True(fsspec/asyn.py:445-448), so_create_executorpicks the thread pool only there.touchis not infsspec.asyn.async_methods, somirror_sync_methodskeeps the explicittouch(). - New offline tests fail on the base for the transaction,
touchand both parallel cases.
Recorded trade-off (not a defect): _put_file_in_transaction repeats ~15 lines of S3FileSystem.put_file (s3.py:1476-1541), because the sync method always writes through the internal filesystem's own open(). Avoiding that would need a hook in S3FileSystem; left as is unless the maintainer prefers that.
Out of scope, pre-existing: compression=/autocommit= kwargs to pipe_file behave differently between the open() and _open() paths; parameter routing is #969.
There was a problem hiding this comment.
Repair (985510d, found in round 2): the starvation premise above was false. The parts of an AioS3File are only waited on in commit(), which a transaction runs in the caller's thread, so the asyncio.to_thread workers never block on them. A mock run of 40 concurrent put() calls inside one transaction completed (120 parts) with both designs. The in-transaction writes now go through self.open() as the issue proposed, and _open_in_transaction (_sync_fs._open + manual registration) is gone. This also brings the compression=/autocommit= handling of fsspec's open() back in line. The "fsspec batches up to 128" note was also imprecise: the default is RLIMIT_NOFILE // 8 (fallback 128). test_s3_async.py: 76 passed live on the new head.
| # As ThreadPoolExecutor does; a semaphore of 0 would never run anything. | ||
| raise ValueError("max_workers must be greater than 0") | ||
| self._loop = loop | ||
| self._semaphore = asyncio.Semaphore(max_workers) |
There was a problem hiding this comment.
Round 1: asyncio.Semaphore binds to a loop only when an acquire() has to wait (3.11+), and every _run() runs on self._loop, so creating it in the thread that calls open() is safe. Cancelling the concurrent future also cancels a task waiting on the semaphore. Exceptions release it via async with. max_workers <= 0 raises like ThreadPoolExecutor instead of hanging on Semaphore(0).
| return | ||
| await asyncio.to_thread(self._sync_fs.pipe_file, path, value, mode=mode, **kwargs) | ||
|
|
||
| def _pipe_file_in_transaction( |
There was a problem hiding this comment.
Self-review round 2 (claims, callers, operations): FINDINGS, repaired.
Scope: full git diff aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe..985510d7f6f85514432bf902949478f71999275a plus PR body, commit messages, docstrings and docs/aio.md.
Claims checked:
- "writes through an
AioS3Filein a transaction could starve the default executor" (my design rationale in the first push): false. See the repair reply on the round-1 thread; a mock run of 40 concurrent transactionalput()calls completed. Repaired by writing throughopen(). - "bounded only by
min(32, os.cpu_count() + 4)": Python 3.13 usesprocess_cpu_count(); PR body narrowed to "the size of the event loop's default executor". asynchronous=True⇒_loop is None(fsspec/asyn.py:445-448): holds for every fsspec version with that constructor; the docs/docstring sentences rely on it.docs/filesystem.md:81-82("writes are deferred … discarded on rollback") now also holds forAioS3FileSystem.pipe_file/put_file(offline test both ways).
Existing callers:
S3AioExecutor(loop=...)withoutmax_workersnow raisesTypeError, a breaking change that the PR body release-notes; its only in-repo caller is_create_executor.AioS3FSCursorbuildsAioS3FileSystem(connection=..., default_block_size=...)(s3fs/result_set.py:134), somax_workersdefaults tocpu_count() * 5, which is at least the default executor size, and its S3 request concurrency is unchanged._touch()now returns a dict; existing callers ignore the value.
Evidence: the offline tests use mocked S3 calls, and the peak counts are fake-provider counts, not measured AWS traffic. The live run covers the S3 write/read/transaction integration tests only. The AioS3FSCursor suites were not run locally (Athena cost) and are left to CI.
| """ | ||
| if mode == "create" and self._sync_fs.exists(path): | ||
| raise FileExistsError(path) | ||
| with self.open(path, "wb", **kwargs) as f: |
There was a problem hiding this comment.
Independent review (relayed): Codex CLI 0.160.0, model gpt-6-astra, codex exec -s read-only --ephemeral, session 01a100cd-1719-7ec2-8cbb-9534f456c49e. Static review only of git diff aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe..985510d7f6f85514432bf902949478f71999275a on a detached snapshot, without the PR text or prior findings. Result: FINDINGS (5). Covered: transaction enrollment/commit/rollback, executor selection, multipart/range reads, shutdown/cancellation, public signatures, sync wrappers, asynchronous=True, AioS3FSCursor callers, touch(), tests, docs.
1 (P2): inside a transaction, pipe_file(path, b"data", ContentType=..., Metadata=...) passes these to open(); only s3_additional_kwargs reaches S3, so the committed object loses them (the delegated small-object path used to forward them).
Author disposition: deferred. S3FileSystem.pipe_file inside a transaction goes through AbstractFileSystem.pipe_file → open() and drops them the same way (s3.py:1350-1354); fixing only the aio side would diverge from the sync filesystem. Parameter routing per operation is #969.
| """ | ||
|
|
||
| def __init__(self, loop: asyncio.AbstractEventLoop | None = None) -> None: | ||
| def __init__(self, loop: asyncio.AbstractEventLoop | None = None, *, max_workers: int) -> None: |
There was a problem hiding this comment.
2 (P2) (relayed from Codex, see the review header on s3_async.py): previously valid S3AioExecutor(loop=running_loop) fails with TypeError because max_workers became mandatory; the class is in the public API docs.
Author disposition: accepted.
There was a problem hiding this comment.
Repair: ef9c594 gives max_workers the S3File default ((cpu_count() or 1) * 5), so S3AioExecutor(loop=...) keeps working; test_s3_executor.py::test_init covers it. Codex follow-up (session 01a100d6-88ec-7481-9bd0-3176bf0b3587, static) confirmed this fix.
| Returns: | ||
| The return value of the function. | ||
| """ | ||
| async with self._semaphore: |
There was a problem hiding this comment.
3 (P2) (relayed from Codex): with max_workers=1, cancel a running function's future while another is queued; async with releases the semaphore though asyncio.to_thread() cannot stop the thread, so the second starts alongside it and exceeds the documented limit.
Author disposition: accepted.
There was a problem hiding this comment.
Repair history:
- ef9c594 kept the permit until the thread returned (
ensure_future+shield+ done callback). - Codex follow-up (session
01a100d6-88ec-7481-9bd0-3176bf0b3587, static) found thatshieldalso stopped the cancellation of functions still queued for a thread, so part uploads could run afterdiscard()aborted the upload. That was a regression from the repair. - f35e70e reverts to
async withand declines this finding: holding permits only for running functions while keeping queued work cancellable needs lock/flag bookkeeping disproportionate to the case.S3Filecancels its futures only indiscard()and after a failed completion, and submits nothing more on those paths unless the abort itself fails (Failed multipart uploads are aborted while parts are in flight, and a failed write-mode open() leaves a half-built S3File #976). - A second follow-up (session
01a100e0-ea82-7e10-a2bf-86d6c0ddac50) showed the docstring overpromised: cancellation reaches the loop asynchronously. 730a00e now documents that a cancelled function may still start or keep running and then no longer counts towards the limit; the PR body states the bound for normal uploads and reads only.
| if self._intrans: | ||
| # The transaction belongs to this filesystem, not to the internal | ||
| # S3FileSystem, so write through open() to defer the commit to it. | ||
| await asyncio.to_thread(self._pipe_file_in_transaction, path, value, mode, **kwargs) |
There was a problem hiding this comment.
4 (P2) (relayed from Codex): with asynchronous=True, cancel/time out await _pipe_file(..., mode="create") inside fs.transaction while its worker checks existence; rollback completes, the worker continues, open() sees _intrans=False, autocommits and publishes after rollback. _put_file() has the same gap.
Author disposition: deferred. It is a general limit of cancelling asyncio.to_thread work: every delegated aio operation (e.g. _rm, _cp_file) keeps running after its coroutine is cancelled, and before this PR these writes were never transactional at all. Binding the write to the transaction captured on the loop would need a design decision; candidate for a follow-up issue.
There was a problem hiding this comment.
Deferred as stated; to be proposed as a follow-up issue together with the other cancellation lifetimes of delegated asyncio.to_thread work, after the maintainer's decision.
| with lock: | ||
| state["active"] += 1 | ||
| state["peak"] = max(state["peak"], state["active"]) | ||
| time.sleep(0.1) |
There was a problem hiding this comment.
5 (P3) (relayed from Codex): the test relies on sleep(0.1) for overlap and asserts peak == 2; slow scheduling can give peak == 1 with a correct limiter, or exactly 2 without one.
Author disposition: accepted.
There was a problem hiding this comment.
Repair: ef9c594 replaced the sleep window with a threading.Condition; follow-up reviews found the bound detection still timing-dependent. 730a00e: each call waits up to 5 s until another overlaps it, then 0.2 s for a third one. Locally, with the semaphore removed, [False] failed 3 of 3 runs; with it, both pass (~0.8 s each).
| # Hold the call until another one overlaps it, then | ||
| # briefly for a third one, which only an executor | ||
| # without the limit runs. | ||
| condition.wait_for(lambda: state["active"] >= 2, timeout=5) |
There was a problem hiding this comment.
Independent follow-up reviews (relayed): Codex CLI 0.160.0, model gpt-6-astra, codex exec -s read-only --ephemeral on detached snapshots of each repair, static only.
985510d7..ef9c5946(session01a100d6-88ec-7481-9bd0-3176bf0b3587): FINDINGS. The optionalmax_workersfix was confirmed. Theshieldrepair stopped the cancellation of queued work (P2, a regression from the repair). The cancellation test could false-pass or hang (P2 ×2). The parallel test was still timing-based (P2). → f35e70e reverted the repair and removed the test.ef9c5946..f35e70e1(session01a100e0-ea82-7e10-a2bf-86d6c0ddac50): FINDINGS.S3Filecan submit after a cancel when the abort itself fails (P2; PR claim narrowed, failure state is Failed multipart uploads are aborted while parts are in flight, and a failed write-mode open() leaves a half-built S3File #976). The docstring overpromised on cancellation (P2). The parallel test was timing-dependent (P2). → 730a00e.f35e70e1..730a00e0(session01a100e9-dbf2-7593-b0e7-1bca65b60881): the docstring finding is resolved; no deadlock or counter issue with the pair ordering, and an unbounded executor "should be detected under ordinary scheduling". One remaining P2: if successive calls are dispatched more than ~5.2 s apart (e.g. a saturated shared pool), a correct limiter givespeak == 1and the test fails; two non-overlapping pairs can still pass without the limiter.
Author disposition of the remaining P2: declined. Some timeout is needed so the test cannot hang, and a 5 s dispatch delay of the loop's default executor would break the other timing-bounded tests too. Locally, with the semaphore removed, [False] failed 3 of 3 runs, and both cases pass in ~0.8 s with it.
| ``as_completed()`` and ``Future.cancel()``. | ||
| ``as_completed()`` and ``Future.cancel()``. At most ``max_workers`` of the | ||
| submitted functions run at once. A function whose future is cancelled | ||
| may still start or keep running in its thread, and then no longer counts |
There was a problem hiding this comment.
Follow-up (session 01a100e9-dbf2-7593-b0e7-1bca65b60881) confirmed that this wording resolves the cancellation docstring finding.
…ouch() - pipe_file() and put_file() inside a transaction of an AioS3FileSystem now defer the commit to that transaction. They used to write through the internal S3FileSystem, which is not in the transaction, so the objects were written immediately and not rolled back. - S3AioExecutor runs at most max_workers functions at once, so AioS3FileSystem.open(max_workers=N) bounds the parallel part uploads and range reads as S3FileSystem.open() does. max_workers is a new required keyword argument of S3AioExecutor. - An AioS3FileSystem created with asynchronous=True has no event loop of its own, so its files use an S3ThreadPoolExecutor instead of an S3AioExecutor that always raised RuntimeError. - touch() delegates to S3FileSystem.touch(): it passes the PutObject parameters, returns the response and raises ValueError for an existing object with truncate=False. _touch() returns the response. Closes #977 Closes #954 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The parts of an AioS3File are only waited on at commit, which a transaction runs in the caller's thread, so writing through open() does not occupy the event loop's default executor and the internal-file detour is unnecessary. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- max_workers defaults to the S3File default, so S3AioExecutor(loop=...) keeps working. - A cancelled function keeps running in its thread, so its permit is released when it returns, not when the coroutine is cancelled. - The parallel request test waits for overlapping calls instead of relying on a sleep window. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Shielding the thread kept the permit of a cancelled running function but also let functions still queued for a thread run after S3File.discard() aborted the upload. Since a future can no longer be cancelled once its function has started, the plain semaphore already holds the permit of a running function until it returns. The parallel request test now holds each call until a third one runs or a timeout passes, so an unbounded executor fails it whatever the scheduling. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The parallel request test waits up to 5 s for two calls to overlap and then 0.2 s for a third one, instead of relying on the scheduling within a single short wait. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
730a00e to
4ddec37
Compare
| future.set_exception(task.exception()) | ||
|
|
||
| task = asyncio.run_coroutine_threadsafe(asyncio.to_thread(run), self._loop) | ||
| task = asyncio.run_coroutine_threadsafe(self._run(run), self._loop) |
There was a problem hiding this comment.
Self-review after rebasing onto #992 (both perspectives): base 6f2585014 (merge-base with master), head 4ddec37. Old series aa0fc914..730a00e0 → new series 6f258501..4ddec371 (conflicts in s3_executor.py, test_s3_executor.py and test_s3_async.py, each resolved by the commit's own delta; the f35e70e/730a00e0 docstring caveats were dropped and their messages reworded because #992 makes running futures non-cancellable).
Round 1 (behavior):
- The semaphore wraps the dispatch of Wait for running parts before aborting multipart uploads, and validate S3File before init #992's
run(). Since the returned future is started and resolved in the function's thread, a running function holds its permit until it returns. A future cancelled before it starts only lets itsrun()take a permit for the instant it returns without callingfn. settle()still owns the future when the task ends while waiting for a permit or a thread (loop shutdown). Covered bytest_loop_shutdown[2-1].discard()waits only for futures whosecancel()returned False, which are always running in a thread, sotest_discard_on_event_loop_threadcannot block on a permit.- Errors raised before
fnstarts (e.g. a shut-down default executor) still reach the future viasettle()(test_executor_shut_down).
Round 2 (claims): the docstring "At most max_workers of the submitted functions run at once" is now strict, and the PR body was updated (executor paragraph, TEST). The AWS runs for 730a00e were cancelled by the concurrency group on push, so current CI is still pending.
Result: CLEAN.
| Returns: | ||
| The PutObject response as a dictionary. | ||
| """ | ||
| return self._sync_fs.touch(path, truncate=truncate, **kwargs) |
There was a problem hiding this comment.
Independent review of the rebased series (relayed): Codex CLI 0.160.0, model gpt-6-astra, codex exec -s read-only --ephemeral, session 01a100f9-9a8a-76f3-83b0-71f70fd27183, static only. Input: git diff 6f258501..4ddec371 plus git range-diff aa0fc914..730a00e0 6f258501..4ddec371, on a detached snapshot of 4ddec37.
Covered: semaphore accounting with cancellation, pre-start exceptions, the claim lock and loop shutdown; discard() on the event loop thread and multipart failure cleanup ("no new deadlock found: permit waiters remain cancellable and are excluded from cleanup waits"); executor selection, transaction helpers, touch(), docstrings and docs/aio.md; the merged tests (both shutdown configurations, the stated 5 s/0.2 s windows); and the full diff and range-diff ("no lost or duplicated changes").
FINDINGS (1): P2, this line. fs.touch(existing_key) inside with fs.transaction: now truncates immediately and is not rolled back. The inherited fsspec touch() went through open() and was deferred. Direct S3FileSystem.touch() and async _touch() already bypassed transactions.
Author disposition: declined, release-noted. #977 asks for parity with S3FileSystem.touch(), which has never been transactional (PutObject directly). The deferral came from fsspec's generic fallback, which also dropped the PutObject parameters. A transactional touch() would be a change to both filesystems and is out of scope; the behavior change is now listed in the PR body's release notes.
WHAT
AioS3FileSystemnow matchesS3FileSystemin three places.pipe_file()andput_file()insidefs.transactiondefer the commit to the transaction of theAioS3FileSystemand are discarded on rollback.In a transaction, they write through
open()of theAioS3FileSystem, so the file joins its transaction as withS3FileSystem.Outside a transaction, both still delegate to the internal
S3FileSystemunchanged.open()executor:S3AioExecutorruns at mostmax_workerssubmitted functions at once (anasyncio.Semaphore), soopen(max_workers=N)bounds the parallelUploadPart/GetObjectrequests asS3FileSystem.open()does (AioS3FileSystem.open() ignores max_workers for parallel uploads and range reads #954).The semaphore wraps the dispatch of Wait for running parts before aborting multipart uploads, and validate S3File before init #992's
run(), which starts and resolves the returned future in the function's thread, so a running function holds its permit until it returns (its future cannot be cancelled once started). A function cancelled before it starts takes a permit only for the instant itsrun()returns without calling it, andsettle()still covers a task cancelled while it waits for a permit (e.g. at loop shutdown).asynchronous=Truehas no event loop of its own (_loopisNone), so its files get anS3ThreadPoolExecutor(max_workers);S3AioExecutorused to raiseRuntimeErroron the first parallel part or range.touch():AioS3FileSystem.touch()delegates toS3FileSystem.touch(), likermdir()and the other sync delegations.It passes the PutObject parameters, returns the response dict and raises
ValueErrorfor an existing object withtruncate=False;_touch()returns the response too.Release-note items (4.0.0):
S3AioExecutor.__init__takes a new optional argumentmax_workers(defaultcpu_count() * 5, asS3File), and raisesValueErrorif it is not positive, asThreadPoolExecutordoes.AioS3FileSystem.open()honorsmax_workersfor the parallel requests of a file. They used to be bounded only by the size of the event loop's default executor.pipe_file()/put_file()inside a transaction of anAioS3FileSystemare now deferred and rolled back.AioS3FileSystem.touch()returns the PutObject response instead ofNone.AioS3FileSystem.touch()inside a transaction now writes immediately, asS3FileSystem.touch()always has; fsspec's defaulttouch()that it inherited deferred the empty file to the transaction.docs/aio.mdand the class docstrings describe the executor choice and themax_workersbound.WHY
Closes #977
Closes #954
TEST
Tested commit: 4ddec37, rebased onto master after #992 (6f25850).
just format,just lint: pass.just docs lint: 0 errors.tests/pyathena/filesystem/test_s3_async.py, run without AWS (pytest --noconftestwith dummy env):test_transaction_pipe_put_file[True|False]: commit uploads both objects at commit time (withContentTypefromput_file); rollback uploads nothing.test_transaction_pipe_file_create_existing:mode="create"in a transaction still raisesFileExistsError.test_touch_sync_wrapper:touch()returns a dict, passesContentType, raisesValueErrorwithtruncate=False.test_open_parallel_requests[False|True]: 4 parts and 4 ranges withmax_workers=2give a peak of exactly 2 concurrent requests, also withasynchronous=True.test_open_invalid_max_workers:max_workers=0raisesValueError.tests/pyathena/filesystem/test_s3_executor.py::test_init:S3AioExecutor(loop=None)still constructs andmax_workers=0raises.threading.Condition); with the semaphore removed,[False]failed in 3 of 3 runs ([True]uses the thread pool, which bounds the requests anyway).git apply -R), 5 of the new tests fail: the transaction cases (2PutObjectcalls on rollback),touch(returnedNone), and both parallel cases (peak 4 withasynchronous=False;RuntimeErrorwithasynchronous=True).test_transaction_pipe_file_create_existingpasses there too, as the old path checkedmode="create"as well.uv run --env-file .env pytest -n 1 tests/pyathena/filesystem/test_s3_async.py tests/pyathena/filesystem/test_s3_executor.py: 85 passed.test_s3_executor.py(Wait for running parts before aborting multipart uploads, and validate S3File before init #992's tests plustest_init, withtest_loop_shutdownparametrized so the pending function waits either for a thread or for a permit) andtest_s3.py::TestS3Fileincludingtest_discard_on_event_loop_thread: all passed.put()calls of 3-part files inside one transaction completed (120 parts), so writing throughAioS3Fileinasyncio.to_threadworkers does not starve the default executor; the parts are waited on at commit in the caller's thread.AioS3FSCursortests (Athena queries), which useS3AioExecutorthroughAioS3FileSystem; left to CI.🤖 Generated with Claude Code