Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions docs/aio.md
Original file line number Diff line number Diff line change
Expand Up @@ -280,6 +280,8 @@ parallel operations. Two implementations are provided:
`AioS3FileSystem` automatically uses `S3AioExecutor` for file handles, so multipart
uploads and parallel range reads are dispatched through the event loop with
`asyncio.to_thread()` instead of a separate `ThreadPoolExecutor` per file.
At most `max_workers` of them run at once.
An instance created with `asynchronous=True` has no event loop of its own, so its file handles use `S3ThreadPoolExecutor`.

### Usage with AioS3FSCursor

Expand Down
110 changes: 106 additions & 4 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,14 +11,16 @@

import asyncio
import logging
import mimetypes
import os
from multiprocessing import cpu_count
from typing import TYPE_CHECKING, Any, cast

from fsspec.asyn import AsyncFileSystem
from fsspec.callbacks import _DEFAULT_CALLBACK

from pyathena.filesystem.s3 import S3File, S3FileSystem
from pyathena.filesystem.s3_executor import S3AioExecutor
from pyathena.filesystem.s3_executor import S3AioExecutor, S3Executor, S3ThreadPoolExecutor
from pyathena.filesystem.s3_object import (
S3Metadata,
S3MultipartUpload,
Expand Down Expand Up @@ -48,6 +50,8 @@ class AioS3FileSystem(AsyncFileSystem):
File handles created by ``_open`` use ``S3AioExecutor`` so that parallel
operations (range reads, multipart uploads) are dispatched through the event
loop with ``asyncio.to_thread`` instead of a ``ThreadPoolExecutor`` per file.
An instance created with ``asynchronous=True`` has no event loop of its own,
so its file handles use a ``ThreadPoolExecutor``.

Attributes:
_sync_fs: The internal synchronous S3FileSystem instance.
Expand Down Expand Up @@ -161,11 +165,74 @@ async def _rm_file(self, path: str, **kwargs) -> None:
async def _pipe_file(
self, path: str, value: bytes | bytearray | memoryview, mode: str = "overwrite", **kwargs
) -> None:
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)

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.

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.

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.

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.

return
await asyncio.to_thread(self._sync_fs.pipe_file, path, value, mode=mode, **kwargs)

def _pipe_file_in_transaction(

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 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 AioS3File in 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 transactional put() calls completed. Repaired by writing through open().
  • "bounded only by min(32, os.cpu_count() + 4)": Python 3.13 uses process_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 for AioS3FileSystem.pipe_file/put_file (offline test both ways).

Existing callers:

  • S3AioExecutor(loop=...) without max_workers now raises TypeError, a breaking change that the PR body release-notes; its only in-repo caller is _create_executor.
  • AioS3FSCursor builds AioS3FileSystem(connection=..., default_block_size=...) (s3fs/result_set.py:134), so max_workers defaults to cpu_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.

self, path: str, value: bytes | bytearray | memoryview, mode: str, **kwargs
) -> None:
"""Write bytes into the path as a file of this filesystem's transaction.

Args:
path: S3 path (s3://bucket/key) to write to.
value: The bytes to write.
mode: "overwrite" or "create". With "create", raise
FileExistsError when the object already exists.
**kwargs: Additional parameters passed to ``open()``.

Raises:
FileExistsError: If the mode is "create" and the path already
exists.
"""
if mode == "create" and self._sync_fs.exists(path):
raise FileExistsError(path)
with self.open(path, "wb", **kwargs) as f:

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): 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.

f.write(value)

