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
20 changes: 20 additions & 0 deletions docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,26 @@ the existing object also count toward the limit. A write with `open` that reache
limit raises `ValueError` and aborts its multipart upload. Multipart copies with `cp`
use parts large enough to stay within the limit.

`cp` copies an object larger than 5 GiB with a multipart upload instead of a single
CopyObject request, with the same result as CopyObject. If HeadObject reports a size
of at most 5 GiB when the copy starts, for example because the size that `cp` found
came from a cached listing, the object is copied with CopyObject instead. CopyObject parameters given as
keyword arguments are sent to the multipart requests that accept them, such as
`CopySourceIfMatch` to each part copy. With the default `COPY` value of
`MetadataDirective`, `TaggingDirective`, and `AnnotationDirective`, the content headers
(such as `ContentType`) and user-defined metadata, the tags, and the annotations of the
source are copied, and the values given for them are ignored, as CopyObject does. A
`REPLACE` directive uses the given values instead, and `AnnotationDirective="EXCLUDE"`
skips the annotations. Copying the tags needs `s3:GetObjectTagging` on the source, and
copying the annotations needs `s3:ListObjectAnnotations` and `s3:GetObjectAnnotation` on

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, AWS effects): base de8cc52ac2a43cdba72bb4571c185883287741fc, head a42788c8ad454869c9cf24a17b33b2520f025ba8. Result: CLEAN after the PR-body corrections (the ObjectIfMatch wording and the added Limits).

Claims checked:

  • "CopyObject ignores given ContentType/Metadata/Tagging under COPY, copies content headers/metadata/tags/annotations, not WebsiteRedirectLocation/storage class/SSE": measured live in the CI bucket on 2026-10-03 (default, explicit COPY, COPY with ContentType/Metadata/Tagging, REPLACE, TaggingDirective=REPLACE without Tagging).
  • "Sync and async multipart copies now match CopyObject": measured live on a 10 MiB + 1 B source with 5 MiB parts, compared with a CopyObject of the same source. Headers, metadata, tags and annotation payloads matched. A failed CopySourceIfMatch left no multipart upload.
  • "Before, CreateMultipartUpload rejected CopySourceIf*/CopySourceSSECustomer*/ExpectedSourceBucketOwner/directives": in the botocore 1.43.102 model these are not members of CreateMultipartUpload; the maintainer's comment on Multipart copies drop the source metadata, and the async multipart copy does not abort on failure #973 reported the same.
  • Existing callers: a ContentType/Metadata/Tagging given for a copy over 5 GiB is now ignored by default. This is a behavior change, release-noted. mv() raises before removing anything if the copy fails, so a failed annotation copy keeps both source and destination.
  • AWS operator: extra requests per copy over 5 GiB are 1 HeadObject + 1 GetObjectTagging + ≥1 ListObjectAnnotations + 2 per annotation, all through _call's retries. Async annotation copies are bounded by max_workers. The new permission requirements are release-noted in the PR body and stated in docs/filesystem.md.
  • Evidence: the offline Stubber tests and the live measurement are kept separate in the PR body, and >5 GiB, SSE-C, directory and requester-pays cases are marked as not tested live.

the source and `s3:PutObjectAnnotation` on the destination. A source without a
`?versionId=` suffix whose HeadObject reports a version ID other than `null` is copied
from that version, which needs `s3:GetObjectVersion` and, to copy the tags,
`s3:GetObjectVersionTagging` on the source. A `null` version is not pinned. The annotations are listed before anything is written and copied after the
upload completes, so the destination exists without them until the last one is
written. If an annotation fails to copy, the error is raised and the destination is
kept. A failed part copy aborts the multipart upload.

Paths are normalized as in fsspec, which drops a trailing slash, so `info`, `isfile`,
and `open` treat `s3://YOUR_S3_BUCKET/dir/` as `s3://YOUR_S3_BUCKET/dir`: the object
`dir` if it exists, and otherwise the directory `dir`. An object whose key ends in a
Expand Down
432 changes: 390 additions & 42 deletions pyathena/filesystem/s3.py

Large diffs are not rendered by default.

195 changes: 149 additions & 46 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -448,29 +448,32 @@ async def _copy_file(self, path1: str, path2: str, **kwargs) -> bool:
# S3FileSystem.cp_file.
return False
size1 = info1.get("size", 0)
if size1 <= S3FileSystem.MULTIPART_UPLOAD_MAX_PART_SIZE:
await asyncio.to_thread(
self._sync_fs._copy_object,
bucket1=bucket1,
key1=key1,
version_id1=version_id1,
bucket2=bucket2,
key2=key2,
**kwargs,
)
else:
await self._copy_object_with_multipart_upload(
bucket1=bucket1,
key1=key1,
version_id1=version_id1,
size1=size1,
bucket2=bucket2,
key2=key2,
max_workers=max_workers,
block_size=block_size,
**kwargs,
)
self._sync_fs.invalidate_cache(path2)
try:
if size1 <= S3FileSystem.MULTIPART_UPLOAD_MAX_PART_SIZE:
await asyncio.to_thread(
self._sync_fs._copy_object,
bucket1=bucket1,
key1=key1,
version_id1=version_id1,
bucket2=bucket2,
key2=key2,
**kwargs,
)
else:
await self._copy_object_with_multipart_upload(
bucket1=bucket1,
key1=key1,
version_id1=version_id1,
size1=size1,
bucket2=bucket2,
key2=key2,
max_workers=max_workers,
block_size=block_size,
**kwargs,
)
finally:
# See S3FileSystem._copy_file.
self._sync_fs.invalidate_cache(path2)
return True

