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
14 changes: 10 additions & 4 deletions docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -227,8 +227,13 @@ fs.clear_multipart_uploads("s3://YOUR_S3_BUCKET/path/to/")

With `version_aware=True`, reads pin the object version observed at open time, so a
file handle keeps returning consistent data even if the object is overwritten while
reading. Explicit versions can always be read with the `?versionId=` suffix or the
`version_id` argument.
reading. The file path carries the pinned version as a `?versionId=` suffix, which
`metadata()`, `getxattr()` and `url()` of the file also use. Explicit versions can
always be read with the `?versionId=` suffix or the `version_id` argument. Only a
`?versionId=` (or `?versionID=`, `?versionid=`, `?version_id=`) query at the end of a
path is a version; any other `?` is part of the key. A path with a version is not a
glob pattern: `copy()`, `mv()` and `get()` copy that version to a destination named
after its key.

```python
fs = S3FileSystem(
Expand All @@ -239,8 +244,9 @@ fs = S3FileSystem(
with fs.open("s3://YOUR_S3_BUCKET/path/to/object", "rb") as f:
data = f.read() # Pinned to the version observed at open time.

# List all versions of the objects under a prefix.
fs.ls("s3://YOUR_S3_BUCKET/path/to/", versions=True, detail=True)
# List all versions of the objects under a prefix. Each version is named
# "bucket/key?versionId=<id>", except the "null" version, which is named by its key.
fs.ls("s3://YOUR_S3_BUCKET/path/to/", versions=True)

# Typed version information, including delete markers if requested.
versions = fs.object_version_info("s3://YOUR_S3_BUCKET/path/to/object")
Expand Down
248 changes: 205 additions & 43 deletions pyathena/filesystem/s3.py

Large diffs are not rendered by default.

125 changes: 122 additions & 3 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@
from fsspec.asyn import AsyncFileSystem, sync
from fsspec.callbacks import _DEFAULT_CALLBACK
from fsspec.core import get_compression
from fsspec.implementations.local import LocalFileSystem, make_path_posix
from fsspec.utils import check_contained

from pyathena.filesystem.s3 import CompressedBuffer, S3File, S3FileSystem
from pyathena.filesystem.s3_executor import S3AioExecutor, S3Executor, S3ThreadPoolExecutor
Expand Down Expand Up @@ -392,6 +394,87 @@ def mv(self, path1, path2, recursive=False, maxdepth=None, **kwargs) -> None:
"""
sync(self.loop, self._mv, path1, path2, recursive=recursive, maxdepth=maxdepth, **kwargs)

async def _copy(
self,
path1,
path2,
recursive=False,
on_error=None,
maxdepth=None,
batch_size=None,
**kwargs,
) -> None:
"""Copy files within S3.

See :meth:`S3FileSystem.copy`. The copies run as in fsspec's
``_copy()``.

Args:
path1: Source S3 path, glob pattern, or list of them.
path2: Destination S3 path, or list of paths when ``path1`` is a
list.
recursive: Whether to copy the directories with their contents.
on_error: ``"raise"`` or ``"ignore"`` for a missing source.
maxdepth: Maximum depth of a recursive copy.
batch_size: Number of copies to run at the same time.
**kwargs: Additional S3 copy parameters passed to ``_cp_file()``.
"""
sources = [path1] if isinstance(path1, (str, os.PathLike)) else path1
if isinstance(path2, str) and any(S3Path.has_version_id(p) for p in sources):
path1, path2 = await asyncio.to_thread(
self._sync_fs._copy_paths, path1, path2, recursive=recursive, maxdepth=maxdepth
)
if not path1:
return
await super()._copy(
path1,
path2,
recursive=recursive,
on_error=on_error,
maxdepth=maxdepth,
batch_size=batch_size,
**kwargs,
)

async def _get(
self, rpath, lpath, recursive=False, callback=_DEFAULT_CALLBACK, maxdepth=None, **kwargs
) -> None:
"""Copy files from S3 to the local filesystem.

See :meth:`S3FileSystem.get`. The downloads run as in fsspec's
``_get()``.

Args:
rpath: Source S3 path, glob pattern, or list of them.
lpath: Local destination path, or list of paths when ``rpath`` is
a list.
recursive: Whether to copy the directories with their contents.
callback: Progress callback.
maxdepth: Maximum depth of a recursive copy.
**kwargs: Additional parameters passed to ``_get_file()``.

Raises:
ValueError: If a source with a version ID is paired, and a
destination lies outside ``lpath``.
"""
sources = [rpath] if isinstance(rpath, (str, os.PathLike)) else rpath
if isinstance(lpath, (str, os.PathLike)) and any(S3Path.has_version_id(p) for p in sources):
root = make_path_posix(lpath)
rpath, lpath = await asyncio.to_thread(
self._sync_fs._copy_paths,
rpath,
root,
recursive=recursive,
maxdepth=maxdepth,
isdir=LocalFileSystem().isdir,
)
check_contained(root, lpath)
if not rpath:
return
await super()._get(
rpath, lpath, recursive=recursive, callback=callback, maxdepth=maxdepth, **kwargs
)

async def _cp_file(self, path1: str, path2: str, **kwargs) -> None:
"""Copy an S3 object, using async parallel multipart upload for large files.

Expand Down Expand Up @@ -423,9 +506,6 @@ async def _copy_file(self, path1: str, path2: str, **kwargs) -> bool:
Raises:
ValueError: If trying to copy to a versioned file or copy buckets.
"""
# fsspec < 2026.6.0 leaks the typo'd "onerror" keyword from mv();
# see S3FileSystem.cp_file.
kwargs.pop("onerror", None)
# Parameters of the multipart copy, not of the S3 requests.
block_size = kwargs.pop("block_size", None)
max_workers = kwargs.pop("max_workers", None)
Expand Down Expand Up @@ -657,6 +737,45 @@ async def _find(
return {f.name: f for f in files}
return [f.name for f in files]

async def _expand_path(self, path, recursive=False, maxdepth=None, **kwargs) -> list[str]:
"""Expand glob patterns and directories into the paths they match.

See :meth:`S3FileSystem.expand_path`. The other paths are expanded by
fsspec's ``_expand_path()``, whose glob lists only the keys under the
literal part of the pattern.

Args:
path: S3 path, glob pattern, or list of them.
recursive: Whether to include the paths below the directories.
maxdepth: Maximum depth of the expansion, at least 1.
**kwargs: Additional arguments passed to fsspec's
``_expand_path()``, such as ``assume_literal``.

Returns:
The sorted matching paths.

Raises:
ValueError: If ``maxdepth`` is less than 1.
FileNotFoundError: If nothing matches.
"""
if maxdepth is not None and maxdepth < 1:
raise ValueError("maxdepth must be at least 1")
versions, others = self._sync_fs._split_version_paths(path)

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), finding 6 (P2): confirmed and repaired. Delegating aio _expand_path() to the sync expand_path() lost fsspec's async glob prefix= (the stem before the first wildcard), so bucket/reports/2026-*.csv listed reports/ instead of reports/2026-.

The aio override now splits off the version paths with the shared _split_version_paths() and passes the others to fsspec's async _expand_path(). test_expand_path_glob_lists_stem_prefix asserts prefix="2026-".

out = {p for p in versions if not recursive or await self._exists(p)}
if others:
try:
out.update(
await super()._expand_path(
others, recursive=recursive, maxdepth=maxdepth, **kwargs
)
)
except FileNotFoundError:
if not out:
raise
if not out:
raise FileNotFoundError(path)
return sorted(out)

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

Expand Down
52 changes: 48 additions & 4 deletions pyathena/filesystem/s3_path.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@

from __future__ import annotations

import os
import re
from dataclasses import dataclass, replace
from re import Pattern
Expand All @@ -24,8 +25,9 @@ class S3Path:
Paths are parsed from strings such as ``s3://bucket/key``,
``s3a://bucket/key`` or ``bucket/key``, optionally followed by a version
ID query (``?versionId=``, ``?versionID=``, ``?versionid=`` or
``?version_id=``). The root path, which names no bucket, is not an
``S3Path``.
``?version_id=``). Only a query at the end of the path is a version ID;
any other ``?`` is part of the key. The root path, which names no bucket,
is not an ``S3Path``.

Attributes:
bucket: The name of the bucket.
Expand All @@ -46,9 +48,16 @@ class S3Path:
's3://bucket/dir/key?versionId=v1'
"""

# Version IDs do not contain "?", so only the last query can be one.
VERSION_QUERY: ClassVar[Pattern[str]] = re.compile(
r"\?version(Id|ID|id|_id)=(?P<version_id>[^?]+)\Z"
)
# Keys may contain any character, including newlines. A bare "/" is tried
# before a key so that "bucket/?versionId=..." names the bucket.
PATTERN: ClassVar[Pattern[str]] = re.compile(
r"(^s3://|^s3a://|^)(?P<bucket>[a-zA-Z0-9.\-_]+)(/(?P<key>[^?]+)|/)?"
r"($|\?version(Id|ID|id|_id)=(?P<version_id>.+)$)"
r"(^s3://|^s3a://|^)(?P<bucket>[a-zA-Z0-9.\-_]+)(/|/(?P<key>.+?))?"
rf"(\Z|{VERSION_QUERY.pattern})",
re.DOTALL,

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), findings 2 and 3 (P2): confirmed and repaired.

  • Finding 2: .+? excluded newlines. bucket/a\nb raised ValueError, and bucket/a\n parsed as key a because $ matched before the final newline. Both keys parsed intact on the base. The pattern now uses re.DOTALL and \Z.
  • Finding 3: bucket/?versionId=v1 parsed as the key ?versionId=v1, whereas the base read it as a version of the bucket path. A bare / is now tried before a key.

Both are covered by new TestS3Path.test_parse cases (a newline inside, at the end, and before a version).

)

bucket: str
Expand All @@ -74,6 +83,41 @@ def parse(cls, path: str) -> S3Path:
raise ValueError(f"Invalid S3 path format {path}.")
return cls(match.group("bucket"), match.group("key"), match.group("version_id"))

@classmethod
def split_version_id(cls, path: str) -> tuple[str, str | None]:
"""Split the version ID query from the end of a path string.

Unlike :meth:`parse`, any string is accepted, such as the names that
fsspec builds.

Args:
path: The path string.

Returns:
Tuple of the string without the version ID query and the version
ID, or of the string itself and None if it ends with no version
ID query.
"""
match = cls.VERSION_QUERY.search(path)
if not match:
return path, None
return path[: match.start()], match.group("version_id")

@classmethod
def has_version_id(cls, path: str | os.PathLike[str]) -> bool:
"""Return whether a path string ends with a version ID query.

Unlike :meth:`parse`, any string is accepted, as with
:meth:`split_version_id`.

Args:
path: The path.

Returns:
Whether the path ends with a version ID query.
"""
return cls.split_version_id(os.fspath(path))[1] is not None

@property
def is_bucket(self) -> bool:
"""Whether the path names the bucket: it has no key, or a key of only slashes."""
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ dependencies = [
"boto3>=1.43.31",
"botocore>=1.43.31",
"tenacity>=4.1.0",
"fsspec",
"fsspec>=2026.9.0",
"python-dateutil",
]
requires-python = ">=3.11"
Expand Down
Loading
Loading