| 746 | } |
| 747 | |
| 748 | Status KrpcDataStreamRecvr::CreateMerger(const TupleRowComparator& less_than, |
| 749 | const CodegenFnPtr<SortedRunMerger::HeapifyHelperFn>& codegend_heapify_helper_fn) { |
| 750 | DCHECK(is_merging_); |
| 751 | DCHECK(TestInfo::is_test() || FragmentInstanceState::IsFragmentExecThread()); |
| 752 | vector<SortedRunMerger::RunBatchSupplierFn> input_batch_suppliers; |
| 753 | input_batch_suppliers.reserve(sender_queues_.size()); |
| 754 | |
| 755 | // Create the merger that will a single stream of sorted rows. |
| 756 | merger_.reset(new SortedRunMerger(less_than, row_desc_, profile_, false, |
| 757 | codegend_heapify_helper_fn)); |
| 758 | |
| 759 | for (SenderQueue* queue: sender_queues_) { |
| 760 | input_batch_suppliers.push_back( |
| 761 | [queue](RowBatch** next_batch) -> Status { |
| 762 | return queue->GetBatch(next_batch); |
| 763 | }); |
| 764 | } |
| 765 | |
| 766 | RETURN_IF_ERROR(merger_->Prepare(input_batch_suppliers)); |
| 767 | return Status::OK(); |
| 768 | } |
| 769 | |
| 770 | void KrpcDataStreamRecvr::TransferAllResources(RowBatch* transfer_batch) { |
| 771 | DCHECK(TestInfo::is_test() || FragmentInstanceState::IsFragmentExecThread()); |