From 484f3f2d4c666f00b43f1ebfb1da928efe94ff51 Mon Sep 17 00:00:00 2001 From: Paul Masurel Date: Tue, 1 Sep 2026 11:05:59 +0200 Subject: [PATCH] Following comment Removing Fused from the element that are private to the module. Re-added optimization in stats for single doc requests to avoid perf regression Renamed Fused -> Flattened --- benches/agg_bench.rs | 8 +- src/aggregation/block_accessor.rs | 6 - src/aggregation/bucket/histogram/histogram.rs | 12 +- ...stogram.rs => flattened_term_histogram.rs} | 151 +++++++++--------- src/aggregation/bucket/term_agg/mod.rs | 9 +- src/aggregation/metric/stats.rs | 17 +- 6 files changed, 107 insertions(+), 96 deletions(-) rename src/aggregation/bucket/term_agg/{term_histogram.rs => flattened_term_histogram.rs} (87%) diff --git a/benches/agg_bench.rs b/benches/agg_bench.rs index b9f55b564..c02a4bc32 100644 --- a/benches/agg_bench.rs +++ b/benches/agg_bench.rs @@ -500,7 +500,7 @@ fn terms_status_with_date_histogram() -> AggregationRequest { }) } -/// Same fused terms × date_histogram, but with `hard_bounds`. The timestamps span 0..120h; the +/// Same flattened terms × date_histogram, but with `hard_bounds`. The timestamps span 0..120h; the /// bounds drop only the first and last hour (ms: 1h=3_600_000, 119h=428_400_000), so almost every /// doc is in-bounds. This exercises the collector's hard-bounds path: `bounds.contains` runs per /// doc (the `all_docs_in_bounds` short-circuit is off) and the rare out-of-bounds doc takes the @@ -562,9 +562,9 @@ fn terms_status_with_date_histogram_hard_bounds() -> AggregationRequest { }) } -/// Same fused terms × date_histogram, but with a sibling terms aggregation next to it. The fused -/// fast path should still trigger for `my_texts` (sibling aggregations are independent top-level -/// aggregations, so they don't change its eligibility). +/// Same flattened terms × date_histogram, but with a sibling terms aggregation next to it. The +/// flattened fast path should still trigger for `my_texts` (sibling aggregations are independent +/// top-level aggregations, so they don't change its eligibility). fn terms_status_with_date_histogram_and_sibling_terms() -> AggregationRequest { json!({ "my_texts": { diff --git a/src/aggregation/block_accessor.rs b/src/aggregation/block_accessor.rs index 1d41e5e3e..2eb592f1b 100644 --- a/src/aggregation/block_accessor.rs +++ b/src/aggregation/block_accessor.rs @@ -55,12 +55,6 @@ impl BlockValueSource for Column { let cardinality = self.index.get_cardinality(); if cardinality.is_full() { load_full_column_values(docs, self, values); - } else if docs.len() == 1 { - values.clear(); - values.extend(self.values_for_doc(docs[0])); - docids.clear(); - docids.resize(values.len(), docs[0]); - row_ids.clear(); } else { docids.clear(); row_ids.clear(); diff --git a/src/aggregation/bucket/histogram/histogram.rs b/src/aggregation/bucket/histogram/histogram.rs index ce5590644..e97a9ca5c 100644 --- a/src/aggregation/bucket/histogram/histogram.rs +++ b/src/aggregation/bucket/histogram/histogram.rs @@ -630,9 +630,9 @@ impl SegmentHistogramCollector { impl SegmentHistogramCollector<()> { /// Builds a histogram collector whose parent `t` is a dense histogram filled from /// `counts[t * num_time_buckets .. (t + 1) * num_time_buckets]` (row-major), consolidating each - /// cell's count lanes. Used by the fused terms×histogram collector to turn its flat 2D counters - /// into the regular intermediate result, so cross-segment merging is shared with the general - /// path. + /// cell's count lanes. Used by the flattened terms×histogram collector to turn its flat 2D + /// counters into the regular intermediate result, so cross-segment merging is shared with the + /// general path. pub(crate) fn from_dense_rows( req_data: HistogramAggReqData, base_pos: i64, @@ -685,7 +685,7 @@ fn normalize_histogram_req(req_data: &mut HistogramAggReqData) -> crate::Result< req_data.offset = req_data.req.offset.unwrap_or(0.0); // Drop `hard_bounds` that can't exclude any value (the column's range already sits inside // them): the per-doc `bounds.contains` check is then a no-op, so collapsing to the unbounded - // sentinel lets the histogram hot loop skip it and the fused term×histogram path derive + // sentinel lets the histogram hot loop skip it and the flattened term×histogram path derive // per-term counts from the grid. Only this collect-time filter is touched — empty-bucket // emission reads `req.hard_bounds` directly (see `get_req_min_max`), and `hard_bounds` only // ever clips that range, so a wider-than-data bound leaves the result unchanged. @@ -704,7 +704,7 @@ fn normalize_histogram_req(req_data: &mut HistogramAggReqData) -> crate::Result< /// Clones and normalizes (resolving interval/offset/bounds) the histogram request at `node`, and /// returns it together with its dense bucket range — or `None` if the column has no usable range. -/// Used by the fused terms×histogram collector, which then owns the normalized request. +/// Used by the flattened terms×histogram collector, which then owns the normalized request. pub(crate) fn prepare_histogram_dense_range( agg_data: &AggregationsSegmentCtx, node: &AggRefNode, @@ -1437,7 +1437,7 @@ mod tests { /// `hard_bounds` wider than the data (here with mid-interval edges, to cover the "bound cuts a /// bucket" case) can't exclude any value, so the result must be identical to the same request /// without bounds. Guards the normalization that collapses such bounds to the unbounded - /// sentinel so the hot loop / fused path can skip the per-doc bounds check. + /// sentinel so the hot loop / flattened path can skip the per-doc bounds check. fn histogram_non_binding_hard_bounds_test_with_opt(merge_segments: bool) -> crate::Result<()> { let values = vec![10.0, 12.0, 14.0, 16.0, 10.0, 13.0, 10.0, 12.0]; let index = get_test_index_from_values(merge_segments, &values)?; diff --git a/src/aggregation/bucket/term_agg/term_histogram.rs b/src/aggregation/bucket/term_agg/flattened_term_histogram.rs similarity index 87% rename from src/aggregation/bucket/term_agg/term_histogram.rs rename to src/aggregation/bucket/term_agg/flattened_term_histogram.rs index 2d1d201b7..675171448 100644 --- a/src/aggregation/bucket/term_agg/term_histogram.rs +++ b/src/aggregation/bucket/term_agg/flattened_term_histogram.rs @@ -1,8 +1,8 @@ -//! Fused collector for the very common shape `terms` (low cardinality) × a single +//! Flattened collector for the very common shape `terms` (low cardinality) × a single //! `histogram`/`date_histogram` sub-aggregation with nothing nested below it. //! -//! See [`FusedTermHistogramCollector`] for the approach and [`maybe_build_fused_collector`] for -//! the conditions under which it is used. +//! See [`FlattenedTermHistogramCollector`] for the approach and +//! [`maybe_build_flattened_collector`] for the conditions under which it is used. use std::fmt::Debug; @@ -24,12 +24,12 @@ use crate::aggregation::intermediate_agg_result::{ use crate::aggregation::segment_agg_result::{BucketIdProvider, SegmentAggregationCollector}; use crate::aggregation::{f64_from_fastfield_u64, BucketId, ColumnBlockAccessor}; -/// Maximum number of physical counters in the fused flat grid. Above this the grid would be too +/// Maximum number of physical counters in the flattened flat grid. Above this the grid would be too /// large/cache-unfriendly, so we fall back to the general buffered path. Count lanes are included /// in this limit: `num_terms × num_time_buckets × LANES` may not exceed it. /// /// Since we are only at the top-level, this won't be multiplied by any parent buckets. -const MAX_FUSED_GRID_COUNTERS: usize = 16_384; +const MAX_FLATTENED_GRID_COUNTERS: usize = 16_384; /// Scalar storage for grids whose term cardinality is too high for count lanes to pay off. const SINGLE_COUNT_LANE: usize = 1; @@ -39,9 +39,9 @@ const SINGLE_COUNT_LANE: usize = 1; const NUM_SMALL_LINEAR_BUCKETS: usize = 4; const NUM_LARGE_LINEAR_BUCKETS: usize = 8; -trait FusedBucketResolver: Debug + 'static { +trait BucketResolver: Debug + 'static { /// Fetches the histogram values needed for this block. Resolvers that do not inspect the - /// histogram column (notably [`FusedSingleBucketResolver`]) leave this as a no-op. + /// histogram column (notably [`FlattenedSingleBucketResolver`]) leave this as a no-op. fn prepare_block(&mut self, docs: &[crate::DocId]); /// Number of logical histogram buckets produced by this resolver. @@ -82,21 +82,21 @@ fn increment_grid_count( /// Resolver for a histogram whose entire value range maps to one bucket. It deliberately owns no /// block accessor: collecting this shape does not read or decode the histogram column at all. #[derive(Debug)] -struct FusedSingleBucketResolver { +struct SingleBucketResolver { next_count_lane: usize, } -impl FusedSingleBucketResolver { +impl SingleBucketResolver { fn new(hist_req_data: &HistogramAggReqData) -> Self { assert!( hist_req_data.accessor.get_cardinality().is_full(), - "FusedSingleBucketResolver requires a full histogram column" + "SingleBucketResolver requires a full histogram column" ); Self { next_count_lane: 0 } } } -impl FusedBucketResolver for FusedSingleBucketResolver { +impl BucketResolver for SingleBucketResolver { #[inline] fn prepare_block(&mut self, _docs: &[crate::DocId]) {} @@ -123,14 +123,14 @@ impl FusedBucketResolver for FusedSingleBucketResolver { _counts: &mut [[u32; LANES]], _term_counts: &mut [[u32; LANES]], ) { - unreachable!("FusedSingleBucketResolver is only constructed without hard bounds"); + unreachable!("SingleBucketResolver is only constructed without hard bounds"); } } /// The general resolver. It preserves the existing field conversion and floating-point bucket /// calculation for histograms that do not use a specialized resolver. #[derive(Debug)] -struct FusedComputedBucketResolver { +struct ComputedBucketResolver { hist_block: ColumnBlockAccessor, next_count_lane: usize, accessor: Column, @@ -142,11 +142,11 @@ struct FusedComputedBucketResolver { bounds: crate::aggregation::bucket::HistogramBounds, } -impl FusedComputedBucketResolver { +impl ComputedBucketResolver { fn new(hist_req_data: &HistogramAggReqData, base_pos: i64, num_buckets: usize) -> Self { assert!( hist_req_data.accessor.get_cardinality().is_full(), - "FusedComputedBucketResolver requires a full histogram column" + "ComputedBucketResolver requires a full histogram column" ); Self { hist_block: ColumnBlockAccessor::default(), @@ -162,7 +162,7 @@ impl FusedComputedBucketResolver { } } -impl FusedBucketResolver for FusedComputedBucketResolver { +impl BucketResolver for ComputedBucketResolver { #[inline] fn prepare_block(&mut self, docs: &[crate::DocId]) { self.hist_block @@ -224,7 +224,7 @@ impl FusedBucketResolver for FusedComputedBucketResolver { /// Resolver for a small histogram grid. Bucket starts are precomputed in monotonic fast-field /// `u64` space, then scanned linearly. `NUM_BUCKETS` is fixed so the optimizer can unroll the scan. #[derive(Debug)] -struct FusedLinearBucketResolver { +struct LinearBucketResolver { hist_block: ColumnBlockAccessor, next_count_lane: usize, accessor: Column, @@ -232,7 +232,7 @@ struct FusedLinearBucketResolver { num_buckets: usize, } -impl FusedLinearBucketResolver { +impl LinearBucketResolver { fn new( hist_req_data: &HistogramAggReqData, base_pos: i64, @@ -241,7 +241,7 @@ impl FusedLinearBucketResolver { assert!(num_time_buckets > 1 && num_time_buckets <= NUM_BUCKETS); assert!( hist_req_data.accessor.get_cardinality().is_full(), - "FusedLinearBucketResolver requires a full histogram column" + "LinearBucketResolver requires a full histogram column" ); let max_encoded_value = hist_req_data.accessor.max_value(); // Padding must compare false for every column value. There is no such `u64` sentinel when @@ -278,7 +278,7 @@ impl FusedLinearBucketResolver { } } -impl FusedBucketResolver for FusedLinearBucketResolver { +impl BucketResolver for LinearBucketResolver { #[inline] fn prepare_block(&mut self, docs: &[crate::DocId]) { self.hist_block @@ -312,8 +312,8 @@ impl FusedBucketResolver for FusedLinearBucketResolver _term_counts: &mut [[u32; LANES]], ) { panic!( - "FusedLinearBucketResolver does not support hard bounds and should not be constructed \ - with them" + "LinearBucketResolver does not support hard bounds and should not be constructed with \ + them" ); } } @@ -344,9 +344,9 @@ fn first_encoded_value_for_bucket( encoded_lower_bound } -/// Fused collector for `terms` (low cardinality) × a single `histogram`/`date_histogram` leaf with -/// nothing nested below it, when the resulting counter grid is small (see -/// [`MAX_FUSED_GRID_COUNTERS`]). +/// Flattened collector for `terms` (low cardinality) × a single `histogram`/`date_histogram` leaf +/// with nothing nested below it, when the resulting counter grid is small (see +/// [`MAX_FLATTENED_GRID_COUNTERS`]). /// /// It keeps a flat, fully dense 2D counter grid /// (`counts[term * num_time_buckets + bucket][lane]`) and a per-term total. Cycling writes through @@ -360,7 +360,7 @@ fn first_encoded_value_for_bucket( /// handed to the shared intermediate-result builders, so cross-segment merging is identical to the /// general path. #[derive(Debug)] -struct FusedTermHistogramCollector { +struct FlattenedTermHistogramCollector { /// Per-term count of docs *outside* `hard_bounds` (still in `doc_count`, but in no bucket). /// Per-term total = this + the term's `counts` row-sum; left empty when there are no hard /// bounds (every doc is in-bounds, so there's no remainder to track). @@ -385,8 +385,8 @@ struct FusedTermHistogramCollector { all_docs_in_bounds: bool, } -impl SegmentAggregationCollector - for FusedTermHistogramCollector +impl SegmentAggregationCollector + for FlattenedTermHistogramCollector { fn add_intermediate_aggregation_result( &mut self, @@ -396,7 +396,7 @@ impl SegmentAggregationCollector ) -> crate::Result<()> { debug_assert_eq!( parent_bucket_id, 0, - "fused term-histogram collector is top-level only" + "flattened term-histogram collector is top-level only" ); // Expand the flat grid back into the regular structures and reuse the shared builders, so // ordering/cut-off/dict handling and cross-segment merging match the general path exactly. @@ -449,11 +449,11 @@ impl SegmentAggregationCollector ) -> crate::Result<()> { debug_assert_eq!( parent_bucket_id, 0, - "fused term-histogram collector is top-level only" + "flattened term-histogram collector is top-level only" ); // The term column is always needed. The resolver fetches the histogram column only when - // bucket selection depends on its values; `FusedSingleBucketResolver` makes this a no-op. + // bucket selection depends on its values; `SingleBucketResolver` makes this a no-op. self.term_block .fetch_full_column_block(docs, &self.terms_req_data.accessor); self.bucket_resolver.prepare_block(docs); @@ -499,13 +499,13 @@ impl SegmentAggregationCollector } } -/// Builds the fused terms×histogram collector for a single top-level parent, when the shape is +/// Builds the flattened terms×histogram collector for a single top-level parent, when the shape is /// eligible. Returns `Ok(None)` to fall back to the general buffered terms path. /// /// Eligibility: top-level, low-cardinality terms over a full column with no missing/include-exclude /// handling; a single `histogram`/`date_histogram` leaf (no nesting below it) over a full column; -/// and a physical counter grid no larger than [`MAX_FUSED_GRID_COUNTERS`]. -pub(super) fn maybe_build_fused_collector( +/// and a physical counter grid no larger than [`MAX_FLATTENED_GRID_COUNTERS`]. +pub(super) fn maybe_build_flattened_collector( agg_data: &mut AggregationsSegmentCtx, node: &AggRefNode, terms_req_data: &TermsAggReqData, @@ -517,7 +517,7 @@ pub(super) fn maybe_build_fused_collector( // no-op (`fetch_block_with_missing` early-returns on full columns), so we needn't check for it. // // We don't cap the term cardinality here: the flat grid is bounded by the total physical - // counter count (`num_terms * num_time_buckets * LANES <= MAX_FUSED_GRID_COUNTERS`) checked + // counter count (`num_terms * num_time_buckets * LANES <= MAX_FLATTENED_GRID_COUNTERS`) checked // below, which subsumes it. // // We only allow this at the top-level, since we don't know how many buckets are created. We @@ -546,10 +546,10 @@ pub(super) fn maybe_build_fused_collector( return Ok(None); } - // Clone + normalize the histogram request and get its dense bucket range; only take the fused - // path when the physical counter grid is small enough. Very small logical grids use multiple - // counters per cell; larger grids retain scalar cells to avoid paying for lanes when writes are - // already spread across many locations. + // Clone + normalize the histogram request and get its dense bucket range; only take the + // flattened path when the physical counter grid is small enough. Very small logical grids use + // multiple counters per cell; larger grids retain scalar cells to avoid paying for lanes when + // writes are already spread across many locations. let Some((hist_req_data, range)) = prepare_histogram_dense_range(agg_data, &node.children[0])? else { return Ok(None); @@ -562,12 +562,12 @@ pub(super) fn maybe_build_fused_collector( } else { SINGLE_COUNT_LANE }; - if num_grid_cells.saturating_mul(num_count_lanes) > MAX_FUSED_GRID_COUNTERS { + if num_grid_cells.saturating_mul(num_count_lanes) > MAX_FLATTENED_GRID_COUNTERS { return Ok(None); } let collector = if use_count_lanes { - build_fused_collector::( + build_flattened_collector::( agg_data, terms_req_data, hist_req_data, @@ -576,7 +576,7 @@ pub(super) fn maybe_build_fused_collector( range.base_pos, )? } else { - build_fused_collector::( + build_flattened_collector::( agg_data, terms_req_data, hist_req_data, @@ -588,7 +588,7 @@ pub(super) fn maybe_build_fused_collector( Ok(Some(collector)) } -fn build_fused_collector( +fn build_flattened_collector( agg_data: &mut AggregationsSegmentCtx, terms_req_data: &TermsAggReqData, hist_req_data: HistogramAggReqData, @@ -596,13 +596,13 @@ fn build_fused_collector( num_time_buckets: usize, base_pos: i64, ) -> crate::Result> { - const { assert!(LANES > 0, "a fused grid needs at least one count lane") }; + const { assert!(LANES > 0, "a flattened grid needs at least one count lane") }; let all_docs_in_bounds = hist_req_data.bounds.min == f64::MIN && hist_req_data.bounds.max == f64::MAX; if all_docs_in_bounds && num_time_buckets == 1 { - let resolver = FusedSingleBucketResolver::new(&hist_req_data); - return build_fused_collector_with_resolver::( + let resolver = SingleBucketResolver::new(&hist_req_data); + return build_flattened_collector_with_resolver::( agg_data, terms_req_data, hist_req_data, @@ -612,12 +612,12 @@ fn build_fused_collector( ); } if all_docs_in_bounds && num_time_buckets <= NUM_SMALL_LINEAR_BUCKETS { - if let Some(resolver) = FusedLinearBucketResolver::::new( + if let Some(resolver) = LinearBucketResolver::::new( &hist_req_data, base_pos, num_time_buckets, ) { - return build_fused_collector_with_resolver::<_, LANES>( + return build_flattened_collector_with_resolver::<_, LANES>( agg_data, terms_req_data, hist_req_data, @@ -627,12 +627,12 @@ fn build_fused_collector( ); } } else if all_docs_in_bounds && num_time_buckets <= NUM_LARGE_LINEAR_BUCKETS { - if let Some(resolver) = FusedLinearBucketResolver::::new( + if let Some(resolver) = LinearBucketResolver::::new( &hist_req_data, base_pos, num_time_buckets, ) { - return build_fused_collector_with_resolver::<_, LANES>( + return build_flattened_collector_with_resolver::<_, LANES>( agg_data, terms_req_data, hist_req_data, @@ -643,8 +643,8 @@ fn build_fused_collector( } } - let resolver = FusedComputedBucketResolver::new(&hist_req_data, base_pos, num_time_buckets); - build_fused_collector_with_resolver::<_, LANES>( + let resolver = ComputedBucketResolver::new(&hist_req_data, base_pos, num_time_buckets); + build_flattened_collector_with_resolver::<_, LANES>( agg_data, terms_req_data, hist_req_data, @@ -654,7 +654,7 @@ fn build_fused_collector( ) } -fn build_fused_collector_with_resolver( +fn build_flattened_collector_with_resolver( agg_data: &mut AggregationsSegmentCtx, terms_req_data: &TermsAggReqData, hist_req_data: HistogramAggReqData, @@ -681,7 +681,7 @@ fn build_fused_collector_with_resolver { + Ok(Box::new(FlattenedTermHistogramCollector:: { term_counts, counts, base_pos, @@ -702,17 +702,17 @@ mod tests { }; use crate::aggregation::AggregationLimitsGuard; - /// Hand-computed correctness check for the fused terms×histogram fast path - /// ([`super::FusedTermHistogramCollector`]): low-cardinality terms × a histogram leaf over + /// Hand-computed correctness check for the flattened terms×histogram fast path + /// ([`super::FlattenedTermHistogramCollector`]): low-cardinality terms × a histogram leaf over /// full columns, exercised single- and multi-segment. #[test] - fn fused_term_histogram_test() -> crate::Result<()> { - fused_term_histogram_with_opt(false)?; - fused_term_histogram_with_opt(true)?; + fn flattened_term_histogram_test() -> crate::Result<()> { + flattened_term_histogram_with_opt(false)?; + flattened_term_histogram_with_opt(true)?; Ok(()) } - fn fused_term_histogram_with_opt(merge_segments: bool) -> crate::Result<()> { + fn flattened_term_histogram_with_opt(merge_segments: bool) -> crate::Result<()> { // 300 docs: term = {a, b, c} by i % 3, histogram value = i % 20 (interval 1 => buckets // 0..19). gcd(3, 20) = 1, so every (term, bucket) pair occurs exactly 300 / 60 = 5 times. let docs: Vec<(f64, String)> = (0..300u64) @@ -723,7 +723,8 @@ mod tests { ) }) .collect(); - // Two segments, to also exercise cross-segment merging of the fused per-term histograms. + // Two segments, to also exercise cross-segment merging of the flattened per-term + // histograms. let segments = vec![docs[..150].to_vec(), docs[150..].to_vec()]; let index = get_test_index_from_values_and_terms(merge_segments, &segments)?; @@ -757,7 +758,7 @@ mod tests { /// A histogram whose values all map to one bucket uses the resolver that counts terms without /// reading the histogram column. #[test] - fn fused_term_histogram_single_bucket_resolver() -> crate::Result<()> { + fn flattened_term_histogram_single_bucket_resolver() -> crate::Result<()> { // The ten distinct values all map to bucket 0 at interval 10. With three terms coprime to // the value count, each term occurs 30 times across the 90 documents. let docs: Vec<(f64, String)> = (0..90usize) @@ -791,7 +792,7 @@ mod tests { /// Four and eight histogram buckets take their corresponding fixed-size linear resolvers. /// Negative values also exercise monotonic `f64` fast-field boundaries on both sides of zero. #[test] - fn fused_term_histogram_linear_bucket_resolver() -> crate::Result<()> { + fn flattened_term_histogram_linear_bucket_resolver() -> crate::Result<()> { for num_buckets in [4usize, 8] { // Three terms are coprime with both bucket counts, so every pair occurs 10 times. let docs: Vec<(f64, String)> = (0..3 * num_buckets * 10) @@ -833,11 +834,11 @@ mod tests { Ok(()) } - /// A `missing` config on a *full* term column still takes the fused path (the string sentinel - /// is just `col_max + 1`, so the column stays low-cardinality). Since no doc is missing, the - /// real term buckets must be exactly as without `missing`. + /// A `missing` config on a *full* term column still takes the flattened path (the string + /// sentinel is just `col_max + 1`, so the column stays low-cardinality). Since no doc is + /// missing, the real term buckets must be exactly as without `missing`. #[test] - fn fused_term_histogram_with_missing_on_full_column() -> crate::Result<()> { + fn flattened_term_histogram_with_missing_on_full_column() -> crate::Result<()> { let docs: Vec<(f64, String)> = (0..300u64) .map(|i| { ( @@ -877,7 +878,7 @@ mod tests { /// Term cardinality is not what gates fusing: the flat grid is bounded by the total cell count /// (`num_terms * num_time_buckets`), not the term count, so many terms still fuse. #[test] - fn fused_term_histogram_many_terms() -> crate::Result<()> { + fn flattened_term_histogram_many_terms() -> crate::Result<()> { let num_terms = 150usize; let docs_per_term = 2usize; // All docs share histogram value 0 (a single bucket), so the grid is 150 x 1 = 150 cells. @@ -919,7 +920,7 @@ mod tests { /// is the case where the per-doc `term_counts` increment cannot be replaced by the grid /// row-sum. #[test] - fn fused_term_histogram_with_hard_bounds() -> crate::Result<()> { + fn flattened_term_histogram_with_hard_bounds() -> crate::Result<()> { // 300 docs: term = {a, b, c} by i % 3, value = i % 20. Per term: 100 docs, each value in // 0..=19 occurring 5 times. let docs: Vec<(f64, String)> = (0..300u64) @@ -975,7 +976,7 @@ mod tests { /// histogram row-sum. This is the case that previously fell back to the per-doc counter only /// because `bounds != [MIN, MAX]`. #[test] - fn fused_term_histogram_with_non_binding_hard_bounds() -> crate::Result<()> { + fn flattened_term_histogram_with_non_binding_hard_bounds() -> crate::Result<()> { // 300 docs: term = {a, b, c} by i % 3, value = i % 20. Data values span [0, 19]. let docs: Vec<(f64, String)> = (0..300u64) .map(|i| { @@ -1021,12 +1022,12 @@ mod tests { Ok(()) } - /// Regression: with hard bounds the fused path allocates `term_counts` (one `u32`/term) on top - /// of the grid, and that allocation must be charged to the memory limit. With many terms and a - /// single time bucket the two are equal in size, so a limit admitting the grid alone but not - /// grid + `term_counts` must fail. + /// Regression: with hard bounds the flattened path allocates `term_counts` (one `u32`/term) on + /// top of the grid, and that allocation must be charged to the memory limit. With many terms + /// and a single time bucket the two are equal in size, so a limit admitting the grid alone but + /// not grid + `term_counts` must fail. #[test] - fn fused_term_histogram_hard_bounds_charges_term_counts() -> crate::Result<()> { + fn flattened_term_histogram_hard_bounds_charges_term_counts() -> crate::Result<()> { // 16k distinct terms, one doc each; values alternate in/out of the single-bucket bounds // [5, 5] so the bounds bind and `term_counts` is allocated. num_terms=16000, // num_time_buckets=1 => `counts` and `term_counts` are ~64 KB each. diff --git a/src/aggregation/bucket/term_agg/mod.rs b/src/aggregation/bucket/term_agg/mod.rs index cd9df9755..f59588762 100644 --- a/src/aggregation/bucket/term_agg/mod.rs +++ b/src/aggregation/bucket/term_agg/mod.rs @@ -30,7 +30,7 @@ use crate::aggregation::{format_date, BucketId, Key}; use crate::error::DataCorruption; use crate::TantivyError; -mod term_histogram; +mod flattened_term_histogram; /// Contains all information required by the SegmentTermCollector to perform the /// terms aggregation on a segment. @@ -427,9 +427,10 @@ pub(crate) fn build_segment_term_collector( let max_column_val: u64 = col_max_value.max(terms_req_data.missing_value_for_accessor.unwrap_or(0u64)); - // Fused fast path: low-cardinality terms × a single `histogram`/`date_histogram` leaf over full - // columns with a small enough bucket grid. Anything else falls through to the general path. - if let Some(collector) = term_histogram::maybe_build_fused_collector( + // Flattened fast path: low-cardinality terms × a single `histogram`/`date_histogram` leaf over + // full columns with a small enough bucket grid. Anything else falls through to the general + // path. + if let Some(collector) = flattened_term_histogram::maybe_build_flattened_collector( req_data, node, &terms_req_data, diff --git a/src/aggregation/metric/stats.rs b/src/aggregation/metric/stats.rs index 8b9f4f09c..06a39f6d1 100644 --- a/src/aggregation/metric/stats.rs +++ b/src/aggregation/metric/stats.rs @@ -266,7 +266,7 @@ impl SegmentAggregationCollector return Err(TantivyError::InvalidArgument(format!( "Unsupported stats type for stats aggregation: {:?}", self.collecting_for - ))) + ))); } }; @@ -285,6 +285,21 @@ impl SegmentAggregationCollector docs: &[crate::DocId], agg_data: &mut AggregationsSegmentCtx, ) -> crate::Result<()> { + // Fast path: the caller hands us a single doc, which is the dominant case when a metric + // sub-agg sits under a high-cardinality bucket agg. Streaming straight from the column + // skips the block accessor's buffers entirely. + // Only valid without a missing value: `values_for_doc` yields nothing for a doc without a + // value, so the substitute would be silently dropped. + // TODO: remove once we fetch all values for all bucket ids in one go + if docs.len() == 1 && self.missing_u64.is_none() { + collect_stats::( + &mut self.buckets[parent_bucket_id as usize], + self.accessor.values_for_doc(docs[0]), + self.is_number_or_date_type, + )?; + + return Ok(()); + } agg_data.column_block_accessor.fetch_block_with_missing( docs, &self.accessor,