Skip to content
13 changes: 13 additions & 0 deletions docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,19 @@ 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.

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
Expand Down
98 changes: 87 additions & 11 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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))

Copy link
Copy Markdown
Member Author

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_SIZE because the validated block_size is 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 through self._sync_fs._get_copy_ranges.

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:

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review rounds one and two of the pre-check extension — base f3ef301d9fec92f1fc53b147bebfef87bfc919ca, range 8b99abbb7612b1594263d626b75f7b79a51d5cc7..6901304bea3cb5112903407b4b7d5d78eb50ed38 (a contract expansion requested by the maintainer after consulting claude-fable-5-1). Result: CLEAN.

Round one (behavior): the check runs before parse_path-dependent work, exists() (mode="create"), any API call, and the transaction branch of pipe_file, so nothing is opened or added to transaction.files; sizes up to one block never reach it; an empty file passes; AioS3FileSystem._put_file/_pipe_file delegate to the sync methods; put_file passes block_size to open(), where S3File.__init__ still rejects a write block size below 5 MiB. New tests fail without the source change (8 cases) and pass with it; offline tests/pyathena/filesystem/: 171 passed, the same 108 credential-less integration failures as the base.

Round two (claims, callers): size > block_size * MULTIPART_UPLOAD_MAX_PARTS is exact for put_file (it writes remote.blocksize chunks, so parts = ceil(size / block)) and is the documented rule for pipe_file, which therefore rejects the sizes in (10,000·B, 10,000·B + 5 MiB) that the tail merge would have fit (release-note). The minimum block size in the message is max(ceil(size / 10,000), 5 MiB), so it never suggests a size open() rejects. put(…, block_size=…) previously reached the S3 API as a parameter and failed, so accepting it breaks no caller. Pre-existing and out of scope: put_file still forwards max_workers to the S3 API.

"""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:
Expand All @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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):
Expand Down Expand Up @@ -2153,6 +2196,7 @@ class S3File(AbstractBufferedFile):
"""

fs: S3FileSystem
buffer: BytesIO | None

def __init__(
self,
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand All @@ -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

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The 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 gpt-6-astra, reasoning effort max, sandbox read-only, session 01a100db-fc47-7d52-93c8-beff25916d64. Base f3ef301d9fec92f1fc53b147bebfef87bfc919ca, head 6901304bea3cb5112903407b4b7d5d78eb50ed38 (full diff, detached snapshot; intended behavior and conventions only). Static review only. Result: FINDINGS (2). Covered (reviewer): the supplied diff, sync/async upload paths, fsspec 2026.9.0 lifecycle and transactions, append accounting, merged/split tails, copy-range bounds through 5 TiB, pre-check placement, typing, documentation, and tests.

1. P2 — Mid-stream overflow leaves the executor running (introduced). Outside a context manager, writing past the limit sets closed=True, drops the buffer, and aborts, but never shuts down the executor; commit()/discard() do not, and fsspec's __del__() skips closed files. The tests masked it because their helper always exits a context manager.

Author verification: confirmed (fsspec __del__: if not self.closed: self.close()).

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Repaired in da55411: the error path calls self._executor.shutdown() after the abort (idempotent with close(); S3AioExecutor.shutdown() is a no-op). New test_write_exceeding_max_parts_without_close (no context manager) fails without the repair and passes with it; the existing failure tests now assert shutdown was called rather than called once. Self-review of the repair (behavior: waits for in-flight parts like close(); claims: fsspec __del__ skips closed files) found nothing further.

self.closed = True

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed) — reviewer: Codex CLI 0.160.0, model gpt-6-astra, reasoning effort max, sandbox read-only, session 01a100b3-237e-7d81-81b1-935b1d6ce6d7. Base f3ef301d9fec92f1fc53b147bebfef87bfc919ca, head 41abb14fde412a6a81dd20d54201db3875cac369, detached snapshot; given the diff, the intended behavior, and repository conventions only (no PR text, commit messages, or prior findings). Static review only: no builds, tests, or network access. Result: FINDINGS (3).

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 None through the cast.

1. P2 — Limit failures can leave the executor running. (reviewer text, condensed) With autocommit=False, write 10,000 full blocks followed by one byte, then call close(). The final flush raises ValueError, so execution never reaches _executor.shutdown() at pyathena/filesystem/s3.py:2286; the file is already marked closed, so __del__ skips cleanup. A direct mid-stream write() failure also closes the file without shutting down its executor. The new failure path introduces this exposure; close() lacking exception-safe shutdown is pre-existing.

Author verification: confirmed for the final-flush path (super().close() raises before shutdown()); a mid-stream failure followed by close() reaches shutdown() because super().close() returns early on a closed file. Pre-existing for other final-flush errors; folded as a contained fix.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Repaired in b8224b1: S3File.close() shuts the executor down in a finally block, so a failing final flush (the part limit error or any other) no longer skips it. The new executor.shutdown.assert_called_once() in test_write_exceeding_max_parts fails for the close()-path cases without the repair and passes with it.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review of the repair b8224b1 (both perspectives, range 41abb14f..b8224b19, same base). Behavior: the try/finally in close() does not change which exception propagates; S3AioExecutor.shutdown() is a no-op and S3ThreadPoolExecutor.shutdown() behaves as on a normal close; a mid-stream error followed by close() reaches shutdown() as before. Claims: put_file writes remote.blocksize chunks (one part per block) and pipe writes once (blocks plus a merged tail), so the docs formula holds for them; the open sentence says "can" because a remainder below the minimum part size is merged. Offline tests/pyathena/filesystem/: 163 passed, the same 108 credential-less integration failures as the base. Result: CLEAN.

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(
Expand Down
Loading
Loading