| 347 | } |
| 348 | |
| 349 | void RemoteTabletNode::WriteTablet(google::protobuf::RpcController* controller, |
| 350 | const WriteTabletRequest* request, WriteTabletResponse* response, |
| 351 | google::protobuf::Closure* done) { |
| 352 | int64_t start_micros = get_micros(); |
| 353 | done = WriteDoneWrapper::NewInstance(start_micros, request, response, done); |
| 354 | VLOG(8) << "accept RPC (WriteTablet): [" << request->tablet_name() << "] " |
| 355 | << tera::utils::GetRemoteAddress(controller); |
| 356 | static uint32_t last_print = time(NULL); |
| 357 | int32_t row_num = request->row_list_size(); |
| 358 | write_request_counter.Add(row_num); |
| 359 | if (write_pending_counter.Get() > FLAGS_tera_request_pending_limit) { |
| 360 | response->set_sequence_id(request->sequence_id()); |
| 361 | response->set_status(kTabletNodeIsBusy); |
| 362 | write_reject_counter.Add(row_num); |
| 363 | done->Run(); |
| 364 | uint32_t now_time = time(NULL); |
| 365 | if (now_time > last_print) { |
| 366 | LOG(WARNING) << "Too many pending write requests, return TabletNode Is Busy!"; |
| 367 | last_print = now_time; |
| 368 | } |
| 369 | VLOG(8) << "finish RPC (WriteTablet)"; |
| 370 | } else { |
| 371 | // check user identification & access |
| 372 | if (!access_entry_->VerifyAndAuthorize(request, response)) { |
| 373 | response->set_sequence_id(request->sequence_id()); |
| 374 | VLOG(20) << "Access VerifyAndAuthorize failed for WriteTablet"; |
| 375 | done->Run(); |
| 376 | return; |
| 377 | } |
| 378 | |
| 379 | // sum write bytes |
| 380 | int64_t sum_write_bytes = 0; |
| 381 | for (int32_t row_index = 0; row_index < row_num; ++row_index) { |
| 382 | sum_write_bytes += request->row_list(row_index).ByteSize(); |
| 383 | } |
| 384 | if (!TsWriteFlowController::Instance().TryWrite(sum_write_bytes)) { |
| 385 | response->set_sequence_id(request->sequence_id()); |
| 386 | response->set_status(kFlowControlLimited); |
| 387 | write_reject_counter.Add(row_num); |
| 388 | VLOG(20) << "Reject write request due to write flow controller"; |
| 389 | done->Run(); |
| 390 | return; |
| 391 | } |
| 392 | if (write_pending_counter.Get() >= |
| 393 | FLAGS_tera_request_pending_limit * FLAGS_tera_quota_unlimited_pending_ratio) { |
| 394 | if (!quota_entry_->CheckAndConsume( |
| 395 | request->tablet_name(), |
| 396 | quota::OpTypeAmountList{std::make_pair(kQuotaWriteReqs, row_num), |
| 397 | std::make_pair(kQuotaWriteBytes, sum_write_bytes)})) { |
| 398 | response->set_sequence_id(request->sequence_id()); |
| 399 | response->set_status(kQuotaLimited); |
| 400 | write_quota_reject_counter.Add(row_num); |
| 401 | VLOG(20) << "quota_entry check failed for WriteTablet"; |
| 402 | done->Run(); |
| 403 | return; |
| 404 | } |
| 405 | } |
| 406 | write_pending_counter.Add(row_num); |
no test coverage detected