S3 filesystem¶
PyAthena ships its own fsspec-compatible
filesystem implementation for Amazon S3 (S3FileSystem), built on boto3, with an API
surface compatible with s3fs for users migrating from it.
The filesystem is used internally by the pandas, Polars, and S3FS result sets to read query results from S3, and can also be used independently for S3 file operations.
fsspec registration¶
Importing pyathena.pandas or pyathena.polars registers S3FileSystem as the fsspec
s3 / s3a protocols via pyathena.filesystem.register_s3_filesystem. This replaces
fsspec’s default lazy mapping of the s3 protocol to s3fs, which means
fsspec.filesystem("s3") returns PyAthena’s implementation and s3fs-specific settings
(such as the S3FS_LOGGING_LEVEL environment variable) have no effect.
A filesystem class that has already been registered explicitly is also overwritten, with a warning log, so that the replacement is diagnosable. To restore another implementation, re-register it afterwards:
import fsspec
import s3fs
import pyathena.pandas # Registers PyAthena's S3FileSystem.
fsspec.register_implementation("s3", s3fs.S3FileSystem, clobber=True)
Basic usage¶
The filesystem can be constructed from a PyAthena connection, or directly with s3fs-compatible credential arguments:
from pyathena import connect
from pyathena.filesystem.s3 import S3FileSystem
fs = S3FileSystem(connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
region_name="us-west-2"))
# Or with direct credentials (s3fs-compatible arguments).
fs = S3FileSystem(key="YOUR_ACCESS_KEY", secret="YOUR_SECRET_KEY")
# Or anonymously for public buckets.
fs = S3FileSystem(anon=True)
Standard fsspec operations work as expected:
fs.ls("s3://YOUR_S3_BUCKET/path/to/")
fs.find("s3://YOUR_S3_BUCKET/path/to/")
fs.exists("s3://YOUR_S3_BUCKET/path/to/object")
fs.info("s3://YOUR_S3_BUCKET/path/to/object")
with fs.open("s3://YOUR_S3_BUCKET/path/to/object", "rb") as f:
data = f.read()
fs.pipe("s3://YOUR_S3_BUCKET/path/to/object", b"data")
fs.cat("s3://YOUR_S3_BUCKET/path/to/object")
fs.cp("s3://YOUR_S3_BUCKET/src", "s3://YOUR_S3_BUCKET/dst")
fs.rm("s3://YOUR_S3_BUCKET/path/to/", recursive=True)
Writes with pipe/pipe_file issue a single PutObject request for data up to the
block size (5 MiB by default); larger data is uploaded as a parallel multipart upload
through the buffered file path. Inside an
fsspec transaction,
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.
A multipart upload consists of at most 10,000 parts. put and pipe upload one part
per block, so with the default block size they can upload up to about 48.8 GiB
(10,000 × 5 MiB). To upload a larger object, use a block size of at least its size
divided by 10,000, either with the block_size argument of put, pipe, and open
or with the default_block_size argument of S3FileSystem. put and pipe check
the size before uploading anything and raise ValueError with the minimum block size
if the data needs more parts. A file written with open can take more parts, because
each write that fills the buffer uploads the data beyond its last full block as a
separate part when that data is at least 5 MiB. In an append, the parts copied from
the existing object also count toward the limit. A write with open that reaches the
limit raises ValueError and aborts its multipart upload. Multipart copies with cp
use parts large enough to stay within the limit.
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
slash, such as a folder marker, is therefore not a file for these methods. Opening
dir/ for reading reads the object dir or raises FileNotFoundError, and opening it
for writing writes the object dir. A path with a ?versionId= suffix keeps the slash
and refers to the object. find, and ls of the directory, list the object as a file
entry. cat_file uses the key as written. Without a ?versionId= suffix, it reads
such an object without a range, with a non-empty range of non-negative offsets, or with
a negative start and no end, and raises FileNotFoundError for other ranges.
S3 request parameters, such as ContentType, ServerSideEncryption, or RequestPayer,
can be given to open, pipe, and put as keyword arguments or in
s3_additional_kwargs, and to all of them through the s3_additional_kwargs argument
of S3FileSystem. The parameters of a call take precedence over those of the
filesystem. A file sends each of its requests only the parameters that the S3 operation
accepts, so, for example, ServerSideEncryption for writes is not sent with reads. A
pipe of data up to the block size sends its parameters with a single PutObject
request as given. put sets ContentType from the file extension unless the call or
the filesystem gives one.
fs = S3FileSystem(s3_additional_kwargs={"ServerSideEncryption": "AES256"})
with fs.open("s3://YOUR_S3_BUCKET/path/to/data.csv", "wb", ContentType="text/csv") as f:
f.write(b"col1\n1\n")
Error translation¶
S3 error responses are translated into standard Python exceptions, so filesystem
operations raise natural errors instead of botocore’s ClientError:
S3 error |
Python exception |
|---|---|
|
|
|
|
|
|
|
|
|
|
Others |
|
Multipart upload management¶
Incomplete multipart uploads continue to accrue storage costs until they are completed or aborted. The filesystem can discover and abort them:
uploads = fs.list_multipart_uploads("s3://YOUR_S3_BUCKET")
for upload in uploads:
print(upload.key, upload.upload_id, upload.initiated)
# Abort all incomplete uploads to a key and the keys under it.
fs.clear_multipart_uploads("s3://YOUR_S3_BUCKET/path/to/")
Versioning¶
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.
fs = S3FileSystem(
connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/", region_name="us-west-2"),
version_aware=True,
)
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)
# Typed version information, including delete markers if requested.
versions = fs.object_version_info("s3://YOUR_S3_BUCKET/path/to/object")
for version in versions:
print(version.version_id, version.is_latest, version.last_modified)
Version-aware operations require the s3:GetObjectVersion and
s3:ListBucketVersions permissions.
Bucket lifecycle¶
Bucket creation and deletion are infrastructure-level changes and are disabled by
default: mkdir/makedirs and rmdir raise PermissionError when they would
create or delete a bucket. Pass the opt-in flags to enable them:
fs = S3FileSystem(
connect(s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/", region_name="us-west-2"),
allow_bucket_creation=True,
allow_bucket_deletion=True,
)
fs.mkdir("s3://YOUR_NEW_BUCKET")
fs.rmdir("s3://YOUR_NEW_BUCKET") # The bucket must be empty.
Creating a key prefix under an existing bucket requires no operation (S3 has no real directories below the bucket level) and is always a no-op.
Async filesystem¶
AioS3FileSystem provides the same functionality on top of fsspec’s
AsyncFileSystem, dispatching parallel operations through the asyncio event loop.
See AioS3FileSystem for details.