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..62094ee6b8ac 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,98 @@ 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) { + 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(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. + 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) + << "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) {