From 9d17000a0c8e390a8e5f3c6037158ed17be6791e Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 16:04:57 +0900 Subject: [PATCH 1/8] Keep multipart uploads within the 10,000-part limit S3 multipart uploads accept at most 10,000 parts. S3File numbered write parts without an upper bound, so a write of more than 10,000 blocks (about 48.8 GiB with the default 5 MiB block size) submitted part 10,001 and failed only when S3 rejected it. Multipart copies with a small block_size had the same gap. Writes now raise ValueError before submitting a part beyond the limit, naming the block size to use. The multipart upload is aborted and the file is closed without its buffered data, so that neither close() nor a deferred commit() uploads a partial object. Parts copied from the existing object in an append count toward the limit. Multipart copies know the source size, so _get_copy_ranges raises the split size to the size divided by 10,000 instead. Closes #953 Co-Authored-By: Claude Opus 5.5 --- docs/filesystem.md | 8 +++ pyathena/filesystem/s3.py | 40 +++++++++--- tests/pyathena/filesystem/test_s3.py | 93 +++++++++++++++++++++++++++- 3 files changed, 133 insertions(+), 8 deletions(-) diff --git a/docs/filesystem.md b/docs/filesystem.md index 4f8135cf0..2fee0a199 100644 --- a/docs/filesystem.md +++ b/docs/filesystem.md @@ -85,6 +85,14 @@ The block size for writing, given by the `block_size` argument of `open` or by t 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, one per block, so a write with +the default block size can upload up to about 48.8 GiB (10,000 × 5 MiB). To write a +larger object, use a block size of at least its size divided by 10,000, either with +the `block_size` argument of `open` and `pipe` or with the `default_block_size` +argument of `S3FileSystem`. A write that would need more parts raises `ValueError` +and aborts its multipart upload. Multipart copies with `cp` use parts large enough +to stay within the limit. + ## Error translation S3 error responses are translated into standard Python exceptions, so filesystem diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index 2bbda039c..252bd227e 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -3,6 +3,7 @@ from __future__ import annotations import logging +import math import mimetypes import os.path import re @@ -10,6 +11,7 @@ from concurrent.futures import Future, as_completed, wait from copy import deepcopy from datetime import datetime +from io import BytesIO from multiprocessing import cpu_count from re import Pattern from typing import Any, cast @@ -104,6 +106,9 @@ class S3FileSystem(AbstractFileSystem): # https://docs.aws.amazon.com/AmazonS3/latest/userguide/qfacts.html # The maximum size of a part in a multipart upload is 5GiB. MULTIPART_UPLOAD_MAX_PART_SIZE: int = 5 * 2**30 # 5GiB + # https://docs.aws.amazon.com/AmazonS3/latest/userguide/qfacts.html + # The maximum number of parts per multipart upload is 10,000. + MULTIPART_UPLOAD_MAX_PARTS: int = 10_000 # https://docs.aws.amazon.com/AmazonS3/latest/API/API_DeleteObjects.html DELETE_OBJECTS_MAX_KEYS: int = 1000 DEFAULT_BLOCK_SIZE: int = 5 * 2**20 # 5MiB @@ -1318,7 +1323,8 @@ def _get_copy_ranges(self, size: int, block_size: int) -> list[tuple[int, int]]: """Split an object into the source ranges of a multipart copy. The object is split into ranges of ``block_size`` bytes, whatever the - number of workers. A last range shorter than + number of workers, or of a larger size that splits it into at most + ``MULTIPART_UPLOAD_MAX_PARTS`` ranges. A last range shorter than ``MULTIPART_UPLOAD_MIN_PART_SIZE`` is merged into the previous one, which is split in half if the result exceeds ``MULTIPART_UPLOAD_MAX_PART_SIZE``. Every range is then within the @@ -1330,14 +1336,16 @@ def _get_copy_ranges(self, size: int, block_size: int) -> list[tuple[int, int]]: size: The size of the source object in bytes. block_size: The size in bytes to split the object by, between ``MULTIPART_UPLOAD_MIN_PART_SIZE`` and - ``MULTIPART_UPLOAD_MAX_PART_SIZE``. The range that a short - last range is merged into can be longer, up to - ``MULTIPART_UPLOAD_MAX_PART_SIZE``. + ``MULTIPART_UPLOAD_MAX_PART_SIZE``. It is raised to + ``size`` divided by ``MULTIPART_UPLOAD_MAX_PARTS``, rounded + up, if smaller. The range that a short last range is merged + into can be longer, up to ``MULTIPART_UPLOAD_MAX_PART_SIZE``. Returns: The ``(start, end)`` byte ranges, with an exclusive end, that cover the whole object in order. """ + block_size = max(block_size, math.ceil(size / self.MULTIPART_UPLOAD_MAX_PARTS)) starts = list(range(0, size, block_size)) if len(starts) > 1 and size - starts[-1] < self.MULTIPART_UPLOAD_MIN_PART_SIZE: starts.pop() @@ -2219,6 +2227,7 @@ class S3File(AbstractBufferedFile): """ fs: S3FileSystem + buffer: BytesIO | None def __init__( self, @@ -2429,15 +2438,18 @@ def _upload_chunk(self, final: bool = False) -> bool: if not self.multipart_upload: raise RuntimeError("Multipart upload is not initialized.") + # fsspec's flush() never calls this on a closed file, whose buffer + # may have been dropped. + buffer = cast(BytesIO, self.buffer) part_number = len(self.multipart_upload_parts) - self.buffer.seek(0) - data = self.buffer.read(self.blocksize) + buffer.seek(0) + data = buffer.read(self.blocksize) while data: # Only the last part of a multipart upload may be smaller than the # minimum part size, and more data may follow a mid-stream chunk. # A single write() can leave several blocks in the buffer, so look # ahead one block and merge a short last block into this one. - next_data = self.buffer.read(self.blocksize) + next_data = buffer.read(self.blocksize) next_data_size = len(next_data) if 0 < next_data_size < self.fs.MULTIPART_UPLOAD_MIN_PART_SIZE: upload_data = data + next_data @@ -2452,6 +2464,20 @@ def _upload_chunk(self, final: bool = False) -> bool: uploads = [data] for upload in uploads: + if part_number >= self.fs.MULTIPART_UPLOAD_MAX_PARTS: + # Abort the upload and close the file without the + # buffered data, so that neither close() nor commit() + # uploads it. + self.discard() + self.buffer = None + self.closed = True + raise ValueError( + f"Cannot upload more than {self.fs.MULTIPART_UPLOAD_MAX_PARTS} " + f"parts to s3://{self.bucket}/{self.key} with a block size of " + f"{self.blocksize} bytes. Write the file with a block_size, or " + "a default_block_size of the filesystem, of at least its total " + f"size divided by {self.fs.MULTIPART_UPLOAD_MAX_PARTS}." + ) part_number += 1 self.multipart_upload_parts.append( self._executor.submit( diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 95b47d016..b68acfc47 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -12,7 +12,7 @@ import uuid from concurrent.futures import Future, ThreadPoolExecutor, wait from datetime import UTC, datetime -from itertools import chain +from itertools import chain, pairwise from pathlib import Path from types import SimpleNamespace from unittest import mock @@ -806,6 +806,32 @@ def finish(): def test_get_copy_ranges(self, size, block_size, ranges): assert self._make_fs()._get_copy_ranges(size, block_size) == ranges + @pytest.mark.parametrize( + ("size", "num_ranges"), + [ + # The block size splits the object into the maximum number of parts. + (10_000 * 5 * 2**20, 10_000), + # GH-953: a larger object is split by a larger size instead of + # into more parts than the maximum, + (10_000 * 5 * 2**20 + 1, 9_999), + (50 * 2**30, 10_000), + # including the maximum object size. + (5 * 2**40, 10_000), + ], + ) + def test_get_copy_ranges_max_parts(self, size, num_ranges): + fs = self._make_fs() + ranges = fs._get_copy_ranges(size, 5 * 2**20) + + assert len(ranges) == num_ranges + assert ranges[0][0] == 0 + assert ranges[-1][1] == size + assert all(end == start for (_, end), (start, _) in pairwise(ranges)) + assert all( + fs.MULTIPART_UPLOAD_MIN_PART_SIZE <= end - start <= fs.MULTIPART_UPLOAD_MAX_PART_SIZE + for start, end in ranges + ) + @pytest.mark.parametrize("max_workers", [1, 4]) def test_copy_object_with_multipart_upload_part_sizes(self, max_workers): # GH-951: the parts are within the S3 part size limits whatever the @@ -2197,6 +2223,7 @@ def _make_write_file(data: bytes, autocommit: bool): file.s3_additional_kwargs = {} file.autocommit = autocommit file.blocksize = S3FileSystem.MULTIPART_UPLOAD_MIN_PART_SIZE + file.fs.MULTIPART_UPLOAD_MAX_PARTS = S3FileSystem.MULTIPART_UPLOAD_MAX_PARTS file.append_block = False file.multipart_upload = None file.multipart_upload_parts = [] @@ -2229,6 +2256,7 @@ def _make_append_fs(existing: bytes): fs = mock.MagicMock(spec=S3FileSystem) fs.MULTIPART_UPLOAD_MIN_PART_SIZE = 4 fs.MULTIPART_UPLOAD_MAX_PART_SIZE = 64 + fs.MULTIPART_UPLOAD_MAX_PARTS = S3FileSystem.MULTIPART_UPLOAD_MAX_PARTS fs.exists.return_value = True fs.info.return_value = S3Object( init={"ContentLength": len(existing)}, @@ -2342,6 +2370,69 @@ def test_write_part_sizes(self, writes, block_size): assert all(size >= fs.MULTIPART_UPLOAD_MIN_PART_SIZE for size in sizes[:-1]) assert all(size <= fs.MULTIPART_UPLOAD_MAX_PART_SIZE for size in sizes) + @staticmethod + def _write_and_close(f, writes: list[bytes]) -> None: + with f: + for data in writes: + f.write(data) + + @pytest.mark.parametrize( + ("existing", "mode", "writes"), + [ + # The data fills the maximum number of parts. + (b"", "wb", [b"a" * 4] * 3), + (b"", "wb", [b"a" * 14]), + # The parts copied from the existing object in an append count + # toward the maximum. + (b"a" * 6, "ab", [b"b" * 4] * 2), + ], + ) + @pytest.mark.parametrize("autocommit", [True, False]) + def test_write_max_parts(self, existing, mode, writes, autocommit): + fs = self._make_append_fs(existing) + fs.MULTIPART_UPLOAD_MAX_PARTS = 3 + + f = S3File(fs, "s3://bucket/key.txt", mode=mode, block_size=4, autocommit=autocommit) + self._write_and_close(f, writes) + if not autocommit: + f.commit() + + assert self._uploaded_object(fs, existing) == existing + b"".join(writes) + assert fs._upload_part_copy.call_count + fs._upload_part.call_count == 3 + + @pytest.mark.parametrize( + ("existing", "mode", "writes"), + [ + # GH-953: the part after the maximum is not uploaded, whether it is + # flushed by a write() + (b"", "wb", [b"a" * 4] * 4), + # or by close(), + (b"", "wb", [b"a" * 12, b"b" * 3]), + # including after the parts copied in an append. + (b"a" * 6, "ab", [b"b" * 4] * 3), + ], + ) + @pytest.mark.parametrize("autocommit", [True, False]) + def test_write_exceeding_max_parts(self, existing, mode, writes, autocommit): + fs = self._make_append_fs(existing) + fs.MULTIPART_UPLOAD_MAX_PARTS = 3 + + f = S3File(fs, "s3://bucket/key.txt", mode=mode, block_size=4, autocommit=autocommit) + with pytest.raises(ValueError, match="block_size"): + self._write_and_close(f, writes) + # The upload is aborted, and committing a deferred write afterwards + # uploads nothing. + if not autocommit: + f.commit() + + assert f.closed + assert fs._upload_part_copy.call_count + fs._upload_part.call_count == 3 + fs._call.assert_called_once_with( + "abort_multipart_upload", Bucket="bucket", Key="key.txt", UploadId="uploadid" + ) + fs._finish_multipart_upload.assert_not_called() + fs._put_object.assert_not_called() + def test_append_discard(self): # Rolling back an append aborts its multipart upload without the # existing object's metadata, which AbortMultipartUpload rejects, From de7b2dbc0156edff77ab6c98b28f73a1f69b1f70 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 16:34:54 +0900 Subject: [PATCH 2/8] Keep the part limit error when the abort fails If AbortMultipartUpload failed, its error replaced the part limit error, the file stayed open so close() retried the abort, and a deferred commit() completed the parts uploaded so far. Close the file before the abort, log an abort failure, clear the upload state, and always raise the part limit error. Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 16 ++++++++++++---- tests/pyathena/filesystem/test_s3.py | 21 +++++++++++++++++++++ 2 files changed, 33 insertions(+), 4 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index 252bd227e..c25fb2856 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -2465,12 +2465,20 @@ def _upload_chunk(self, final: bool = False) -> bool: for upload in uploads: if part_number >= self.fs.MULTIPART_UPLOAD_MAX_PARTS: - # Abort the upload and close the file without the - # buffered data, so that neither close() nor commit() - # uploads it. - self.discard() + # Close the file without the buffered data, so that + # neither close() nor commit() uploads it, and abort the + # upload. An abort failure does not mask this error, and + # commit() does not complete the upload afterwards. self.buffer = None self.closed = True + try: + self.discard() + except Exception: + _logger.exception( + f"Failed to abort multipart upload to s3://{self.bucket}/{self.key}." + ) + self.multipart_upload = None + self.multipart_upload_parts = [] raise ValueError( f"Cannot upload more than {self.fs.MULTIPART_UPLOAD_MAX_PARTS} " f"parts to s3://{self.bucket}/{self.key} with a block size of " diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index b68acfc47..9a9d2f581 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -2433,6 +2433,27 @@ def test_write_exceeding_max_parts(self, existing, mode, writes, autocommit): fs._finish_multipart_upload.assert_not_called() fs._put_object.assert_not_called() + @pytest.mark.parametrize("autocommit", [True, False]) + def test_write_exceeding_max_parts_abort_failure(self, autocommit): + # An abort failure is logged; the part limit error propagates, and + # neither closing the file nor committing a deferred write retries + # the upload or completes it. + fs = self._make_append_fs(b"") + fs.MULTIPART_UPLOAD_MAX_PARTS = 3 + fs._call.side_effect = PermissionError("abort failed") + + f = S3File(fs, "s3://bucket/key.txt", mode="wb", block_size=4, autocommit=autocommit) + with pytest.raises(ValueError, match="block_size"): + self._write_and_close(f, [b"a" * 4] * 4) + if not autocommit: + f.commit() + + assert f.closed + assert fs._upload_part.call_count == 3 + fs._call.assert_called_once() + fs._finish_multipart_upload.assert_not_called() + fs._put_object.assert_not_called() + def test_append_discard(self): # Rolling back an append aborts its multipart upload without the # existing object's metadata, which AbortMultipartUpload rejects, From fdc7bf4c94bd64738cc0d1a413d1fa455fda3076 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 16:36:59 +0900 Subject: [PATCH 3/8] Do not give an exact block size for appends In an append to an object smaller than the maximum part size, the copy takes one part, so a block size of the total size divided by 10,000 can still need 10,001 parts. The error now asks for a block size that fits the parts, including the copied ones, and the docs note that the copied parts count toward the limit. Co-Authored-By: Claude Opus 5.5 --- docs/filesystem.md | 3 ++- pyathena/filesystem/s3.py | 5 +++-- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/docs/filesystem.md b/docs/filesystem.md index 2fee0a199..3cca78a88 100644 --- a/docs/filesystem.md +++ b/docs/filesystem.md @@ -89,7 +89,8 @@ A multipart upload consists of at most 10,000 parts, one per block, so a write w the default block size can upload up to about 48.8 GiB (10,000 × 5 MiB). To write a larger object, use a block size of at least its size divided by 10,000, either with the `block_size` argument of `open` and `pipe` or with the `default_block_size` -argument of `S3FileSystem`. A write that would need more parts raises `ValueError` +argument of `S3FileSystem`. In an append, the parts copied from the existing object +also count toward the limit. A write that would need more parts raises `ValueError` and aborts its multipart upload. Multipart copies with `cp` use parts large enough to stay within the limit. diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index c25fb2856..0ed3b9fca 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -2483,8 +2483,9 @@ def _upload_chunk(self, final: bool = False) -> bool: f"Cannot upload more than {self.fs.MULTIPART_UPLOAD_MAX_PARTS} " f"parts to s3://{self.bucket}/{self.key} with a block size of " f"{self.blocksize} bytes. Write the file with a block_size, or " - "a default_block_size of the filesystem, of at least its total " - f"size divided by {self.fs.MULTIPART_UPLOAD_MAX_PARTS}." + "a default_block_size of the filesystem, large enough for it to " + f"fit in {self.fs.MULTIPART_UPLOAD_MAX_PARTS} parts, including " + "the parts copied from the existing object in an append." ) part_number += 1 self.multipart_upload_parts.append( From 1015c1005c9276e7660d10427377a69151971702 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 16:46:09 +0900 Subject: [PATCH 4/8] Shut down the executor when closing fails S3File.close() skipped the executor shutdown when the final flush raised, as the part limit error does, and the file was already marked closed. Shut it down in a finally block. The part limit failure tests asserted the number of executed part uploads, which depends on whether the abort cancels queued uploads. They now record the submitted parts through a wrapped executor and check its shutdown. The docs claimed one part per block for every write, but each write through open() that fills the buffer uploads its data beyond the last full block as a separate part. Narrow the sizing guidance to put and pipe. Co-Authored-By: Claude Opus 5.5 --- docs/filesystem.md | 16 ++++++++------- pyathena/filesystem/s3.py | 7 +++++-- tests/pyathena/filesystem/test_s3.py | 29 +++++++++++++++++++++++----- 3 files changed, 38 insertions(+), 14 deletions(-) diff --git a/docs/filesystem.md b/docs/filesystem.md index 3cca78a88..57a999e87 100644 --- a/docs/filesystem.md +++ b/docs/filesystem.md @@ -85,13 +85,15 @@ The block size for writing, given by the `block_size` argument of `open` or by t 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, one per block, so a write with -the default block size can upload up to about 48.8 GiB (10,000 × 5 MiB). To write a -larger object, use a block size of at least its size divided by 10,000, either with -the `block_size` argument of `open` and `pipe` or with the `default_block_size` -argument of `S3FileSystem`. In an append, the parts copied from the existing object -also count toward the limit. A write that would need more parts raises `ValueError` -and aborts its multipart upload. Multipart copies with `cp` use parts large enough +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 `pipe` and `open` or with +the `default_block_size` argument of `S3FileSystem`. 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. In an append, the parts copied from the existing +object also count toward the limit. A write that would need more parts raises +`ValueError` and aborts its multipart upload. Multipart copies with `cp` use parts large enough to stay within the limit. ## Error translation diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index 0ed3b9fca..0086898c0 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -2374,8 +2374,11 @@ def __init__( def close(self) -> None: """Close the file, flushing any written data, and shut down its executor.""" - super().close() - self._executor.shutdown() + try: + super().close() + finally: + # The executor is shut down even if the final flush fails. + self._executor.shutdown() def _initiate_upload(self) -> None: if not self.append_block and self.tell() < self.blocksize: diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 9a9d2f581..d172f2a9c 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -23,7 +23,7 @@ import pyathena from pyathena.filesystem import register_s3_filesystem from pyathena.filesystem.s3 import S3File, S3FileSystem -from pyathena.filesystem.s3_executor import S3AioExecutor +from pyathena.filesystem.s3_executor import S3AioExecutor, S3ThreadPoolExecutor from pyathena.filesystem.s3_object import S3Object, S3ObjectType, S3StorageClass from pyathena.util import RetryConfig from tests import ENV @@ -2417,7 +2417,15 @@ def test_write_exceeding_max_parts(self, existing, mode, writes, autocommit): fs = self._make_append_fs(existing) fs.MULTIPART_UPLOAD_MAX_PARTS = 3 - f = S3File(fs, "s3://bucket/key.txt", mode=mode, block_size=4, autocommit=autocommit) + executor = mock.MagicMock(wraps=S3ThreadPoolExecutor(max_workers=1)) + f = S3File( + fs, + "s3://bucket/key.txt", + mode=mode, + block_size=4, + autocommit=autocommit, + executor=executor, + ) with pytest.raises(ValueError, match="block_size"): self._write_and_close(f, writes) # The upload is aborted, and committing a deferred write afterwards @@ -2426,7 +2434,9 @@ def test_write_exceeding_max_parts(self, existing, mode, writes, autocommit): f.commit() assert f.closed - assert fs._upload_part_copy.call_count + fs._upload_part.call_count == 3 + # The submitted parts, some of which the abort may have cancelled. + assert [c.kwargs["part_number"] for c in executor.submit.call_args_list] == [1, 2, 3] + executor.shutdown.assert_called_once() fs._call.assert_called_once_with( "abort_multipart_upload", Bucket="bucket", Key="key.txt", UploadId="uploadid" ) @@ -2442,14 +2452,23 @@ def test_write_exceeding_max_parts_abort_failure(self, autocommit): fs.MULTIPART_UPLOAD_MAX_PARTS = 3 fs._call.side_effect = PermissionError("abort failed") - f = S3File(fs, "s3://bucket/key.txt", mode="wb", block_size=4, autocommit=autocommit) + executor = mock.MagicMock(wraps=S3ThreadPoolExecutor(max_workers=1)) + f = S3File( + fs, + "s3://bucket/key.txt", + mode="wb", + block_size=4, + autocommit=autocommit, + executor=executor, + ) with pytest.raises(ValueError, match="block_size"): self._write_and_close(f, [b"a" * 4] * 4) if not autocommit: f.commit() assert f.closed - assert fs._upload_part.call_count == 3 + assert [c.kwargs["part_number"] for c in executor.submit.call_args_list] == [1, 2, 3] + executor.shutdown.assert_called_once() fs._call.assert_called_once() fs._finish_multipart_upload.assert_not_called() fs._put_object.assert_not_called() From d44602afad4793fb9616d9ffea016ffdcc5c02cd Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 16:52:04 +0900 Subject: [PATCH 5/8] Note the minimum size of separate remainder parts A remainder smaller than the minimum part size is merged into the preceding part, so only a remainder of at least 5 MiB adds a part. Co-Authored-By: Claude Opus 5.5 --- docs/filesystem.md | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/docs/filesystem.md b/docs/filesystem.md index 57a999e87..a921cfb7f 100644 --- a/docs/filesystem.md +++ b/docs/filesystem.md @@ -91,10 +91,10 @@ per block, so with the default block size they can upload up to about 48.8 GiB divided by 10,000, either with the `block_size` argument of `pipe` and `open` or with the `default_block_size` argument of `S3FileSystem`. 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. In an append, the parts copied from the existing -object also count toward the limit. A write that would need more parts raises -`ValueError` and aborts its multipart upload. Multipart copies with `cp` use parts large enough -to stay within the limit. +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 that would +need more parts raises `ValueError` and aborts its multipart upload. Multipart copies +with `cp` use parts large enough to stay within the limit. ## Error translation From fdf0a70a47f2eb30509b2aa5ab4d5fc15f0545f4 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 17:22:27 +0900 Subject: [PATCH 6/8] Reject put and pipe data that cannot fit before uploading put_file() and pipe_file() know the size of the data, so they now raise ValueError before uploading anything when it needs more than 10,000 blocks, naming the minimum block size, instead of failing after about 48.8 GiB has been uploaded at the default block size. The check follows the documented rule (a block size of at least the size divided by 10,000), so pipe_file() also rejects data up to 5 MiB over 10,000 blocks, whose short last block the multipart upload would merge. put_file() now accepts block_size and passes it to open(); before, it went to the S3 API as a request parameter. Co-Authored-By: Claude Opus 5.5 --- docs/filesystem.md | 16 ++++--- pyathena/filesystem/s3.py | 38 +++++++++++++++- tests/pyathena/filesystem/test_s3.py | 68 ++++++++++++++++++++++++++++ 3 files changed, 113 insertions(+), 9 deletions(-) diff --git a/docs/filesystem.md b/docs/filesystem.md index a921cfb7f..0096ee03b 100644 --- a/docs/filesystem.md +++ b/docs/filesystem.md @@ -88,13 +88,15 @@ 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 `pipe` and `open` or with -the `default_block_size` argument of `S3FileSystem`. 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 that would -need more parts raises `ValueError` and aborts its multipart upload. Multipart copies -with `cp` use parts large enough to stay within the limit. +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. ## Error translation diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index 0086898c0..3d43122f1 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -1353,6 +1353,31 @@ def _get_copy_ranges(self, size: int, block_size: int) -> list[tuple[int, int]]: starts.append(starts[-1] + (size - starts[-1]) // 2) return list(zip(starts, [*starts[1:], size], strict=True)) + def _check_multipart_upload_size(self, path: str, size: int, block_size: int) -> None: + """Check that data fits in a multipart upload before uploading it. + + Args: + path: The path that the data is written to. + size: The size of the data in bytes. + block_size: The block size of the write in bytes. + + Raises: + ValueError: If the data takes more than + ``MULTIPART_UPLOAD_MAX_PARTS`` blocks. + """ + if size > block_size * self.MULTIPART_UPLOAD_MAX_PARTS: + min_block_size = max( + math.ceil(size / self.MULTIPART_UPLOAD_MAX_PARTS), + self.MULTIPART_UPLOAD_MIN_PART_SIZE, + ) + raise ValueError( + f"Cannot upload {size} bytes to {path} in " + f"{self.MULTIPART_UPLOAD_MAX_PARTS} parts with a block size of " + f"{block_size} bytes. Write the file with a block_size, or a " + "default_block_size of the filesystem, of at least " + f"{min_block_size} bytes." + ) + def pipe_file( self, path: str, value: bytes | bytearray | memoryview, mode: str = "overwrite", **kwargs ) -> None: @@ -1379,9 +1404,11 @@ def pipe_file( FileExistsError: If the mode is "create" and the path already exists. ValueError: If the path does not contain a key or specifies a - version. + version, or if the data takes more than + ``MULTIPART_UPLOAD_MAX_PARTS`` blocks. """ block_size = kwargs.get("block_size") or self.default_block_size + self._check_multipart_upload_size(path, len(value), block_size) if self._intrans or len(value) > 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 @@ -1523,6 +1550,11 @@ def put_file(self, lpath: str, rpath: str, callback=_DEFAULT_CALLBACK, **kwargs) rpath: S3 destination path (s3://bucket/key). callback: Progress callback for tracking upload progress. **kwargs: Additional S3 parameters (e.g., ContentType, StorageClass). + The ``block_size`` parameter of ``open()`` is also accepted. + + Raises: + ValueError: If the file takes more than + ``MULTIPART_UPLOAD_MAX_PARTS`` blocks. Note: Directories are not supported for upload. If lpath is a directory, @@ -1539,6 +1571,8 @@ def put_file(self, lpath: str, rpath: str, callback=_DEFAULT_CALLBACK, **kwargs) return size = os.path.getsize(lpath) + block_size = kwargs.pop("block_size", None) or self.default_block_size + self._check_multipart_upload_size(rpath, size, block_size) callback.set_size(size) if "ContentType" not in kwargs: content_type, _ = mimetypes.guess_type(lpath) @@ -1546,7 +1580,7 @@ def put_file(self, lpath: str, rpath: str, callback=_DEFAULT_CALLBACK, **kwargs) kwargs["ContentType"] = content_type with ( - self.open(rpath, "wb", s3_additional_kwargs=kwargs) as remote, + self.open(rpath, "wb", block_size=block_size, s3_additional_kwargs=kwargs) as remote, open(lpath, "rb") as local, ): while data := local.read(remote.blocksize): diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index d172f2a9c..0392a0e4f 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -514,6 +514,74 @@ def test_pipe_file_small_drops_max_workers(self): fs.pipe_file("s3://bucket/key", b"data", max_workers=2) fs._put_object.assert_called_once_with(bucket="bucket", key="key", body=b"data") + @pytest.mark.parametrize( + ("size", "block_size", "min_block_size"), + [ + # The data fits in the maximum number of parts. + (12, 4, None), + # GH-953: more data is rejected with the minimum block size, + (13, 4, 5), + # which is at least the minimum part size. + (5, 1, 4), + ], + ) + def test_check_multipart_upload_size(self, size, block_size, min_block_size): + fs = self._make_fs() + fs.MULTIPART_UPLOAD_MIN_PART_SIZE = 4 + fs.MULTIPART_UPLOAD_MAX_PARTS = 3 + + if min_block_size is None: + fs._check_multipart_upload_size("s3://bucket/key", size, block_size) + else: + with pytest.raises(ValueError, match=f"at least {min_block_size} bytes"): + fs._check_multipart_upload_size("s3://bucket/key", size, block_size) + + @pytest.mark.parametrize("kwargs", [{"block_size": 4}, {}]) + def test_put_file_exceeding_max_parts(self, tmp_path, kwargs): + # GH-953: a file that does not fit in the maximum number of parts is + # rejected before anything is uploaded. + fs = self._make_fs() + fs.MULTIPART_UPLOAD_MAX_PARTS = 3 + fs.default_block_size = 4 + fs.open = mock.MagicMock() + lpath = tmp_path / "data" + lpath.write_bytes(b"a" * 13) + + with pytest.raises(ValueError, match="block_size"): + fs.put_file(str(lpath), "s3://bucket/key", **kwargs) + fs.open.assert_not_called() + fs._call.assert_not_called() + + def test_put_file_block_size(self, tmp_path): + # block_size is passed to open() instead of the S3 API. + fs = self._make_fs() + fs.open = mock.MagicMock() + fs.open.return_value.__enter__.return_value.blocksize = 8 + lpath = tmp_path / "data" + lpath.write_bytes(b"a" * 13) + + fs.put_file(str(lpath), "s3://bucket/key", block_size=8) + + fs.open.assert_called_once_with( + "s3://bucket/key", "wb", block_size=8, s3_additional_kwargs={} + ) + + @pytest.mark.parametrize("kwargs", [{"block_size": 4}, {}]) + def test_pipe_file_exceeding_max_parts(self, kwargs): + # GH-953: data that does not fit in the maximum number of parts is + # rejected before anything is uploaded. + fs = self._make_fs() + fs.MULTIPART_UPLOAD_MAX_PARTS = 3 + fs.default_block_size = 4 + fs.open = mock.MagicMock() + fs._put_object = mock.MagicMock() + + with pytest.raises(ValueError, match="block_size"): + fs.pipe_file("s3://bucket/key", b"a" * 13, **kwargs) + fs.open.assert_not_called() + fs._put_object.assert_not_called() + fs._call.assert_not_called() + def test_open_max_workers(self): fs = self._make_fs() fs.default_cache_type = "bytes" From f13adf5da5011f9bca65f4dd3403c9df984dfda4 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 17:31:36 +0900 Subject: [PATCH 7/8] Shut down the executor on the part limit error and count bytes A write() that reached the part limit closed the file without shutting down its executor, and fsspec does not close a closed file again when it is garbage collected, so the executor was left running unless close() was called. Shut it down in the error path as well. pipe_file() measured a memoryview with len(), which counts its items, so a view with multi-byte items passed the part limit check and the PutObject dispatch with a fraction of its size. Use its size in bytes. Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 11 ++++++--- tests/pyathena/filesystem/test_s3.py | 34 ++++++++++++++++++++++++---- 2 files changed, 37 insertions(+), 8 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index 3d43122f1..1cf5a820b 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -1408,8 +1408,10 @@ def pipe_file( ``MULTIPART_UPLOAD_MAX_PARTS`` blocks. """ block_size = kwargs.get("block_size") or self.default_block_size - self._check_multipart_upload_size(path, len(value), block_size) - if self._intrans or len(value) > min(block_size, self.MULTIPART_UPLOAD_MAX_PART_SIZE): + # The size in bytes; the length of a memoryview counts its items. + 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. @@ -2505,7 +2507,9 @@ def _upload_chunk(self, final: bool = False) -> bool: # Close the file without the buffered data, so that # neither close() nor commit() uploads it, and abort the # upload. An abort failure does not mask this error, and - # commit() does not complete the upload afterwards. + # commit() does not complete the upload afterwards. The + # executor is shut down here, as fsspec does not close a + # closed file again when it is garbage collected. self.buffer = None self.closed = True try: @@ -2516,6 +2520,7 @@ def _upload_chunk(self, final: bool = False) -> bool: ) self.multipart_upload = None self.multipart_upload_parts = [] + self._executor.shutdown() raise ValueError( f"Cannot upload more than {self.fs.MULTIPART_UPLOAD_MAX_PARTS} " f"parts to s3://{self.bucket}/{self.key} with a block size of " diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 0392a0e4f..00f976805 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -566,8 +566,16 @@ def test_put_file_block_size(self, tmp_path): "s3://bucket/key", "wb", block_size=8, s3_additional_kwargs={} ) - @pytest.mark.parametrize("kwargs", [{"block_size": 4}, {}]) - def test_pipe_file_exceeding_max_parts(self, kwargs): + @pytest.mark.parametrize( + ("value", "kwargs"), + [ + (b"a" * 13, {"block_size": 4}), + (b"a" * 13, {}), + # The size of a memoryview is counted in bytes, not items. + (memoryview(b"a" * 16).cast("I"), {}), + ], + ) + def test_pipe_file_exceeding_max_parts(self, value, kwargs): # GH-953: data that does not fit in the maximum number of parts is # rejected before anything is uploaded. fs = self._make_fs() @@ -577,7 +585,7 @@ def test_pipe_file_exceeding_max_parts(self, kwargs): fs._put_object = mock.MagicMock() with pytest.raises(ValueError, match="block_size"): - fs.pipe_file("s3://bucket/key", b"a" * 13, **kwargs) + fs.pipe_file("s3://bucket/key", value, **kwargs) fs.open.assert_not_called() fs._put_object.assert_not_called() fs._call.assert_not_called() @@ -2504,7 +2512,7 @@ def test_write_exceeding_max_parts(self, existing, mode, writes, autocommit): assert f.closed # The submitted parts, some of which the abort may have cancelled. assert [c.kwargs["part_number"] for c in executor.submit.call_args_list] == [1, 2, 3] - executor.shutdown.assert_called_once() + executor.shutdown.assert_called() fs._call.assert_called_once_with( "abort_multipart_upload", Bucket="bucket", Key="key.txt", UploadId="uploadid" ) @@ -2536,11 +2544,27 @@ def test_write_exceeding_max_parts_abort_failure(self, autocommit): assert f.closed assert [c.kwargs["part_number"] for c in executor.submit.call_args_list] == [1, 2, 3] - executor.shutdown.assert_called_once() + executor.shutdown.assert_called() fs._call.assert_called_once() fs._finish_multipart_upload.assert_not_called() fs._put_object.assert_not_called() + def test_write_exceeding_max_parts_without_close(self): + # The executor of the closed file is shut down, as fsspec does not + # close it again when it is garbage collected. + fs = self._make_append_fs(b"") + fs.MULTIPART_UPLOAD_MAX_PARTS = 3 + executor = mock.MagicMock(wraps=S3ThreadPoolExecutor(max_workers=1)) + f = S3File(fs, "s3://bucket/key.txt", mode="wb", block_size=4, executor=executor) + + for _ in range(3): + f.write(b"a" * 4) + with pytest.raises(ValueError, match="block_size"): + f.write(b"a" * 4) + + assert f.closed + executor.shutdown.assert_called_once() + def test_append_discard(self): # Rolling back an append aborts its multipart upload without the # existing object's metadata, which AbortMultipartUpload rejects, From a5f0ec7ab45619d563865e6616bb333d6640511e Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sat, 3 Oct 2026 17:38:44 +0900 Subject: [PATCH 8/8] Keep the PutObject dispatch of pipe_file on the item count Measuring the dispatch in bytes sent a non-contiguous memoryview with multi-byte items to the buffered path, whose BytesIO cannot write it, when its item count fit in a block. Only the part limit check needs the size in bytes; the dispatch keeps using len() as before. Co-Authored-By: Claude Opus 5.5 --- pyathena/filesystem/s3.py | 5 ++--- tests/pyathena/filesystem/test_s3.py | 12 ++++++++++++ 2 files changed, 14 insertions(+), 3 deletions(-) diff --git a/pyathena/filesystem/s3.py b/pyathena/filesystem/s3.py index 1cf5a820b..e5cf41035 100644 --- a/pyathena/filesystem/s3.py +++ b/pyathena/filesystem/s3.py @@ -1409,9 +1409,8 @@ def pipe_file( """ block_size = kwargs.get("block_size") or self.default_block_size # The size in bytes; the length of a memoryview counts its items. - 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): + 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): # Defer to the buffered open() path, which keeps the # deferred-commit semantics of fsspec transactions and uploads # large data as a parallel multipart upload. diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 00f976805..c17a9bef7 100644 --- a/tests/pyathena/filesystem/test_s3.py +++ b/tests/pyathena/filesystem/test_s3.py @@ -506,6 +506,18 @@ def test_pipe_file_invalid_path_raises(self): with pytest.raises(ValueError, match="version"): fs.pipe_file("s3://bucket/key?versionId=12345abcde", b"data") + def test_pipe_file_non_contiguous_memoryview(self): + # A non-contiguous memoryview within the block size in items, 4 items + # of 8 bytes here, is uploaded with PutObject, as the buffered path + # cannot write it. + fs = self._make_fs() + fs._put_object = mock.MagicMock() + value = memoryview(b"ab" * 8).cast("H")[::2] + + fs.pipe_file("s3://bucket/key", value, block_size=6) + + fs._put_object.assert_called_once_with(bucket="bucket", key="key", body=b"ab" * 4) + def test_pipe_file_small_drops_max_workers(self): fs = self._make_fs() fs._put_object = mock.MagicMock()