| 2249 | } |
| 2250 | |
| 2251 | void 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(); |
no test coverage detected