Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion cpp/src/arrow/dataset/file_ipc.cc
Original file line number Diff line number Diff line change
Expand Up @@ -147,7 +147,7 @@ Result<RecordBatchGenerator> IpcFileFormat::ScanBatchesAsync(
-> Future<std::shared_ptr<ipc::RecordBatchFileReader>> {
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;
Expand Down
97 changes: 97 additions & 0 deletions cpp/src/arrow/dataset/file_ipc_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,11 @@

#include "arrow/dataset/file_ipc.h"

#include <atomic>
#include <chrono>
#include <future>
#include <memory>
#include <thread>
#include <utility>
#include <vector>

Expand All @@ -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 {

Expand Down Expand Up @@ -62,6 +67,98 @@ class IpcFormatHelper {

class TestIpcFileFormat : public FileFormatFixtureMixin<IpcFormatHelper> {};

class GatedIpcBufferReader : public io::BufferReader {
public:
explicit GatedIpcBufferReader(std::shared_ptr<Buffer> buffer)
: io::BufferReader(std::move(buffer)),
pending_(Future<std::shared_ptr<Buffer>>::Make()) {}

Future<std::shared_ptr<Buffer>> 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<void> Started() { return started_promise_.get_future(); }

void Release() {
pending_.MarkFinished(io::BufferReader::ReadAsync(io::IOContext(), position_, nbytes_,
allow_short_read_)
.result());
}

private:
std::atomic<bool> started_{false};
std::promise<void> started_promise_;
Future<std::shared_ptr<Buffer>> 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<GatedIpcBufferReader>(buffer);
auto read_started = delayed_reader->Started();
std::atomic<int> opens{0};
FileSource source(
[&]() -> Result<std::shared_ptr<io::RandomAccessFile>> {
if (opens.fetch_add(1) == 0) {
return std::make_shared<io::BufferReader>(buffer);
}
return delayed_reader;
},
buffer->size());
SetSchema(schema_->fields());
auto fragment = MakeFragment(source);

std::promise<Result<RecordBatchGenerator>> 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<std::promise<void>>();
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) {
Expand Down
Loading