| 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) { |
nothing calls this directly
no test coverage detected