From 59b6859c16a7ab3012f58867fe76863b63701d35 Mon Sep 17 00:00:00 2001 From: Pascal Seitz Date: Mon, 20 Jul 2026 11:07:41 +0200 Subject: [PATCH] move term req data to collector --- src/aggregation/bucket/multi_terms/mod.rs | 217 +++++++++++++++------- src/aggregation/bucket/term_agg/mod.rs | 6 +- 2 files changed, 154 insertions(+), 69 deletions(-) diff --git a/src/aggregation/bucket/multi_terms/mod.rs b/src/aggregation/bucket/multi_terms/mod.rs index c2a4457e7..11cddb466 100644 --- a/src/aggregation/bucket/multi_terms/mod.rs +++ b/src/aggregation/bucket/multi_terms/mod.rs @@ -16,13 +16,14 @@ use crate::aggregation::agg_data::{ }; use crate::aggregation::agg_req::Aggregations; use crate::aggregation::bucket::term_agg::{ - cut_off_buckets, get_agg_name_and_property, GetDocCount, HashMapTermBuckets, PagedTermMap, - TermAggregationMap, VecTermBuckets, VecTermBucketsNoAgg, + add_memory_consumption, cut_off_buckets, get_agg_name_and_property, GetDocCount, + HashMapTermBuckets, PagedTermMap, TermAggregationMap, VecTermBuckets, + LAZY_BUCKET_ID_GENERATION_THRESHOLD, MAX_NUM_TERMS_FOR_LOWCARD_SUBAGG, MAX_NUM_TERMS_FOR_PAGED_MAP, MAX_NUM_TERMS_FOR_VEC, }; use crate::aggregation::buffered_sub_aggs::{ - BufferedSubAggs, HighCardSubAggBuffer, LowCardBufferedSubAggs, LowCardSubAggBuffer, - SubAggBuffer, + BufferedSubAggs, HighCardBufferedSubAggs, HighCardSubAggBuffer, LowCardBufferedSubAggs, + LowCardSubAggBuffer, SubAggBuffer, }; use crate::aggregation::intermediate_agg_result::{ IntermediateAggregationResult, IntermediateAggregationResults, IntermediateBucketResult, @@ -256,7 +257,7 @@ pub struct SegmentMultiTermsCollector { parent_buckets: Vec>, sub_agg: Option>, bucket_id_provider: BucketIdProvider, - req_data_idx: usize, + req_data: MultiTermsAggReqData, } /// Shared validation for both the fast and recursive-fallback collectors: checks the @@ -314,14 +315,18 @@ pub(crate) fn build_segment_multi_terms_collector( } Ok(Box::new(SegmentMultiTermsCollector::from_validated( - req, node, + req, node, req_data, )?)) } impl SegmentMultiTermsCollector { /// Constructs the recursive fallback collector. Assumes `validate_multi_terms` has /// already been run by [`build_segment_multi_terms_collector`]. - fn from_validated(req: &mut AggregationsSegmentCtx, node: &AggRefNode) -> crate::Result { + fn from_validated( + req: &mut AggregationsSegmentCtx, + node: &AggRefNode, + req_data: MultiTermsAggReqData, + ) -> crate::Result { let sub_agg = if !node.children.is_empty() { let sub_collector = build_segment_agg_collectors(req, &node.children)?; Some(BufferedSubAggs::new(sub_collector)) @@ -333,7 +338,7 @@ impl SegmentMultiTermsCollector { parent_buckets: vec![FxHashMap::default()], sub_agg, bucket_id_provider: BucketIdProvider::default(), - req_data_idx: node.idx_in_req_data, + req_data, }) } @@ -549,7 +554,7 @@ impl SegmentAggregationCollector for SegmentMultiTermsCollector { ) -> crate::Result<()> { self.prepare_max_bucket(bucket, agg_data)?; let bucket_map = std::mem::take(&mut self.parent_buckets[bucket as usize]); - let req_data = &agg_data.per_request.multi_terms_req_data[self.req_data_idx]; + let req_data = &self.req_data; let name = req_data.name.clone(); let result = Self::into_intermediate_bucket_result_inner( req_data, @@ -570,7 +575,7 @@ impl SegmentAggregationCollector for SegmentMultiTermsCollector { docs: &[crate::DocId], agg_data: &mut AggregationsSegmentCtx, ) -> crate::Result<()> { - let req_data = &agg_data.per_request.multi_terms_req_data[self.req_data_idx]; + let req_data = &self.req_data; let num_fields = req_data.fields.len(); let mem_pre = self.get_memory_consumption(parent_bucket_id, num_fields); @@ -785,67 +790,143 @@ fn build_fast_multi_terms_collector( }; let mut bucket_id_provider = BucketIdProvider::default(); - if is_top_level && max_packed < MAX_NUM_TERMS_FOR_VEC && !has_sub_aggregations { - let term_buckets = VecTermBucketsNoAgg::new(max_packed + 1, &mut bucket_id_provider); - let collector: SegmentMultiTermsFastCollector<_, HighCardSubAggBuffer> = - SegmentMultiTermsFastCollector { - parent_buckets: vec![term_buckets], - sub_agg: None, - bucket_id_provider, - req_data_idx: node.idx_in_req_data, - packs, - max_packed, - field_block_accessors: vec![ColumnBlockAccessor::default(); num_fields], - packed_keys_buf: Vec::new(), - }; - Ok(Box::new(collector)) - } else if is_top_level && max_packed < MAX_NUM_TERMS_FOR_VEC { - let term_buckets = - VecTermBuckets::::new(max_packed + 1, &mut bucket_id_provider); - let sub_agg = sub_agg_collector.map(LowCardBufferedSubAggs::new); + let num_terms = max_packed.saturating_add(1); + + // Keep low-cardinality sub-aggregations on the per-bucket document buffer. At higher + // cardinalities, mirror the terms collector's dense/paged/hash storage selection while using + // the partitioned high-cardinality sub-aggregation buffer. + if is_top_level + && has_sub_aggregations + && max_packed < MAX_NUM_TERMS_FOR_LOWCARD_SUBAGG + { + let term_buckets = VecTermBuckets::::new(num_terms, &mut bucket_id_provider); + add_memory_consumption(&term_buckets, req)?; let collector: SegmentMultiTermsFastCollector<_, LowCardSubAggBuffer> = SegmentMultiTermsFastCollector { parent_buckets: vec![term_buckets], - sub_agg, + sub_agg: sub_agg_collector.map(LowCardBufferedSubAggs::new), bucket_id_provider, - req_data_idx: node.idx_in_req_data, + req_data, packs, max_packed, field_block_accessors: vec![ColumnBlockAccessor::default(); num_fields], packed_keys_buf: Vec::new(), }; - Ok(Box::new(collector)) - } else if max_packed < MAX_NUM_TERMS_FOR_PAGED_MAP && is_top_level { - let term_buckets = PagedTermMap::new(max_packed + 1, &mut bucket_id_provider); - let sub_agg = sub_agg_collector.map(BufferedSubAggs::new); - let collector: SegmentMultiTermsFastCollector = - SegmentMultiTermsFastCollector { - parent_buckets: vec![term_buckets], - sub_agg, - bucket_id_provider, - req_data_idx: node.idx_in_req_data, - packs, - max_packed, - field_block_accessors: vec![ColumnBlockAccessor::default(); num_fields], - packed_keys_buf: Vec::new(), - }; - Ok(Box::new(collector)) - } else { - let term_buckets = HashMapTermBuckets::default(); - let sub_agg = sub_agg_collector.map(BufferedSubAggs::new); - let collector: SegmentMultiTermsFastCollector = - SegmentMultiTermsFastCollector { - parent_buckets: vec![term_buckets], - sub_agg, - bucket_id_provider, - req_data_idx: node.idx_in_req_data, - packs, - max_packed, - field_block_accessors: vec![ColumnBlockAccessor::default(); num_fields], - packed_keys_buf: Vec::new(), - }; - Ok(Box::new(collector)) + return Ok(Box::new(collector)); } + + let sub_agg = sub_agg_collector.map(BufferedSubAggs::new); + if is_top_level && max_packed < MAX_NUM_TERMS_FOR_VEC { + if has_sub_aggregations { + if num_terms > LAZY_BUCKET_ID_GENERATION_THRESHOLD { + let term_buckets = + VecTermBuckets::::new(num_terms, &mut bucket_id_provider); + add_memory_consumption(&term_buckets, req)?; + Ok(boxed_high_card_multi_terms_collector( + term_buckets, + sub_agg, + bucket_id_provider, + req_data, + packs, + max_packed, + num_fields, + )) + } else { + let term_buckets = + VecTermBuckets::::new(num_terms, &mut bucket_id_provider); + add_memory_consumption(&term_buckets, req)?; + Ok(boxed_high_card_multi_terms_collector( + term_buckets, + sub_agg, + bucket_id_provider, + req_data, + packs, + max_packed, + num_fields, + )) + } + } else { + let term_buckets = VecTermBuckets::<()>::new(num_terms, &mut bucket_id_provider); + add_memory_consumption(&term_buckets, req)?; + Ok(boxed_high_card_multi_terms_collector( + term_buckets, + sub_agg, + bucket_id_provider, + req_data, + packs, + max_packed, + num_fields, + )) + } + } else if is_top_level && max_packed < MAX_NUM_TERMS_FOR_PAGED_MAP { + if has_sub_aggregations { + let term_buckets = PagedTermMap::::new(num_terms, &mut bucket_id_provider); + Ok(boxed_high_card_multi_terms_collector( + term_buckets, + sub_agg, + bucket_id_provider, + req_data, + packs, + max_packed, + num_fields, + )) + } else { + let term_buckets = PagedTermMap::<()>::new(num_terms, &mut bucket_id_provider); + Ok(boxed_high_card_multi_terms_collector( + term_buckets, + sub_agg, + bucket_id_provider, + req_data, + packs, + max_packed, + num_fields, + )) + } + } else if has_sub_aggregations { + let term_buckets = HashMapTermBuckets::::default(); + Ok(boxed_high_card_multi_terms_collector( + term_buckets, + sub_agg, + bucket_id_provider, + req_data, + packs, + max_packed, + num_fields, + )) + } else { + let term_buckets = HashMapTermBuckets::<()>::default(); + Ok(boxed_high_card_multi_terms_collector( + term_buckets, + sub_agg, + bucket_id_provider, + req_data, + packs, + max_packed, + num_fields, + )) + } +} + +/// Boxes a fast-path collector backed by the high-cardinality sub-aggregation buffer. +fn boxed_high_card_multi_terms_collector( + term_buckets: M, + sub_agg: Option, + bucket_id_provider: BucketIdProvider, + req_data: MultiTermsAggReqData, + packs: Vec, + max_packed: u64, + num_fields: usize, +) -> Box { + Box::new(SegmentMultiTermsFastCollector:: { + parent_buckets: vec![term_buckets], + sub_agg, + bucket_id_provider, + req_data, + packs, + max_packed, + field_block_accessors: vec![ColumnBlockAccessor::default(); num_fields], + packed_keys_buf: Vec::new(), + }) } /// Fast-path segment collector for [`MultiTermsAggregation`]: packs each field's raw @@ -862,7 +943,7 @@ struct SegmentMultiTermsFastCollector, sub_agg: Option>, bucket_id_provider: BucketIdProvider, - req_data_idx: usize, + req_data: MultiTermsAggReqData, packs: Vec, /// Upper bound on the packed key space; used to initialize each parent bucket's /// `TermMap`. @@ -898,7 +979,7 @@ impl SegmentAggregationCollector &mut self.parent_buckets[bucket as usize], TermMap::new(0, &mut self.bucket_id_provider), ); - let req_data = &agg_data.per_request.multi_terms_req_data[self.req_data_idx]; + let req_data = &self.req_data; let name = req_data.name.clone(); let result = into_intermediate_bucket_result_fast( req_data, @@ -922,7 +1003,7 @@ impl SegmentAggregationCollector ) -> crate::Result<()> { let mem_pre = self.get_memory_consumption(parent_bucket_id); - let req_data = &agg_data.per_request.multi_terms_req_data[self.req_data_idx]; + let req_data = &self.req_data; self.packed_keys_buf.clear(); self.packed_keys_buf.resize(docs.len(), 0u64); @@ -1671,7 +1752,13 @@ mod tests { parent_buckets: vec![FxHashMap::default()], sub_agg: None, bucket_id_provider: BucketIdProvider::default(), - req_data_idx: 0, + req_data: MultiTermsAggReqData { + name: String::new(), + req: MultiTermsAggregation::default(), + fields: Vec::new(), + sub_aggregations: Aggregations::default(), + is_top_level: true, + }, }; for i in 0..64u64 { let mut key: MultiTermsKey = SmallVec::new(); diff --git a/src/aggregation/bucket/term_agg/mod.rs b/src/aggregation/bucket/term_agg/mod.rs index f8f4ea945..d81752608 100644 --- a/src/aggregation/bucket/term_agg/mod.rs +++ b/src/aggregation/bucket/term_agg/mod.rs @@ -358,7 +358,7 @@ pub const MAX_NUM_TERMS_FOR_PAGED_MAP: u64 = 8_000_000; /// terms. /// /// TODO: Benchmark this threshold. -const LAZY_BUCKET_ID_GENERATION_THRESHOLD: u64 = 4_096; +pub(crate) const LAZY_BUCKET_ID_GENERATION_THRESHOLD: u64 = 4_096; /// Average docs-per-bucket below which term counts cluster too tightly (mostly 1s and 2s) for /// `select_nth_unstable` to beat `sort_unstable`'s adaptive paths, so we fall back to a full sort. @@ -371,7 +371,7 @@ const DOCS_PER_BUCKET_QUICKSELECT_THRESHOLD: u64 = 2; /// the dense term maps preallocate all `num_terms` slots at construction, so their memory would /// otherwise escape the limit entirely (relevant now that the `Vec` threshold reaches /// [`MAX_NUM_TERMS_FOR_VEC`]). -fn add_memory_consumption( +pub(crate) fn add_memory_consumption( term_buckets: &M, req_data: &mut AggregationsSegmentCtx, ) -> crate::Result<()> { @@ -912,8 +912,6 @@ pub(crate) struct VecTermBuckets< buckets: Vec>, } -/// Dense term storage without sub-aggregation bucket ids. -pub(crate) type VecTermBucketsNoAgg = VecTermBuckets<()>; impl TermAggregationMap for VecTermBuckets