Compress pipe_file() values up front, and route and key them consistently - #1038
Conversation
pipe_file() passed fsspec's compression argument to PutObject on the single-request path, which botocore rejects, and to open() on the buffered path, which returns a compression wrapper. A failed write on the buffered path then raised AttributeError from _close_without_commit() on the wrapper, hiding the original error, and still uploaded an empty compressed object. Compress the value in memory with the codec that open() uses, including "infer" from the path, and upload the compressed bytes through the usual paths. The same applies to the transaction path of AioS3FileSystem. Closes #1037 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The zstandard stream writer, fsspec's zstd codec without compression.zstd or backports.zstd, closes the file that it writes to, so reading the compressed bytes from the BytesIO afterwards raised ValueError. Write to a BytesIO whose close() keeps the data. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
| _logger = logging.getLogger(__name__) | ||
|
|
||
|
|
||
| class _CompressedBuffer(BytesIO): |
There was a problem hiding this comment.
Self-review round 1 (behavior and implementation): FINDINGS (repaired)
Base fe21250a3e228e3d38e61898334c5baacc1285c2 (merge-base with master). Initial pass on e25850c869662e5e938b486cbfdf923d3784894f; repair ebb855d7cc6759eec2981f00f5d2601eb4bd27c7, rechecked on 846227079ea492e3edddb55ea7fd4e0911d1bee8.
Covered: _compress(), S3FileSystem.pipe_file(), AioS3FileSystem._pipe_file_in_transaction() (and _pipe_file() outside a transaction, which delegates to the sync pipe_file()), the new tests.
Finding (repaired in ebb855d7): _compress() wrote to a plain BytesIO. fsspec's zstd codec falls back to the zstandard stream writer when neither compression.zstd nor backports.zstd is available (fsspec 2026.9.0 compression.py:158-177). That writer closes the file it writes to, so buffer.getvalue() raised ValueError: I/O operation on closed file. Reproduced with uv run --with zstandard. _CompressedBuffer keeps the data on close(). test_compress_codec_closing_its_file fails with a plain BytesIO. With zstandard and lz4 installed, every registered codec (zip, bz2, gzip, lzma, xz, lz4, zstd) round-trips through _compress() and the codec's reader.
Checked, no finding:
compressionis popped before the size check, so the routing and_check_multipart_upload_size()use the compressed size, andcompressionreaches neitheropen()nor PutObject.compression=Nonegiven explicitly is popped and ignored, asopen()ignored it."infer"without a codec returns the value unchanged, so it takes the uncompressed paths as before.- A non-contiguous memoryview is converted before compressing (codecs need a contiguous buffer).
- Memory: the value is already in memory; the compressed copy is one more buffer, the same as the wrapper's buffering on the old path.
put_file()never passedcompressiontoopen()and is unchanged.
Validation: just lint passed; tests/pyathena/filesystem/ 545 passed against live S3 on ebb855d7.
| 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 |
There was a problem hiding this comment.
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 holdscompression(_make_api_callcaptured), and botocorevalidate_parametersagainst 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 onfe21250a. - "7 cases fail with the source reverted": rerun with only
pyathena/reverted. - Live round trip (ad hoc script, objects deleted):
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. - 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.
fsspec's open() infers the codec from the path without the protocol, which also drops a trailing slash. pipe_file() inferred it from the path as given, so "key.gz/" was uploaded uncompressed. Strip the protocol first, also in the transaction path of AioS3FileSystem. Cover compressed multipart uploads and non-contiguous memoryviews too. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
| Raises: | ||
| ValueError: If the codec is not supported. | ||
| """ | ||
| compression = get_compression(path, compression) |
There was a problem hiding this comment.
Independent review (relayed): FINDINGS — regression
Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only, model reported gpt-6-astra), session 01a10244-528d-7e50-8ab1-84e0bb7d40a2. Static review only. Base fe21250a3e228e3d38e61898334c5baacc1285c2, head 846227079ea492e3edddb55ea7fd4e0911d1bee8, detached snapshot (unchanged afterwards). The prompt held the diff, the intent, and conventions, without the PR number, PR text, or prior findings.
Covered: sync and async pipe/pipe_file; single-request and multipart uploads; transaction commit/rollback; mode="create"; block sizes and failure cleanup; bytes, bytearray, contiguous/non-contiguous and non-byte-format memoryviews; compression None, "infer", unsupported names; all codec registration branches incl. writers that close their buffer; changed tests, docstrings, docs, conventions.
P2 — Compression inference differs from
open()for trailing-slash paths.s3.py:84infers compression from the unnormalized path. fsspecopen()first strips trailing slashes, then infers compression. Consequently, inside a transaction,fs.pipe_file("s3://bucket/key.gz/", b"data", compression="infer")now stores uncompressed bytes atkey.gz; previously it stored gzip data. Reading that object with inferred compression fails. Both sync and async transaction paths are affected. Infer from the same normalized path thatopen()uses, and add a regression case for it.The added tests assert observable uploaded bytes, keyword removal, and absence of uploads after write failures. [...] However, every successful compression case fits within one compressed block, so successful compressed multipart uploads and compressed memoryview variants remain uncovered.
Author verification: confirmed. fsspec 2026.9.0 AbstractFileSystem.open() calls self._strip_protocol(path) before get_compression(path, compression), and _strip_protocol("s3://bucket/key.gz/") is bucket/key.gz; on 84622707 the transaction write sent b"data" uncompressed to key.gz.
There was a problem hiding this comment.
Repaired in 135e64872fad15cc949a5f659f4a71f6fec49cc4: pipe_file() and AioS3FileSystem._pipe_file_in_transaction() pass self._strip_protocol(path) to _compress(), and its docstring says so. test_pipe_file_compression gains a "s3://bucket/key.gz/" + "infer" case (4 parametrizations, which fail without the repair). The coverage gap is closed by test_pipe_file_compression_multipart (compressed data above the block size, in and outside a transaction) and test_pipe_file_compression_non_contiguous_memoryview.
Self-review of the repair:
- Round 1 (behavior):
_strip_protocol()is fsspec's default for both filesystems (neither overrides it) and is whatopen()calls; a?versionId=suffix stays in the path, as inopen(), so inference matches there too. Only"infer"uses the path. - Round 2 (claims): PR body updated (inference path, tests, tested commit, 552 passed).
- Validation:
just lintpassed;tests/pyathena/filesystem/552 passed against live S3 on135e6487.
There was a problem hiding this comment.
Independent follow-up review (relayed): FINDINGS (test coverage only)
Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only, model reported gpt-6-astra), session 01a1024b-b4ed-7260-9acb-50ada0979845. Static review only; scope git range-diff fe21250a..84622707 fe21250a..135e6487 and git diff 84622707..135e6487; snapshot at 135e64872fad15cc949a5f659f4a71f6fec49cc4, unchanged afterwards. Reported scope deviation: the reviewer read local review guidance and searched the agent memory index outside the given roots (no matches); it changed nothing.
Earlier point (1): resolved across all three routes. Both compression callers now use the same
_strip_protocol()normalization asopen(); nontransaction aio delegates to sync. This covers bare paths,s3,s3a, and trailing slashes. Query suffixes remain in the inference input, matchingopen(); version-qualified writes are rejected.
Earlier point (2): addressed for sync. [...] These assert observable behavior.P3 — Remaining test-coverage gap.
test_s3_async.py:281still tests only small bytes with explicitgzip. [...] Reverting justs3_async.py:203would leave these compression tests passing while an aio transaction writings3://bucket/key.gz/withcompression="infer"again stores uncompressed data atkey.gz. Add corresponding aio cases.
Repaired in 6d90268d601bd27f631049e1d72eb0e23325f451 (tests only): aio test_pipe_file_compression gains the "s3://bucket/key.gz/" + "infer" case, in and outside a transaction. With only the aio call reverted to the unstripped path, the transaction case fails. Self-review: round 1, test-only change, no production code touched; round 2, PR body updated (aio test cases, tested commit). just lint passed; tests/pyathena/filesystem/ 554 passed against live S3 on 6d90268d.
There was a problem hiding this comment.
Independent follow-up review of the test-only repair (relayed): CLEAN
Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only), session 01a10253-7f50-71c0-a134-43a7e597c557. Static review only; scope git diff 135e6487..6d90268d; snapshot at 6d90268d601bd27f631049e1d72eb0e23325f451, unchanged afterwards.
CLEAN. The trailing-slash inference case runs both inside and outside a transaction. Using the unstripped path in the aio helper would leave
b"data"uncompressed, causing the transaction case'sgzip.decompressassertion to fail. Assertions verify exactly one upload, no forwardedcompressionkeyword, and an uploaded gzip payload that decompresses tob"data".
| version, if the compression is not supported, or if the data | ||
| takes more than ``MULTIPART_UPLOAD_MAX_PARTS`` blocks. | ||
| """ | ||
| compression = kwargs.pop("compression", None) |
There was a problem hiding this comment.
Independent review (relayed): pre-existing issue (reviewer and scope: see the comment on pyathena/filesystem/s3.py:84)
P2 — Typed memoryviews can bypass multipart routing.
s3.py:1828compareslen(value)—an item count—with byte thresholds. Outside transactions, a 6 GiB memoryview cast to"I", withblock_size=5 GiB, passes the size guard but takes the single PutObject path, exceeding its size limit. This remains reachable with compression omitted,None, or"infer"resolving to no codec. Use the byte count for routing too.
Author: pre-existing and unchanged by this PR: the routing by len(value) predates it (kept in #1003), and a compressed value is bytes, so it is not affected. Deferred out of #1037's scope; to be reported to the maintainer.
There was a problem hiding this comment.
Fixed here at the maintainer's request, in 407ec3edcd953f420b05c877eccda1abd7b7444d: pipe_file() computes size = memoryview(value).nbytes once and uses it for both _check_multipart_upload_size() and the routing. The PR was returned to Draft for this change.
test_pipe_file_memoryview_routed_by_bytes: acast("I")memoryview of block size + 4 bytes goes to the multipart path, and the uploaded parts join to the data. It fails with thelen(value)routing.test_pipe_file_non_contiguous_memoryviewpinned the item count (4 items of 2 bytes,block_size=6→ PutObject). It now usesblock_size=8, keeping its purpose: a non-contiguous memoryview within the block size is converted to bytes for PutObject.
Self-review of the change:
- Round 1 (behavior): for
bytesandbytearray,nbytes == len(), so their routing is unchanged; only multi-byte memoryviews change. A non-contiguous view'snbytesis the size of itstobytes(). On the buffered path,BytesIO.write()takes any buffer and returns the byte count, as the new test shows. The aio transaction path is always buffered and has no routing; aio outside a transaction delegates to this method. - Round 2 (claims): PR title, WHAT, release notes, and TEST updated; issue pipe_file() rejects compression for small values and uploads an empty object when a compressed write fails #1037 notes the addition. The 5 GiB single-PUT limit is the documented S3 PutObject maximum; the over-5-GiB case was not uploaded (offline test only, stated in the PR).
- Validation:
just lintpassed;tests/pyathena/filesystem/555 passed against live S3 on407ec3ed.
There was a problem hiding this comment.
Independent follow-up review of 407ec3ed (relayed): CLEAN
Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only), session 01a10268-a4e0-7930-9efa-75d21c2765b4. Static review only; scope git diff 6d90268d..407ec3ed with the surrounding pipe_file(), _write_and_close(), _upload_chunk(), aio _pipe_file()/_pipe_file_in_transaction(), and the changed tests; snapshot at 407ec3edcd953f420b05c877eccda1abd7b7444d, unchanged afterwards.
nbyteshandles bytes, bytearray, and memoryviews independently of format, shape, or contiguity. Compression precedes sizing. No remaining raw-value item-count routing was found in the local upload paths. [...] The adjusted test preserves non-contiguous PutObject conversion coverage with an eight-byte boundary. The new test would fail under the old condition and checks observable behavior: no PutObject, exact uploaded payload, and multipart completion. The changed comments are accurate.CLEAN — the reported pre-existing routing defect is resolved. No regressions or additional pre-existing defects identified in scope.
S3File.writeis inherited from fsspec, whose implementation is absent from this tree, so that dependency implementation could not be inspected.
Author: test_pipe_file_memoryview_routed_by_bytes runs the real fsspec write() (only the S3 requests are mocked), and the joined parts equal the data.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
pipe_file() chose between PutObject and the buffered path by len(value), which counts the items of a memoryview. A typed memoryview larger than the block size in bytes, but not in items, went to PutObject, which accepts at most 5 GiB. Compare the size in bytes, as the multipart size check already does. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
_compress() took the path only to resolve "infer" and returned the value unchanged when no codec was inferred. Resolve the codec in the callers with get_compression(), as open() does, and compress with the public CompressedBuffer.compress(), which takes a codec name and always returns bytes. s3_async no longer imports a private name from s3. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
| pass | ||
|
|
||
| @classmethod | ||
| def compress(cls, value: bytes | bytearray | memoryview, compression: str) -> bytes: |
There was a problem hiding this comment.
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()callget_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 raisesValueError("Compression type … not supported")fromget_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 intocls(), and always returnsbytes.s3_asyncimportsCompressedBufferinstead of the private_compress.- While verifying the closing-codec test, a
git checkout -- pyathena/filesystem/s3.pyreverted the uncommitteds3.pypart 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().
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
| """ | ||
| # The codec is taken as open() takes it, also from the extension of | ||
| # the path without the protocol for "infer". | ||
| compression = get_compression(self._strip_protocol(path), kwargs.pop("compression", None)) |
There was a problem hiding this comment.
Independent review of the full PR at edc338db (relayed): FINDINGS — finding 1 of 2
Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only, model reported gpt-6-astra), session 01a1027c-e4b5-74e3-aed0-7a2ba6dafcc4. Static review only. Scope git diff fe21250a..edc338db (full PR, as the public CompressedBuffer API expanded the contract); snapshot at edc338db70f70ce38d2267953be9d898119c1c33, unchanged afterwards. The prompt held the diff, intent, and conventions, without the PR number, PR text, or prior findings, and limited reading to the snapshot and the fsspec source.
Covered: sync/aio pipe/pipe_file, transactions, create mode, failures, block sizes, multipart limits; bytes/bytearray/memoryview variants; None/"infer"/unknown codecs, fsspec codec adapters, sink-closing finalization; CompressedBuffer's public contract, conventions, docs, tests.
P2 — Compression-based routing can change the destination key.
s3.py:1825. Outside a transaction, callpipe_file("s3://bucket/key.gz/", b"a" * (5 * 2**20 + 1), compression="gzip"). Previously, the uncompressed size selectedopen(), which strips the trailing slash and targetskey.gz. Now compression reduces the size below the threshold, selecting PutObject with the original path and targetingkey.gz/. [...] Both routes should use the same destination normalization. The new trailing-slash tests check decompression but omit destination-key assertions.Pre-existing versus introduced: the underlying mismatch between raw-path parsing at
s3.py:1833and normalization throughopen()already existed for uncompressed writes. Finding 1 identifies the newly introduced destination change when compression switches an existing call between those routes.
Author verification: the pre-existing mismatch is confirmed on master fe21250a without compression: pipe_file("s3://bucket/key/", …) writes the key key/ for 10 bytes and key for block size + 1 bytes. A fix means choosing one destination for a trailing-slash path on both routes, which changes uncompressed writes too; that is a design decision for the maintainer, so this PR is kept in Draft until it is decided.
There was a problem hiding this comment.
Maintainer's decision: write a path with a trailing slash to the key without it on both paths, as open() does (s3fs was considered: its pipe_file() keeps the slash on both paths while its open() drops it, so it is inconsistent itself).
Repaired in 73cdb62537e951d8aa58a9baaf75cfe04599a9e1: pipe_file() starts with path = self._strip_protocol(path), and the routing, the codec of "infer", the single-request key, the mode="create" exists() check, and the cache invalidation all use that path. Docstring and the trailing-slash paragraph of docs/filesystem.md updated; release note added to the PR body.
test_pipe_file_trailing_slashasserts the key written fors3://bucket/dir/key/(up to and above the block size, in and outside a transaction), and the trailing-slash cases oftest_pipe_file_compressionnow assert the key (key.gz). Without the normalization, the 3 single-request cases outside a transaction fail, including the scenario above (compressible block size + 1 bytes).- Live check (ad hoc, objects deleted):
pipe_file()of 1 byte and of 6 MiB, andput_file(), to paths ending in/wrotesmall,big, andput.
Self-review of the repair:
- Round 1 (behavior):
"s3://bucket"still raises "Cannot write to a bucket" and a?versionId=path still raises (the query stays in the normalized path;test_pipe_file_invalid_path_raisespasses).s3a://is stripped too. Aio outside a transaction delegates to this method; aio in a transaction goes throughopen(), which already normalizes. Error messages now show the path without the protocol, like the errors raised fromopen(). - Round 2 (claims): "
touch()still writes the key as given" is from code reading (touch()parses the raw path); the docs sentence onpipe/putis backed by the live check. - Validation:
just lintandjust docs lintpassed;tests/pyathena/filesystem/565 passed against live S3 on73cdb625.
There was a problem hiding this comment.
Independent follow-up review of 73cdb625 (relayed): FINDINGS (docs only)
Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only), session 01a10290-bde7-7370-9a5f-25f9a2291403. Static review only; scope git diff edc338db..73cdb625 and git range-diff fe21250a..edc338db fe21250a..73cdb625; snapshot at 73cdb62537e951d8aa58a9baaf75cfe04599a9e1, unchanged afterwards.
Both original review points are resolved: Normalization makes the upload destination independent of size and compression. Aio delegates to this implementation outside transactions and already normalizes through
open()inside transactions. No implementation regression found in the traced paths. Theclose()docstring accurately describes the retained no-op behavior. The changed tests assert observable upload keys and decompressed payloads. [...]P3 (low) —
put()does not write the documented destination.docs/filesystem.md:112newly states thatputtargetingdir/writes objectdir. For an existing local file,fs.put("local.txt", "s3://bucket/dir/")instead writesdir/local.txt. Inherited fsspecput()detects the directory destination before normalization and appends the source basename throughother_paths; async_put()behaves likewise. Replaceputwithput_filehere and distinguishput's directory-target behavior.
Repaired in 67a3db4c921a4843cb8f37a137be9deb630b4107 (docs only): "Opening it for writing, pipe, pipe_file, and put_file write the object dir, and put writes the file into the directory dir." Live check (ad hoc, objects deleted): put() to dir/ wrote dir/local.txt; pipe() to piped/ wrote piped. Also verified while fixing: fsspec pipe() strips the path before calling pipe_file(), so the release note now limits the trailing-slash change to direct pipe_file() calls. Self-review: docs/PR text only, no behavior change; just docs lint passed.
There was a problem hiding this comment.
Independent follow-up review of the docs repair 67a3db4c (relayed): FINDINGS (docs only)
Reviewer: Codex CLI 0.160.0 (codex exec --sandbox read-only), session 01a10296-b8f5-7321-a823-9216e8e2a9e3. Static review only; scope git diff 73cdb625..67a3db4c, checked against open/_open, pipe_file, put_file, fsspec pipe/put/open, and s3_async.py; snapshot at 67a3db4c921a4843cb8f37a137be9deb630b4107, unchanged afterwards.
docs/filesystem.md:112: "putwrites the file into the directorydir" needs qualification to a string destination. With a readable localfile.txt,fs.put(["file.txt"], ["s3://bucket/dir/"])writes objectdir[...] Both fsspec implementations bypass directory expansion for paired lists (spec.py:1119,asyn.py:688), then PyAthena'sput_filestrips the slash throughopen. [...]The remaining paragraph statements match the inspected sync and async paths.
Repaired in c6297d35b3552af208fb931ac8a60201178f47d2 by removing the put statement: the paragraph now says only "Opening it for writing, pipe, pipe_file, and put_file write the object dir", the statements that this review found accurate. fsspec's put() resolves directory destinations itself and is not described. No further independent pass was run for this deletion, as it leaves only the reviewed statements. just docs lint passed.
| """ | ||
|
|
||
| @override | ||
| def close(self) -> None: |
There was a problem hiding this comment.
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:closedremains false and subsequent writes remain possible. Method-level help can inherit the incompatibleBytesIO.close()description.docs/contributing.md:85requires 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.
The single PutObject request wrote a path with a trailing slash to the key with the slash, while the buffered path, through open(), wrote it to the key without it, so the key depended on the size of the value. With the compression argument, compressing could now switch a call between the two. Normalize the path as open() does before choosing, so both paths write the same key. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
fsspec's put() treats a destination ending in a slash as a directory and writes the file under it, unlike put_file(). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
fsspec's put() resolves a destination ending in a slash itself, as a directory for a single source but not for paired lists, so the paragraph no longer describes it. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
WHAT
S3FileSystem.pipe_file()andAioS3FileSystem._pipe_file_in_transaction()compress the value up front when fsspec'scompressionargument is given, and upload the compressed bytes through the usual paths.pipe_file()normalizes the path asopen()does (self._strip_protocol(path), which also drops a trailing slash) before anything else, so the single-request path writes the same key as the buffered path. A path with a trailing slash used to be written to the key with the slash by PutObject and without it throughopen(), depending on the size of the value; compression could switch a call between the two.The callers resolve the codec as
open()does:fsspec.core.get_compression()on that normalized path, including"infer"from its extension. When"infer"finds no codec, the value is uploaded as it is, asopen()does.A new public
CompressedBufferinpyathena/filesystem/s3.py(not added to the API reference, like the other support classes):CompressedBuffer.compress(value, compression)compresses a value with a codec name offsspec.compression.comprand returnsbytes; an unknown name (also"infer") raisesValueError. The buffer is aBytesIOwhoseclose()keeps the data, because some codecs close the file they write to (thezstandardstream writer, fsspec'szstdcodec when neithercompression.zstdnorbackports.zstdis available).compressionis no longer passed toopen()or to S3. The routing between the single-request and buffered paths, and the multipart size check, use the compressed size.pipe_file()chooses between the single-request and buffered paths by the size of the value in bytes (memoryview(value).nbytes, already used for the multipart size check) instead oflen(value), which counts the items of a memoryview.test_pipe_file_non_contiguous_memoryviewpinned the item count (4 items, 8 bytes,block_size=6→ PutObject); it now usesblock_size=8for the same PutObject case.docs/filesystem.md: the trailing-slash paragraph now says thatpipe,pipe_file, andput_filealso write the objectdirfordir/(fsspec'sput()resolves directory destinations itself and is not described); and it states thatpipe/pipe_filecompress the data before uploading it with thecompressionargument, and that the sizes of its write section apply to the compressed data.Behavior changes (release note):
pipe_file(..., compression=...)for a value up to the block size outside a transaction used to fail withParamValidationError: Unknown parameter in input: "compression"; it now uploads the compressed value.AttributeError: 'GzipFile' object has no attribute '_close_without_commit', hiding the original error, and to upload a 10-byte object holding only the gzip header, in and outside a transaction (checked offline on masterfe21250a); it now raises the original error and leaves the existing object unchanged, as Write pipe_file() data without committing a failed write #1003 does for uncompressed values.ValueErrorbefore any request.memoryview(data).cast("I")) larger than the block size in bytes, but not in items, used to go to PutObject, which fails above 5 GiB; it is now uploaded through the buffered (multipart) path.pipe_file()called directly with a value up to the block size, outside a transaction, and a path with a trailing slash (s3://bucket/dir/) used to write the keydir/; it now writesdir, asopen(),put_file(), larger values, and fsspec'spipe()(which strips the path before callingpipe_file()) already did.touch()still writes the key as given, e.g. to create a folder marker.WHY
Closes #1037.
The routing by bytes is a pre-existing defect that the independent review of this PR found (reported in its review thread); the maintainer asked to fix it here.
Found in the review of PR #1035 (#1034): moving the write-and-close helper to
S3Filebreaks compressed buffered writes whileopen()returns a compression wrapper. With this changeopen()always returns theS3File, and #1035 will be rebased onto it.TEST
Tested commit: 73cdb62 (code); 67a3db4 and c6297d3 change only docs/filesystem.md (
just docs lintpassed)just lint: passed.just docs lintandjust docs build: passed on 8462270 (docs unchanged since).uv run --env-file .env pytest -n 8 tests/pyathena/filesystem/: 565 passed (live S3).test_pipe_file_compression(gzip,"infer", and"infer"with a trailing slash, up to and above the block size, in and outside a transaction),test_pipe_file_compression_multipart(compressed data above the block size, in and outside a transaction),test_pipe_file_compression_non_contiguous_memoryview,test_pipe_file_compression_inferred_none,test_pipe_file_unsupported_compression,test_pipe_file_compression_failed_write(in and outside a transaction), aiotest_pipe_file_compression(gzip and"infer"with a trailing slash, in and outside a transaction),TestCompressedBuffer(gzip/bz2/xz round trips, a non-contiguous memoryview, unsupported names incl."infer", and a codec that closes its file, which fails without theclose()override),test_pipe_file_trailing_slash(the key written fors3://bucket/dir/key/, up to and above the block size, in and outside a transaction; the single-request case fails without the normalization, as do the trailing-slash cases oftest_pipe_file_compression, which now also assert the key), andtest_pipe_file_memoryview_routed_by_bytes(acast("I")memoryview of block size + 4 bytes goes to the multipart path; fails withlen(value)routing)."infer"without a codec, the unsupported codec, the failed write (2), and aio outside a transaction. The 7 buffered-path cases pass on master too, which already compressed those.zstandardandlz4installed in a temporary environment (zip, bz2, gzip, lzma, xz, lz4, zstd) round-trips throughCompressedBuffer.compress()and the codec's reader;zstdfailed with a plainBytesIO.pipe_file()of 1 byte and of 6 MiB, andput_file(), to paths ending in/wrote the keyssmall,big, andputwithout the slash; fsspec'sput()todir/wrotedir/local.txt, andpipe()topiped/wrotepiped.pipe_file()withcompression="infer"on a.gzkey (7 bytes, single request) andcompression="gzip"(6 MiB of random bytes, multipart), each in and outside a transaction; every object read back withcat_file()decompressed to the value. The objects were deleted.🤖 Generated with Claude Code