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
4 changes: 4 additions & 0 deletions docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,10 @@ through the buffered file path. Inside an
[fsspec transaction](https://filesystem-spec.readthedocs.io/en/latest/features.html#transactions),
writes are deferred until the transaction commits and are discarded on rollback.

The block size for writing, given by the `block_size` argument of `open` or by the
filesystem's `default_block_size`, must be between 5 MiB and 5 GiB, inclusive, the part
size limits of a multipart upload. Otherwise, `open` raises `ValueError`.

## Error translation

S3 error responses are translated into standard Python exceptions, so filesystem
Expand Down
91 changes: 59 additions & 32 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
import os.path
import re
from collections.abc import Callable, Iterator
from concurrent.futures import Future, as_completed
from concurrent.futures import Future, as_completed, wait
from copy import deepcopy
from datetime import datetime
from multiprocessing import cpu_count
Expand Down Expand Up @@ -1275,7 +1275,11 @@ def _copy_object_with_multipart_upload(
block_size < self.MULTIPART_UPLOAD_MIN_PART_SIZE
or block_size > self.MULTIPART_UPLOAD_MAX_PART_SIZE
):
raise ValueError("Block size must be greater than 5MiB and less than 5GiB.")
raise ValueError(
"Block size must be between "
f"5 MiB ({self.MULTIPART_UPLOAD_MIN_PART_SIZE} bytes) and "
f"5 GiB ({self.MULTIPART_UPLOAD_MAX_PART_SIZE} bytes), inclusive: {block_size}."
)

copy_source = {
"Bucket": bucket1,
Expand Down Expand Up @@ -1407,9 +1411,10 @@ def _finish_multipart_upload(
) -> S3CompleteMultipartUpload:
"""Collect the uploaded parts and complete the multipart upload.

When any part fails, the remaining parts are cancelled and the
multipart upload is aborted so that no incomplete upload is left
behind, then the original error is re-raised.
When any part or the completion fails, the parts that have not
started are cancelled, the running ones are waited for, and the
multipart upload is aborted so that no incomplete upload or part is
left behind. The original error is then re-raised.

Args:
bucket: S3 bucket name.
Expand All @@ -1431,8 +1436,10 @@ def _finish_multipart_upload(
parts=parts,
)
except Exception:
for future in futures:
future.cancel()
# A part that is still uploading when the upload is aborted may
# be stored after the abort, so wait for the parts that could not
# be cancelled first.
wait([future for future in futures if not future.cancel()])
try:
self._call(
self._client.abort_multipart_upload,
Expand Down Expand Up @@ -2249,8 +2256,9 @@ def __init__(
part copies.
executor: The executor for parallel operations. If None, a new
``S3ThreadPoolExecutor`` is created.
block_size: The block size for reads and writes. Must be at least
``MULTIPART_UPLOAD_MIN_PART_SIZE`` unless reading.
block_size: The block size for reads and writes. Must be between
``MULTIPART_UPLOAD_MIN_PART_SIZE`` and
``MULTIPART_UPLOAD_MAX_PART_SIZE``, inclusive, unless reading.
cache_type: The fsspec cache type for reads.
autocommit: Whether to commit the written data when the file is
closed. If False, :meth:`commit` must be called.
Expand All @@ -2263,13 +2271,16 @@ def __init__(

Raises:
ValueError: If the path has no key, the version IDs do not match,
a version is given for writing, or the block size is too small
for writing.
a version is given for writing, or the block size is not
between ``MULTIPART_UPLOAD_MIN_PART_SIZE`` and
``MULTIPART_UPLOAD_MAX_PART_SIZE`` for writing.
"""
self.max_workers = max_workers
self._executor: S3Executor = executor or S3ThreadPoolExecutor(max_workers=max_workers)
self.s3_additional_kwargs = s3_additional_kwargs if s3_additional_kwargs else {}

# The arguments are validated, and the objects looked up, before the
# base class initializer: a file that fails here is never opened, so
# its garbage collection does not close (flush and commit) it.
bucket, key, path_version_id = S3FileSystem.parse_path(path)
self.bucket = bucket
if not key:
Expand All @@ -2292,8 +2303,20 @@ def __init__(
# Carry the version in the path, as with the ?versionId= suffix,
# so that a reopened (e.g., unpickled) file reads the same version.
path = f"{path}?versionId={self.version_id}"
if "r" not in mode and not (

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 (PR description corrected; no code change)

Scope: git diff aa0fc9146f4e683d6cdf96b894f963f9dd8f7abe..0191c1f683ff6ee0283d4cbdf431ecb8f876b87b, the PR body, commit messages, changed docstrings/comments, and docs/filesystem.md.

Claims checked:

  • AbortMultipartUpload: "parts currently in progress might or might not succeed … it might be necessary to abort a given multipart upload multiple times" — confirmed in the botocore S3 service model shipped in uv.lock.
  • 5 MiB to 5 GiB part size: confirmed in the S3 user guide's multipart upload limits table; PutObject single-operation limit of 5 GB confirmed in the "Uploading objects" guide. The PR body now says the >5 GiB part rejection is documented, not reproduced.
  • Issue premise "missing-key check raises after the base initializer": true at e0e85da (issue base), but PR Accept version_id in S3FileSystem.open() and cat_file() #958 (bbf3759) moved it before super().__init__. Reproduced on aa0fc91: the issue script prints one "Exception ignored" line (the block size file), not two. Corrected the WHY section; the missing-key test cases pass on the base and are kept for the no-request contract.
  • Caller impact of the new upper bound: open()/put_file() with default_block_size > 5 GiB now fail at open() even for small writes that used to succeed with a single PutObject; pipe_file() of data that fits in one PutObject never opens a file and is unaffected (s3.py pipe_file routes on min(block_size, MULTIPART_UPLOAD_MAX_PART_SIZE)). Added this to the release-note items. Read-mode block sizes (e.g. S3FSResultSet's default_block_size) are not checked.
  • Operator: the added wait is bounded per part by botocore timeouts and _call's retry configuration, as the part request itself already was; stated in the release note. No new requests are added.
  • docs/filesystem.md: no other statement about block size limits exists in docs/ or the README.
  • Evidence: offline results are local; the live filesystem run (298 passed) is at b2c0fe7, and 0191c1f changes only an offline test. The in-flight abort is not exercised against real S3 (stated in TEST).

fs.MULTIPART_UPLOAD_MIN_PART_SIZE <= block_size <= fs.MULTIPART_UPLOAD_MAX_PART_SIZE
):
# When writing, every full block is uploaded as a part of a
# multipart upload.
raise ValueError(
"Block size for writing must be between "
f"5 MiB ({fs.MULTIPART_UPLOAD_MIN_PART_SIZE} bytes) and "
f"5 GiB ({fs.MULTIPART_UPLOAD_MAX_PART_SIZE} bytes), inclusive: {block_size}."
)

self._details: S3Object | dict[str, Any] = {}
append_info: S3Object | None = None
append_data: bytes | None = None
if "r" in mode:
# Looked up before the base class initializer, which would
# otherwise take the size from the latest version of the object.
Expand All @@ -2308,7 +2331,14 @@ def __init__(
self._details = info
if size is None:
size = info.get("size")
elif "a" in mode and fs.exists(path):
append_info = fs.info(path)
if append_info.get("size", 0) < fs.MULTIPART_UPLOAD_MIN_PART_SIZE:
# Too small to be a part of a multipart upload: rewritten
# from the buffer.
append_data = fs.cat(path)

self._executor: S3Executor = executor or S3ThreadPoolExecutor(max_workers=max_workers)
super().__init__(
fs=fs,
path=path,
Expand All @@ -2319,28 +2349,19 @@ def __init__(
cache_options=cache_options,
size=size,
)
if "r" not in mode and block_size < self.fs.MULTIPART_UPLOAD_MIN_PART_SIZE:
# When writing occurs, the block size should not be smaller
# than the minimum size of a part in a multipart upload.
raise ValueError(f"Block size must be >= {self.fs.MULTIPART_UPLOAD_MIN_PART_SIZE}MB.")

self.append_block = False
if "a" in mode and self.fs.exists(path):
info = self.fs.info(self.path, version_id=self.version_id)
loc = info.get("size", 0)
if loc < self.fs.MULTIPART_UPLOAD_MIN_PART_SIZE:
# Too small to be a part of a multipart upload: rewrite it
# from the buffer.
self.write(self.fs.cat(self.path))
self.multipart_upload: S3MultipartUpload | None = None
self.multipart_upload_parts: list[Future[S3MultipartUploadPart]] = []
if append_info is not None:
if append_data is not None:
self.write(append_data)
else:
# Copied with UploadPartCopy as the leading part(s).
self.append_block = True
self.loc = loc
self.s3_additional_kwargs.update(info.to_api_repr())
self._details = info

self.multipart_upload: S3MultipartUpload | None = None
self.multipart_upload_parts: list[Future[S3MultipartUploadPart]] = []
self.loc = append_info.get("size", 0)
self.s3_additional_kwargs.update(append_info.to_api_repr())
self._details = append_info

def close(self) -> None:
"""Close the file, flushing any written data, and shut down its executor."""
Expand Down Expand Up @@ -2497,10 +2518,16 @@ def commit(self) -> None:
self.fs.invalidate_cache(self.path)

def discard(self) -> None:
"""Cancel pending part uploads and abort the multipart upload, if any."""
"""Abort the multipart upload, if any.

The part uploads that have not started are cancelled, and the
running ones are waited for before the abort.
"""
if self.multipart_upload:
for f in self.multipart_upload_parts:
f.cancel()
# A part that is still uploading when the upload is aborted may
# be stored after the abort, so wait for the parts that could not
# be cancelled first.
wait([f for f in self.multipart_upload_parts if not f.cancel()])
# s3_additional_kwargs also holds object parameters (e.g., the
# existing object's metadata in append mode) that
# AbortMultipartUpload rejects.
Expand Down
7 changes: 6 additions & 1 deletion pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -281,7 +281,12 @@ async def _copy_object_with_multipart_upload(
block_size < S3FileSystem.MULTIPART_UPLOAD_MIN_PART_SIZE
or block_size > S3FileSystem.MULTIPART_UPLOAD_MAX_PART_SIZE
):
raise ValueError("Block size must be greater than 5MiB and less than 5GiB.")
raise ValueError(
"Block size must be between "
f"5 MiB ({S3FileSystem.MULTIPART_UPLOAD_MIN_PART_SIZE} bytes) and "
f"5 GiB ({S3FileSystem.MULTIPART_UPLOAD_MAX_PART_SIZE} bytes), "
f"inclusive: {block_size}."
)

copy_source: dict[str, Any] = {
"Bucket": bucket1,
Expand Down
49 changes: 45 additions & 4 deletions pyathena/filesystem/s3_executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
from __future__ import annotations

import asyncio
import threading
from abc import ABCMeta, abstractmethod
from collections.abc import Callable
from concurrent.futures import Future
Expand Down Expand Up @@ -77,7 +78,9 @@ class S3AioExecutor(S3Executor):
Uses ``asyncio.run_coroutine_threadsafe(asyncio.to_thread(fn), loop)`` to
dispatch blocking functions onto the event loop's thread pool, returning
``concurrent.futures.Future`` objects that are compatible with
``as_completed()`` and ``Future.cancel()``.
``as_completed()``, ``wait()`` and ``Future.cancel()``. As with
``ThreadPoolExecutor``, a future cannot be cancelled once its function has
started.

This avoids thread-in-thread nesting when ``S3File`` is used from within
``asyncio.to_thread()`` calls (the pattern used by ``AioS3FileSystem``).
Expand All @@ -101,9 +104,47 @@ def __init__(self, loop: asyncio.AbstractEventLoop | None = None) -> None:
@override
def submit(self, fn: Callable[..., T], *args: Any, **kwargs: Any) -> Future[T]:
if self._loop is not None and self._loop.is_running():
return asyncio.run_coroutine_threadsafe(
asyncio.to_thread(fn, *args, **kwargs), self._loop
)
# The future of run_coroutine_threadsafe can be cancelled while
# the function keeps running in its thread, so the returned future
# is started and resolved by the function's thread instead.
future: Future[T] = Future()
# Acquired once, by run() or by settle(), whichever comes first,
# so that the future is started or settled exactly once.
claim = threading.Lock()

def run() -> None:
"""Run the function and resolve the future unless it was cancelled."""
if not claim.acquire(blocking=False) or not future.set_running_or_notify_cancel():
return
try:
result = fn(*args, **kwargs)
except BaseException as e:
future.set_exception(e)
else:
future.set_result(result)

def settle(task: Future[None]) -> None:
"""Resolve the future if the task ended before the function started.

This happens, for example, when the event loop shuts down.

Args:
task: The finished future of the task that runs the function.
"""
if not claim.acquire(blocking=False):
# run() has claimed the future and resolves it.
return
if task.cancelled():
future.cancel()
# Notify the waiters of the cancellation, as an executor
# does when it drops a cancelled function.
future.set_running_or_notify_cancel()
elif future.set_running_or_notify_cancel():
future.set_exception(task.exception())

task = asyncio.run_coroutine_threadsafe(asyncio.to_thread(run), self._loop)
task.add_done_callback(settle)
return future
raise RuntimeError(
"S3AioExecutor requires a running event loop. "
"Use S3ThreadPoolExecutor for synchronous usage."
Expand Down
Loading
Loading