async def _put_file(self, lpath: str, rpath: str, callback=_DEFAULT_CALLBACK, **kwargs) -> None:
if self._intrans:
# See _pipe_file.
await asyncio.to_thread(self._put_file_in_transaction, lpath, rpath, callback, **kwargs)
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:

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 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(): same mode="create" check, path stripping, autocommit=False, registration in transaction.files; Transaction.complete commits/discards the S3File (its thread pool is shut down on close() after the parts finish, as in the sync path).
  • put_file helper vs S3FileSystem.put_file (s3.py:1476): same directory/bucket-only early returns, size callback, ContentType guess, s3_additional_kwargs, cache invalidation.
  • Starvation: an AioS3File here would wait in a default-executor thread on parts queued to the same executor; the internal S3File uses its own pool, so concurrent put() (fsspec batches up to 128) cannot deadlock.
  • fsspec sets _loop = None exactly when asynchronous=True (fsspec/asyn.py:445-448), so _create_executor picks the thread pool only there. touch is not in fsspec.asyn.async_methods, so mirror_sync_methods keeps the explicit touch().
  • New offline tests fail on the base for the transaction, touch and 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.

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.

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.

"""Upload a local file as a file of this filesystem's transaction.

Mirrors :meth:`S3FileSystem.put_file`, but writes through ``open()``
of this filesystem.

Args:
lpath: Local file path to upload.
rpath: S3 destination path (s3://bucket/key).
callback: Progress callback for tracking upload progress.
**kwargs: Additional S3 parameters (e.g., ContentType, StorageClass).
"""
if os.path.isdir(lpath):
return
_, key, _ = self.parse_path(rpath)
if not key:
return

callback.set_size(os.path.getsize(lpath))
if "ContentType" not in kwargs:
content_type, _ = mimetypes.guess_type(lpath)
if content_type is not None:
kwargs["ContentType"] = content_type

with (
self.open(rpath, "wb", s3_additional_kwargs=kwargs) as remote,
open(lpath, "rb") as local,
):
while data := local.read(remote.blocksize):
remote.write(data)
callback.relative_update(len(data))
self.invalidate_cache(rpath)

async def _get_file(self, rpath: str, lpath: str, callback=_DEFAULT_CALLBACK, **kwargs) -> None:
await asyncio.to_thread(self._sync_fs.get_file, rpath, lpath, callback=callback, **kwargs)

