-
Notifications
You must be signed in to change notification settings - Fork 116
Align AioS3FileSystem with S3FileSystem in transactions, open() and touch() #988
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
9ae2bd6
2f65a46
4e95bc5
cbe4e96
71a24b6
dec6e12
4ddec37
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 |
|---|---|---|
|
|
@@ -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, | ||
|
|
@@ -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. | ||
|
|
@@ -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) | ||
| return | ||
| await asyncio.to_thread(self._sync_fs.pipe_file, path, value, mode=mode, **kwargs) | ||
|
|
||
| def _pipe_file_in_transaction( | ||
|
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, repaired. Scope: full Claims checked:
Existing callers:
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 |
||
| 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: | ||
|
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. Independent review (relayed): Codex CLI 0.160.0, model gpt-6-astra, 1 (P2): inside a transaction, Author disposition: deferred. |
||
| 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: | ||
|
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 (implementation behavior): CLEAN, no repairs. Scope: Checked:
Recorded trade-off (not a defect):
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. Repair (985510d, found in round 2): the starvation premise above was false. The parts of an |
||
| """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) | ||
|
|
||
|
|
@@ -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, | ||
|
|
@@ -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, | ||
|
|
@@ -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) | ||
|
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. Independent review of the rebased series (relayed): Codex CLI 0.160.0, model gpt-6-astra, Covered: semaphore accounting with cancellation, pre-start exceptions, the claim lock and loop shutdown; FINDINGS (1): P2, this line. Author disposition: declined, release-noted. #977 asks for parity with |
||
|
|
||
|
|
||
| class AioS3File(S3File): | ||
|
|
@@ -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``. | ||
| """ | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 | ||
|
|
@@ -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) | ||
|
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. Round 1: |
||
|
|
||
| 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: | ||
|
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. 3 (P2) (relayed from Codex): with Author disposition: accepted.
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. Repair history:
|
||
| return await asyncio.to_thread(fn, *args, **kwargs) | ||
|
|
||
| @override | ||
| def submit(self, fn: Callable[..., T], *args: Any, **kwargs: Any) -> Future[T]: | ||
|
|
@@ -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) | ||
|
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 after rebasing onto #992 (both perspectives): base 6f2585014 (merge-base with master), head 4ddec37. Old series Round 1 (behavior):
Round 2 (claims): the docstring "At most Result: CLEAN. |
||
| task.add_done_callback(settle) | ||
| return future | ||
| raise RuntimeError( | ||
|
|
||
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.
4 (P2) (relayed from Codex): with
asynchronous=True, cancel/time outawait _pipe_file(..., mode="create")insidefs.transactionwhile 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_threadwork: 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.
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_threadwork, after the maintainer's decision.