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
21 changes: 15 additions & 6 deletions docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -280,10 +280,13 @@ directories below the bucket level) and is always a no-op.
## Typed S3 operations

`S3FileSystem.core` is an `S3Core`, the typed operations that the filesystem sends
its listing, lookup and delete requests with. It can also be built on a boto3 S3
client. Each operation sends one request (one per page for the iterators) with the
retry policy, raises `FileNotFoundError` for a missing bucket, or for a missing object
or version that it reads, and caches nothing.
its listing, lookup, delete and multipart upload requests with. It can also be built
on a boto3 S3 client. Each operation sends one request (one per page for the
iterators) with the retry policy, raises `FileNotFoundError` for a missing bucket or
multipart upload, or for a missing object or version that it reads, and caches
nothing. Requests sent
through `fs.core` do not invalidate the filesystem's cache: call
`fs.invalidate_cache()` after a change, or make it through the filesystem.

```python
import boto3
Expand All @@ -306,8 +309,7 @@ for page in core.list_objects("YOUR_S3_BUCKET", prefix="path/to/", delimiter="/"
DeleteObjects request accepts. `S3DeleteBatch.from_paths()` groups paths into batches
of up to 1,000 objects per bucket. The objects that S3 could not delete are in the
`errors` of the returned `S3DeleteResult`, not raised; as in S3, deleting a key that
does not exist is not an error. Requests sent through `fs.core` do not invalidate the
filesystem's cache: call `fs.invalidate_cache()` after them, or delete with `fs.rm()`.
does not exist is not an error.

```python
from pyathena.filesystem.s3_core import S3DeleteBatch
Expand All @@ -322,6 +324,13 @@ for batch in S3DeleteBatch.from_paths(paths):
print(error) # path (code: message)
```

`create_multipart_upload()`, `upload_part()`, `upload_part_copy()`,
`complete_multipart_upload()` and `abort_multipart_upload()` send the requests of a
multipart upload. `part_ranges()` sends no request: it splits an object into the byte

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, AWS operation, docs). Base 911492c3e63ae6ba8c8892be2543fd1b51a75447, head 3825953dfc51b842229f1eb563931176f3797880. Result: CLEAN after PR-text corrections.

Claims checked:

  • "No request change on the wire except one edge case": every primitive keeps {**params, **request} and S3Core.call's request_kwargs merge; discard() and clear_multipart_uploads() send the same AbortMultipartUpload fields. The append's CopySource changed from "bucket/key" to a dict: botocore.handlers.handle_copy_source_param (botocore 1.43.102) encodes both identically for keys with spaces, +, %, non-ASCII and ?; only a key like a?versionId=v1?b differs, where the string form was misread as a version. Stated in the PR.
  • Docs: abort_multipart_upload of a missing upload → FileNotFoundError (measured live); part_ranges() sends no request; the cache note now covers all core requests. The other filesystem docs (5 MiB, 5 GiB, 10,000 parts) still hold.
  • Existing callers: the removed MULTIPART_UPLOAD_* attributes raise AttributeError when read and are ignored when a subclass sets them — added to the PR's breaking note. The removed private helpers have no remaining references in pyathena/, tests/, docs/ or benchmarks/.
  • AWS operation: retries and error translation still go through S3Core.call; the aio copy keeps to_thread + wait/shield and its cleanup unchanged. Pre-existing Multipart uploads with ChecksumAlgorithm cannot complete, and clear_multipart_uploads() fails on finished uploads #1070 (checksum completion, clear race) measured live and left unchanged.
  • Evidence: offline failure set identical to master at 3825953 (664 passed / 119 AWS-dependent failed); the live runs were on f730dd1 (PR body corrected to say so). Full AWS suite pending CI on Ready.

ranges of the parts that copy it, by the part limits `MULTIPART_UPLOAD_MIN_PART_SIZE` (5 MiB),
`MULTIPART_UPLOAD_MAX_PART_SIZE` (5 GiB) and `MULTIPART_UPLOAD_MAX_PARTS` (10,000) of
`S3Core`.

## Async filesystem

`AioS3FileSystem` provides the same functionality on top of fsspec's
Expand Down
301 changes: 73 additions & 228 deletions pyathena/filesystem/s3.py

Large diffs are not rendered by default.

66 changes: 26 additions & 40 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
S3CompleteMultipartUpload,
S3Metadata,
S3MultipartUpload,
S3MultipartUploadPart,
S3Object,
S3ObjectType,
S3ObjectVersion,
Expand Down Expand Up @@ -205,7 +206,7 @@ def _pipe_file_in_transaction(
FileExistsError: If the mode is "create" and the path already
exists.
ValueError: If the compression is not supported, or if the data
takes more than ``MULTIPART_UPLOAD_MAX_PARTS`` blocks.
takes more than ``S3Core.MULTIPART_UPLOAD_MAX_PARTS`` blocks.
"""
# See S3FileSystem.pipe_file.
compression = get_compression(self._strip_protocol(path), kwargs.pop("compression", None))
Expand Down Expand Up @@ -258,7 +259,7 @@ def _put_file_in_transaction(
FileExistsError: If the mode is "create" and the path already
exists.
ValueError: If the file takes more than
``MULTIPART_UPLOAD_MAX_PARTS`` blocks.
``S3Core.MULTIPART_UPLOAD_MAX_PARTS`` blocks.
"""
if os.path.isdir(lpath):
return
Expand Down Expand Up @@ -530,7 +531,7 @@ async def _copy_file(self, path1: str, path2: str, **kwargs) -> bool:
return False
size1 = info1.get("size", 0)
try:
if size1 <= S3FileSystem.MULTIPART_UPLOAD_MAX_PART_SIZE:
if size1 <= self.core.MULTIPART_UPLOAD_MAX_PART_SIZE:
await asyncio.to_thread(
self._sync_fs._copy_object,
bucket1=source.bucket,
Expand Down Expand Up @@ -596,22 +597,22 @@ async def _copy_object_with_multipart_upload(
directive has an invalid value.
"""
max_workers = max_workers if max_workers else self._sync_fs.max_workers
block_size = block_size if block_size else S3FileSystem.MULTIPART_UPLOAD_MAX_PART_SIZE
block_size = block_size if block_size else self.core.MULTIPART_UPLOAD_MAX_PART_SIZE
if (
block_size < S3FileSystem.MULTIPART_UPLOAD_MIN_PART_SIZE
or block_size > S3FileSystem.MULTIPART_UPLOAD_MAX_PART_SIZE
block_size < self.core.MULTIPART_UPLOAD_MIN_PART_SIZE
or block_size > self.core.MULTIPART_UPLOAD_MAX_PART_SIZE
):
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"5 MiB ({self.core.MULTIPART_UPLOAD_MIN_PART_SIZE} bytes) and "
f"5 GiB ({self.core.MULTIPART_UPLOAD_MAX_PART_SIZE} bytes), "
f"inclusive: {block_size}."
)

create_kwargs, version_id1, head_size = await asyncio.to_thread(
self._sync_fs._get_multipart_copy_kwargs, bucket1, key1, version_id1, kwargs
)
if head_size is not None and head_size <= S3FileSystem.MULTIPART_UPLOAD_MAX_PART_SIZE:
if head_size is not None and head_size <= self.core.MULTIPART_UPLOAD_MAX_PART_SIZE:
# See S3FileSystem._copy_object_with_multipart_upload.
await asyncio.to_thread(
self._sync_fs._copy_object,
Expand All @@ -624,15 +625,9 @@ async def _copy_object_with_multipart_upload(
)
return
# The size of the copied version; see S3FileSystem.
ranges = self._sync_fs._get_copy_ranges(
size1 if head_size is None else head_size, block_size
)
copy_source: dict[str, Any] = {
"Bucket": bucket1,
"Key": key1,
}
if version_id1:
copy_source["VersionId"] = version_id1
ranges = self.core.part_ranges(size1 if head_size is None else head_size, block_size)
source = S3Path(bucket1, key1, version_id1)
destination = S3Path(bucket2, key2)
# Listed before anything is written; see S3FileSystem.
annotations = (
await asyncio.to_thread(
Expand All @@ -642,42 +637,34 @@ async def _copy_object_with_multipart_upload(
else []
)
multipart_upload = await asyncio.to_thread(
self._sync_fs._create_multipart_upload,
bucket=bucket2,
key=key2,
**create_kwargs,
self.core.create_multipart_upload, destination, **create_kwargs
)
upload_id = cast(str, multipart_upload.upload_id)

semaphore = asyncio.Semaphore(max_workers)
part_kwargs = self._sync_fs.core.operation_params("upload_part_copy", kwargs)
part_kwargs = self.core.operation_params("upload_part_copy", kwargs)
failed = False

async def _upload_part(i: int, range_: tuple[int, int]) -> dict[str, Any] | None:
async def _upload_part(i: int, range_: tuple[int, int]) -> S3MultipartUploadPart | None:
nonlocal failed
async with semaphore:
if failed:
# The upload is being aborted; do not start more parts.
return None
try:
result = await asyncio.to_thread(
self._sync_fs._upload_part_copy,
bucket=bucket2,
key=key2,
copy_source=copy_source,
return await asyncio.to_thread(
self.core.upload_part_copy,
path=destination,
upload_id=upload_id,
part_number=i + 1,
copy_source_ranges=range_,
source=source,
range_=range_,
**part_kwargs,
)
except Exception:
# Set before the semaphore lets a waiting part start.
failed = True
raise
return {
"ETag": result.etag,
"PartNumber": result.part_number,
}

tasks = [asyncio.ensure_future(_upload_part(i, r)) for i, r in enumerate(ranges)]
completion: asyncio.Task[S3CompleteMultipartUpload] | None = None
Expand Down Expand Up @@ -708,12 +695,11 @@ async def _abort() -> None:
parts = [task.result() for task in tasks]
completion = asyncio.ensure_future(
asyncio.to_thread(
self._sync_fs._complete_multipart_upload,
bucket=bucket2,
key=key2,
upload_id=upload_id,
parts=cast(list[dict[str, Any]], parts),
**self._sync_fs.core.operation_params("complete_multipart_upload", kwargs),
self.core.complete_multipart_upload,
destination,
upload_id,
cast(list[S3MultipartUploadPart], parts),
**self.core.operation_params("complete_multipart_upload", kwargs),
)
)
# shield keeps a cancellation from cancelling the completion, whose
Expand Down
Loading
Loading