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

Method CollectPendingDeferredRpcs

be/src/runtime/krpc-data-stream-recvr.cc:712–740  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

710}
711
712int KrpcDataStreamRecvr::SenderQueue::CollectPendingDeferredRpcs() {
713 lock_.DCheckLocked();
714 DCHECK(!deferred_rpcs_.empty());
715
716 // Try dequeuing multiple entries from 'deferred_rpcs_' to parallelize the CPU
717 // bound deserialization work. No point in dequeuing more than number of
718 // deserialization threads available.
719 DCHECK_EQ(pending_deferred_rpcs_.size(), num_deserialize_tasks_pending_);
720 int free_deserialization_threads =
721 recvr_->mgr_->num_deserialization_threads() - num_deserialize_tasks_pending_;
722 DCHECK_GE(free_deserialization_threads, 0);
723 int num_to_try_dequeue =
724 min(free_deserialization_threads, static_cast<int>(deferred_rpcs_.size()));
725 int num_to_dequeue = 0;
726 for(; num_to_dequeue < num_to_try_dequeue; num_to_dequeue++) {
727 int64_t size = deferred_rpcs_.front()->deserialized_size;
728 if (!CanEnqueue(size)) break;
729 pending_deserialized_size_ += size;
730 if (pending_deferred_rpcs_.empty()) {
731 has_pending_deferred_rpcs_start_time_ns_ = MonotonicNanos();
732 }
733 pending_deferred_rpcs_.push(DequeueDeferredRpc());
734 num_deserialize_tasks_pending_++;
735 }
736 DCHECK_EQ(pending_deferred_rpcs_.size(), num_deserialize_tasks_pending_);
737 DCHECK(!batch_queue_.empty() || !pending_deferred_rpcs_.empty()
738 || num_pending_enqueue_ > 0);
739 return num_to_dequeue;
740}
741
742void KrpcDataStreamRecvr::SenderQueue::RespondClosed(TransmitDataCtx* ctx) {
743 Status cancel_status = Status::Expected(TErrorCode::DATASTREAM_RECVR_CLOSED,

Callers

nothing calls this directly

Calls 8

minFunction · 0.85
MonotonicNanosFunction · 0.85
DCheckLockedMethod · 0.80
frontMethod · 0.80
pushMethod · 0.80
emptyMethod · 0.45
sizeMethod · 0.45

Tested by

no test coverage detected