| 284 | } |
| 285 | |
| 286 | void 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 | |
| 322 | Status KrpcDataStreamMgr::DeregisterRecvr( |
| 323 | const TUniqueId& finst_id, PlanNodeId dest_node_id) { |