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

Method UpdateAggregates

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

Source from the content-addressed store, hash-verified

2249}
2250
2251void AdmissionController::PoolStats::UpdateAggregates(HostMemMap* host_mem_reserved) {
2252 const string& coord_id = parent_->host_id_;
2253 int64_t num_running = 0;
2254 int64_t num_queued = 0;
2255 int64_t mem_reserved = 0;
2256 AggregatedUserLoads new_agg_user_loads;
2257 for (const PoolStats::RemoteStatsMap::value_type& remote_entry : remote_stats_) {
2258 const string& host = remote_entry.first;
2259 // Skip an update from this subscriber as the information may be outdated.
2260 // The stats from this coordinator will be added below.
2261 if (host == coord_id) continue;
2262 const TPoolStats& remote_pool_stats = remote_entry.second;
2263 DCHECK_GE(remote_pool_stats.num_admitted_running, 0);
2264 DCHECK_GE(remote_pool_stats.num_queued, 0);
2265 DCHECK_GE(remote_pool_stats.backend_mem_reserved, 0);
2266 num_running += remote_pool_stats.num_admitted_running;
2267 num_queued += remote_pool_stats.num_queued;
2268
2269 new_agg_user_loads.add_loads(remote_pool_stats.user_loads);
2270
2271 // Update the per-pool and per-host aggregates with the mem reserved by this host in
2272 // this pool.
2273 mem_reserved += remote_pool_stats.backend_mem_reserved;
2274 (*host_mem_reserved)[host] += remote_pool_stats.backend_mem_reserved;
2275 // TODO(IMPALA-8762): For multiple coordinators, need to track the number of running
2276 // queries per executor, i.e. every admission controller needs to send the full map to
2277 // everyone else.
2278 }
2279 num_running += local_stats_.num_admitted_running;
2280 num_queued += local_stats_.num_queued;
2281 mem_reserved += local_stats_.backend_mem_reserved;
2282 (*host_mem_reserved)[coord_id] += local_stats_.backend_mem_reserved;
2283 new_agg_user_loads.add_loads(local_stats_.user_loads);
2284
2285 DCHECK_GE(num_running, 0);
2286 DCHECK_GE(num_queued, 0);
2287 DCHECK_GE(mem_reserved, 0);
2288 DCHECK_GE(num_running, local_stats_.num_admitted_running);
2289 DCHECK_GE(num_queued, local_stats_.num_queued);
2290
2291 if (agg_num_running_ == num_running && agg_num_queued_ == num_queued
2292 && agg_mem_reserved_ == mem_reserved
2293 && agg_user_loads_.get_user_loads() == new_agg_user_loads.get_user_loads()) {
2294 // Nothing changed
2295 DCHECK_EQ(num_running, metrics_.agg_num_running->GetValue());
2296 DCHECK_EQ(num_queued, metrics_.agg_num_queued->GetValue());
2297 DCHECK_EQ(mem_reserved, metrics_.agg_mem_reserved->GetValue());
2298 return;
2299 }
2300 VLOG_ROW << "Recomputed agg stats, previous: " << DebugString();
2301 agg_num_running_ = num_running;
2302 agg_num_queued_ = num_queued;
2303 agg_mem_reserved_ = mem_reserved;
2304 metrics_.agg_num_running->SetValue(num_running);
2305 metrics_.agg_num_queued->SetValue(num_queued);
2306 metrics_.agg_mem_reserved->SetValue(mem_reserved);
2307
2308 agg_user_loads_.clear();

Callers 1

Calls 7

add_loadsMethod · 0.80
export_usersMethod · 0.80
clearMethod · 0.65
DebugStringFunction · 0.50
GetValueMethod · 0.45
SetValueMethod · 0.45
ResetMethod · 0.45

Tested by

no test coverage detected