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

Method TakeOverEarlySender

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

Source from the content-addressed store, hash-verified

599}
600
601void KrpcDataStreamRecvr::SenderQueue::TakeOverEarlySender(
602 unique_ptr<TransmitDataCtx> ctx) {
603 // TakeOverEarlySender() is called by the same thread which calls Close().
604 // The receiver cannot be closed while this function is in progress so
605 // 'recvr_->mgr_' shouldn't be NULL.
606 DCHECK(TestInfo::is_test() || FragmentInstanceState::IsFragmentExecThread());
607 DCHECK(!recvr_->closed_ && recvr_->mgr_ != nullptr);
608 COUNTER_ADD(recvr_->total_received_batches_counter_, 1);
609 int64_t serialized_size;
610 Status status = GetBatchSize(ctx->request, ctx->rpc_context, &serialized_size);
611 if (UNLIKELY(!status.ok())) {
612 {
613 unique_lock<SpinLock> l(lock_);
614 if (!is_cancelled_) MarkErrorStatus(status, l);
615 }
616 TRACE_TO(
617 ctx->rpc_context->trace(), "Error unpacking request: $0", status.GetDetail());
618 DataStreamService::RespondRpc(status, ctx->response, ctx->rpc_context);
619 return;
620 }
621 COUNTER_ADD(recvr_->bytes_received_counter_, serialized_size);
622 DCHECK_GE(ctx->deserialized_size, 0);
623 int sender_id = ctx->request->sender_id();
624 int num_to_dequeue = 0;
625 {
626 unique_lock<SpinLock> l(lock_);
627 // Only enqueue a deferred RPC if the sender queue is not yet cancelled.
628 if (UNLIKELY(is_cancelled_)) {
629 l.unlock();
630 TRACE_TO(ctx->rpc_context->trace(), "Recvr closed");
631 RespondClosed(ctx.get());
632 return;
633 }
634 recvr_->deferred_rpc_tracker()->Consume(ctx->rpc_context->GetTransferSize());
635 EnqueueDeferredRpc(move(ctx));
636 num_to_dequeue = CollectPendingDeferredRpcs();
637 }
638
639 if (num_to_dequeue > 0) {
640 recvr_->mgr_->EnqueueDeserializeTask(recvr_->fragment_instance_id(),
641 recvr_->dest_node_id(), sender_id, num_to_dequeue);
642 }
643}
644
645bool KrpcDataStreamRecvr::SenderQueue::DecrementSenders() {
646 lock_guard<SpinLock> l(lock_);

Callers 1

CreateRecvrMethod · 0.80

Calls 11

moveFunction · 0.85
GetDetailMethod · 0.80
deferred_rpc_trackerMethod · 0.80
dest_node_idMethod · 0.80
getMethod · 0.65
okMethod · 0.45
traceMethod · 0.45
unlockMethod · 0.45
ConsumeMethod · 0.45
GetTransferSizeMethod · 0.45

Tested by

no test coverage detected