| 597 | } |
| 598 | |
| 599 | void ExecEnv::SetImpalaServer(ImpalaServer* server) { |
| 600 | DCHECK(impala_server_ == nullptr) << "ImpalaServer already registered"; |
| 601 | DCHECK(server != nullptr); |
| 602 | impala_server_ = server; |
| 603 | // Register the ImpalaServer with the cluster membership manager |
| 604 | cluster_membership_mgr_->SetLocalBeDescFn([server]() { |
| 605 | return server->GetLocalBackendDescriptor(); |
| 606 | }); |
| 607 | if (FLAGS_is_coordinator) { |
| 608 | cluster_membership_mgr_->RegisterUpdateCallbackFn( |
| 609 | [server](const ClusterMembershipMgr::SnapshotPtr& snapshot) { |
| 610 | std::unordered_set<BackendIdPB> current_backend_set; |
| 611 | for (const auto& it : snapshot->current_backends) { |
| 612 | current_backend_set.insert(it.second.backend_id()); |
| 613 | } |
| 614 | server->CancelQueriesOnFailedBackends(current_backend_set); |
| 615 | }); |
| 616 | } |
| 617 | if (FLAGS_is_executor && !TestInfo::is_test()) { |
| 618 | cluster_membership_mgr_->RegisterUpdateCallbackFn( |
| 619 | [](const ClusterMembershipMgr::SnapshotPtr& snapshot) { |
| 620 | std::unordered_set<BackendIdPB> current_backend_set; |
| 621 | for (const auto& it : snapshot->current_backends) { |
| 622 | current_backend_set.insert(it.second.backend_id()); |
| 623 | } |
| 624 | ExecEnv::GetInstance()->query_exec_mgr()->CancelQueriesForFailedCoordinators( |
| 625 | current_backend_set); |
| 626 | }); |
| 627 | } |
| 628 | } |
| 629 | |
| 630 | void ExecEnv::InitBufferPool(int64_t min_buffer_size, int64_t capacity, |
| 631 | int64_t clean_pages_limit) { |
no test coverage detected