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

Method AddPoolAndPerHostStatsUpdates

be/src/scheduling/admission-controller.cc:2624–2670  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

2622}
2623
2624void AdmissionController::AddPoolAndPerHostStatsUpdates(
2625 vector<TTopicDelta>* topic_updates) {
2626 // local_stats_ are updated eagerly except for backend_mem_reserved (which isn't used
2627 // for local admission control decisions). Update that now before sending local_stats_.
2628 for (auto& entry : pool_stats_) {
2629 entry.second.UpdateMemTrackerStats();
2630 }
2631 if (pools_for_updates_.empty()) {
2632 // No pool updates means no changes to host stats as well, so just return.
2633 return;
2634 }
2635 topic_updates->push_back(TTopicDelta());
2636 TTopicDelta& topic_delta = topic_updates->back();
2637 topic_delta.topic_name = request_queue_topic_name_;
2638 for (const string& pool_name: pools_for_updates_) {
2639 DCHECK(pool_stats_.find(pool_name) != pool_stats_.end());
2640 PoolStats* stats = GetPoolStats(pool_name);
2641 VLOG_ROW << "Sending topic update " << stats->DebugString();
2642 topic_delta.topic_entries.push_back(TTopicItem());
2643 TTopicItem& topic_item = topic_delta.topic_entries.back();
2644 topic_item.key = MakePoolTopicKey(pool_name, host_id_);
2645 Status status =
2646 thrift_serializer_.SerializeToString(&stats->local_stats(), &topic_item.value);
2647 if (!status.ok()) {
2648 LOG(WARNING) << "Failed to serialize query pool stats: " << status.GetDetail();
2649 topic_delta.topic_entries.pop_back();
2650 }
2651 }
2652 pools_for_updates_.clear();
2653
2654 // Now add the host stats
2655 topic_delta.topic_entries.push_back(TTopicItem());
2656 TTopicItem& topic_item = topic_delta.topic_entries.back();
2657 topic_item.key = Substitute("$0$1", TOPIC_KEY_STAT_PREFIX, host_id_);
2658 TPerHostStatsUpdate update;
2659 for (const auto& elem : host_stats_) {
2660 update.per_host_stats.emplace_back();
2661 TPerHostStatsUpdateElement& inserted_elem = update.per_host_stats.back();
2662 inserted_elem.__set_host_addr(elem.first);
2663 inserted_elem.__set_stats(elem.second);
2664 }
2665 Status status = thrift_serializer_.SerializeToString(&update, &topic_item.value);
2666 if (!status.ok()) {
2667 LOG(WARNING) << "Failed to serialize host stats: " << status.GetDetail();
2668 topic_delta.topic_entries.pop_back();
2669 }
2670}
2671
2672void AdmissionController::DequeueLoop() {
2673 unique_lock<mutex> lock(admission_ctrl_lock_);

Callers

nothing calls this directly

Calls 12

TTopicDeltaClass · 0.85
SubstituteFunction · 0.85
UpdateMemTrackerStatsMethod · 0.80
push_backMethod · 0.80
SerializeToStringMethod · 0.80
GetDetailMethod · 0.80
clearMethod · 0.65
emptyMethod · 0.45
findMethod · 0.45
endMethod · 0.45
DebugStringMethod · 0.45
okMethod · 0.45

Tested by

no test coverage detected