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

Method AddData

be/src/runtime/krpc-data-stream-mgr.cc:218–261  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

216}
217
218void KrpcDataStreamMgr::AddData(const TransmitDataRequestPB* request,
219 TransmitDataResponsePB* response, kudu::rpc::RpcContext* rpc_context) {
220 TUniqueId finst_id;
221 finst_id.__set_lo(request->dest_fragment_instance_id().lo());
222 finst_id.__set_hi(request->dest_fragment_instance_id().hi());
223 TPlanNodeId dest_node_id = request->dest_node_id();
224 VLOG_ROW << "AddData(): fragment_instance_id=" << PrintId(finst_id)
225 << " node_id=" << request->dest_node_id()
226 << " #rows=" << request->row_batch_header().num_rows()
227 << " sender_id=" << request->sender_id();
228 bool already_unregistered = false;
229 shared_ptr<KrpcDataStreamRecvr> recvr;
230 {
231 lock_guard<mutex> l(lock_);
232 recvr = FindRecvr(finst_id, request->dest_node_id(), &already_unregistered);
233 // If no receiver is found and it's not in the closed stream cache, best guess is
234 // that it is still preparing, so add payload to per-receiver early senders' list.
235 // If the receiver doesn't show up after FLAGS_datastream_sender_timeout_ms ms
236 // (e.g. if the receiver was closed and has already been retired from the
237 // closed_stream_cache_), the sender is timed out by the maintenance thread.
238 if (!already_unregistered && recvr == nullptr) {
239 AddEarlySender(finst_id, request, response, rpc_context);
240 TRACE_TO(rpc_context->trace(), "Added early sender");
241 return;
242 }
243 }
244 if (already_unregistered) {
245 TRACE_TO(rpc_context->trace(), "Sender already unregistered");
246 // The receiver may remove itself from the receiver map via DeregisterRecvr() at any
247 // time without considering the remaining number of senders. As a consequence,
248 // FindRecvr() may return nullptr even though the receiver was once present. We
249 // detect this case by checking already_unregistered - if true then the receiver was
250 // already closed deliberately, and there's no unexpected error here.
251 ErrorMsg msg(TErrorCode::DATASTREAM_RECVR_CLOSED, PrintId(finst_id), dest_node_id);
252 DataStreamService::RespondAndReleaseRpc(Status::Expected(msg), response, rpc_context,
253 service_mem_tracker_);
254 return;
255 }
256 DCHECK(recvr != nullptr);
257 int64_t transfer_size = rpc_context->GetTransferSize();
258 recvr->AddBatch(request, response, rpc_context);
259 // Release memory. The receiver already tracks it in its instance tracker.
260 service_mem_tracker_->Release(transfer_size);
261}
262
263void KrpcDataStreamMgr::EnqueueDeserializeTask(const TUniqueId& finst_id,
264 PlanNodeId dest_node_id, int sender_id, int num_requests) {

Callers 2

TransmitDataMethod · 0.45
TransmitDataMethod · 0.45

Calls 7

PrintIdFunction · 0.85
dest_node_idMethod · 0.80
num_rowsMethod · 0.45
traceMethod · 0.45
GetTransferSizeMethod · 0.45
AddBatchMethod · 0.45
ReleaseMethod · 0.45

Tested by 1

TransmitDataMethod · 0.36