From 19f9a99cb9ede39b888993919b04c869eb46b0d4 Mon Sep 17 00:00:00 2001 From: Edmond Abraham <54119952+gitedmond@users.noreply.github.com> Date: Sun, 27 Sep 2026 23:31:19 -0400 Subject: [PATCH 1/2] GH-51547: Reopen IPC dataset readers asynchronously --- cpp/src/arrow/dataset/file_ipc.cc | 2 +- cpp/src/arrow/dataset/file_ipc_test.cc | 92 ++++++++++++++++++++++++++ 2 files changed, 93 insertions(+), 1 deletion(-) diff --git a/cpp/src/arrow/dataset/file_ipc.cc b/cpp/src/arrow/dataset/file_ipc.cc index 8e19b2bbee2a..f1906eaa7d58 100644 --- a/cpp/src/arrow/dataset/file_ipc.cc +++ b/cpp/src/arrow/dataset/file_ipc.cc @@ -147,7 +147,7 @@ Result IpcFileFormat::ScanBatchesAsync( -> Future> { ARROW_ASSIGN_OR_RAISE(auto options, GetReadOptions(*reader->schema(), *self, *options)); - return OpenReader(source, options); + return OpenReaderAsync(source, options); }; auto readahead_level = options->batch_readahead; auto default_fragment_scan_options = this->default_fragment_scan_options; diff --git a/cpp/src/arrow/dataset/file_ipc_test.cc b/cpp/src/arrow/dataset/file_ipc_test.cc index 370c6d978200..556faac297e5 100644 --- a/cpp/src/arrow/dataset/file_ipc_test.cc +++ b/cpp/src/arrow/dataset/file_ipc_test.cc @@ -17,7 +17,11 @@ #include "arrow/dataset/file_ipc.h" +#include +#include +#include #include +#include #include #include @@ -34,6 +38,7 @@ #include "arrow/testing/gtest_util.h" #include "arrow/testing/util.h" #include "arrow/util/key_value_metadata.h" +#include "arrow/util/thread_pool.h" namespace arrow { @@ -62,6 +67,93 @@ class IpcFormatHelper { class TestIpcFileFormat : public FileFormatFixtureMixin {}; +class GatedIpcBufferReader : public io::BufferReader { + public: + explicit GatedIpcBufferReader(std::shared_ptr buffer) + : io::BufferReader(std::move(buffer)), + pending_(Future>::Make()) {} + + Future> ReadAsync(const io::IOContext& context, + int64_t position, int64_t nbytes, + bool allow_short_read) override { + if (!started_.exchange(true)) { + position_ = position; + nbytes_ = nbytes; + allow_short_read_ = allow_short_read; + started_promise_.set_value(); + return pending_; + } + return io::BufferReader::ReadAsync(context, position, nbytes, allow_short_read); + } + + std::future Started() { return started_promise_.get_future(); } + + void Release() { + pending_.MarkFinished(io::BufferReader::ReadAsync(io::IOContext(), position_, nbytes_, + allow_short_read_) + .result()); + } + + private: + std::atomic started_{false}; + std::promise started_promise_; + Future> pending_; + int64_t position_ = 0; + int64_t nbytes_ = 0; + bool allow_short_read_ = false; +}; + +TEST_F(TestIpcFileFormat, ReopeningReaderDoesNotWaitForIO) { + using namespace std::chrono_literals; + const int original_capacity = GetCpuThreadPoolCapacity(); + ASSERT_OK(SetCpuThreadPoolCapacity(1)); + auto schema_ = schema({field("f64", float64())}); + auto reader = GetRecordBatchReader(schema_); + auto buffer = GetFileSource(reader.get())->buffer(); + auto delayed_reader = std::make_shared(buffer); + auto read_started = delayed_reader->Started(); + std::atomic opens{0}; + FileSource source( + [&]() -> Result> { + if (opens.fetch_add(1) == 0) { + return std::make_shared(buffer); + } + return delayed_reader; + }, + buffer->size()); + SetSchema(schema_->fields()); + auto fragment = MakeFragment(source); + + std::promise> scan_promise; + auto scan_result = scan_promise.get_future(); + std::thread scan_thread( + [&] { scan_promise.set_value(fragment->ScanBatchesAsync(opts_)); }); + + // The second open has reached an outstanding IO operation. The scan setup must + // be able to return without waiting for that operation to finish. + auto started = read_started.wait_for(5s); + auto returned = scan_result.wait_for(1s); + std::promise pool_progress; + auto progress = pool_progress.get_future(); + ASSERT_OK(::arrow::internal::GetCpuThreadPool()->Spawn( + [&pool_progress] { pool_progress.set_value(); })); + auto progressed = progress.wait_for(1s); + delayed_reader->Release(); + scan_thread.join(); + + // All callbacks have been unblocked before restoring the shared executor. + ASSERT_EQ(progress.wait_for(5s), std::future_status::ready); + ASSERT_OK(SetCpuThreadPoolCapacity(original_capacity)); + + ASSERT_EQ(started, std::future_status::ready); + ASSERT_EQ(returned, std::future_status::ready) + << "IPC scan setup blocked waiting for the second reader's IO"; + ASSERT_EQ(progressed, std::future_status::ready) + << "IPC reader reopening blocked the only CPU worker waiting for IO"; + ASSERT_OK_AND_ASSIGN(auto batches, scan_result.get()); + ASSERT_FINISHES_OK(CollectAsyncGenerator(std::move(batches))); +} + TEST_F(TestIpcFileFormat, WriteRecordBatchReader) { TestWrite(); } TEST_F(TestIpcFileFormat, WriteRecordBatchReaderCustomOptions) { From 8729183bb2d5a61c2129f6b014ddc5e89f5033d5 Mon Sep 17 00:00:00 2001 From: Edmond Abraham <54119952+gitedmond@users.noreply.github.com> Date: Sun, 27 Sep 2026 23:45:37 -0400 Subject: [PATCH 2/2] GH-51547: Make IPC regression test cleanup safe --- cpp/src/arrow/dataset/file_ipc_test.cc | 25 +++++++++++++++---------- 1 file changed, 15 insertions(+), 10 deletions(-) diff --git a/cpp/src/arrow/dataset/file_ipc_test.cc b/cpp/src/arrow/dataset/file_ipc_test.cc index 556faac297e5..62094ee6b8ac 100644 --- a/cpp/src/arrow/dataset/file_ipc_test.cc +++ b/cpp/src/arrow/dataset/file_ipc_test.cc @@ -104,7 +104,6 @@ class GatedIpcBufferReader : public io::BufferReader { }; TEST_F(TestIpcFileFormat, ReopeningReaderDoesNotWaitForIO) { - using namespace std::chrono_literals; const int original_capacity = GetCpuThreadPoolCapacity(); ASSERT_OK(SetCpuThreadPoolCapacity(1)); auto schema_ = schema({field("f64", float64())}); @@ -131,19 +130,25 @@ TEST_F(TestIpcFileFormat, ReopeningReaderDoesNotWaitForIO) { // The second open has reached an outstanding IO operation. The scan setup must // be able to return without waiting for that operation to finish. - auto started = read_started.wait_for(5s); - auto returned = scan_result.wait_for(1s); - std::promise pool_progress; - auto progress = pool_progress.get_future(); - ASSERT_OK(::arrow::internal::GetCpuThreadPool()->Spawn( - [&pool_progress] { pool_progress.set_value(); })); - auto progressed = progress.wait_for(1s); + auto started = read_started.wait_for(std::chrono::seconds(5)); + auto returned = scan_result.wait_for(std::chrono::seconds(1)); + auto pool_progress = std::make_shared>(); + auto progress = pool_progress->get_future(); + auto spawn_status = ::arrow::internal::GetCpuThreadPool()->Spawn( + [pool_progress] { pool_progress->set_value(); }); + auto progressed = spawn_status.ok() ? progress.wait_for(std::chrono::seconds(1)) + : std::future_status::timeout; delayed_reader->Release(); scan_thread.join(); // All callbacks have been unblocked before restoring the shared executor. - ASSERT_EQ(progress.wait_for(5s), std::future_status::ready); - ASSERT_OK(SetCpuThreadPoolCapacity(original_capacity)); + auto eventually_progressed = spawn_status.ok() + ? progress.wait_for(std::chrono::seconds(5)) + : std::future_status::timeout; + auto restore_status = SetCpuThreadPoolCapacity(original_capacity); + ASSERT_OK(spawn_status); + ASSERT_EQ(eventually_progressed, std::future_status::ready); + ASSERT_OK(restore_status); ASSERT_EQ(started, std::future_status::ready); ASSERT_EQ(returned, std::future_status::ready)