move term req data to collector

This commit is contained in:
Pascal Seitz
2026-07-20 11:07:41 +02:00
committed by PSeitz
parent 5c1937a1ec
commit 59b6859c16
2 changed files with 154 additions and 69 deletions
+152 -65
View File
@@ -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<FxHashMap<MultiTermsKey, MultiTermsBucket>>,
sub_agg: Option<BufferedSubAggs<HighCardSubAggBuffer>>,
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<Self> {
fn from_validated(
req: &mut AggregationsSegmentCtx,
node: &AggRefNode,
req_data: MultiTermsAggReqData,
) -> crate::Result<Self> {
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::<BucketId>::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::<BucketId>::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<PagedTermMap, HighCardSubAggBuffer> =
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<HashMapTermBuckets, HighCardSubAggBuffer> =
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::<BucketId, true>::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::<BucketId, false>::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::<BucketId>::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::<BucketId>::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<M: TermAggregationMap>(
term_buckets: M,
sub_agg: Option<HighCardBufferedSubAggs>,
bucket_id_provider: BucketIdProvider,
req_data: MultiTermsAggReqData,
packs: Vec<FieldPack>,
max_packed: u64,
num_fields: usize,
) -> Box<dyn SegmentAggregationCollector> {
Box::new(SegmentMultiTermsFastCollector::<M, HighCardSubAggBuffer> {
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<TermMap: TermAggregationMap, B: SubAggBuff
parent_buckets: Vec<TermMap>,
sub_agg: Option<BufferedSubAggs<B>>,
bucket_id_provider: BucketIdProvider,
req_data_idx: usize,
req_data: MultiTermsAggReqData,
packs: Vec<FieldPack>,
/// Upper bound on the packed key space; used to initialize each parent bucket's
/// `TermMap`.
@@ -898,7 +979,7 @@ impl<TermMap: TermAggregationMap, B: SubAggBuffer> 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<TermMap: TermAggregationMap, B: SubAggBuffer> 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();
+2 -4
View File
@@ -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<M: TermAggregationMap>(
pub(crate) fn add_memory_consumption<M: TermAggregationMap>(
term_buckets: &M,
req_data: &mut AggregationsSegmentCtx,
) -> crate::Result<()> {
@@ -912,8 +912,6 @@ pub(crate) struct VecTermBuckets<
buckets: Vec<Bucket<B>>,
}
/// Dense term storage without sub-aggregation bucket ids.
pub(crate) type VecTermBucketsNoAgg = VecTermBuckets<()>;
impl<B: BucketIdSlot, const LAZY_BUCKET_ID_GENERATION: bool> TermAggregationMap
for VecTermBuckets<B, LAZY_BUCKET_ID_GENERATION>