| 216 | } |
| 217 | |
| 218 | void 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 | |
| 263 | void KrpcDataStreamMgr::EnqueueDeserializeTask(const TUniqueId& finst_id, |
| 264 | PlanNodeId dest_node_id, int sender_id, int num_requests) { |