From 7e0af10f71bb1ccc7795d46a018df91aa9bf18b7 Mon Sep 17 00:00:00 2001 From: Pascal Seitz Date: Wed, 7 Oct 2026 23:08:54 +0200 Subject: [PATCH] Report warmup downloads per field and segment component Reads issued within count_reads_for_field are counted per (field, component) by CountingStorage. The component is the extension of the tantivy segment file read by the outermost CountingStorage (term, idx, pos, fast, fieldnorm) and is propagated to inner ones, which only see the .split file. Two counting layers in leaf search: reads missing the ephemeral cache (requested) and reads reaching object storage (download). The breakdown is returned in SplitResourceStats.field_download_stats and summed by (field_name, component). --- .../protos/quickwit/search.proto | 24 +++ .../src/codegen/quickwit/quickwit.search.rs | 37 +++- quickwit/quickwit-search/src/fetch_docs.rs | 2 +- quickwit/quickwit-search/src/leaf.rs | 111 +++++++++--- quickwit/quickwit-search/src/lib.rs | 69 ++++++-- quickwit/quickwit-search/src/list_terms.rs | 2 +- quickwit/quickwit-search/src/root.rs | 14 +- quickwit/quickwit-search/src/tests.rs | 58 +++++- .../quickwit-storage/src/counting_storage.rs | 166 ++++++++++++++++-- quickwit/quickwit-storage/src/lib.rs | 4 +- 10 files changed, 424 insertions(+), 63 deletions(-) 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};