-
Notifications
You must be signed in to change notification settings - Fork 116
Compress pipe_file() values up front, and route and key them consistently #1038
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
e25850c
ebb855d
8462270
135e648
6d90268
407ec3e
edc338d
d8ec220
73cdb62
67a3db4
c6297d3
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 | ||
|
|
@@ -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: | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent review of the full PR at
Repaired in |
||
| """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: | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Self-review (rounds 1 and 2) of Scope: Round 1 (behavior):
Round 2 (claims): "public, not in the API reference like the other support classes": Validation on |
||
| """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. | ||
|
|
||
|
|
@@ -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. | ||
|
|
@@ -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. | ||
|
|
||
There was a problem hiding this comment.
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, head846227079ea492e3edddb55ea7fd4e0911d1bee8.Claims checked:
ParamValidationError: Unknown parameter in input: "compression": the PutObject request on master holdscompression(_make_api_callcaptured), and botocorevalidate_parametersagainst the PutObject input shape rejects it. Not sent to live S3.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 onfe21250a.pyathena/reverted.compression="infer"on a.gzkey (single request) andcompression="gzip"on 6 MiB of random bytes (multipart), in and outside a transaction, decompress to the values.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()passescompressionper file throughpipe_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.