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
9 changes: 6 additions & 3 deletions docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,9 @@ block size (5 MiB by default); larger data is uploaded as a parallel multipart u
through the buffered file path. Inside an
[fsspec transaction](https://filesystem-spec.readthedocs.io/en/latest/features.html#transactions),
writes are deferred until the transaction commits and are discarded on rollback.
With the `compression` argument (a codec of fsspec, or `"infer"` from the extension of

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, operations): FINDINGS (text only, repaired)

Base fe21250a3e228e3d38e61898334c5baacc1285c2, head 846227079ea492e3edddb55ea7fd4e0911d1bee8.

Claims checked:

  • Single-request path on master raises ParamValidationError: Unknown parameter in input: "compression": the PutObject request on master holds compression (_make_api_call captured), and botocore validate_parameters against the PutObject input shape rejects it. Not sent to live S3.
  • Failed compressed write on master: raises AttributeError: 'GzipFile' object has no attribute '_close_without_commit' with the original error as context, and uploads an object, both outside a transaction (incompressible-size path) and inside one, checked offline on fe21250a.
  • "7 cases fail with the source reverted": rerun with only pyathena/ reverted.
  • Live round trip (ad hoc script, objects deleted): compression="infer" on a .gz key (single request) and compression="gzip" on 6 MiB of random bytes (multipart), in and outside a transaction, decompress to the values.
  • Docstrings of pipe_file() and _pipe_file_in_transaction() match the code; this docs sentence matches the routing on the compressed size.

Finding (repaired): the issue, the PR body, and the relayed review on #1035 said a failed write uploaded "an empty gzip object (10 bytes)". The object holds only the 10-byte gzip header, which is not a valid gzip file; an empty gzip file is 20 bytes. All three now say "a 10-byte object holding only the gzip header".

Callers: fsspec pipe() passes compression per file through pipe_file(); AioS3FileSystem.pipe_file() outside a transaction delegates to the sync method. Operationally: no extra S3 requests; small compressed values now take one PutObject instead of failing.

the path), `pipe`/`pipe_file` compress the data before uploading it, and the sizes in
this section apply to the compressed data.

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
Expand All @@ -105,9 +108,9 @@ Paths are normalized as in fsspec, which drops a trailing slash, so `info`, `isf
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
`dir/` for reading reads the object `dir` or raises `FileNotFoundError`. Opening it for
writing, `pipe`, `pipe_file`, and `put_file` write 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.
Expand Down
66 changes: 59 additions & 7 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,8 @@
from botocore.client import BaseClient, Config
from fsspec import AbstractFileSystem
from fsspec.callbacks import _DEFAULT_CALLBACK, Callback
from fsspec.compression import compr
from fsspec.core import get_compression
from fsspec.implementations.local import trailing_sep
from fsspec.spec import AbstractBufferedFile
from fsspec.utils import isfilelike, other_paths, tokenize
Expand All @@ -44,11 +46,47 @@
S3PutObject,
S3StorageClass,
)
from pyathena.util import RetryConfig, retry_api_call
from pyathena.util import RetryConfig, override, retry_api_call

_logger = logging.getLogger(__name__)


class CompressedBuffer(BytesIO):
"""An in-memory buffer of data compressed with a codec of fsspec.

The buffer keeps its data when it is closed, as some codecs, such as
the ``zstandard`` stream writer, close the file that they write to.
"""

@override
def close(self) -> None:

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 of the full PR at edc338db (relayed): finding 2 of 2 (reviewer and scope: see the comment on s3.py:1816)

P3 — The public close() override lacks its changed-behavior docstring. CompressedBuffer.close() leaves the buffer open: closed remains false and subsequent writes remain possible. Method-level help can inherit the incompatible BytesIO.close() description. docs/contributing.md:85 requires a docstring when an override changes behavior. Document the no-op semantics directly on this method.

Repaired in d8ec2207678815b2dee85386d379edd523fd804d: """Do nothing, so that the data can still be read after a codec closes the buffer.""". Self-review: docstring-only change; round 1 no behavior change, round 2 the sentence matches the code. just lint passed; TestCompressedBuffer 7 passed.

"""Do nothing, so that the data can still be read after a codec closes the buffer."""

@classmethod
def compress(cls, value: bytes | bytearray | memoryview, compression: str) -> bytes:

@laughingman7743 laughingman7743 Oct 3, 2026 •

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 (rounds 1 and 2) of edc338db — public CompressedBuffer API: CLEAN

Scope: git diff 407ec3ed..edc338db (base fe21250a3e228e3d38e61898334c5baacc1285c2). The maintainer asked to redesign the private _compress() / _CompressedBuffer as a public API: CompressedBuffer.compress(value, compression) -> bytes takes a codec name only, and the callers resolve the codec as open() does. The PR was returned to Draft for this change.

Round 1 (behavior):

  • pipe_file() and aio _pipe_file_in_transaction() call get_compression(self._strip_protocol(path), kwargs.pop("compression", None)). None → None, and "infer" without a codec → None, so nothing is compressed, as before. An unknown name raises ValueError("Compression type … not supported") from get_compression() before any request, as before.
  • compress() validates the name itself for direct callers ("infer" is not a codec name and raises), converts a non-contiguous memoryview, writes through the codec into cls(), and always returns bytes.
  • s3_async imports CompressedBuffer instead of the private _compress.
  • While verifying the closing-codec test, a git checkout -- pyathena/filesystem/s3.py reverted the uncommitted s3.py part of this change; it was re-applied by the same scripted edit, the diff was re-read, and the suites were rerun before the commit.

Round 2 (claims): "public, not in the API reference like the other support classes": docs/api/filesystem.rst lists none of S3ClientError, S3Metadata, S3ObjectVersion. Docstrings match the code; PR body updated (WHAT, tests, tested commit).

Validation on edc338db: just lint passed; tests/pyathena/filesystem/ 561 passed against live S3; TestCompressedBuffer 7 cases, of which test_compress_codec_closing_its_file fails without the close() override; with zstandard and lz4 installed, all 7 registered codecs round-trip through CompressedBuffer.compress().

"""Compress a value with a codec of fsspec.

Args:
value: The bytes to compress.
compression: Name of a codec in ``fsspec.compression.compr``.

Returns:
The compressed bytes.

Raises:
ValueError: If the codec is not supported.
"""
if compression not in compr:
raise ValueError(f"Compression type {compression} not supported")
if isinstance(value, memoryview) and not value.c_contiguous:
# Codecs cannot compress a non-contiguous memoryview.
value = value.tobytes()
buffer = cls()
with compr[compression](buffer, mode="w") as f:
f.write(value)
return buffer.getvalue()


class S3FileSystem(AbstractFileSystem):
"""A filesystem interface for Amazon S3 that implements the fsspec protocol.

Expand Down Expand Up @@ -1748,7 +1786,8 @@ def pipe_file(
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. A write that
fails on that path leaves the existing object unchanged.
fails on that path leaves the existing object unchanged. Both paths
write to the path without a trailing slash, as ``open()`` does.

Args:
path: S3 path (s3://bucket/key) to write to.
Expand All @@ -1760,20 +1799,33 @@ def pipe_file(
(e.g., ContentType, StorageClass) on the single-request
path. The ``block_size``, ``max_workers``, and
``s3_additional_kwargs`` parameters of the ``open()`` path
are also accepted.
are also accepted, and so is ``compression``: the codec of
``open()`` to compress the value with before it is
uploaded, or ``"infer"`` to take it from the extension of
the path.

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_PARTS`` blocks.
version, if the compression is not supported, or if the data
takes more than ``MULTIPART_UPLOAD_MAX_PARTS`` blocks.
"""
# Normalized as open() normalizes it, so that the key written, and
# the codec that "infer" takes from its extension, do not depend on
# the path that the size of the value selects.
path = self._strip_protocol(path)
compression = get_compression(path, kwargs.pop("compression", None))
if compression is not None:
# Compressed up front, so that every path uploads the compressed
# bytes, and open() returns the file instead of a wrapper.
value = CompressedBuffer.compress(value, compression)
block_size = kwargs.get("block_size") or self.default_block_size
# The size in bytes; the length of a memoryview counts its items.
self._check_multipart_upload_size(path, memoryview(value).nbytes, block_size)
if self._intrans or len(value) > min(block_size, self.MULTIPART_UPLOAD_MAX_PART_SIZE):
size = memoryview(value).nbytes
self._check_multipart_upload_size(path, size, block_size)
if self._intrans or size > min(block_size, self.MULTIPART_UPLOAD_MAX_PART_SIZE):
# Defer to the buffered open() path, which keeps the
# deferred-commit semantics of fsspec transactions and uploads
# large data as a parallel multipart upload.
Expand Down
15 changes: 11 additions & 4 deletions pyathena/filesystem/s3_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,9 @@

from fsspec.asyn import AsyncFileSystem, sync
from fsspec.callbacks import _DEFAULT_CALLBACK
from fsspec.core import get_compression

from pyathena.filesystem.s3 import S3File, S3FileSystem
from pyathena.filesystem.s3 import CompressedBuffer, S3File, S3FileSystem
from pyathena.filesystem.s3_executor import S3AioExecutor, S3Executor, S3ThreadPoolExecutor
from pyathena.filesystem.s3_object import (
S3Metadata,
Expand Down Expand Up @@ -187,14 +188,20 @@ def _pipe_file_in_transaction(
opened in ``xb`` mode: raise FileExistsError when the object
already exists, including one created before the
transaction is committed, which is not replaced.
**kwargs: Additional parameters passed to ``open()``.
**kwargs: Additional parameters passed to ``open()``, except
``compression``, with which the value is compressed before
it is written, as in :meth:`S3FileSystem.pipe_file`.

Raises:
FileExistsError: If the mode is "create" and the path already
exists.
ValueError: If the data takes more than
``MULTIPART_UPLOAD_MAX_PARTS`` blocks.
ValueError: If the compression is not supported, or if the data
takes more than ``MULTIPART_UPLOAD_MAX_PARTS`` blocks.
"""
# See S3FileSystem.pipe_file.
compression = get_compression(self._strip_protocol(path), kwargs.pop("compression", None))
if compression is not None:
value = CompressedBuffer.compress(value, compression)
block_size = kwargs.get("block_size") or self._sync_fs.default_block_size
# The size in bytes; the length of a memoryview counts its items.
self._sync_fs._check_multipart_upload_size(path, memoryview(value).nbytes, block_size)
Expand Down
Loading
Loading