MCPcopy Create free account
hub / github.com/apache/impala / Sender

Method Sender

be/src/runtime/data-stream-test.cc:656–705  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

654 }
655
656 void Sender(int sender_num, int channel_buffer_size,
657 TPartitionType::type partition_type, SenderInfo* info, bool reset_hash_seed,
658 int num_batches) {
659 RuntimeState state(TQueryCtx(), exec_env_.get(), desc_tbl_);
660 VLOG_QUERY << "create sender " << sender_num;
661 const TDataSink sink = GetSink(partition_type);
662 TPlanFragment fragment;
663 fragment.output_sink = sink;
664 PlanFragmentCtxPB fragment_ctx;
665 *fragment_ctx.mutable_destinations() = dest_;
666 FragmentState fragment_state(state.query_state(), fragment, fragment_ctx);
667 DataSinkConfig* data_sink = nullptr;
668 EXPECT_OK(DataSinkConfig::CreateConfig(sink, row_desc_, &fragment_state, &data_sink));
669
670 // We create an object of the base class DataSink and cast to the appropriate sender
671 // according to the 'is_thrift' option.
672 scoped_ptr<DataSink> sender;
673
674 KrpcDataStreamSenderConfig& config =
675 *(static_cast<KrpcDataStreamSenderConfig*>(data_sink));
676
677 // Reset the hash seed with a new query id. Useful for testing hash exchanges are
678 // random.
679 if (reset_hash_seed) {
680 config.exchange_hash_seed_ =
681 GetExchangeHashSeed(UuidToQueryId(random_generator()()));
682 }
683
684 sender.reset(new KrpcDataStreamSender(-1, sender_num, config,
685 data_sink->tsink_->stream_sink, dest_, channel_buffer_size, &state));
686 EXPECT_OK(sender->Prepare(&state, &tracker_));
687 EXPECT_OK(sender->Open(&state));
688 scoped_ptr<RowBatch> batch(CreateRowBatch());
689 int next_val = 0;
690 for (int i = 0; i < num_batches; ++i) {
691 GetNextBatch(batch.get(), &next_val);
692 VLOG_QUERY << "sender " << sender_num << ": #rows=" << batch->num_rows();
693 info->status = sender->Send(&state, batch.get());
694 if (!info->status.ok()) break;
695 }
696 VLOG_QUERY << "closing sender" << sender_num;
697 info->status.MergeStatus(sender->FlushFinal(&state));
698 sender->Close(&state);
699 info->num_bytes_sent = static_cast<KrpcDataStreamSender*>(
700 sender.get())->GetNumDataBytesSent();
701 data_sink->Close();
702
703 batch->Reset();
704 state.ReleaseResources();
705 }
706
707 void TestStream(TPartitionType::type stream_type, int num_senders, int num_receivers,
708 int buffer_size, bool is_merging) {

Callers

nothing calls this directly

Calls 15

UuidToQueryIdFunction · 0.85
MergeStatusMethod · 0.80
GetNumDataBytesSentMethod · 0.80
TQueryCtxClass · 0.70
getMethod · 0.65
resetMethod · 0.65
query_stateMethod · 0.45
PrepareMethod · 0.45
OpenMethod · 0.45
num_rowsMethod · 0.45
SendMethod · 0.45
okMethod · 0.45

Tested by

no test coverage detected