diff --git a/docs/filesystem.md b/docs/filesystem.md index 4f8135cf0..0096ee03b 100644 --- a/docs/filesystem.md +++ b/docs/filesystem.md @@ -85,6 +85,19 @@ 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. `put` and `pipe` upload one part +per block, so with the default block size they can upload up to about 48.8 GiB +(10,000 × 5 MiB). To upload a larger object, use a block size of at least its size +divided by 10,000, either with the `block_size` argument of `put`, `pipe`, and `open` +or with the `default_block_size` argument of `S3FileSystem`. `put` and `pipe` check +the size before uploading anything and raise `ValueError` with the minimum block size +if the data needs more parts. A file written with `open` can take more parts, because +each write that fills the buffer uploads the data beyond its last full block as a +separate part when that data is at least 5 MiB. In an append, the parts copied from +the existing object also count toward the limit. A write with `open` that reaches the +limit raises `ValueError` and aborts its multipart upload. Multipart copies with `cp` +use parts large enough to stay within the limit. + ## 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..e5cf41035 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() @@ -1345,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: @@ -1371,9 +1404,12 @@ 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 + # 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): # Defer to the buffered open() path, which keeps the # deferred-commit semantics of fsspec transactions and uploads @@ -1515,6 +1551,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, @@ -1531,6 +1572,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) @@ -1538,7 +1581,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): @@ -2219,6 +2262,7 @@ class S3File(AbstractBufferedFile): """ fs: S3FileSystem + buffer: BytesIO | None def __init__( self, @@ -2365,8 +2409,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: @@ -2429,15 +2476,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 +2502,32 @@ def _upload_chunk(self, final: bool = False) -> bool: uploads = [data] for upload in uploads: + if part_number >= self.fs.MULTIPART_UPLOAD_MAX_PARTS: + # 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. 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: + 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 = [] + 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 " + f"{self.blocksize} bytes. Write the file with a block_size, or " + "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( self._executor.submit( diff --git a/tests/pyathena/filesystem/test_s3.py b/tests/pyathena/filesystem/test_s3.py index 95b47d016..c17a9bef7 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 @@ -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 @@ -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() @@ -514,6 +526,82 @@ 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( + ("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() + 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", value, **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" @@ -806,6 +894,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 +2311,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 +2344,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 +2458,125 @@ 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 + + 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 + # uploads nothing. + if not autocommit: + f.commit() + + 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() + 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() + + @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") + + 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 [c.kwargs["part_number"] for c in executor.submit.call_args_list] == [1, 2, 3] + 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,