diff --git a/quickwit/quickwit-proto/build.rs b/quickwit/quickwit-proto/build.rs index 5d51c8cd3f9..b7d08264920 100644 --- a/quickwit/quickwit-proto/build.rs +++ b/quickwit/quickwit-proto/build.rs @@ -222,6 +222,7 @@ fn main() -> Result<(), Box> { prost_config .file_descriptor_set_path("src/codegen/quickwit/search_descriptor.bin") .protoc_arg("--experimental_allow_proto3_optional") + .field_attribute("SearchRequest.priority", "#[serde(default)]") // Box the large `LeafSearchResponse` variant so the oneof stays small // (the `Error` variant only carries a `String`). .boxed("LambdaSingleSplitResult.outcome.response"); diff --git a/quickwit/quickwit-proto/protos/quickwit/search.proto b/quickwit/quickwit-proto/protos/quickwit/search.proto index 18136f6b3f3..d266f889508 100644 --- a/quickwit/quickwit-proto/protos/quickwit/search.proto +++ b/quickwit/quickwit-proto/protos/quickwit/search.proto @@ -273,6 +273,13 @@ message SearchRequest { // When true, skip finalization of aggregation results and return // the raw IntermediateAggregationResults bytes instead. bool skip_aggregation_finalization = 19; + + // Field 20 is intentionally reserved for wire compatibility with a downstream fork. + reserved 20; + + // Scheduling priority for leaf search execution. Negative values are allowed, + // and lower values have higher priority. Callers that omit it get priority 0. + int32 priority = 21; } enum CountHits { diff --git a/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.search.rs b/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.search.rs index 1c1ed0b1b03..347892fb124 100644 --- a/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.search.rs +++ b/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.search.rs @@ -211,6 +211,11 @@ pub struct SearchRequest { /// the raw IntermediateAggregationResults bytes instead. #[prost(bool, tag = "19")] pub skip_aggregation_finalization: bool, + /// Scheduling priority for leaf search execution. Negative values are allowed, + /// and lower values have higher priority. Callers that omit it get priority 0. + #[prost(int32, tag = "21")] + #[serde(default)] + pub priority: i32, } #[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] diff --git a/quickwit/quickwit-search/src/leaf.rs b/quickwit/quickwit-search/src/leaf.rs index 99ff06c312b..5b99c9512ff 100644 --- a/quickwit/quickwit-search/src/leaf.rs +++ b/quickwit/quickwit-search/src/leaf.rs @@ -1891,7 +1891,11 @@ async fn schedule_search_tasks( mut splits: Vec<(SplitIdAndFooterOffsets, SearchRequest)>, searcher_context: &SearcherContext, ) -> ScheduleSearchTaskResult { - let task_metadata: Vec = splits + let priority = splits + .first() + .map(|(_, search_request)| search_request.priority) + .unwrap_or_default(); + let split_metadatas = splits .iter() .map(|(split, _)| { let memory_allocation = compute_initial_memory_allocation( @@ -1907,6 +1911,10 @@ async fn schedule_search_tasks( } }) .collect(); + let task_metadata = crate::search_permit_provider::LeafSearchTaskMetadata { + priority, + splits: split_metadatas, + }; let offload_threshold: usize = if searcher_context.lambda_invoker.is_some() && let Some(lambda_config) = &searcher_context.searcher_config.lambda diff --git a/quickwit/quickwit-search/src/leaf_cache.rs b/quickwit/quickwit-search/src/leaf_cache.rs index 4cb9f37d6c5..39758b08956 100644 --- a/quickwit/quickwit-search/src/leaf_cache.rs +++ b/quickwit/quickwit-search/src/leaf_cache.rs @@ -107,6 +107,8 @@ impl CacheKey { // it doesn't matter whether or not we count all hits at the scale of a // single split: either we did process it and got everything, or we didn't. search_request.count_hits = CountHits::CountAll.into(); + // Priority only affects scheduling, not search results. + search_request.priority = 0; CacheKey { split_id: split_info.split_id, @@ -436,4 +438,30 @@ mod tests { assert!(cache.get(split_3.clone(), query_2).is_none()); assert!(cache.get(split_3, query_2bis).is_some()); } + + #[test] + fn test_leaf_search_cache_ignores_priority() { + let cache = LeafSearchCache::new(&ByteSize::mb(64).into()); + let split = SplitIdAndFooterOffsets { + split_id: "split".to_string(), + ..Default::default() + }; + let high_priority_request = SearchRequest { + priority: -10, + ..Default::default() + }; + let low_priority_request = SearchRequest { + priority: 10, + ..Default::default() + }; + + cache.put( + split.clone(), + high_priority_request, + LeafSearchResponse::default(), + ); + + // the two requests differ only in priority, so it should be a cache-hit + assert!(cache.get(split, low_priority_request).is_some()); + } } diff --git a/quickwit/quickwit-search/src/list_terms.rs b/quickwit/quickwit-search/src/list_terms.rs index 9c6e025ac92..78da5d88b6b 100644 --- a/quickwit/quickwit-search/src/list_terms.rs +++ b/quickwit/quickwit-search/src/list_terms.rs @@ -326,7 +326,7 @@ pub async fn leaf_list_terms( splits: &[SplitIdAndFooterOffsets], ) -> Result { info!(split_offsets = ?PrettySample::new(splits, 5)); - let task_metadata: Vec = splits + let split_metadatas = splits .iter() .map(|split| { let memory_allocation = compute_initial_memory_allocation( @@ -342,6 +342,10 @@ pub async fn leaf_list_terms( } }) .collect(); + let task_metadata = crate::search_permit_provider::LeafSearchTaskMetadata { + priority: 0, + splits: split_metadatas, + }; // We have added offloading leaf search to lambdas, but not for list_terms yet. // TODO (Add it) // https://github.com/quickwit-oss/quickwit/issues/6150 diff --git a/quickwit/quickwit-search/src/root.rs b/quickwit/quickwit-search/src/root.rs index 896a763f508..c31a848f8ea 100644 --- a/quickwit/quickwit-search/src/root.rs +++ b/quickwit/quickwit-search/src/root.rs @@ -376,6 +376,7 @@ fn simplify_search_request_for_scroll_api(req: &SearchRequest) -> crate::Result< count_hits: quickwit_proto::search::CountHits::Underestimate as i32, ignore_missing_indexes: req.ignore_missing_indexes, skip_aggregation_finalization: false, + priority: req.priority, }) } @@ -1320,6 +1321,7 @@ pub async fn root_search( count_required = search_request.count_hits().as_str_name(), num_docs = num_docs, num_splits = num_splits, + priority = search_request.priority, "root_search" ); diff --git a/quickwit/quickwit-search/src/search_permit_provider.rs b/quickwit/quickwit-search/src/search_permit_provider.rs index becc3ff6d18..60c85e04f32 100644 --- a/quickwit/quickwit-search/src/search_permit_provider.rs +++ b/quickwit/quickwit-search/src/search_permit_provider.rs @@ -34,11 +34,11 @@ use crate::metrics::{ /// Distributor of permits to perform split search operation. /// -/// Requests are served in order. Each permit initially reserves a slot for the -/// warmup (limit concurrent downloads) and a pessimistic amount of memory. Once -/// the warmup is completed, the actual memory usage is set and the warmup slot -/// is released. Once the search is completed and the permit is dropped, the -/// remaining memory is also released. +/// Requests are served by priority, then by fewest remaining splits. Each permit initially +/// reserves a slot for the warmup (limit concurrent downloads) and a pessimistic amount of +/// memory. Once the warmup is completed, the actual memory usage is set and the warmup slot is +/// released. Once the search is completed and the permit is dropped, the remaining memory is also +/// released. #[derive(Clone)] pub struct SearchPermitProvider { message_sender: mpsc::UnboundedSender, @@ -58,9 +58,16 @@ pub(crate) struct SplitSearchTaskMetadata { pub job_cost: usize, } +/// Metadata for a leaf search task. +pub(crate) struct LeafSearchTaskMetadata { + /// Lower values have higher priority. + pub priority: i32, + pub splits: Vec, +} + pub enum SearchPermitMessage { RequestWithOffload { - task_metadata: Vec, + task_metadata: LeafSearchTaskMetadata, /// Maximum number of pending requests. If granting all /// requested permits would cause the number of pending requests to exceed this threshold, /// some permits will be offloaded to Lambda. @@ -153,9 +160,10 @@ impl SearchPermitProvider { /// The returned futures are guaranteed to resolve in order. pub(crate) async fn get_permits( &self, - splits: Vec, + task_metadata: LeafSearchTaskMetadata, ) -> Vec { - self.get_permits_with_offload(splits, usize::MAX).await + self.get_permits_with_offload(task_metadata, usize::MAX) + .await } /// Returns permits for local splits and a list of split indices to offload. @@ -170,17 +178,17 @@ impl SearchPermitProvider { /// If `offload_threshold` is usize::MAX, all splits are processed locally. pub(crate) async fn get_permits_with_offload( &self, - splits: Vec, + task_metadata: LeafSearchTaskMetadata, offload_threshold: usize, ) -> Vec { - if splits.is_empty() { + if task_metadata.splits.is_empty() { return Vec::new(); } let (permit_sender, permit_receiver) = oneshot::channel(); self.message_sender .send(SearchPermitMessage::RequestWithOffload { permit_resp_tx: permit_sender, - task_metadata: splits, + task_metadata, offload_threshold, }) .expect("Receiver lives longer than sender"); @@ -219,19 +227,23 @@ struct SingleSplitPermitRequest { } struct LeafPermitRequest { + /// Lower values have higher priority. + priority: i32, /// Single split permit requests for this leaf search. single_split_permit_requests: std::vec::IntoIter, } impl Ord for LeafPermitRequest { fn cmp(&self, other: &Self) -> std::cmp::Ordering { - // we compare other with self and not the other way arround because we want a min-heap and + // we compare other with self and not the other way around because we want a min-heap and // Rust's is a max-heap - other - .single_split_permit_requests - .as_slice() - .len() - .cmp(&self.single_split_permit_requests.as_slice().len()) + other.priority.cmp(&self.priority).then_with(|| { + other + .single_split_permit_requests + .as_slice() + .len() + .cmp(&self.single_split_permit_requests.as_slice().len()) + }) } } @@ -252,30 +264,32 @@ impl Eq for LeafPermitRequest {} impl LeafPermitRequest { // `task_metadata` must not be empty. fn from_task_metadata( - task_metadata: Vec, + task_metadata: LeafSearchTaskMetadata, ) -> (Self, Vec) { - assert!(!task_metadata.is_empty(), "task_metadata must not be empty"); + let LeafSearchTaskMetadata { priority, splits } = task_metadata; + assert!(!splits.is_empty(), "task_metadata must not be empty"); // Stamped on every `SingleSplitPermitRequest` we're about to enqueue. // The actor will compute `requested_at.elapsed()` at grant time to // report the permit's acquisition latency. let requested_at = Instant::now(); - let mut permits = Vec::with_capacity(task_metadata.len()); - let mut single_split_permit_requests = Vec::with_capacity(task_metadata.len()); - for meta in task_metadata { + let mut permits = Vec::with_capacity(splits.len()); + let mut single_split_permit_requests = Vec::with_capacity(splits.len()); + for split in splits { let (tx, rx) = oneshot::channel(); // we keep our internal list of permits and the returned wait handles in the // same order to make sure we emit each permit in the right order. Doing otherwise // may cause deadlocks single_split_permit_requests.push(SingleSplitPermitRequest { permit_sender: tx, - permit_size: meta.memory_allocation.as_u64(), - job_cost: meta.job_cost, + permit_size: split.memory_allocation.as_u64(), + job_cost: split.job_cost, requested_at, }); permits.push(SearchPermitFuture(rx)); } ( LeafPermitRequest { + priority, single_split_permit_requests: single_split_permit_requests.into_iter(), }, permits, @@ -321,18 +335,18 @@ impl SearchPermitActor { // If this indeed truncates the task_metadata vector, other splits will be // offloaded to lambdas. - task_metadata.truncate(local_capacity); + task_metadata.splits.truncate(local_capacity); // We special case here in order to avoid pushing empty request in the queue. // (they would never be removed) - if task_metadata.is_empty() { + if task_metadata.splits.is_empty() { let _ = permit_sender.send(Vec::new()); return; } // Increment total_job_cost for all tasks entering the queue (both pending and // those that will be granted immediately by assign_available_permits below). - let added: usize = task_metadata.iter().map(|m| m.job_cost).sum(); + let added: usize = task_metadata.splits.iter().map(|m| m.job_cost).sum(); // fetch_add returns the previous value, so add `added` to get the new total. let new_load = self.total_job_cost.fetch_add(added, Ordering::Relaxed) + added; SEARCHER_NODE_LOAD.set(new_load as f64); @@ -523,13 +537,75 @@ mod tests { use super::*; - fn make_splits(memory_mb: u64, count: usize) -> Vec { - (0..count) - .map(|_| SplitSearchTaskMetadata { - memory_allocation: ByteSize::mb(memory_mb), - job_cost: 5, - }) - .collect() + fn make_splits(memory_mb: u64, count: usize) -> LeafSearchTaskMetadata { + make_splits_with_priority(memory_mb, count, 0) + } + + fn make_splits_with_priority( + memory_mb: u64, + count: usize, + priority: i32, + ) -> LeafSearchTaskMetadata { + LeafSearchTaskMetadata { + priority, + splits: (0..count) + .map(|_| SplitSearchTaskMetadata { + memory_allocation: ByteSize::mb(memory_mb), + job_cost: 5, + }) + .collect(), + } + } + + #[tokio::test] + async fn test_search_permit_priority_precedes_remaining_splits() { + let permit_provider = SearchPermitProvider::new(1, ByteSize::mb(100)); + let blocker = permit_provider + .get_permits(make_splits(10, 1)) + .await + .pop() + .unwrap() + .await; + + let negative_priority = permit_provider + .get_permits(make_splits_with_priority(10, 3, -10)) + .await; + let default_priority = permit_provider.get_permits(make_splits(10, 1)).await; + let positive_priority = permit_provider + .get_permits(make_splits_with_priority(10, 1, 10)) + .await; + + let mut join_set = JoinSet::new(); + for (request, permit_futures) in [ + ("negative", negative_priority), + ("default", default_priority), + ("positive", positive_priority), + ] { + for (split_idx, permit_future) in permit_futures.into_iter().enumerate() { + join_set.spawn(async move { + let permit = permit_future.await; + (request, split_idx, permit) + }); + } + } + + drop(blocker); + + let mut execution_order = Vec::new(); + while let Some(result) = join_set.join_next().await { + let (request, split_idx, _permit) = result.unwrap(); + execution_order.push((request, split_idx)); + } + assert_eq!( + execution_order, + vec![ + ("negative", 0), + ("negative", 1), + ("negative", 2), + ("default", 0), + ("positive", 0), + ] + ); } #[tokio::test] @@ -829,21 +905,24 @@ mod tests { assert_eq!(permit_provider.get_load(), 0); - let splits = vec![ - SplitSearchTaskMetadata { - memory_allocation: ByteSize::mb(10), - job_cost: 7, - }, - SplitSearchTaskMetadata { - memory_allocation: ByteSize::mb(10), - job_cost: 3, - }, - SplitSearchTaskMetadata { - memory_allocation: ByteSize::mb(10), - job_cost: 5, - }, - ]; - let mut permit_futs = permit_provider.get_permits(splits).await; + let task_metadata = LeafSearchTaskMetadata { + priority: 0, + splits: vec![ + SplitSearchTaskMetadata { + memory_allocation: ByteSize::mb(10), + job_cost: 7, + }, + SplitSearchTaskMetadata { + memory_allocation: ByteSize::mb(10), + job_cost: 3, + }, + SplitSearchTaskMetadata { + memory_allocation: ByteSize::mb(10), + job_cost: 5, + }, + ], + }; + let mut permit_futs = permit_provider.get_permits(task_metadata).await; assert_eq!(permit_provider.get_load(), 15); diff --git a/quickwit/quickwit-serve/src/elasticsearch_api/rest_handler.rs b/quickwit/quickwit-serve/src/elasticsearch_api/rest_handler.rs index a5795de6b4e..cec582d934d 100644 --- a/quickwit/quickwit-serve/src/elasticsearch_api/rest_handler.rs +++ b/quickwit/quickwit-serve/src/elasticsearch_api/rest_handler.rs @@ -581,6 +581,7 @@ fn build_request_for_es_api( count_hits, ignore_missing_indexes, skip_aggregation_finalization: false, + ..Default::default() }, has_doc_id_field, )) diff --git a/quickwit/quickwit-serve/src/search_api/rest_handler.rs b/quickwit/quickwit-serve/src/search_api/rest_handler.rs index b1400fa12c0..01ad87d78e9 100644 --- a/quickwit/quickwit-serve/src/search_api/rest_handler.rs +++ b/quickwit/quickwit-serve/src/search_api/rest_handler.rs @@ -266,6 +266,7 @@ pub fn search_request_from_api_request( count_hits: search_request.count_all.into(), ignore_missing_indexes: false, skip_aggregation_finalization: false, + ..Default::default() }; Ok(search_request) }