File System Integration¶
This section covers S3 filesystem integration and object management.
S3 FileSystem¶
- class pyathena.filesystem.s3.S3FileSystem(*args, **kwargs)[source]¶
A filesystem interface for Amazon S3 that implements the fsspec protocol.
This class provides a file-system like interface to Amazon S3, allowing you to use familiar file operations (ls, open, cp, rm, etc.) with S3 objects. It’s designed to be compatible with s3fs while offering PyAthena-specific optimizations.
The filesystem supports standard S3 operations including:
Listing objects and directories
Reading and writing files
Copying and moving objects
Reading and writing object metadata, tags, and canned ACLs
Multipart uploads for large files, including management of incomplete uploads
Version-aware reads and object version listing (see
version_aware)Creating and removing buckets (disabled by default; see
allow_bucket_creation/allow_bucket_deletion)Various S3 storage classes and encryption options
Translating S3 error responses into standard Python exceptions (e.g.,
404->FileNotFoundError,403->PermissionError)
- allow_bucket_creation¶
Whether mkdir/makedirs may create buckets. Defaults to False.
- allow_bucket_deletion¶
Whether rmdir may delete buckets. Defaults to False.
- version_aware¶
Whether reads pin the object version observed at open time and ls may list all versions. Requires the s3:GetObjectVersion / s3:ListBucketVersions permissions. Defaults to False.
Example
>>> from pyathena.filesystem.s3 import S3FileSystem >>> fs = S3FileSystem() >>> >>> # List objects in a bucket >>> files = fs.ls('s3://my-bucket/data/') >>> >>> # Read a file >>> with fs.open('s3://my-bucket/data/file.csv', 'r') as f: ... content = f.read() >>> >>> # Write a file >>> with fs.open('s3://my-bucket/output/result.txt', 'w') as f: ... f.write('Hello, S3!') >>> >>> # Copy files >>> fs.cp('s3://source-bucket/file.txt', 's3://dest-bucket/file.txt')
Note
This filesystem is used internally by PyAthena for handling query results stored in S3, but can also be used independently for S3 file operations.
- OBJECT_ACLS: frozenset[str] = frozenset({'authenticated-read', 'aws-exec-read', 'bucket-owner-full-control', 'bucket-owner-read', 'private', 'public-read', 'public-read-write'})¶
- BUCKET_ACLS: frozenset[str] = frozenset({'authenticated-read', 'private', 'public-read', 'public-read-write'})¶
- PATTERN_PATH: Pattern[str] = re.compile('(^s3://|^s3a://|^)(?P<bucket>[a-zA-Z0-9.\\-_]+)(/(?P<key>[^?]+)|/)?($|\\?version(Id|ID|id|_id)=(?P<version_id>.+)$)')¶
- __init__(connection: Connection[Any] | None = None, default_block_size: int | None = None, default_cache_type: str | None = None, max_workers: int = 20, s3_additional_kwargs=None, allow_bucket_creation: bool = False, allow_bucket_deletion: bool = False, version_aware: bool = False, *args, **kwargs) None[source]¶
Create a filesystem for Amazon S3.
- Parameters:
connection – A PyAthena connection whose session, region, config and retry policy the S3 client uses. Without one, the client is built from s3fs-compatible arguments in
kwargs.default_block_size – The block size for reads and writes; defaults to
DEFAULT_BLOCK_SIZE.default_cache_type – The fsspec cache type for reads; defaults to
"bytes".max_workers – The number of threads for parallel transfers.
s3_additional_kwargs – Extra arguments for the object requests of
open()andpipe_file(); listings and other requests do not use them. Each request receives those that its operation accepts, and the parameters of a call take precedence.allow_bucket_creation – Whether
mkdir/makedirsmay create a bucket.allow_bucket_deletion – Whether
rmdirmay delete a bucket.version_aware – Whether reads pin the object version observed at open time.
*args – Passed to
fsspec.AbstractFileSystem.**kwargs – Passed to
fsspec.AbstractFileSystem; without aconnection, also s3fs-compatible client arguments.requester_pays=Truesends requester-pays requests with the operations that acceptRequestPayer.
- static parse_path(path: str) tuple[str, str | None, str | None][source]¶
Parse an S3 path into its bucket, key and version ID.
The path may have an
s3://ors3a://scheme and a version ID query (?versionId=,?versionID=,?versionid=or?version_id=).- Parameters:
path – The S3 path (e.g., “s3://bucket/key?versionId=…”).
- Returns:
Tuple of the bucket, the key (None for a bucket path) and the version ID (None if the path has none).
- Raises:
ValueError – If the path is not a valid S3 path.
- ls(path: str, detail: bool = False, refresh: bool = False, **kwargs) list[S3Object] | list[str][source]¶
List contents of an S3 path.
Lists buckets (when path is root) or objects within a bucket/prefix. Compatible with fsspec interface for filesystem operations.
- Parameters:
path – S3 path to list (e.g., “s3://bucket” or “s3://bucket/prefix”).
detail – If True, return S3Object instances; if False, return paths as strings.
refresh – If True, bypass cache and fetch fresh results from S3.
**kwargs –
Additional arguments including: versions: If True, list all versions of the objects. Requires
the filesystem to be constructed with
version_aware=True.
- Returns:
List of S3Object instances (if detail=True) or paths as strings (if detail=False).
Example
>>> fs = S3FileSystem() >>> fs.ls("s3://my-bucket") # List objects in bucket >>> fs.ls("s3://my-bucket/", detail=True) # Get detailed object info
- info(path: str, **kwargs) S3Object[source]¶
Return information about an S3 path.
Uses the directory cache first: a cached entry for the path is returned, a cached listing of the path itself makes it a directory, and a cached listing of its parent without it means it does not exist. Otherwise, a key path is looked up with HeadObject and, if no object exists, with a ListObjectsV2 request (
Delimiter="/",MaxKeys=1) that checks whether it is a key prefix; a bucket path is looked up with HeadBucket. Withversion_aware, a cached file entry without a version ID is looked up again. With an explicit version, the cached entries of the path are skipped, and the HeadObject result is cached under the version-qualified path apart from other versions, except for thenullversion, which an overwrite replaces.- Parameters:
path – S3 path (e.g., “s3://bucket” or “s3://bucket/key”).
**kwargs – Additional arguments including: refresh: If True, bypass the cache and query S3. version_id: The version ID to look up when the path has none.
- Returns:
S3Object describing the bucket, directory, or file.
- Raises:
FileNotFoundError – If the path does not exist.
- find(path: str, maxdepth: int | None = None, withdirs: bool | None = None, detail: bool = False, **kwargs) dict[str, S3Object] | list[str][source]¶
Find all files below a given S3 path.
Recursively searches for files under the specified path, with optional depth limiting and directory inclusion. Uses efficient S3 list operations with delimiter handling for performance.
- Parameters:
path – S3 path to search under (e.g., “s3://bucket/prefix”).
maxdepth – Maximum number of levels to descend, at least 1 (None for unlimited). With 1, only the entries directly under the path are listed.
withdirs – Whether to include directories in results (None = default behavior).
detail – If True, return dict of {path: S3Object}; if False, return list of paths.
**kwargs –
Additional arguments including: prefix: Key prefix, relative to the path, to filter the listed keys
by. Without maxdepth, if nothing is listed and the path itself is an object, that object is returned regardless of the prefix.
refresh: If True, bypass the cache and list from S3.
- Returns:
Dictionary mapping paths to S3Objects (if detail=True) or list of paths (if detail=False).
- Raises:
ValueError – If
maxdepthis less than 1 or the path is the root.
Example
>>> fs = S3FileSystem() >>> fs.find("s3://bucket/data/", maxdepth=2) # Limit depth >>> fs.find("s3://bucket/", withdirs=True) # Include directories
- exists(path: str, **kwargs) bool[source]¶
Check if an S3 path exists.
Determines whether a bucket, object, or prefix exists in S3. Uses caching and efficient head operations to minimize API calls.
- Parameters:
path – S3 path to check (e.g., “s3://bucket” or “s3://bucket/key”).
**kwargs – Additional arguments including: refresh: If True, bypass the cache and query S3.
- Returns:
True if the path exists, False otherwise.
Example
>>> fs = S3FileSystem() >>> fs.exists("s3://my-bucket/file.txt") >>> fs.exists("s3://my-bucket/")
- rm_file(path: str, **kwargs) None[source]¶
Delete an S3 object with DeleteObject.
Does nothing for a bucket path. If the path has a version ID, that version is deleted.
- Parameters:
path – S3 path (s3://bucket/key) of the object to delete.
**kwargs – Accepted for fsspec compatibility; not used in the request.
- rm(path, recursive=False, maxdepth=None, **kwargs) None[source]¶
Delete objects with DeleteObjects requests.
Expands the paths with
expand_pathand deletes the matched objects in parallel requests of up toDELETE_OBJECTS_MAX_KEYSkeys each, one set of requests per bucket. A path with a version ID deletes that version without expansion.- Parameters:
path – S3 path (s3://bucket/key) or list of paths to delete.
recursive – Whether to delete all objects below the paths.
maxdepth – Maximum depth to expand when
recursiveis True.**kwargs – Additional parameters passed to the DeleteObjects API.
Quiet(default True) sets the quiet mode of the requests.
- Raises:
ValueError – If a path is a bucket.
OSError – If S3 could not delete some of the objects.
- mkdir(path: str, create_parents: bool = True, **kwargs) None[source]¶
Create an S3 bucket.
S3 has no real directories below the bucket level; creating a key prefix requires no operation. This method creates the bucket when the path points at a bucket (or when
create_parentsis True and the bucket does not exist yet), and does nothing for key prefixes under an existing bucket.Bucket lifecycle operations are disabled by default because they are infrastructure-level changes; pass
allow_bucket_creation=Trueto the filesystem constructor to enable bucket creation.- Parameters:
path – S3 path (e.g., “s3://bucket” or “s3://bucket/prefix”).
create_parents – If True, create the bucket when it does not exist, even if the path contains a key prefix.
**kwargs –
Additional arguments including: acl: Canned ACL to apply to the bucket. region_name: Region to create the bucket in. Defaults to the
client’s region.
- Raises:
FileExistsError – If the path is a bucket that already exists.
FileNotFoundError – If the bucket does not exist and
create_parentsis False.PermissionError – If the bucket would be created but bucket creation is not enabled on this filesystem instance.
ValueError – If the ACL is invalid or the path is empty.
- makedirs(path: str, exist_ok: bool = False) None[source]¶
Recursively create a directory, creating the bucket if necessary.
Creating the bucket requires
allow_bucket_creation=Trueon the filesystem constructor; seemkdir().- Parameters:
path – S3 path (e.g., “s3://bucket” or “s3://bucket/prefix”).
exist_ok – If False, raise FileExistsError when the path is a bucket that already exists.
- Raises:
FileExistsError – If the path is a bucket that already exists and
exist_okis False.PermissionError – If the bucket would be created but bucket creation is not enabled on this filesystem instance.
- rmdir(path: str) None[source]¶
Remove an S3 bucket, which must be empty.
S3 has no real directories below the bucket level, so only bucket paths can be removed.
Bucket lifecycle operations are disabled by default because they are infrastructure-level changes; pass
allow_bucket_deletion=Trueto the filesystem constructor to enable bucket deletion.- Parameters:
path – S3 bucket path (e.g., “s3://bucket”).
- Raises:
FileExistsError – If the path contains a key that exists. The user may have meant
rm(path, recursive=True).FileNotFoundError – If the path contains a key that does not exist, or the bucket does not exist.
PermissionError – If bucket deletion is not enabled on this filesystem instance.
OSError – If the bucket is not empty.
- touch(path: str, truncate: bool = True, **kwargs) dict[str, Any][source]¶
Create an empty object with PutObject.
- Parameters:
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 (see
S3PutObject.to_dict()).- Raises:
ValueError – If the path has a version ID, is a bucket, or exists while
truncateis False.
- cp_file(path1: str, path2: str, recursive=False, maxdepth=None, on_error=None, **kwargs)[source]¶
Copy an S3 object to another S3 location.
Performs server-side copy of S3 objects, which is more efficient than downloading and re-uploading. Automatically chooses between simple copy and multipart copy based on object size.
- Parameters:
path1 – Source S3 path (s3://bucket/key).
path2 – Destination S3 path (s3://bucket/key).
recursive – Unused parameter for fsspec compatibility.
maxdepth – Unused parameter for fsspec compatibility.
on_error – Unused parameter for fsspec compatibility.
**kwargs – Additional S3 copy parameters (e.g., metadata, storage class). The
block_sizeandmax_workersparameters control a multipart copy and are not sent to S3.
- Raises:
ValueError – If trying to copy to a versioned file or copy buckets.
Note
Uses multipart copy for objects larger than the maximum part size to optimize performance for large files. The copy operation is performed entirely on the S3 service without data transfer.
- pipe_file(path: str, value: bytes | bytearray | memoryview, mode: str = 'overwrite', **kwargs) None[source]¶
Write bytes into the path.
Writes data up to the block size with a single PutObject request instead of the inherited
open()+write()path. Larger data and writes inside an fsspec transaction go through the buffered path, which uploads the data as a parallel multipart upload and keeps the deferred-commit semantics of transactions.- Parameters:
path – S3 path (s3://bucket/key) to write to.
value – The bytes to write.
mode – “overwrite” (default) or “create”. With “create”, raise FileExistsError when the object already exists, including one created during the write, which is not replaced.
**kwargs – Additional parameters passed to the PutObject API (e.g., ContentType, StorageClass) on the single-request path. The
block_size,max_workers, ands3_additional_kwargsparameters of theopen()path are also accepted.
- Raises:
FileExistsError – If the mode is “create” and the path already exists, or an object is created at it before the write is committed.
ValueError – If the path does not contain a key or specifies a version, or if the data takes more than
MULTIPART_UPLOAD_MAX_PARTSblocks.
- cat_file(path: str, start: int | None = None, end: int | None = None, **kwargs) bytes[source]¶
Read the contents of an S3 object with GetObject.
startandendselect bytes like a slice of the object: an empty range, or one that starts at or past the end of the object, returnsb"", and an end past the object reads up to its end. Non-negative offsets are sent to S3 as they are, and so is a negativestartwithout anend, as a suffix range of the last bytes. Other negative offsets are resolved against the size frominfo(), which also checks that the object exists for an empty range.- Parameters:
path – S3 path (s3://bucket/key) of the object.
start – Byte offset to start reading at. A negative value counts from the end of the object.
end – Byte offset to stop reading at (exclusive). A negative value counts from the end of the object.
**kwargs – Additional parameters passed to the GetObject API, except
version_id: the version ID to read when the path has none.
- Returns:
The bytes read from the object.
- Raises:
FileNotFoundError – If the path has no key or the key does not exist.
- put_file(lpath: str, rpath: str, callback=<fsspec.callbacks.NoOpCallback object>, mode: str = 'overwrite', **kwargs)[source]¶
Upload a local file to S3.
Uploads a file from the local filesystem to an S3 location. Supports automatic content type detection based on file extension and provides progress callback functionality.
- Parameters:
lpath – Local file path to upload.
rpath – S3 destination path (s3://bucket/key).
callback – Progress callback for tracking upload progress.
mode – “overwrite” (default) or “create”. With “create”, the file is written as with
open()inxbmode: raise FileExistsError when the object already exists, including one created during the upload, which is not replaced.**kwargs – Additional S3 parameters (e.g., ContentType, StorageClass). The
block_size,max_workers, ands3_additional_kwargsparameters ofopen()are also accepted.
- Raises:
FileExistsError – If the mode is “create” and the path already exists, or an object is created at it before the upload is committed.
ValueError – If the file takes more than
MULTIPART_UPLOAD_MAX_PARTSblocks.
Note
Directories are not supported for upload. If lpath is a directory, the method returns without performing any operation. Bucket-only destinations (without key) are also not supported.
- get_file(rpath: str, lpath: str, callback=<fsspec.callbacks.NoOpCallback object>, outfile=None, **kwargs)[source]¶
Download an S3 file to local filesystem.
Downloads a file from S3 to the local filesystem with progress tracking. Reads the file in chunks to handle large files efficiently.
- Parameters:
rpath – S3 source path (s3://bucket/key).
lpath – Local destination file path.
callback – Progress callback for tracking download progress.
outfile – Unused parameter for fsspec compatibility.
**kwargs – Additional S3 parameters passed to open().
Note
If lpath is a directory, the method returns without performing any operation.
- checksum(path: str, **kwargs)[source]¶
Get checksum for S3 object or directory.
Computes a checksum for the specified S3 path. For individual objects, returns the ETag converted to an integer. For directories, returns a checksum based on the directory’s tokenized representation.
- Parameters:
path – S3 path (s3://bucket/key) to get checksum for.
**kwargs – Additional arguments including: refresh: If True, refresh cached info before computing checksum.
- Returns:
Integer checksum value derived from S3 ETag or directory token.
Note
For multipart uploads, ETag format is different and only the first part before the dash is used for checksum calculation.
- sign(path: str, expiration: int = 3600, **kwargs)[source]¶
Generate a presigned URL for S3 object access.
Creates a presigned URL that allows temporary access to an S3 object without requiring AWS credentials. Useful for sharing files or providing time-limited access to resources.
- Parameters:
path – S3 path (s3://bucket/key) to generate URL for.
expiration – URL expiration time in seconds. Defaults to 3600 (1 hour).
**kwargs –
Additional parameters including: client_method: S3 operation (‘get_object’, ‘put_object’, etc.).
Defaults to ‘get_object’.
Additional parameters passed to the S3 operation.
- Returns:
Presigned URL string that provides temporary access to the S3 object.
Example
>>> fs = S3FileSystem() >>> url = fs.sign("s3://my-bucket/file.txt", expiration=7200) >>> # URL valid for 2 hours >>> >>> # Generate upload URL >>> upload_url = fs.sign( ... "s3://my-bucket/upload.txt", ... client_method="put_object" ... )
- metadata(path: str, **kwargs) S3Metadata[source]¶
Return the metadata of the path.
- Parameters:
path – S3 path (s3://bucket/key) to get metadata for.
**kwargs – Additional parameters passed to the HeadObject API.
- Returns:
S3Metadata, which behaves as a read-only mapping of the user-defined metadata (
x-amz-meta-*) and exposes the system-defined metadata (content type, encryption settings, etc.) as typed properties.
- getxattr(path: str, attr_name: str, **kwargs) str | None[source]¶
Get an attribute from the user-defined metadata of the path.
- Parameters:
path – S3 path (s3://bucket/key) to get the attribute for.
attr_name – The name of the attribute.
**kwargs – Additional parameters passed to
metadata().
- Returns:
The value of the attribute, or None if the attribute is not set.
- setxattr(path: str, copy_kwargs: dict[str, Any] | None = None, **kw_args) None[source]¶
Set the user-defined metadata of the path.
S3 does not allow updating the metadata of an existing object in place, so the object is copied onto itself with the REPLACE metadata directive. Note that this rewrites the object and updates its last-modified time.
- Parameters:
path – S3 path (s3://bucket/key) to set metadata for.
copy_kwargs – Additional parameters to use for the underlying CopyObject API call.
**kw_args – Key-value pairs to set, where the values must be strings. The keys are used as-is; names that are not valid Python identifiers (e.g., containing hyphens) can be passed by unpacking a dictionary. Does not alter existing fields, unless the field appears here - if the value is None, delete the field.
Example
>>> fs = S3FileSystem() >>> fs.setxattr("s3://bucket/key", attribute1="value1") >>> fs.setxattr("s3://bucket/key", **{"attribute-2": "value2"})
- get_tags(path: str) dict[str, str][source]¶
Retrieve the tag key/values for the given path.
- Parameters:
path – S3 path (s3://bucket/key) to get tags for.
- Returns:
Dictionary mapping tag keys to tag values.
- put_tags(path: str, tags: dict[str, str], mode: str = 'o') None[source]¶
Set the tags for the given existing key.
Tags are a str:str mapping that can be attached to any key, distinct from the user-defined metadata, which is usually set at key creation time. See https://docs.aws.amazon.com/AmazonS3/latest/userguide/object-tagging.html
- Parameters:
path – S3 path (s3://bucket/key) of the existing key to attach tags to.
tags – Tags to apply.
mode – One of ‘o’ or ‘m’. ‘o’ will over-write any existing tags. ‘m’ will merge in new tags with existing tags, which incurs two remote calls.
- chmod(path: str, acl: str, recursive: bool = False, **kwargs) None[source]¶
Set the Access Control on a bucket/key.
See https://docs.aws.amazon.com/AmazonS3/latest/userguide/acl-overview.html#canned-acl
- Parameters:
path – S3 path (s3://bucket or s3://bucket/key) to set the ACL on.
acl – The value of the canned ACL to apply.
recursive – Whether to apply the ACL to all keys below the given path too.
**kwargs – Additional parameters passed to the PutObjectAcl or PutBucketAcl API.
- list_multipart_uploads(path: str) list[S3MultipartUpload][source]¶
List in-progress (incomplete) multipart uploads in a bucket.
Incomplete multipart uploads continue to accrue storage costs until they are completed or aborted. Use
clear_multipart_uploads()to abort all of them.- Parameters:
path – S3 bucket or key path (e.g., “bucket”, “s3://bucket” or “s3://bucket/prefix”). If the path contains a key, only the uploads to that key and to the keys under
key/are listed, not those to sibling keys that merely start with the same characters (e.g.,prefix2/a).- Returns:
List of S3MultipartUpload instances describing the in-progress multipart uploads.
- object_version_info(path: str, delete_markers: bool = False, **kwargs) list[S3ObjectVersion][source]¶
List the versions of the object or of the objects under the path.
A key path without a trailing slash selects that key if it has any versions or delete markers, and otherwise the keys under
key/. The choice does not depend ondelete_markers, so a key that has only delete markers yields no versions without them. A key path with a trailing slash selects the keys under it, and a bucket path selects all the keys in the bucket. Sibling keys that merely start with the same characters (e.g.,key.bak) are never included.- Parameters:
path – S3 path (s3://bucket/key or a key prefix) to list the versions for.
delete_markers – Whether to include delete markers in the result.
**kwargs – Additional parameters passed to the ListObjectVersions API.
- Returns:
List of S3ObjectVersion instances describing the versions.
- clear_multipart_uploads(path: str) None[source]¶
Abort any incomplete multipart uploads in the bucket.
- Parameters:
path – S3 bucket or key path (e.g., “bucket”, “s3://bucket” or “s3://bucket/prefix”). If the path contains a key, only the uploads to that key and to the keys under
key/are aborted, as listed bylist_multipart_uploads().
- created(path: str) datetime[source]¶
Return the creation time of the path.
Returns the same value as
modified().- Parameters:
path – S3 path (s3://bucket/key).
- Returns:
The last-modified time of the object.
- modified(path: str) datetime[source]¶
Return the last-modified time of the path.
- Parameters:
path – S3 path (s3://bucket/key).
- Returns:
The
last_modifiedfield frominfo(), which is None for buckets and directories.
- invalidate_cache(path: str | None = None) None[source]¶
Remove the cached entries of the path and its parent paths.
A version-qualified path invalidates the version under every query spelling that
parse_pathaccepts, and also the object path without the version, because deleting or copying a version can change the current version of the object.- Parameters:
path – The path to invalidate. If None, clear the whole cache.
- class pyathena.filesystem.s3.S3File(fs: S3FileSystem, path: str, mode: str = 'rb', version_id: str | None = None, max_workers: int = 20, executor: S3Executor | None = None, block_size: int = 5242880, cache_type: str = 'bytes', autocommit: bool = True, cache_options: dict[Any, Any] | None = None, size: int | None = None, s3_additional_kwargs: dict[str, Any] | None = None, **kwargs)[source]¶
A buffered file object for reading and writing an S3 object.
Instances are returned by
S3FileSystem.open().- fs: S3FileSystem¶
- __init__(fs: S3FileSystem, path: str, mode: str = 'rb', version_id: str | None = None, max_workers: int = 20, executor: S3Executor | None = None, block_size: int = 5242880, cache_type: str = 'bytes', autocommit: bool = True, cache_options: dict[Any, Any] | None = None, size: int | None = None, s3_additional_kwargs: dict[str, Any] | None = None, **kwargs) None[source]¶
Initialize the file for the path and mode.
In read mode, the object is looked up with
info()and the reads are made conditional on its ETag (IfMatch). In append mode, an existing object smaller thanMULTIPART_UPLOAD_MIN_PART_SIZEis read into the write buffer; a larger one is copied withUploadPartCopyas the first parts of a multipart upload, whatever the block size. In exclusive-create mode, the object must not exist when the file is opened, and the upload is committed withIfNoneMatch="*"so that it does not replace an object created in the meantime.- Parameters:
fs – The filesystem that the file belongs to.
path – S3 path (s3://bucket/key) of the file.
mode – The file mode:
rb,wb,ab, orxb.version_id – The version ID to read. Must match the version ID in the path if both are given. A version cannot be given, in either form, for writing or appending.
max_workers – The number of parallel workers for range reads and part copies.
executor – The executor for parallel operations. If None, a new
S3ThreadPoolExecutoris created.block_size – The block size for reads and writes. Must be between
MULTIPART_UPLOAD_MIN_PART_SIZEandMULTIPART_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,
commit()must be called.cache_options – Options for the fsspec cache.
size – The size of the object, if known. Passed to
fsspec.spec.AbstractBufferedFile.s3_additional_kwargs – Additional parameters for the S3 requests of the file, such as
ContentTypeorRequestPayer. Each request receives those that its operation accepts.**kwargs – Additional parameters for the S3 requests of the file, which take precedence over
s3_additional_kwargs.
- Raises:
FileExistsError – If an object exists at the path in exclusive-create mode.
FileNotFoundError – If no object exists at the path when reading, including when the path is a prefix.
ValueError – If the path has no key, the version IDs do not match, a version is given for writing, or the block size is not between
MULTIPART_UPLOAD_MIN_PART_SIZEandMULTIPART_UPLOAD_MAX_PART_SIZEfor writing.
- commit() None[source]¶
Complete the upload of the written data.
Creates an empty object if nothing was written, uploads the buffered data with PutObject if no multipart upload part was submitted, and otherwise completes the multipart upload, which is aborted if the completion fails. Invalidates the cache of the path afterwards.
- Raises:
FileExistsError – If an object was created at the path after the file was opened in exclusive-create mode.
RuntimeError – If parts were submitted but no multipart upload is initialized.
- discard() None[source]¶
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.
- url(expiration: int = 3600, **kwargs) str[source]¶
Generate a presigned HTTP URL to read this file (if it already exists).
- Parameters:
expiration – URL expiration time in seconds. Defaults to 3600 (1 hour).
**kwargs – Additional parameters passed to
S3FileSystem.sign().
- Returns:
Presigned URL string that provides temporary access to the S3 object.
- metadata(**kwargs) S3Metadata[source]¶
Return the metadata of the file.
- Parameters:
**kwargs – Additional parameters passed to the HeadObject API.
- Returns:
S3Metadata, which behaves as a read-only mapping of the user-defined metadata and exposes the system-defined metadata as typed properties.
- getxattr(xattr_name: str, **kwargs) str | None[source]¶
Get an attribute from the user-defined metadata of the file.
- Parameters:
xattr_name – The name of the attribute.
**kwargs – Additional parameters passed to the HeadObject API.
- Returns:
The value of the attribute, or None if the attribute is not set.
Async S3 FileSystem¶
- class pyathena.filesystem.s3_async.AioS3FileSystem(*args, **kwargs)[source]¶
An async filesystem interface for Amazon S3 using fsspec’s AsyncFileSystem.
This class wraps
S3FileSystemto provide native asyncio support. Instead of usingThreadPoolExecutorfor parallel operations, it usesasyncio.gatherwithasyncio.to_threadfor natural integration with the asyncio event loop.The implementation uses composition: an internal
S3FileSysteminstance handles all boto3 calls, while this class delegates to it viaasyncio.to_thread(). This avoids diamond inheritance issues and keeps all boto3 logic in one place.File handles created by
_openuseS3AioExecutorso that parallel operations (range reads, multipart uploads) are dispatched through the event loop withasyncio.to_threadinstead of aThreadPoolExecutorper file. An instance created withasynchronous=Truehas no event loop of its own, so its file handles use aThreadPoolExecutor.- _sync_fs¶
The internal synchronous S3FileSystem instance.
Example
>>> from pyathena.filesystem.s3_async import AioS3FileSystem >>> fs = AioS3FileSystem(asynchronous=True) >>> >>> # Use in async context >>> files = await fs._ls('s3://my-bucket/data/') >>> >>> # Sync wrappers (auto-generated by fsspec) need an instance created >>> # without asynchronous=True; they block the caller until done >>> files = AioS3FileSystem().ls('s3://my-bucket/data/')
- mirror_sync_methods = True¶
- async_impl = True¶
- __init__(connection: Connection[Any] | None = None, default_block_size: int | None = None, default_cache_type: str | None = None, max_workers: int = 20, s3_additional_kwargs: dict[str, Any] | None = None, allow_bucket_creation: bool = False, allow_bucket_deletion: bool = False, version_aware: bool = False, asynchronous: bool = False, loop: Any | None = None, batch_size: int | None = None, **kwargs) None[source]¶
Initialize the filesystem and its internal
S3FileSystem.- Parameters:
connection – Passed to the internal
S3FileSystem.default_block_size – Passed to the internal
S3FileSystem.default_cache_type – Passed to the internal
S3FileSystem.max_workers – Passed to the internal
S3FileSystem.s3_additional_kwargs – Passed to the internal
S3FileSystem.allow_bucket_creation – Passed to the internal
S3FileSystem.allow_bucket_deletion – Passed to the internal
S3FileSystem.version_aware – Passed to the internal
S3FileSystem.asynchronous – Passed to
fsspec.asyn.AsyncFileSystem.loop – Passed to
fsspec.asyn.AsyncFileSystem.batch_size – Passed to
fsspec.asyn.AsyncFileSystem.**kwargs – Passed to both
fsspec.asyn.AsyncFileSystemand the internalS3FileSystem.
- static parse_path(path: str) tuple[str, str | None, str | None][source]¶
Parse an S3 path into its bucket, key and version ID.
See
S3FileSystem.parse_path().- Parameters:
path – The S3 path.
- Returns:
Tuple of the bucket, the key and the version ID.
- Raises:
ValueError – If the path is not a valid S3 path.
- rmdir(path: str) None[source]¶
Remove an S3 bucket, which must be empty.
See
S3FileSystem.rmdir().- Parameters:
path – S3 bucket path (e.g., “s3://bucket”).
- sign(path: str, expiration: int = 3600, **kwargs) str[source]¶
Generate a presigned URL for S3 object access.
See
S3FileSystem.sign().- Parameters:
path – S3 path (s3://bucket/key) to generate the URL for.
expiration – URL expiration time in seconds.
**kwargs – Additional parameters passed to
S3FileSystem.sign().
- Returns:
The presigned URL.
- metadata(path: str, **kwargs) S3Metadata[source]¶
Return the metadata of the path.
See
S3FileSystem.metadata().- Parameters:
path – S3 path (s3://bucket/key) to get metadata for.
**kwargs – Additional parameters passed to the HeadObject API.
- Returns:
S3Metadata of the object.
- getxattr(path: str, attr_name: str, **kwargs) str | None[source]¶
Get an attribute from the user-defined metadata of the path.
See
S3FileSystem.getxattr().- Parameters:
path – S3 path (s3://bucket/key) to get the attribute for.
attr_name – The name of the attribute.
**kwargs – Additional parameters passed to the HeadObject API.
- Returns:
The value of the attribute, or None if the attribute is not set.
- setxattr(path: str, copy_kwargs: dict[str, Any] | None = None, **kwargs) None[source]¶
Set the user-defined metadata of the path.
See
S3FileSystem.setxattr().- Parameters:
path – S3 path (s3://bucket/key) to set metadata for.
copy_kwargs – Additional parameters to use for the underlying CopyObject API call.
**kwargs – Key-value pairs of metadata to set; a None value deletes the key.
- get_tags(path: str) dict[str, str][source]¶
Retrieve the tag key/values for the given path.
See
S3FileSystem.get_tags().- Parameters:
path – S3 path (s3://bucket/key) to get tags for.
- Returns:
Dictionary mapping tag keys to tag values.
- put_tags(path: str, tags: dict[str, str], mode: str = 'o') None[source]¶
Set the tags for the given existing key.
See
S3FileSystem.put_tags().- Parameters:
path – S3 path (s3://bucket/key) of the existing key.
tags – Tags to apply.
mode –
oto overwrite the existing tags ormto merge with them.
- chmod(path: str, acl: str, recursive: bool = False, **kwargs) None[source]¶
Set the canned ACL of a bucket or key.
See
S3FileSystem.chmod().- Parameters:
path – S3 path (s3://bucket or s3://bucket/key) to set the ACL on.
acl – The canned ACL to apply.
recursive – Whether to apply the ACL to all keys below the path too.
**kwargs – Additional parameters passed to the PutObjectAcl or PutBucketAcl API.
- object_version_info(path: str, delete_markers: bool = False, **kwargs) list[S3ObjectVersion][source]¶
List the versions of the object or of the objects under the path.
See
S3FileSystem.object_version_info().- Parameters:
path – S3 path (s3://bucket/key or a key prefix) to list the versions for.
delete_markers – Whether to include delete markers in the result.
**kwargs – Additional parameters passed to the ListObjectVersions API.
- Returns:
List of S3ObjectVersion instances describing the versions.
- list_multipart_uploads(path: str) list[S3MultipartUpload][source]¶
List in-progress (incomplete) multipart uploads in a bucket.
See
S3FileSystem.list_multipart_uploads().- Parameters:
path – S3 bucket or key path (e.g., “s3://bucket” or “s3://bucket/prefix”).
- Returns:
List of S3MultipartUpload instances describing the uploads.
- clear_multipart_uploads(path: str) None[source]¶
Abort any incomplete multipart uploads in the bucket.
See
S3FileSystem.clear_multipart_uploads().- Parameters:
path – S3 bucket or key path (e.g., “s3://bucket” or “s3://bucket/prefix”).
- checksum(path: str, **kwargs) int[source]¶
Get the checksum of an S3 object or directory.
See
S3FileSystem.checksum().- Parameters:
path – S3 path (s3://bucket/key) to get the checksum for.
**kwargs – Additional arguments passed to
S3FileSystem.checksum().
- Returns:
Integer checksum derived from the ETag or the directory token.
- created(path: str) datetime[source]¶
Return the creation time of the path.
See
S3FileSystem.created().- Parameters:
path – S3 path (s3://bucket/key).
- Returns:
The last-modified time of the object.
- modified(path: str) datetime[source]¶
Return the last-modified time of the path.
See
S3FileSystem.modified().- Parameters:
path – S3 path (s3://bucket/key).
- Returns:
The last-modified time of the object.
- invalidate_cache(path: str | None = None) None[source]¶
Remove the cached entries of the path and its parent paths.
See
S3FileSystem.invalidate_cache().- Parameters:
path – The path to invalidate. If None, clear the whole cache.
- touch(path: str, truncate: bool = True, **kwargs) dict[str, Any][source]¶
Create an empty object with PutObject.
See
S3FileSystem.touch().- Parameters:
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.
- class pyathena.filesystem.s3_async.AioS3File(fs: S3FileSystem, path: str, mode: str = 'rb', version_id: str | None = None, max_workers: int = 20, executor: S3Executor | None = None, block_size: int = 5242880, cache_type: str = 'bytes', autocommit: bool = True, cache_options: dict[Any, Any] | None = None, size: int | None = None, s3_additional_kwargs: dict[str, Any] | None = None, **kwargs)[source]¶
Async-aware S3 file handle using
S3AioExecutor.Functionally identical to
S3File; exists as a distinct type forisinstancechecks and to document the async execution model. All parallel operations (range reads, multipart uploads) are dispatched through theS3Executorinterface — theS3AioExecutorprovided byAioS3FileSystemdispatches them through the event loop withasyncio.to_threadinstead of aThreadPoolExecutorper file. For anAioS3FileSystemcreated withasynchronous=True, it is anS3ThreadPoolExecutor.
S3 Executor¶
- class pyathena.filesystem.s3_executor.S3Executor[source]¶
Abstract executor for parallel S3 operations.
Defines the interface used by
S3FileandS3FileSystemfor submitting work to run in parallel and for shutting down the executor when done. Bothsubmitandshutdownmirror theconcurrent.futures.Executorinterface so thatas_completed()andFuture.cancel()work unchanged.
- class pyathena.filesystem.s3_executor.S3ThreadPoolExecutor(max_workers: int)[source]¶
Executor that delegates to a
ThreadPoolExecutor.This is the default executor used by
S3FileandS3FileSystemfor synchronous parallel operations.- __init__(max_workers: int) None[source]¶
Initialize the executor with a new
ThreadPoolExecutor.- Parameters:
max_workers – The maximum number of threads of the thread pool.
- class pyathena.filesystem.s3_executor.S3AioExecutor(loop: AbstractEventLoop | None = None, max_workers: int = 20)[source]¶
Executor that schedules work on an asyncio event loop.
Uses
asyncio.run_coroutine_threadsafe(asyncio.to_thread(fn), loop)to dispatch blocking functions onto the event loop’s thread pool, returningconcurrent.futures.Futureobjects that are compatible withas_completed(),wait()andFuture.cancel(). As withThreadPoolExecutor, a future cannot be cancelled once its function has started. At mostmax_workersof the submitted functions run at once.This avoids thread-in-thread nesting when
S3Fileis used from withinasyncio.to_thread()calls (the pattern used byAioS3FileSystem).- Parameters:
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
submitis called.
- __init__(loop: AbstractEventLoop | None = None, max_workers: int = 20) None[source]¶
Initialize the executor with the event loop to schedule work on.
- Parameters:
loop – The asyncio event loop.
submitraisesRuntimeErrorif it is None or not running.max_workers – The maximum number of submitted functions that run at once.
- Raises:
ValueError – If
max_workersis not positive.
S3 Objects¶
- class pyathena.filesystem.s3_object.S3Object(init: dict[str, Any], **kwargs)[source]¶
Represents an S3 object with metadata and filesystem-like properties.
This class provides a dictionary-like interface to S3 object metadata, making it easier to work with S3 objects in filesystem operations. It handles the mapping between S3 API field names and more pythonic property names.
The object supports both dictionary-style access and property-style access to metadata fields like content type, storage class, encryption settings, and object lock configurations. Dictionary-style access behaves like a dictionary, so a missing key raises KeyError. Property-style access returns None for a known field that the object does not have, and raises AttributeError for any other missing name.
Example
>>> s3_obj = S3Object({"ContentType": "text/csv", "ContentLength": 1024}) >>> print(s3_obj.content_type) # "text/csv" >>> print(s3_obj["content_length"]) # 1024 >>> s3_obj.storage_class = "STANDARD_IA"
Note
This class is primarily used internally by S3FileSystem for representing S3 objects in filesystem operations.
- __init__(init: dict[str, Any], **kwargs) None[source]¶
Initialize the object from an S3 API response.
Only the fields of
initthat have a property mapping are kept, under their property names (e.g.,ContentType->content_type).storage_classdefaults toSTANDARDwheninithas noStorageClass, andsizeis taken fromSizeorContentLength.nameis set tobucket/key, or to the bucket when there is no key.- Parameters:
init – An S3 API response or listing entry, such as a HeadObject response or a ListObjectsV2
Contentsentry.**kwargs – Additional fields stored as-is, such as
type,bucket,keyandversion_id. S3 API field names are stored under their property names.
- __getattr__(item: str) Any[source]¶
Return None for a known field that the object does not have.
Called only when normal attribute lookup fails, so fields that the object has are returned without reaching this method.
- Parameters:
item – The attribute name.
- Returns:
None, if
itemis a known S3 object field.- Raises:
AttributeError – If
itemis not a known S3 object field.
- copy() S3Object[source]¶
Return a shallow copy of the object.
- Returns:
A new S3Object with the same fields.
- class pyathena.filesystem.s3_object.S3ObjectType[source]¶
Constants for S3 object types in filesystem operations.
These constants are used to distinguish between directories and files when working with S3 paths through the S3FileSystem interface.
- class pyathena.filesystem.s3_object.S3StorageClass[source]¶
Constants for Amazon S3 storage classes.
S3 storage classes determine the availability, durability, and cost characteristics of stored objects. Each class is optimized for different access patterns and use cases.
- Storage classes:
STANDARD: Default storage for frequently accessed data
REDUCED_REDUNDANCY: Lower cost, reduced durability (deprecated)
STANDARD_IA: Infrequently accessed data with rapid retrieval
ONEZONE_IA: Lower cost IA storage in single availability zone
INTELLIGENT_TIERING: Automatic tiering between frequent/infrequent
GLACIER: Archive storage for long-term backup
DEEP_ARCHIVE: Lowest cost archive storage
GLACIER_IR: Archive with faster retrieval than standard Glacier
OUTPOSTS: Storage on AWS Outposts
BUCKET: Pseudo storage class PyAthena assigns to bucket entries
DIRECTORY: Pseudo storage class PyAthena assigns to directory entries
See also
AWS S3 storage classes documentation: https://docs.aws.amazon.com/s3/latest/userguide/storage-class-intro.html
S3 Upload Operations¶
- class pyathena.filesystem.s3_object.S3PutObject(response: dict[str, Any])[source]¶
Represents the response from an S3 PUT object operation.
This class encapsulates the metadata returned when uploading an object to S3, including encryption details, versioning information, and integrity checksums.
- expiration¶
Object expiration time if lifecycle policy applies.
- version_id¶
Version ID if bucket versioning is enabled.
- etag¶
Entity tag for the uploaded object.
- server_side_encryption¶
Server-side encryption method used.
- Various checksum properties
For data integrity verification.
Note
This class is used internally by S3FileSystem operations and typically not instantiated directly by users.
- __init__(response: dict[str, Any]) None[source]¶
Initialize the result from a PutObject response.
- Parameters:
response – The PutObject response.
- property server_side_encryption: str | None¶
The
ServerSideEncryptionalgorithm of the uploaded object.
- property sse_customer_key_md5: str | None¶
The
SSECustomerKeyMD5of the customer-provided key for the uploaded object.
- property sse_kms_encryption_context: str | None¶
The
SSEKMSEncryptionContextof the uploaded object.
- class pyathena.filesystem.s3_object.S3MultipartUpload(response: dict[str, Any])[source]¶
Represents an S3 multipart upload operation.
This class manages the metadata for multipart uploads, which allow uploading large files in chunks for better reliability and performance. It tracks upload identifiers, encryption settings, and lifecycle rules.
- bucket¶
S3 bucket name for the upload.
- key¶
Object key being uploaded.
- upload_id¶
Unique identifier for the multipart upload.
- server_side_encryption¶
Encryption method applied to the upload.
- abort_date/abort_rule_id
Lifecycle rule information for upload cleanup.
- initiated/storage_class/owner/initiator
Fields returned by the ListMultipartUploads API for in-progress uploads.
Note
Used internally by S3FileSystem for large file upload operations, and returned by
S3FileSystem.list_multipart_uploads.- __init__(response: dict[str, Any]) None[source]¶
Initialize the upload from an S3 API response.
- Parameters:
response – A CreateMultipartUpload response or an
Uploadsentry of a ListMultipartUploads response.
- property abort_rule_id: str | None¶
The
AbortRuleIdof the lifecycle rule that applies to the upload.
- property sse_customer_key_md5: str | None¶
The
SSECustomerKeyMD5of the customer-provided key for the upload.
- property bucket_key_enabled: bool | None¶
Whether the upload uses an S3 Bucket Key (
BucketKeyEnabled).
- property initiated: datetime | None¶
The
Initiatedtime of the upload, returned by ListMultipartUploads.
- class pyathena.filesystem.s3_object.S3MultipartUploadPart(part_number: int, response: dict[str, Any])[source]¶
Represents a single part in an S3 multipart upload operation.
Each part in a multipart upload has its own metadata including checksums, encryption details, and part identification. This class manages that metadata and provides methods to convert it to API-compatible formats.
- part_number¶
The sequential part number (1-based).
- etag¶
Entity tag for this specific part.
- checksum_*
Various integrity checksums for the part data.
- server_side_encryption¶
Encryption settings for this part.
Note
Parts must be at least 5MB except for the last part. Used internally by S3FileSystem for chunked upload operations.
- __init__(part_number: int, response: dict[str, Any]) None[source]¶
Initialize the part from an UploadPart or UploadPartCopy response.
For an UploadPartCopy response, the
ETag,LastModifiedand checksums are read from itsCopyPartResult.- Parameters:
part_number – The part number of the part.
response – The UploadPart or UploadPartCopy response.
- property copy_source_version_id: str | None¶
The
CopySourceVersionIdof the source object of a copied part.
- property last_modified: datetime | None¶
The
LastModifiedtime fromCopyPartResult; None for uploaded parts.
- property sse_customer_key_md5: str | None¶
The
SSECustomerKeyMD5of the customer-provided key for the part.
- class pyathena.filesystem.s3_object.S3CompleteMultipartUpload(response: dict[str, Any])[source]¶
Represents the completion of an S3 multipart upload operation.
This class encapsulates the final response when a multipart upload is completed, including the final object location, versioning information, and consolidated metadata from all parts.
- location¶
Final S3 URL of the completed object.
- bucket¶
S3 bucket containing the object.
- key¶
Final object key.
- version_id¶
Version ID if bucket versioning is enabled.
- etag¶
Final entity tag of the complete object.
- server_side_encryption¶
Encryption applied to the final object.
Note
This represents the successful completion of a multipart upload. Used internally by S3FileSystem operations.
- __init__(response: dict[str, Any]) None[source]¶
Initialize the result from a CompleteMultipartUpload response.
- Parameters:
response – The CompleteMultipartUpload response.
- property server_side_encryption: str | None¶
The
ServerSideEncryptionalgorithm of the completed object.