-
Notifications
You must be signed in to change notification settings - Fork 116
Keep multipart uploads within the 10,000-part limit #968
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
692167d
ff6de77
41abb14
b8224b1
8b99abb
6901304
da55411
cf3cddf
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -3,13 +3,15 @@ | |
| from __future__ import annotations | ||
|
|
||
| import logging | ||
| import math | ||
| import mimetypes | ||
| import os.path | ||
| import re | ||
| from collections.abc import Callable, Iterator | ||
| from concurrent.futures import Future, as_completed | ||
| 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 | ||
|
|
@@ -1277,7 +1282,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 | ||
|
|
@@ -1289,21 +1295,48 @@ 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() | ||
| if size - starts[-1] > self.MULTIPART_UPLOAD_MAX_PART_SIZE: | ||
| 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: | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Self-review rounds one and two of the pre-check extension — base Round one (behavior): the check runs before Round two (claims, callers): |
||
| """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: | ||
|
|
@@ -1330,9 +1363,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 | ||
|
|
@@ -1466,6 +1502,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, | ||
|
|
@@ -1482,14 +1523,16 @@ 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) | ||
| if content_type is not None: | ||
| 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): | ||
|
|
@@ -2153,6 +2196,7 @@ class S3File(AbstractBufferedFile): | |
| """ | ||
|
|
||
| fs: S3FileSystem | ||
| buffer: BytesIO | None | ||
|
|
||
| def __init__( | ||
| self, | ||
|
|
@@ -2273,8 +2317,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: | ||
|
|
@@ -2337,15 +2384,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 | ||
|
|
@@ -2360,6 +2410,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 | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent review of the expanded scope (relayed) — reviewer: Codex CLI 0.160.0, model 1. P2 — Mid-stream overflow leaves the executor running (introduced). Outside a context manager, writing past the limit sets Author verification: confirmed (fsspec
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Repaired in da55411: the error path calls |
||
| self.closed = True | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent review (relayed) — reviewer: Codex CLI 0.160.0, model Covered (reviewer): the supplied diff, sync/async write and copy callers, append accounting, fsspec 2026.9.0 buffering and transactions, abort handling, executor lifecycle, typing, documentation, and tests. The reviewer confirmed that the per-part guard covers multi-block writes, merged/split tails, and copied append parts; that clearing the state prevents later publication, including after an abort exception; that copy ranges satisfy count, size, and coverage bounds through 5 TiB; and that no supported flush path exposes 1. P2 — Limit failures can leave the executor running. (reviewer text, condensed) With Author verification: confirmed for the final-flush path (
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Repaired in b8224b1:
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Self-review of the repair b8224b1 (both perspectives, range |
||
| 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( | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Self-review round one — copy split size checked, no finding:
math.ceil(size / 10_000)is exact for sizes within the 5 TiB object limit (below 2**53), stays ≥MULTIPART_UPLOAD_MIN_PART_SIZEbecause the validatedblock_sizeis the lower bound, and stays ≤ about 550 MB, so a merged short tail never needs the half split and the count stays ≤ 10,000. The async copy reaches it throughself._sync_fs._get_copy_ranges.