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

Method CloseSender

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

Source from the content-addressed store, hash-verified

284}
285
286void KrpcDataStreamMgr::CloseSender(const EndDataStreamRequestPB* request,
287 EndDataStreamResponsePB* response, kudu::rpc::RpcContext* rpc_context) {
288 TUniqueId finst_id;
289 finst_id.__set_lo(request->dest_fragment_instance_id().lo());
290 finst_id.__set_hi(request->dest_fragment_instance_id().hi());
291 VLOG_ROW << "CloseSender(): fragment_instance_id=" << PrintId(finst_id)
292 << " node_id=" << request->dest_node_id()
293 << " sender_id=" << request->sender_id();
294 shared_ptr<KrpcDataStreamRecvr> recvr;
295 {
296 lock_guard<mutex> l(lock_);
297 bool already_unregistered;
298 recvr = FindRecvr(finst_id, request->dest_node_id(), &already_unregistered);
299 // If no receiver is found and it's not in the closed stream cache, we still need
300 // to make sure that the close operation is performed so add to per-recvr list of
301 // pending closes. It's possible for a sender to issue EOS RPC without sending any
302 // rows if no rows are materialized at all in the sender side.
303 if (!already_unregistered && recvr == nullptr) {
304 AddEarlyClosedSender(finst_id, request, response, rpc_context);
305 TRACE_TO(rpc_context->trace(), "Added early closed sender");
306 return;
307 }
308 }
309
310 // If we reach this point, either the receiver is found or it has been unregistered
311 // already. In either cases, it's safe to just return an OK status.
312 TRACE_TO(
313 rpc_context->trace(), "Found receiver? $0", recvr != nullptr ? "true" : "false");
314 if (LIKELY(recvr != nullptr)) {
315 recvr->RemoveSender(request->sender_id());
316 TRACE_TO(rpc_context->trace(), "Removed sender from receiver");
317 }
318 DataStreamService::RespondAndReleaseRpc(Status::OK(), response, rpc_context,
319 service_mem_tracker_);
320}
321
322Status KrpcDataStreamMgr::DeregisterRecvr(
323 const TUniqueId& finst_id, PlanNodeId dest_node_id) {

Callers 2

EndDataStreamMethod · 0.80
EndDataStreamMethod · 0.80

Calls 5

PrintIdFunction · 0.85
OKFunction · 0.85
dest_node_idMethod · 0.80
RemoveSenderMethod · 0.80
traceMethod · 0.45

Tested by 1

EndDataStreamMethod · 0.64