async def _copy_object_with_multipart_upload(
Expand All @@ -485,6 +488,28 @@ async def _copy_object_with_multipart_upload(
version_id1: str | None = None,
**kwargs,
) -> None:
"""Copy an object with a multipart upload of its byte ranges.

See :meth:`S3FileSystem._copy_object_with_multipart_upload`. The part
and annotation copies run in parallel with ``asyncio.gather`` and
``asyncio.to_thread``.

Args:
bucket1: Source S3 bucket name.
key1: Source object key.
size1: Size of the source object in bytes.
bucket2: Destination S3 bucket name.
key2: Destination object key.
max_workers: Maximum number of parallel requests.
block_size: Size in bytes of the copied ranges.
version_id1: Source version ID, if any.
**kwargs: The CopyObject parameters of the copy; each request
receives those that it accepts.

Raises:
ValueError: If ``block_size`` is out of the part size limits or a
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
if (
Expand All @@ -498,52 +523,130 @@ async def _copy_object_with_multipart_upload(
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:
# See S3FileSystem._copy_object_with_multipart_upload.
await asyncio.to_thread(
self._sync_fs._copy_object,
bucket1=bucket1,
key1=key1,
version_id1=version_id1,
bucket2=bucket2,
key2=key2,
**kwargs,
)
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._sync_fs._get_copy_ranges(size1, block_size)
# Listed before anything is written; see S3FileSystem.
annotations = (
await asyncio.to_thread(
self._sync_fs._list_object_annotations, bucket1, key1, version_id1, kwargs
)
if self._sync_fs._copies_annotations(bucket1, kwargs)
else []
)
multipart_upload = await asyncio.to_thread(
self._sync_fs._create_multipart_upload,
bucket=bucket2,
key=key2,
**kwargs,
**create_kwargs,
)
upload_id = cast(str, multipart_upload.upload_id)

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

async def _upload_part(i: int, range_: tuple[int, int]) -> dict[str, Any]:
async def _upload_part(i: int, range_: tuple[int, int]) -> dict[str, Any] | None:
nonlocal failed
async with semaphore:
result = await asyncio.to_thread(
self._sync_fs._upload_part_copy,
bucket=bucket2,
key=key2,
copy_source=copy_source,
upload_id=cast(str, multipart_upload.upload_id),
part_number=i + 1,
copy_source_ranges=range_,
**part_kwargs,
)
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,
upload_id=upload_id,
part_number=i + 1,
copy_source_ranges=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,
}

parts = await asyncio.gather(*[_upload_part(i, r) for i, r in enumerate(ranges)])
parts_list = sorted(parts, key=lambda x: x["PartNumber"])
tasks = [asyncio.ensure_future(_upload_part(i, r)) for i, r in enumerate(ranges)]
try:
# gather keeps the part-number order of the tasks.
parts = await asyncio.gather(*tasks)
completed = await 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._get_operation_kwargs("complete_multipart_upload", kwargs),
)
except Exception:
failed = True
# A part that is still copying when the upload is aborted may be
# stored after the abort, so wait for the running parts first.
await asyncio.gather(*tasks, return_exceptions=True)
await asyncio.to_thread(
self._sync_fs._abort_multipart_upload, bucket2, key2, upload_id, kwargs
)
raise

await asyncio.to_thread(
self._sync_fs._complete_multipart_upload,
bucket=bucket2,
key=key2,
upload_id=cast(str, multipart_upload.upload_id),
parts=parts_list,
**self._sync_fs._get_operation_kwargs("complete_multipart_upload", kwargs),
failed = False

async def _copy_annotation(name: str) -> None:
nonlocal failed
async with semaphore:
if failed:
# Do not start more copies after one failed, as
# S3FileSystem does.
return
try:
await asyncio.to_thread(
self._sync_fs._copy_object_annotation,
name,
bucket1,
key1,
version_id1,
bucket2,
key2,
completed,
kwargs,
)
except Exception:
failed = True
raise

results = await asyncio.gather(
*[_copy_annotation(name) for name in annotations], return_exceptions=True
)
for result in results:
if isinstance(result, BaseException):
raise result

async def _find(
self,
Expand Down
4 changes: 2 additions & 2 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,8 @@ maintainers = [
{name = "laughingman7743", email = "laughingman7743@gmail.com"},
]
dependencies = [
"boto3>=1.41.2",
"botocore>=1.41.2",
"boto3>=1.43.31",
"botocore>=1.43.31",
"tenacity>=4.1.0",
"fsspec",
"python-dateutil",
Expand Down
Loading
Loading