MCPcopy Create free account
hub / github.com/apache/arrow / TestDoExchangeTotal

Method TestDoExchangeTotal

cpp/src/arrow/flight/test_definitions.cc:487–539  ·  view source on GitHub ↗

Test interleaved reading/writing

Source from the content-addressed store, hash-verified

485}
486// Test interleaved reading/writing
487void DataTest::TestDoExchangeTotal() {
488 auto descr = FlightDescriptor::Command("total");
489 std::unique_ptr<FlightStreamReader> reader;
490 std::unique_ptr<FlightStreamWriter> writer;
491 {
492 auto a1 = ArrayFromJSON(arrow::int32(), "[4, 5, 6, null]");
493 auto schema = arrow::schema({field("f1", a1->type())});
494 // XXX: as noted in flight/client.cc, Begin() is lazy and the
495 // schema message won't be written until some data is also
496 // written. There's also timing issues; hence we check each status
497 // here.
498 EXPECT_RAISES_WITH_MESSAGE_THAT(
499 Invalid, ::testing::HasSubstr("Field is not INT64: f1"), ([&]() {
500 ARROW_ASSIGN_OR_RAISE(auto exchange, client_->DoExchange(descr));
501 reader = std::move(exchange.reader);
502 writer = std::move(exchange.writer);
503 RETURN_NOT_OK(writer->Begin(schema));
504 auto batch = RecordBatch::Make(schema, /* num_rows */ 4, {a1});
505 RETURN_NOT_OK(writer->WriteRecordBatch(*batch));
506 return writer->Close();
507 })());
508 }
509 {
510 auto a1 = ArrayFromJSON(arrow::int64(), "[1, 2, null, 3]");
511 auto a2 = ArrayFromJSON(arrow::int64(), "[null, 4, 5, 6]");
512 auto schema = arrow::schema({field("f1", a1->type()), field("f2", a2->type())});
513 ASSERT_OK_AND_ASSIGN(auto exchange, client_->DoExchange(descr));
514 reader = std::move(exchange.reader);
515 writer = std::move(exchange.writer);
516 ASSERT_OK(writer->Begin(schema));
517 auto batch = RecordBatch::Make(schema, /* num_rows */ 4, {a1, a2});
518 ASSERT_OK(writer->WriteRecordBatch(*batch));
519 ASSERT_OK_AND_ASSIGN(auto server_schema, reader->GetSchema());
520 AssertSchemaEqual(*schema, *server_schema);
521
522 ASSERT_OK_AND_ASSIGN(auto chunk, reader->Next());
523 ASSERT_NE(nullptr, chunk.data);
524 auto expected1 = RecordBatch::Make(
525 schema, /* num_rows */ 1,
526 {ArrayFromJSON(arrow::int64(), "[6]"), ArrayFromJSON(arrow::int64(), "[15]")});
527 AssertBatchesEqual(*expected1, *chunk.data);
528
529 ASSERT_OK(writer->WriteRecordBatch(*batch));
530 ASSERT_OK_AND_ASSIGN(chunk, reader->Next());
531 ASSERT_NE(nullptr, chunk.data);
532 auto expected2 = RecordBatch::Make(
533 schema, /* num_rows */ 1,
534 {ArrayFromJSON(arrow::int64(), "[12]"), ArrayFromJSON(arrow::int64(), "[30]")});
535 AssertBatchesEqual(*expected2, *chunk.data);
536
537 ASSERT_OK(writer->Close());
538 }
539}
540// Ensure server errors get propagated no matter what we try
541void DataTest::TestDoExchangeError() {
542 auto descr = FlightDescriptor::Command("error");

Callers

nothing calls this directly

Calls 12

ArrayFromJSONFunction · 0.85
AssertBatchesEqualFunction · 0.85
CommandFunction · 0.70
ASSERT_OK_AND_ASSIGNFunction · 0.70
schemaFunction · 0.50
fieldFunction · 0.50
MakeFunction · 0.50
typeMethod · 0.45
BeginMethod · 0.45
WriteRecordBatchMethod · 0.45
CloseMethod · 0.45
NextMethod · 0.45

Tested by

no test coverage detected