Expand Down Expand Up @@ -344,6 +411,23 @@ async def _find(
return {f.name: f for f in files}
return [f.name for f in files]

def _create_executor(self, max_workers: int) -> S3Executor:
"""Create the executor for the parallel operations of a file.

An instance created with ``asynchronous=True`` has no event loop of
its own, so its files run the operations in a thread pool.

Args:
max_workers: The maximum number of operations that run at once.

Returns:
An ``S3AioExecutor`` on the event loop of this filesystem, or an
``S3ThreadPoolExecutor`` if it has none.
"""
if self._loop is None:
return S3ThreadPoolExecutor(max_workers=max_workers)
return S3AioExecutor(loop=self._loop, max_workers=max_workers)

def _open(
self,
path: str,
Expand All @@ -367,7 +451,7 @@ def _open(
path,
mode,
max_workers=max_workers,
executor=S3AioExecutor(loop=self._loop),
executor=self._create_executor(max_workers=max_workers),
block_size=block_size,
cache_type=cache_type,
autocommit=autocommit,
Expand Down Expand Up @@ -576,8 +660,24 @@ def invalidate_cache(self, path: str | None = None) -> None:
"""
self._sync_fs.invalidate_cache(path)

async def _touch(self, path: str, truncate: bool = True, **kwargs) -> None:
await asyncio.to_thread(self._sync_fs.touch, path, truncate=truncate, **kwargs)
async def _touch(self, path: str, truncate: bool = True, **kwargs) -> dict[str, Any]:
return await asyncio.to_thread(self._sync_fs.touch, path, truncate=truncate, **kwargs)

def touch(self, path: str, truncate: bool = True, **kwargs) -> dict[str, Any]:
"""Create an empty object with PutObject.

See :meth:`S3FileSystem.touch`.

Args:
path: S3 path (s3://bucket/key) of the object.
truncate: If True, replace an existing object with an empty one;
if False, raise if the object exists.
**kwargs: Additional parameters passed to the PutObject API.

Returns:
The PutObject response as a dictionary.
"""
return self._sync_fs.touch(path, truncate=truncate, **kwargs)

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 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.



class AioS3File(S3File):
Expand All @@ -589,4 +689,6 @@ class AioS3File(S3File):
through the ``S3Executor`` interface — the ``S3AioExecutor``
provided by ``AioS3FileSystem`` dispatches them through the event loop with
``asyncio.to_thread`` instead of a ``ThreadPoolExecutor`` per file.
For an ``AioS3FileSystem`` created with ``asynchronous=True``, it is an
``S3ThreadPoolExecutor``.
"""
35 changes: 32 additions & 3 deletions pyathena/filesystem/s3_executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
from collections.abc import Callable
from concurrent.futures import Future
from concurrent.futures.thread import ThreadPoolExecutor
from multiprocessing import cpu_count
from typing import Any, TypeVar

from pyathena.util import override
Expand Down Expand Up @@ -80,26 +81,54 @@ class S3AioExecutor(S3Executor):
``concurrent.futures.Future`` objects that are compatible with
``as_completed()``, ``wait()`` and ``Future.cancel()``. As with
``ThreadPoolExecutor``, a future cannot be cancelled once its function has
started.
started. At most ``max_workers`` of the submitted functions run at once.

This avoids thread-in-thread nesting when ``S3File`` is used from within
``asyncio.to_thread()`` calls (the pattern used by ``AioS3FileSystem``).

Args:
loop: A running asyncio event loop.
max_workers: The maximum number of submitted functions that run at once.

Raises:
RuntimeError: If the event loop is not running when ``submit`` is called.
"""

def __init__(self, loop: asyncio.AbstractEventLoop | None = None) -> None:
def __init__(
self,
loop: asyncio.AbstractEventLoop | None = None,
max_workers: int = (cpu_count() or 1) * 5,
) -> None:
"""Initialize the executor with the event loop to schedule work on.

Args:
loop: The asyncio event loop. ``submit`` raises ``RuntimeError``
if it is None or not running.
max_workers: The maximum number of submitted functions that run
at once.

Raises:
ValueError: If ``max_workers`` is not positive.
"""
if max_workers <= 0:
# 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)

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.

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


async def _run(self, fn: Callable[..., T], *args: Any, **kwargs: Any) -> T:
"""Run the function in a thread once fewer than ``max_workers`` run.

Args:
fn: The blocking function to run.
*args: Positional arguments passed to the function.
**kwargs: Keyword arguments passed to the function.

Returns:
The return value of the function.
"""
async with self._semaphore:

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.

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.

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.

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 that shield also stopped the cancellation of functions still queued for a thread, so part uploads could run after discard() aborted the upload. That was a regression from the repair.
  • f35e70e reverts to async with and declines this finding: holding permits only for running functions while keeping queued work cancellable needs lock/flag bookkeeping disproportionate to the case. S3File cancels its futures only in discard() 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.

return await asyncio.to_thread(fn, *args, **kwargs)

@override
def submit(self, fn: Callable[..., T], *args: Any, **kwargs: Any) -> Future[T]:
Expand Down Expand Up @@ -142,7 +171,7 @@ def settle(task: Future[None]) -> None:
elif future.set_running_or_notify_cancel():
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)

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 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 its run() take a permit for the instant it returns without calling fn.
  • settle() still owns the future when the task ends while waiting for a permit or a thread (loop shutdown). Covered by test_loop_shutdown[2-1].
  • discard() waits only for futures whose cancel() returned False, which are always running in a thread, so test_discard_on_event_loop_thread cannot block on a permit.
  • Errors raised before fn starts (e.g. a shut-down default executor) still reach the future via settle() (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.

task.add_done_callback(settle)
return future
raise RuntimeError(
Expand Down
Loading
Loading