From 67c2a153b5b564b2de90e82dd5b8abeec57e5f16 Mon Sep 17 00:00:00 2001 From: Adrien Courtois Date: Tue, 29 Sep 2026 10:14:39 +0200 Subject: [PATCH 1/6] GH-40502: [Python] Expose NativeFile.abort() --- python/pyarrow/includes/libarrow.pxd | 1 + python/pyarrow/io.pxi | 14 ++++++++++++++ python/pyarrow/tests/test_fs.py | 15 +++++++++++++++ python/pyarrow/tests/test_io.py | 10 ++++++++++ 4 files changed, 40 insertions(+) diff --git a/python/pyarrow/includes/libarrow.pxd b/python/pyarrow/includes/libarrow.pxd index 17dcb87a80e9..6f54eb646a76 100644 --- a/python/pyarrow/includes/libarrow.pxd +++ b/python/pyarrow/includes/libarrow.pxd @@ -1633,6 +1633,7 @@ cdef extern from "arrow/io/api.h" namespace "arrow::io" nogil: cdef cppclass FileInterface: CStatus Close() + CStatus Abort() CResult[int64_t] Tell() FileMode mode() c_bool closed() diff --git a/python/pyarrow/io.pxi b/python/pyarrow/io.pxi index dbecd80909c6..aad002afce8b 100644 --- a/python/pyarrow/io.pxi +++ b/python/pyarrow/io.pxi @@ -193,6 +193,20 @@ cdef class NativeFile(_Weakrefable): else: check_status(self.output_stream.get().Close()) + def abort(self): + """ + Close the stream, discarding written data if the stream supports it. + + For example, an S3 output stream aborts its multipart upload, so no + object is written. Other streams are simply closed. + """ + if not self.closed: + with nogil: + if self.is_readable: + check_status(self.input_stream.get().Abort()) + else: + check_status(self.output_stream.get().Abort()) + cdef set_random_access_file(self, shared_ptr[CRandomAccessFile] handle): self.input_stream = handle self.random_access = handle diff --git a/python/pyarrow/tests/test_fs.py b/python/pyarrow/tests/test_fs.py index 5bf1950c0654..d0d42f831f76 100644 --- a/python/pyarrow/tests/test_fs.py +++ b/python/pyarrow/tests/test_fs.py @@ -1143,6 +1143,21 @@ def test_open_output_stream_metadata(fs, pathfn): assert got_metadata == {} +def test_open_output_stream_abort(fs, pathfn): + p = pathfn('open-output-stream-abort') + with fs.open_output_stream(p) as f: + f.write(b'some data') + f.abort() + assert f.closed + + # Other filesystems may keep the written data, so only these are checked + if fs.type_name == 's3': + assert fs.get_file_info(p).type == FileType.NotFound + elif 'mock' in fs.type_name: + with fs.open_input_stream(p) as f: + assert f.read().startswith(b'MockFSOutputStream aborted') + + def test_localfs_options(): # LocalFileSystem instantiation LocalFileSystem(use_mmap=False) diff --git a/python/pyarrow/tests/test_io.py b/python/pyarrow/tests/test_io.py index 1cfebb4936a1..88e14143664a 100644 --- a/python/pyarrow/tests/test_io.py +++ b/python/pyarrow/tests/test_io.py @@ -989,6 +989,16 @@ def test_inmemory_write_after_closed(): f.write(b'not ok') +def test_inmemory_write_after_abort(): + f = pa.BufferOutputStream() + f.write(b'ok') + f.abort() + assert f.closed + + with pytest.raises(ValueError): + f.write(b'not ok') + + def test_buffer_protocol_ref_counting(): def make_buffer(bytes_obj): return bytearray(pa.py_buffer(bytes_obj)) From 6c0fccb2596193b0a68a296c65c2b14de258c2a6 Mon Sep 17 00:00:00 2001 From: Adrien Courtois Date: Tue, 29 Sep 2026 10:52:53 +0200 Subject: [PATCH 2/6] Test abort on readable streams, stream wrappers and uploaded S3 parts --- python/pyarrow/tests/test_fs.py | 27 +++++++++++++++++++++++++-- python/pyarrow/tests/test_io.py | 9 +++++++++ 2 files changed, 34 insertions(+), 2 deletions(-) diff --git a/python/pyarrow/tests/test_fs.py b/python/pyarrow/tests/test_fs.py index d0d42f831f76..195f09fe7295 100644 --- a/python/pyarrow/tests/test_fs.py +++ b/python/pyarrow/tests/test_fs.py @@ -1143,9 +1143,19 @@ def test_open_output_stream_metadata(fs, pathfn): assert got_metadata == {} -def test_open_output_stream_abort(fs, pathfn): +@pytest.mark.gzip +@pytest.mark.parametrize( + ('compression', 'buffer_size'), + [ + (None, None), + (None, 64), + ('gzip', None), + ('gzip', 256), + ] +) +def test_open_output_stream_abort(fs, pathfn, compression, buffer_size): p = pathfn('open-output-stream-abort') - with fs.open_output_stream(p) as f: + with fs.open_output_stream(p, compression, buffer_size) as f: f.write(b'some data') f.abort() assert f.closed @@ -1158,6 +1168,19 @@ def test_open_output_stream_abort(fs, pathfn): assert f.read().startswith(b'MockFSOutputStream aborted') +@pytest.mark.s3 +def test_s3_output_stream_abort_after_part_upload(s3fs): + fs, pathfn = s3fs['fs'], s3fs['pathfn'] + p = pathfn('abort-after-part-upload') + with fs.open_output_stream(p) as f: + # Flushing a full 10 MiB part waits for its upload to complete + f.write(b'x' * 10 * 1024 * 1024) + f.flush() + f.abort() + + assert fs.get_file_info(p).type == FileType.NotFound + + def test_localfs_options(): # LocalFileSystem instantiation LocalFileSystem(use_mmap=False) diff --git a/python/pyarrow/tests/test_io.py b/python/pyarrow/tests/test_io.py index 88e14143664a..bafa89ff99dc 100644 --- a/python/pyarrow/tests/test_io.py +++ b/python/pyarrow/tests/test_io.py @@ -999,6 +999,15 @@ def test_inmemory_write_after_abort(): f.write(b'not ok') +def test_inmemory_read_after_abort(): + f = pa.BufferReader(b'data') + f.abort() + assert f.closed + + with pytest.raises(ValueError): + f.read() + + def test_buffer_protocol_ref_counting(): def make_buffer(bytes_obj): return bytearray(pa.py_buffer(bytes_obj)) From b23ad1bff53a29bd6b3910bc99b56fe7fd8a72be Mon Sep 17 00:00:00 2001 From: Adrien Courtois Date: Tue, 29 Sep 2026 13:08:19 +0200 Subject: [PATCH 3/6] Close the S3 output stream even if aborting fails --- cpp/src/arrow/filesystem/s3fs.cc | 11 ++++++----- python/pyarrow/tests/test_fs.py | 23 +++++++++++++++++++++++ 2 files changed, 29 insertions(+), 5 deletions(-) diff --git a/cpp/src/arrow/filesystem/s3fs.cc b/cpp/src/arrow/filesystem/s3fs.cc index 1c6763a4aee9..181328c8be57 100644 --- a/cpp/src/arrow/filesystem/s3fs.cc +++ b/cpp/src/arrow/filesystem/s3fs.cc @@ -1748,8 +1748,13 @@ class ObjectOutputStream final : public io::OutputStream { return Status::OK(); } + // Close the stream even if aborting fails, so that the upload is never completed. + auto holder = std::move(holder_); + current_part_.reset(); + closed_ = true; + if (IsMultipartCreated()) { - ARROW_ASSIGN_OR_RAISE(auto client_lock, holder_->Lock()); + ARROW_ASSIGN_OR_RAISE(auto client_lock, holder->Lock()); S3Model::AbortMultipartUploadRequest req; req.SetBucket(ToAwsString(path_.bucket)); @@ -1765,10 +1770,6 @@ class ObjectOutputStream final : public io::OutputStream { } } - current_part_.reset(); - holder_ = nullptr; - closed_ = true; - return Status::OK(); } diff --git a/python/pyarrow/tests/test_fs.py b/python/pyarrow/tests/test_fs.py index 195f09fe7295..a0f6a71dc578 100644 --- a/python/pyarrow/tests/test_fs.py +++ b/python/pyarrow/tests/test_fs.py @@ -1181,6 +1181,29 @@ def test_s3_output_stream_abort_after_part_upload(s3fs): assert fs.get_file_info(p).type == FileType.NotFound +@pytest.mark.s3 +def test_s3_output_stream_failed_abort(s3_server): + from pyarrow.fs import S3FileSystem + # The limited user isn't allowed to abort multipart uploads + _configure_s3_limited_user(s3_server, _minio_limited_policy, + 'test_fs_abort_user', 'abort123') + host, port, _, _ = s3_server['connection'] + fs = S3FileSystem( + access_key='test_fs_abort_user', + secret_key='abort123', + endpoint_override=f'{host}:{port}', + scheme='http' + ) + p = 'existing-bucket/failed-abort' + with fs.open_output_stream(p) as f: + f.write(b'some data') + with pytest.raises(OSError, match="AbortMultipartUpload"): + f.abort() + assert f.closed + + assert fs.get_file_info(p).type == FileType.NotFound + + def test_localfs_options(): # LocalFileSystem instantiation LocalFileSystem(use_mmap=False) From ac97c71f2fb56cfff16192e5b4dce741d858f45c Mon Sep 17 00:00:00 2001 From: Adrien Courtois Date: Tue, 29 Sep 2026 13:22:28 +0200 Subject: [PATCH 4/6] Test failed abort through the with block and the destructor --- cpp/src/arrow/filesystem/s3fs.cc | 3 ++- python/pyarrow/tests/test_fs.py | 10 ++++++---- 2 files changed, 8 insertions(+), 5 deletions(-) diff --git a/cpp/src/arrow/filesystem/s3fs.cc b/cpp/src/arrow/filesystem/s3fs.cc index 181328c8be57..ac2731f521c3 100644 --- a/cpp/src/arrow/filesystem/s3fs.cc +++ b/cpp/src/arrow/filesystem/s3fs.cc @@ -1748,7 +1748,8 @@ class ObjectOutputStream final : public io::OutputStream { return Status::OK(); } - // Close the stream even if aborting fails, so that the upload is never completed. + // Close the stream even if aborting fails, so that a later Close() doesn't + // complete the upload. auto holder = std::move(holder_); current_part_.reset(); closed_ = true; diff --git a/python/pyarrow/tests/test_fs.py b/python/pyarrow/tests/test_fs.py index a0f6a71dc578..ebcad7c5305c 100644 --- a/python/pyarrow/tests/test_fs.py +++ b/python/pyarrow/tests/test_fs.py @@ -1195,12 +1195,14 @@ def test_s3_output_stream_failed_abort(s3_server): scheme='http' ) p = 'existing-bucket/failed-abort' - with fs.open_output_stream(p) as f: - f.write(b'some data') - with pytest.raises(OSError, match="AbortMultipartUpload"): + with pytest.raises(OSError, match="AbortMultipartUpload"): + with fs.open_output_stream(p) as f: + f.write(b'some data') f.abort() - assert f.closed + assert f.closed + # Neither closing nor destroying the stream completes the upload + del f assert fs.get_file_info(p).type == FileType.NotFound From 0de1c38390e1fde2ed673cc935242e286f788b27 Mon Sep 17 00:00:00 2001 From: Adrien Courtois Date: Tue, 29 Sep 2026 13:34:22 +0200 Subject: [PATCH 5/6] Tighten comment --- cpp/src/arrow/filesystem/s3fs.cc | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/cpp/src/arrow/filesystem/s3fs.cc b/cpp/src/arrow/filesystem/s3fs.cc index ac2731f521c3..e72695c096d8 100644 --- a/cpp/src/arrow/filesystem/s3fs.cc +++ b/cpp/src/arrow/filesystem/s3fs.cc @@ -1748,8 +1748,7 @@ class ObjectOutputStream final : public io::OutputStream { return Status::OK(); } - // Close the stream even if aborting fails, so that a later Close() doesn't - // complete the upload. + // Close even if the abort request fails auto holder = std::move(holder_); current_part_.reset(); closed_ = true; From 754c66f92eef5d2442dba21ece2d1462756fcfda Mon Sep 17 00:00:00 2001 From: Adrien Courtois Date: Tue, 29 Sep 2026 13:38:21 +0200 Subject: [PATCH 6/6] Reword test comments --- python/pyarrow/tests/test_fs.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/python/pyarrow/tests/test_fs.py b/python/pyarrow/tests/test_fs.py index ebcad7c5305c..4db78a04ad95 100644 --- a/python/pyarrow/tests/test_fs.py +++ b/python/pyarrow/tests/test_fs.py @@ -1160,7 +1160,7 @@ def test_open_output_stream_abort(fs, pathfn, compression, buffer_size): f.abort() assert f.closed - # Other filesystems may keep the written data, so only these are checked + # Note that only S3 discards the written data on abort for now if fs.type_name == 's3': assert fs.get_file_info(p).type == FileType.NotFound elif 'mock' in fs.type_name: @@ -1201,7 +1201,6 @@ def test_s3_output_stream_failed_abort(s3_server): f.abort() assert f.closed - # Neither closing nor destroying the stream completes the upload del f assert fs.get_file_info(p).type == FileType.NotFound