| 2524 | } |
| 2525 | |
| 2526 | void ImpalaServer::WaitForCatalogUpdateTopicPropagation( |
| 2527 | const TUniqueId& catalog_service_id, RuntimeProfile::EventSequence* timeline) { |
| 2528 | unique_lock<mutex> unique_lock(catalog_version_lock_); |
| 2529 | int64_t min_req_subscriber_topic_version = |
| 2530 | catalog_update_info_.catalog_topic_version; |
| 2531 | VLOG_QUERY << "Waiting for min subscriber topic version: " |
| 2532 | << min_req_subscriber_topic_version << " current version: " |
| 2533 | << min_subscriber_catalog_topic_version_; |
| 2534 | while (min_subscriber_catalog_topic_version_ < min_req_subscriber_topic_version && |
| 2535 | catalog_update_info_.catalog_service_id == catalog_service_id) { |
| 2536 | catalog_version_update_cv_.Wait(unique_lock); |
| 2537 | } |
| 2538 | |
| 2539 | if (catalog_update_info_.catalog_service_id != catalog_service_id) { |
| 2540 | MarkTimelineEvent(timeline, "Catalog service ID changed when waiting for " |
| 2541 | "catalog propagation"); |
| 2542 | VLOG_QUERY << "Detected catalog service ID changed when waiting for " |
| 2543 | "catalog propagation"; |
| 2544 | } else { |
| 2545 | MarkTimelineEvent(timeline, |
| 2546 | Substitute("Min catalog topic version of coordinators reached $0", |
| 2547 | min_req_subscriber_topic_version)); |
| 2548 | VLOG_QUERY << "Min catalog topic version of coordinators: " |
| 2549 | << min_req_subscriber_topic_version; |
| 2550 | } |
| 2551 | } |
| 2552 | |
| 2553 | void ImpalaServer::WaitForMinCatalogUpdate(const int64_t min_req_catalog_object_version, |
| 2554 | const TUniqueId& catalog_service_id, RuntimeProfile::EventSequence* timeline) { |
nothing calls this directly
no test coverage detected