diff --git a/quickwit/quickwit-proto/protos/quickwit/search.proto b/quickwit/quickwit-proto/protos/quickwit/search.proto index a1a10efbe09..aa985279de1 100644 --- a/quickwit/quickwit-proto/protos/quickwit/search.proto +++ b/quickwit/quickwit-proto/protos/quickwit/search.proto @@ -389,6 +389,24 @@ message LeafSearchRequest { repeated string index_uris = 9; } +// Warmup reads of one segment component of one field. +message FieldDownloadStats { + // Field name, or JSON path for fast fields. + string field_name = 1; + // Tantivy segment component, as file extension: `term` (term dictionary), + // `idx` (postings), `pos` (positions), `fast` (fast fields), `fieldnorm`. + string component = 2; + // Bytes read that missed the warmup short-lived cache, served by the + // long-term caches or object storage. + uint64 requested_num_bytes = 3; + // Number of read requests counted in `requested_num_bytes`. + uint64 requested_num_requests = 4; + // Bytes downloaded from object storage. Part of `download_num_bytes`. + uint64 download_num_bytes = 5; + // Number of read requests counted in `download_num_bytes`. + uint64 download_num_requests = 6; +} + // Per-split resource statistics. // // All fields are extensive (sum across splits is meaningful) except where noted. @@ -417,6 +435,12 @@ message SplitResourceStats { // Excludes any cross-split finalize work performed outside the single-split // search. uint64 cpu_search_microsecs = 9; + + // Warmup reads per field and segment component, merged by + // `(field_name, component)` when summed. Reads of the split footer and + // hotcache are not attributed to a field: they are only part of + // `download_num_bytes`. + repeated FieldDownloadStats field_download_stats = 10; } // Resource statistics for a single leaf-search call (over one or more splits). diff --git a/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.search.rs b/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.search.rs index ad0c1099cf7..cb1ae0124ff 100644 --- a/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.search.rs +++ b/quickwit/quickwit-proto/src/codegen/quickwit/quickwit.search.rs @@ -313,11 +313,36 @@ pub struct LeafSearchRequest { #[prost(string, repeated, tag = "9")] pub index_uris: ::prost::alloc::vec::Vec<::prost::alloc::string::String>, } +/// Warmup reads of one segment component of one field. +#[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct FieldDownloadStats { + /// Field name, or JSON path for fast fields. + #[prost(string, tag = "1")] + pub field_name: ::prost::alloc::string::String, + /// Tantivy segment component, as file extension: `term` (term dictionary), + /// `idx` (postings), `pos` (positions), `fast` (fast fields), `fieldnorm`. + #[prost(string, tag = "2")] + pub component: ::prost::alloc::string::String, + /// Bytes read that missed the warmup short-lived cache, served by the + /// long-term caches or object storage. + #[prost(uint64, tag = "3")] + pub requested_num_bytes: u64, + /// Number of read requests counted in `requested_num_bytes`. + #[prost(uint64, tag = "4")] + pub requested_num_requests: u64, + /// Bytes downloaded from object storage. Part of `download_num_bytes`. + #[prost(uint64, tag = "5")] + pub download_num_bytes: u64, + /// Number of read requests counted in `download_num_bytes`. + #[prost(uint64, tag = "6")] + pub download_num_requests: u64, +} /// Per-split resource statistics. /// /// All fields are extensive (sum across splits is meaningful) except where noted. #[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] -#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)] +#[derive(Clone, PartialEq, ::prost::Message)] pub struct SplitResourceStats { /// Number of documents in the split. #[prost(uint64, tag = "1")] @@ -351,6 +376,12 @@ pub struct SplitResourceStats { /// search. #[prost(uint64, tag = "9")] pub cpu_search_microsecs: u64, + /// Warmup reads per field and segment component, merged by + /// `(field_name, component)` when summed. Reads of the split footer and + /// hotcache are not attributed to a field: they are only part of + /// `download_num_bytes`. + #[prost(message, repeated, tag = "10")] + pub field_download_stats: ::prost::alloc::vec::Vec, } /// Resource statistics for a single leaf-search call (over one or more splits). /// If the configuration allows it, leaf nodes can offload part of their computation to @@ -367,7 +398,7 @@ pub struct SplitResourceStats { /// are also exceptions — they are merged with `min` rather than `sum` (see /// their per-field doc below). #[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] -#[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)] +#[derive(Clone, PartialEq, ::prost::Message)] pub struct LeafResourceStats { /// Number of splits whose results came from the partial result cache. #[prost(uint64, tag = "1")] @@ -433,7 +464,7 @@ pub struct LeafResourceStats { } /// Resource statistics for a root search. #[derive(serde::Serialize, serde::Deserialize, utoipa::ToSchema)] -#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +#[derive(Clone, PartialEq, ::prost::Message)] pub struct RootResourceStats { /// The leaf with the largest `wall_time_microsecs`. #[prost(message, optional, tag = "1")] diff --git a/quickwit/quickwit-search/src/fetch_docs.rs b/quickwit/quickwit-search/src/fetch_docs.rs index 29be8eb2be2..d2680049422 100644 --- a/quickwit/quickwit-search/src/fetch_docs.rs +++ b/quickwit/quickwit-search/src/fetch_docs.rs @@ -165,7 +165,7 @@ async fn fetch_docs_in_split( global_doc_addrs.sort_by_key(|doc| doc.doc_addr); // Opens the index without the ephemeral unbounded cache, this cache is indeed not useful // when fetching docs as we will fetch them only once. - let (mut index, _) = open_index_with_caches( + let (mut index, _, _) = open_index_with_caches( &searcher_context, index_storage, split, diff --git a/quickwit/quickwit-search/src/leaf.rs b/quickwit/quickwit-search/src/leaf.rs index e60dfad251a..9bc7b5a5721 100644 --- a/quickwit/quickwit-search/src/leaf.rs +++ b/quickwit/quickwit-search/src/leaf.rs @@ -34,8 +34,9 @@ use quickwit_doc_mapper::{Automaton, DocMapper, FastFieldWarmupInfo, TermRange, use quickwit_metrics::{GaugeGuard, HistogramTimer}; use quickwit_proto::search::lambda_single_split_result::Outcome; use quickwit_proto::search::{ - CountHits, LeafResourceStats, LeafSearchRequest, LeafSearchResponse, PartialHit, SearchRequest, - SortOrder, SortValue, SplitIdAndFooterOffsets, SplitResourceStats, SplitSearchError, + CountHits, FieldDownloadStats, LeafResourceStats, LeafSearchRequest, LeafSearchResponse, + PartialHit, SearchRequest, SortOrder, SortValue, SplitIdAndFooterOffsets, SplitResourceStats, + SplitSearchError, }; use quickwit_proto::types::SplitId; use quickwit_query::query_ast::{ @@ -44,8 +45,9 @@ use quickwit_query::query_ast::{ }; use quickwit_query::tokenizers::TokenizerManager; use quickwit_storage::{ - BundleStorage, ByteRangeCache, CountingStorage, MemorySizedCache, OwnedBytes, SearchSplitCache, - Storage, StorageResolver, TimeoutAndRetryStorage, wrap_storage_with_cache, + BundleStorage, ByteRangeCache, CountingStorage, DownloadCounters, FieldComponent, + MemorySizedCache, OwnedBytes, SearchSplitCache, Storage, StorageResolver, + TimeoutAndRetryStorage, count_reads_for_field, wrap_storage_with_cache, }; use tantivy::aggregation::AggContextParams; use tantivy::aggregation::agg_req::{AggregationVariants, Aggregations}; @@ -235,13 +237,17 @@ fn configure_storage_retries( /// - A fast fields cache given by `SearcherContext.storage_long_term_cache`. /// - An ephemeral unbounded cache directory (whose lifetime is tied to the returned `Index` if no /// `ByteRangeCache` is provided). +/// +/// Also returns the counters of the reads that missed the ephemeral cache. Its +/// storage wrapper is the outermost one, so it resolves the field component of +/// the reads (see [`count_reads_for_field`]). pub(crate) async fn open_index_with_caches( searcher_context: &SearcherContext, index_storage: Arc, split_and_footer_offsets: &SplitIdAndFooterOffsets, tokenizer_manager: Option<&TokenizerManager>, ephemeral_unbounded_cache: Option, -) -> anyhow::Result<(Index, HotDirectory)> { +) -> anyhow::Result<(Index, HotDirectory, Arc)> { let index_storage_with_retry_on_timeout = configure_storage_retries(searcher_context, index_storage); @@ -257,7 +263,9 @@ pub(crate) async fn open_index_with_caches( Arc::new(bundle_storage), ); - let directory = StorageDirectory::new(bundle_storage_with_cache); + let (requested_storage, requested_counters) = + CountingStorage::instrument_storage(bundle_storage_with_cache); + let directory = StorageDirectory::new(requested_storage); let hot_directory = if let Some(cache) = ephemeral_unbounded_cache { let caching_directory = CachingDirectory::new(Arc::new(directory), cache); @@ -276,7 +284,7 @@ pub(crate) async fn open_index_with_caches( .clone(), ); set_expr_compilation_cache(&mut index); - Ok((index, hot_directory)) + Ok((index, hot_directory, requested_counters)) } /// Runs `fut`, racing it against `cancel`. If cancellation fires first, the @@ -410,12 +418,14 @@ async fn warm_up_full_term_dictionaries( ) -> anyhow::Result<()> { let mut warm_up_futures = Vec::new(); for field in term_dict_fields { + let field_name = searcher.schema().get_field_name(*field); for segment_reader in searcher.segment_readers() { let inverted_index = segment_reader.inverted_index(*field)?.clone(); - warm_up_futures.push(async move { + let warm_up_fut = async move { let dict = inverted_index.terms(); dict.warm_up_dictionary().await - }); + }; + warm_up_futures.push(count_reads_for_field(field_name.to_string(), warm_up_fut)); } } try_join_all(warm_up_futures).await?; @@ -426,9 +436,11 @@ async fn warm_up_full_term_dictionaries( async fn warm_up_all_postings(searcher: &Searcher, fields: &HashSet) -> anyhow::Result<()> { let mut warm_up_futures = Vec::new(); for field in fields { + let field_name = searcher.schema().get_field_name(*field); for segment_reader in searcher.segment_readers() { let inverted_index = segment_reader.inverted_index(*field)?.clone(); - warm_up_futures.push(async move { inverted_index.warm_postings_full(false).await }); + let warm_up_fut = async move { inverted_index.warm_postings_full(false).await }; + warm_up_futures.push(count_reads_for_field(field_name.to_string(), warm_up_fut)); } } try_join_all(warm_up_futures).await?; @@ -468,7 +480,10 @@ async fn warm_up_fastfields( let fast_field_reader = segment_reader.fast_fields(); for fast_field in fast_fields { let warm_up_fut = warm_up_fastfield(fast_field_reader, fast_field); - warm_up_futures.push(Box::pin(warm_up_fut)); + warm_up_futures.push(Box::pin(count_reads_for_field( + fast_field.name.clone(), + warm_up_fut, + ))); } } futures::future::try_join_all(warm_up_futures).await?; @@ -484,6 +499,7 @@ async fn warm_up_terms( ) -> anyhow::Result<()> { let mut warm_up_futures = Vec::new(); for (field, terms) in terms_grouped_by_field { + let field_name = searcher.schema().get_field_name(*field); for segment_reader in searcher.segment_readers() { let inv_idx = segment_reader.inverted_index(*field)?; let segment_id = segment_reader.segment_id(); @@ -496,7 +512,7 @@ async fn warm_up_terms( Some(abort_token) if required_terms.contains(term) => Some(abort_token), _ => None, }; - warm_up_futures.push(async move { + let warm_up_fut = async move { let found = inv_idx_clone.warm_postings(term, *position_needed).await?; if !found && let Some(abort_token) = cancel_on_empty { // Report the absence and fire the abort token. Both are synchronous, so @@ -505,7 +521,8 @@ async fn warm_up_terms( abort_token.cancel(); } anyhow::Ok(()) - }); + }; + warm_up_futures.push(count_reads_for_field(field_name.to_string(), warm_up_fut)); } } } @@ -524,16 +541,18 @@ async fn warm_up_term_ranges( ) -> anyhow::Result<()> { let mut warm_up_futures = Vec::new(); for (field, terms) in terms_grouped_by_field { + let field_name = searcher.schema().get_field_name(*field); for segment_reader in searcher.segment_readers() { let inv_idx = segment_reader.inverted_index(*field)?; for (term_range, position_needed) in terms.iter() { let inv_idx_clone = inv_idx.clone(); let range = (term_range.start.as_ref(), term_range.end.as_ref()); - warm_up_futures.push(async move { + let warm_up_fut = async move { inv_idx_clone .warm_postings_range(range, term_range.limit, *position_needed) .await - }); + }; + warm_up_futures.push(count_reads_for_field(field_name.to_string(), warm_up_fut)); } } } @@ -554,11 +573,12 @@ async fn warm_up_automatons( .map_err(|_| std::io::Error::other("task panicked"))? }; for (field, automatons) in terms_grouped_by_field { + let field_name = searcher.schema().get_field_name(*field); for segment_reader in searcher.segment_readers() { let inv_idx = segment_reader.inverted_index(*field)?; for automaton in automatons { let inv_idx_clone = inv_idx.clone(); - warm_up_futures.push(async move { + let warm_up_fut = async move { match automaton { Automaton::Regex(path, regex_str) => { let regex = get_or_compile_cached_fst_regex(regex_str) @@ -575,7 +595,8 @@ async fn warm_up_automatons( .context("failed to load automaton") } } - }); + }; + warm_up_futures.push(count_reads_for_field(field_name.to_string(), warm_up_fut)); } } } @@ -588,12 +609,16 @@ async fn warm_up_fieldnorms(searcher: &Searcher, requires_scoring: bool) -> anyh return Ok(()); } let mut warm_up_futures = Vec::new(); - for field in searcher.schema().fields() { + for (field, field_entry) in searcher.schema().fields() { for segment_reader in searcher.segment_readers() { let fieldnorm_readers = segment_reader.fieldnorms_readers(); - let file_handle_opt = fieldnorm_readers.get_inner_file().open_read(field.0); + let file_handle_opt = fieldnorm_readers.get_inner_file().open_read(field); if let Some(file_handle) = file_handle_opt { - warm_up_futures.push(async move { file_handle.read_bytes_async().await }) + let warm_up_fut = async move { file_handle.read_bytes_async().await }; + warm_up_futures.push(count_reads_for_field( + field_entry.name().to_string(), + warm_up_fut, + )); } } } @@ -601,6 +626,43 @@ async fn warm_up_fieldnorms(searcher: &Searcher, requires_scoring: bool) -> anyh Ok(()) } +/// Returns the warmup reads per field component: `requested` counts the reads +/// that missed the ephemeral cache, `downloaded` the reads that reached object +/// storage. +fn field_download_stats( + requested_counters: &DownloadCounters, + download_counters: &DownloadCounters, +) -> Vec { + fn stats_entry( + stats_per_field_component: &mut HashMap, + field_component: FieldComponent, + ) -> &mut FieldDownloadStats { + stats_per_field_component + .entry(field_component) + .or_insert_with_key(|field_component| FieldDownloadStats { + field_name: field_component.field_name.clone(), + component: field_component.component.clone(), + ..Default::default() + }) + } + let mut stats_per_field_component: HashMap = HashMap::new(); + for (field_component, (num_bytes, num_requests)) in + requested_counters.per_field_component_snapshot() + { + let stats = stats_entry(&mut stats_per_field_component, field_component); + stats.requested_num_bytes = num_bytes; + stats.requested_num_requests = num_requests; + } + for (field_component, (num_bytes, num_requests)) in + download_counters.per_field_component_snapshot() + { + let stats = stats_entry(&mut stats_per_field_component, field_component); + stats.download_num_bytes = num_bytes; + stats.download_num_requests = num_requests; + } + stats_per_field_component.into_values().collect() +} + fn get_leaf_resp_from_count(count: u64) -> LeafSearchResponse { LeafSearchResponse { num_hits: count, @@ -632,12 +694,12 @@ fn leaf_resource_stats_for_split(split_stats: SplitResourceStats) -> LeafResourc LeafResourceStats { localexec_num_splits: 1, localexec_num_docs: split_stats.split_num_docs, - split_resources_sum: Some(split_stats), - split_resources_worst: Some(split_stats), min_wait_for_search_permit_microsecs: Some( split_stats.wait_for_search_permit_microsecs, ), min_wait_for_cpu_pool_microsecs: Some(split_stats.wait_for_cpu_pool_microsecs), + split_resources_worst: Some(split_stats.clone()), + split_resources_sum: Some(split_stats), ..Default::default() } } @@ -721,7 +783,7 @@ async fn leaf_search_single_split( // ABOVE this wrapper, so reads served from cache do not contribute to the // counters — that is the desired "downloaded from object storage" semantics. let (storage, download_counters) = CountingStorage::instrument_storage(storage); - let (index, hot_directory) = open_index_with_caches( + let (index, hot_directory, requested_counters) = open_index_with_caches( &ctx.searcher_context, storage, &split, @@ -859,6 +921,7 @@ async fn leaf_search_single_split( provably_empty }; let warmup_end = Instant::now(); + let field_download_stats = field_download_stats(&requested_counters, &download_counters); let warmup_duration: Duration = warmup_end.duration_since(warmup_start); let warmup_size = ByteSize(byte_range_cache.get_num_bytes()); if warmup_size > search_permit.memory_allocation() { @@ -902,6 +965,7 @@ async fn leaf_search_single_split( warmup_microsecs: warmup_duration.as_micros() as u64, wait_for_cpu_pool_microsecs: 0, cpu_search_microsecs: 0, + field_download_stats, }; let mut leaf_search_response = get_leaf_resp_from_count(0); leaf_search_response.resource_stats = Some(leaf_resource_stats_for_split(split_stats)); @@ -981,6 +1045,7 @@ async fn leaf_search_single_split( warmup_microsecs: warmup_duration.as_micros() as u64, wait_for_cpu_pool_microsecs: cpu_thread_pool_wait_microsecs.as_micros() as u64, cpu_search_microsecs: cpu_start.elapsed().as_micros() as u64, + field_download_stats, }; leaf_search_response.resource_stats = Some(leaf_resource_stats_for_split(split_stats)); diff --git a/quickwit/quickwit-search/src/lib.rs b/quickwit/quickwit-search/src/lib.rs index 64b28776bd5..a745ac87d06 100644 --- a/quickwit/quickwit-search/src/lib.rs +++ b/quickwit/quickwit-search/src/lib.rs @@ -403,6 +403,7 @@ pub(crate) fn split_phase_sum_microsecs(stats: &SplitResourceStats) -> u64 { } /// Field-wise sum of two `SplitResourceStats` (every field is extensive). +/// `field_download_stats` are summed by `(field_name, component)`. pub(crate) fn add_split_stats(acc: &mut SplitResourceStats, other: &SplitResourceStats) { acc.split_num_docs += other.split_num_docs; acc.input_memory_bytes += other.input_memory_bytes; @@ -413,6 +414,21 @@ pub(crate) fn add_split_stats(acc: &mut SplitResourceStats, other: &SplitResourc acc.warmup_microsecs += other.warmup_microsecs; acc.wait_for_cpu_pool_microsecs += other.wait_for_cpu_pool_microsecs; acc.cpu_search_microsecs += other.cpu_search_microsecs; + for other_field_stats in &other.field_download_stats { + // Linear scan: a query warms up a handful of field components. + let acc_field_stats_opt = acc.field_download_stats.iter_mut().find(|acc_field_stats| { + acc_field_stats.field_name == other_field_stats.field_name + && acc_field_stats.component == other_field_stats.component + }); + let Some(acc_field_stats) = acc_field_stats_opt else { + acc.field_download_stats.push(other_field_stats.clone()); + continue; + }; + acc_field_stats.requested_num_bytes += other_field_stats.requested_num_bytes; + acc_field_stats.requested_num_requests += other_field_stats.requested_num_requests; + acc_field_stats.download_num_bytes += other_field_stats.download_num_bytes; + acc_field_stats.download_num_requests += other_field_stats.download_num_requests; + } } /// Min of two `Option`, treating `None` as "no contribution": @@ -462,10 +478,13 @@ pub(crate) fn add_leaf_stats(acc: &mut LeafResourceStats, other: &LeafResourceSt .get_or_insert_with(SplitResourceStats::default); add_split_stats(acc_split, other_split); } - acc.split_resources_worst = [acc.split_resources_worst, other.split_resources_worst] - .into_iter() - .flatten() - .max_by_key(split_phase_sum_microsecs); + acc.split_resources_worst = [ + acc.split_resources_worst.take(), + other.split_resources_worst.clone(), + ] + .into_iter() + .flatten() + .max_by_key(split_phase_sum_microsecs); } /// Merge an iterator of `Option` into a single `Option`. @@ -488,6 +507,8 @@ pub(crate) fn merge_leaf_stats_it<'a>( #[cfg(test)] mod stats_merge_tests { + use quickwit_proto::search::FieldDownloadStats; + use super::*; fn split_stats(num_docs: u64, warmup: u64, search: u64) -> SplitResourceStats { @@ -503,7 +524,7 @@ mod stats_merge_tests { LeafResourceStats { localexec_num_splits: 1, localexec_num_docs: split.split_num_docs, - split_resources_sum: Some(split), + split_resources_sum: Some(split.clone()), split_resources_worst: Some(split), ..Default::default() } @@ -515,6 +536,18 @@ mod stats_merge_tests { /// destructure forces the test to be updated whenever a field is added. #[test] fn test_add_split_stats_sums_every_field() { + let term_stats = FieldDownloadStats { + field_name: "body".to_string(), + component: "term".to_string(), + requested_num_bytes: 1, + requested_num_requests: 10, + download_num_bytes: 100, + download_num_requests: 1_000, + }; + let postings_stats = FieldDownloadStats { + component: "idx".to_string(), + ..term_stats.clone() + }; let mut acc = SplitResourceStats { split_num_docs: 1, input_memory_bytes: 10, @@ -525,6 +558,7 @@ mod stats_merge_tests { warmup_microsecs: 1_000_000, wait_for_cpu_pool_microsecs: 10_000_000, cpu_search_microsecs: 100_000_000, + field_download_stats: vec![term_stats.clone()], }; let other = SplitResourceStats { split_num_docs: 2, @@ -536,6 +570,7 @@ mod stats_merge_tests { warmup_microsecs: 2_000_000, wait_for_cpu_pool_microsecs: 20_000_000, cpu_search_microsecs: 200_000_000, + field_download_stats: vec![term_stats.clone(), postings_stats.clone()], }; // Destructure on the proto type itself so a newly-added field forces // an update to this test (otherwise the assertion below would silently @@ -550,7 +585,8 @@ mod stats_merge_tests { warmup_microsecs: _, wait_for_cpu_pool_microsecs: _, cpu_search_microsecs: _, - } = other; + field_download_stats: _, + } = &other; add_split_stats(&mut acc, &other); @@ -563,6 +599,17 @@ mod stats_merge_tests { assert_eq!(acc.warmup_microsecs, 3_000_000); assert_eq!(acc.wait_for_cpu_pool_microsecs, 30_000_000); assert_eq!(acc.cpu_search_microsecs, 300_000_000); + let summed_term_stats = FieldDownloadStats { + requested_num_bytes: 2, + requested_num_requests: 20, + download_num_bytes: 200, + download_num_requests: 2_000, + ..term_stats + }; + assert_eq!( + acc.field_download_stats, + vec![summed_term_stats, postings_stats] + ); } #[test] @@ -573,7 +620,7 @@ mod stats_merge_tests { cpu_search_microsecs: 1_000, ..Default::default() }; - let snapshot = acc; + let snapshot = acc.clone(); add_split_stats(&mut acc, &SplitResourceStats::default()); assert_eq!(acc, snapshot); } @@ -607,7 +654,7 @@ mod stats_merge_tests { lambda_bottleneck: 0, localexec_num_splits: 1_000_000, localexec_num_docs: 10_000_000, - split_resources_worst: Some(split_a), + split_resources_worst: Some(split_a.clone()), split_resources_sum: Some(split_a), min_wait_for_search_permit_microsecs: Some(100), min_wait_for_cpu_pool_microsecs: Some(2_000), @@ -624,7 +671,7 @@ mod stats_merge_tests { lambda_bottleneck: 1, localexec_num_splits: 2_000_000, localexec_num_docs: 20_000_000, - split_resources_worst: Some(split_b), + split_resources_worst: Some(split_b.clone()), split_resources_sum: Some(split_b), min_wait_for_search_permit_microsecs: Some(50), min_wait_for_cpu_pool_microsecs: Some(1_000), @@ -709,7 +756,7 @@ mod stats_merge_tests { lambda_bottleneck: 1, localexec_num_splits: 1, localexec_num_docs: 3, - split_resources_sum: Some(split), + split_resources_sum: Some(split.clone()), split_resources_worst: Some(split), // Set the min fields explicitly so this test verifies that // `min_opt(Some(x), None) = Some(x)` keeps them unchanged. @@ -718,7 +765,7 @@ mod stats_merge_tests { wall_time_microsecs: 42, ..Default::default() }; - let snapshot = acc; + let snapshot = acc.clone(); add_leaf_stats(&mut acc, &LeafResourceStats::default()); assert_eq!(acc, snapshot); } diff --git a/quickwit/quickwit-search/src/list_terms.rs b/quickwit/quickwit-search/src/list_terms.rs index 087c0749bd1..4efe7cd7b47 100644 --- a/quickwit/quickwit-search/src/list_terms.rs +++ b/quickwit/quickwit-search/src/list_terms.rs @@ -220,7 +220,7 @@ async fn leaf_list_terms_single_split( split: SplitIdAndFooterOffsets, ) -> crate::Result { let cache = ByteRangeCache::with_infinite_capacity(); - let (index, _) = + let (index, _, _) = open_index_with_caches(searcher_context, storage, &split, None, Some(cache)).await?; let split_schema = index.schema(); let reader = index diff --git a/quickwit/quickwit-search/src/root.rs b/quickwit/quickwit-search/src/root.rs index e316c60bf85..0ef1034e34d 100644 --- a/quickwit/quickwit-search/src/root.rs +++ b/quickwit/quickwit-search/src/root.rs @@ -2990,7 +2990,7 @@ mod tests { localexec_num_splits: 1, localexec_num_docs: 10, lambda_bottleneck: 0, - split_resources_sum: Some(split_a), + split_resources_sum: Some(split_a.clone()), split_resources_worst: Some(split_a), wall_time_microsecs: 1_000, ..Default::default() @@ -3002,7 +3002,7 @@ mod tests { localexec_num_splits: 1, localexec_num_docs: 20, lambda_bottleneck: 1, - split_resources_sum: Some(split_b), + split_resources_sum: Some(split_b.clone()), split_resources_worst: Some(split_b), wall_time_microsecs: 2_500, ..Default::default() @@ -3071,7 +3071,7 @@ mod tests { resource_stats: Some(LeafResourceStats { localexec_num_splits: 1, localexec_num_docs: 7, - split_resources_sum: Some(split), + split_resources_sum: Some(split.clone()), split_resources_worst: Some(split), wall_time_microsecs: 500, ..Default::default() @@ -3156,7 +3156,7 @@ mod tests { let leaf_stats_1 = LeafResourceStats { localexec_num_splits: 1, localexec_num_docs: 10, - split_resources_sum: Some(split1_stats), + split_resources_sum: Some(split1_stats.clone()), split_resources_worst: Some(split1_stats), wall_time_microsecs: 1_000, ..Default::default() @@ -3164,7 +3164,7 @@ mod tests { let leaf_stats_2 = LeafResourceStats { localexec_num_splits: 1, localexec_num_docs: 20, - split_resources_sum: Some(split2_stats), + split_resources_sum: Some(split2_stats.clone()), split_resources_worst: Some(split2_stats), wall_time_microsecs: 2_500, ..Default::default() @@ -3182,7 +3182,7 @@ mod tests { failed_splits: Vec::new(), num_attempted_splits: 1, num_successful_splits: 1, - resource_stats: Some(leaf_stats_1), + resource_stats: Some(leaf_stats_1.clone()), ..Default::default() }) }, @@ -3203,7 +3203,7 @@ mod tests { failed_splits: Vec::new(), num_attempted_splits: 1, num_successful_splits: 1, - resource_stats: Some(leaf_stats_2), + resource_stats: Some(leaf_stats_2.clone()), ..Default::default() }) }, diff --git a/quickwit/quickwit-search/src/tests.rs b/quickwit/quickwit-search/src/tests.rs index d15c14a3cec..d7080780fa6 100644 --- a/quickwit/quickwit-search/src/tests.rs +++ b/quickwit/quickwit-search/src/tests.rs @@ -2638,7 +2638,7 @@ async fn test_time_bounded_query_populates_and_reuses_complete_predicate_cache() let first_input_memory_bytes = first_response .resource_stats .as_ref() - .and_then(|stats| stats.split_resources_sum) + .and_then(|stats| stats.split_resources_sum.as_ref()) .expect("the split should report resource stats") .input_memory_bytes; @@ -2669,7 +2669,7 @@ async fn test_time_bounded_query_populates_and_reuses_complete_predicate_cache() let full_split_input_memory_bytes = full_split_response .resource_stats .as_ref() - .and_then(|stats| stats.split_resources_sum) + .and_then(|stats| stats.split_resources_sum.as_ref()) .expect("the split should report resource stats") .input_memory_bytes; assert!( @@ -2735,6 +2735,60 @@ async fn test_time_bounded_query_populates_and_reuses_complete_predicate_cache() test_sandbox.assert_quit().await; } +#[tokio::test] +async fn test_leaf_search_reports_field_download_stats() { + let (test_sandbox, searcher_context, storage, splits, doc_mapper, _start_timestamp) = + negative_cache_ts_test_setup().await; + let search_request = SearchRequest { + index_id_patterns: vec!["negative-cache-ts-index".to_string()], + query_ast: qast_json_helper("info", &["body"]), + max_hits: 10, + sort_fields: vec![SortField { + field_name: "ts".to_string(), + sort_order: SortOrder::Desc as i32, + sort_datetime_format: None, + }], + ..Default::default() + }; + let response = single_doc_mapping_leaf_search( + searcher_context, + std::sync::Arc::new(search_request), + storage, + splits, + doc_mapper, + ) + .await + .unwrap(); + test_sandbox.assert_quit().await; + assert_eq!(response.num_hits, 10); + let split_stats = response + .resource_stats + .and_then(|stats| stats.split_resources_sum) + .expect("the split should report resource stats"); + let mut field_components: Vec<(&str, &str)> = split_stats + .field_download_stats + .iter() + .map(|stats| (stats.field_name.as_str(), stats.component.as_str())) + .collect(); + field_components.sort(); + assert_eq!( + field_components, + [("body", "idx"), ("body", "term"), ("ts", "fast")] + ); + for stats in &split_stats.field_download_stats { + assert!(stats.requested_num_bytes > 0); + assert!(stats.requested_num_requests > 0); + } + let field_download_num_bytes: u64 = split_stats + .field_download_stats + .iter() + .map(|stats| stats.download_num_bytes) + .sum(); + // The footer and hotcache are downloaded without a field. + assert!(field_download_num_bytes > 0); + assert!(field_download_num_bytes < split_stats.download_num_bytes); +} + #[tokio::test] async fn test_negative_cache_short_circuits_across_time_windows() { // A per-term absence does not depend on the time window: the required term diff --git a/quickwit/quickwit-storage/src/counting_storage.rs b/quickwit/quickwit-storage/src/counting_storage.rs index 14aba79573b..29601484059 100644 --- a/quickwit/quickwit-storage/src/counting_storage.rs +++ b/quickwit/quickwit-storage/src/counting_storage.rs @@ -12,10 +12,12 @@ // See the License for the specific language governing permissions and // limitations under the License. +use std::collections::HashMap; +use std::future::Future; use std::ops::Range; use std::path::Path; -use std::sync::Arc; use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Arc, Mutex}; use async_trait::async_trait; use quickwit_common::uri::Uri; @@ -25,6 +27,58 @@ use tokio::io::AsyncRead; use crate::storage::SendableAsync; use crate::{BulkDeleteError, ListObjectsStream, PutPayload, Storage, StorageResult}; +/// Field and tantivy segment component a read is counted for. +#[derive(Debug, Clone, PartialEq, Eq, Hash)] +pub struct FieldComponent { + /// Name of the field, or JSON path of the fast field. + pub field_name: String, + /// Extension of the file read by the outermost [`CountingStorage`]: `term` + /// (term dictionary), `idx` (postings), `pos` (positions), `fast` (fast + /// fields), `fieldnorm`, ... + pub component: String, +} + +tokio::task_local! { + /// Field set by [`count_reads_for_field`], and the component resolved by the + /// outermost [`CountingStorage`], if any. + static READ_CONTEXT: (String, Option); +} + +/// Runs `fut` so that the reads it issues through a [`CountingStorage`] are also +/// counted per [`FieldComponent`] of `field_name`. +/// +/// The field is set for the duration of each `poll` of `fut`, so concurrently +/// joined futures count their reads for their own field. Reads issued from a +/// task spawned by `fut` are not counted per field. +pub fn count_reads_for_field( + field_name: String, + fut: F, +) -> impl Future { + READ_CONTEXT.scope((field_name, None), fut) +} + +/// Returns the field component a read of `path` is counted for, if it happens +/// within [`count_reads_for_field`]. +/// +/// The component is taken from the extension of the file read by the outermost +/// `CountingStorage`. Inner storages may read a different file (a +/// `BundleStorage` maps segment files to ranges of the `.split` file), so the +/// outermost component is propagated to them by [`CountingStorage::count_read`]. +fn current_field_component(path: &Path) -> Option { + READ_CONTEXT + .try_with(|(field_name, component)| { + let component = component.clone().unwrap_or_else(|| { + let extension = path.extension().unwrap_or_default(); + extension.to_string_lossy().into_owned() + }); + FieldComponent { + field_name: field_name.clone(), + component, + } + }) + .ok() +} + /// Per-request download counters tracked by [`CountingStorage`]. /// /// `bytes` accumulates the size of every successfully fulfilled read; `requests` @@ -35,6 +89,8 @@ use crate::{BulkDeleteError, ListObjectsStream, PutPayload, Storage, StorageResu pub struct DownloadCounters { bytes: AtomicU64, requests: AtomicU64, + /// `(bytes, requests)` of the reads issued within [`count_reads_for_field`]. + per_field_component: Mutex>, } impl DownloadCounters { @@ -46,9 +102,22 @@ impl DownloadCounters { ) } - fn record_read(&self, num_bytes: u64) { + /// Snapshots the counters of the reads issued within [`count_reads_for_field`] + /// as `field component -> (bytes, requests)`. + pub fn per_field_component_snapshot(&self) -> HashMap { + self.per_field_component.lock().unwrap().clone() + } + + fn record_read(&self, num_bytes: u64, field_component_opt: Option) { self.bytes.fetch_add(num_bytes, Ordering::Relaxed); self.requests.fetch_add(1, Ordering::Relaxed); + let Some(field_component) = field_component_opt else { + return; + }; + let mut per_field_component = self.per_field_component.lock().unwrap(); + let (bytes, requests) = per_field_component.entry(field_component).or_default(); + *bytes += num_bytes; + *requests += 1; } } @@ -81,6 +150,29 @@ impl CountingStorage { }; (Arc::new(instrumented_storage), counters) } + + /// Awaits `read` of `path` and counts it, propagating the field component to + /// the reads of the inner storage. + async fn count_read( + &self, + path: &Path, + read: impl Future>, + num_bytes: impl FnOnce(&T) -> u64, + ) -> StorageResult { + let Some(field_component) = current_field_component(path) else { + let output = read.await?; + self.counters.record_read(num_bytes(&output), None); + return Ok(output); + }; + let read_context = ( + field_component.field_name.clone(), + Some(field_component.component.clone()), + ); + let output = READ_CONTEXT.scope(read_context, read).await?; + self.counters + .record_read(num_bytes(&output), Some(field_component)); + Ok(output) + } } #[async_trait] @@ -103,15 +195,14 @@ impl Storage for CountingStorage { } async fn copy_to_file(&self, path: &Path, output_path: &Path) -> StorageResult { - let num_bytes = self.inner.copy_to_file(path, output_path).await?; - self.counters.record_read(num_bytes); - Ok(num_bytes) + let read = self.inner.copy_to_file(path, output_path); + self.count_read(path, read, |num_bytes| *num_bytes).await } async fn get_slice(&self, path: &Path, range: Range) -> StorageResult { - let bytes = self.inner.get_slice(path, range).await?; - self.counters.record_read(bytes.len() as u64); - Ok(bytes) + let read = self.inner.get_slice(path, range); + self.count_read(path, read, |bytes| bytes.len() as u64) + .await } async fn get_slice_stream( @@ -123,15 +214,14 @@ impl Storage for CountingStorage { // The stream may yield fewer bytes if the caller drops it early, but // that is rare and the over-count is bounded by the requested range. let range_len = range.len() as u64; - let stream = self.inner.get_slice_stream(path, range).await?; - self.counters.record_read(range_len); - Ok(stream) + let read = self.inner.get_slice_stream(path, range); + self.count_read(path, read, |_| range_len).await } async fn get_all(&self, path: &Path) -> StorageResult { - let bytes = self.inner.get_all(path).await?; - self.counters.record_read(bytes.len() as u64); - Ok(bytes) + let read = self.inner.get_all(path); + self.count_read(path, read, |bytes| bytes.len() as u64) + .await } async fn delete(&self, path: &Path) -> StorageResult<()> { @@ -164,6 +254,54 @@ mod tests { use super::*; use crate::RamStorageBuilder; + fn field_component(field_name: &str, component: &str) -> FieldComponent { + FieldComponent { + field_name: field_name.to_string(), + component: component.to_string(), + } + } + + #[tokio::test] + async fn test_counting_storage_counts_reads_per_field_component() { + let inner = RamStorageBuilder::default() + .put("seg.fast", b"hello world") + .put("seg.term", b"hello world") + .build(); + let (inner_storage, inner_counters) = CountingStorage::instrument_storage(Arc::new(inner)); + let (storage, counters) = CountingStorage::instrument_storage(inner_storage); + let status_read = count_reads_for_field("status".to_string(), async { + tokio::task::yield_now().await; + storage + .get_slice(Path::new("seg.fast"), 0..5) + .await + .unwrap(); + }); + let body_read = count_reads_for_field("body".to_string(), async { + tokio::task::yield_now().await; + storage + .get_slice(Path::new("seg.term"), 5..11) + .await + .unwrap(); + }); + tokio::join!(status_read, body_read); + storage + .get_slice(Path::new("seg.fast"), 0..1) + .await + .unwrap(); + + let expected_per_field_component = HashMap::from([ + (field_component("status", "fast"), (5, 1)), + (field_component("body", "term"), (6, 1)), + ]); + for counters in [counters, inner_counters] { + assert_eq!(counters.snapshot(), (12, 3)); + assert_eq!( + counters.per_field_component_snapshot(), + expected_per_field_component + ); + } + } + #[tokio::test] async fn test_counting_storage_counts_get_slice() { let inner = RamStorageBuilder::default() diff --git a/quickwit/quickwit-storage/src/lib.rs b/quickwit/quickwit-storage/src/lib.rs index 73ef173eee2..5b535519eef 100644 --- a/quickwit/quickwit-storage/src/lib.rs +++ b/quickwit/quickwit-storage/src/lib.rs @@ -69,7 +69,9 @@ pub use self::cache::{ ByteRangeCache, FileByteRangeCache, MemorySizedCache, QuickwitCache, StorageCache, wrap_storage_with_cache, }; -pub use self::counting_storage::{CountingStorage, DownloadCounters}; +pub use self::counting_storage::{ + CountingStorage, DownloadCounters, FieldComponent, count_reads_for_field, +}; pub use self::local_file_storage::{LocalFileStorage, LocalFileStorageFactory}; #[cfg(feature = "azure")] pub use self::object_storage::{AzureBlobStorage, AzureBlobStorageFactory};