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
This commit is contained in:
Paul Masurel
2026-09-01 12:26:23 +02:00
parent 3552d1d769
commit 484f3f2d4c
6 changed files with 107 additions and 96 deletions
+4 -4
View File
@@ -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": {
-6
View File
@@ -55,12 +55,6 @@ impl BlockValueSource for Column<u64> {
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();
@@ -630,9 +630,9 @@ impl<B: BucketIdSlot> SegmentHistogramCollector<B> {
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<const LANES: usize>(
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)?;
@@ -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<const LANES: usize>(
/// 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<u64>,
@@ -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<const NUM_BUCKETS: usize> {
struct LinearBucketResolver<const NUM_BUCKETS: usize> {
hist_block: ColumnBlockAccessor,
next_count_lane: usize,
accessor: Column<u64>,
@@ -232,7 +232,7 @@ struct FusedLinearBucketResolver<const NUM_BUCKETS: usize> {
num_buckets: usize,
}
impl<const NUM_BUCKETS: usize> FusedLinearBucketResolver<NUM_BUCKETS> {
impl<const NUM_BUCKETS: usize> LinearBucketResolver<NUM_BUCKETS> {
fn new(
hist_req_data: &HistogramAggReqData,
base_pos: i64,
@@ -241,7 +241,7 @@ impl<const NUM_BUCKETS: usize> FusedLinearBucketResolver<NUM_BUCKETS> {
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<const NUM_BUCKETS: usize> FusedLinearBucketResolver<NUM_BUCKETS> {
}
}
impl<const NUM_BUCKETS: usize> FusedBucketResolver for FusedLinearBucketResolver<NUM_BUCKETS> {
impl<const NUM_BUCKETS: usize> BucketResolver for LinearBucketResolver<NUM_BUCKETS> {
#[inline]
fn prepare_block(&mut self, docs: &[crate::DocId]) {
self.hist_block
@@ -312,8 +312,8 @@ impl<const NUM_BUCKETS: usize> 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<R: FusedBucketResolver, const LANES: usize> {
struct FlattenedTermHistogramCollector<R: BucketResolver, const LANES: usize> {
/// 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<R: FusedBucketResolver, const LANES: usize> {
all_docs_in_bounds: bool,
}
impl<R: FusedBucketResolver, const LANES: usize> SegmentAggregationCollector
for FusedTermHistogramCollector<R, LANES>
impl<R: BucketResolver, const LANES: usize> SegmentAggregationCollector
for FlattenedTermHistogramCollector<R, LANES>
{
fn add_intermediate_aggregation_result(
&mut self,
@@ -396,7 +396,7 @@ impl<R: FusedBucketResolver, const LANES: usize> 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<R: FusedBucketResolver, const LANES: usize> 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<R: FusedBucketResolver, const LANES: usize> 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::<NUM_LOW_CARD_COUNT_LANES>(
build_flattened_collector::<NUM_LOW_CARD_COUNT_LANES>(
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::<SINGLE_COUNT_LANE>(
build_flattened_collector::<SINGLE_COUNT_LANE>(
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<const LANES: usize>(
fn build_flattened_collector<const LANES: usize>(
agg_data: &mut AggregationsSegmentCtx,
terms_req_data: &TermsAggReqData,
hist_req_data: HistogramAggReqData,
@@ -596,13 +596,13 @@ fn build_fused_collector<const LANES: usize>(
num_time_buckets: usize,
base_pos: i64,
) -> crate::Result<Box<dyn SegmentAggregationCollector>> {
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::<FusedSingleBucketResolver, LANES>(
let resolver = SingleBucketResolver::new(&hist_req_data);
return build_flattened_collector_with_resolver::<SingleBucketResolver, LANES>(
agg_data,
terms_req_data,
hist_req_data,
@@ -612,12 +612,12 @@ fn build_fused_collector<const LANES: usize>(
);
}
if all_docs_in_bounds && num_time_buckets <= NUM_SMALL_LINEAR_BUCKETS {
if let Some(resolver) = FusedLinearBucketResolver::<NUM_SMALL_LINEAR_BUCKETS>::new(
if let Some(resolver) = LinearBucketResolver::<NUM_SMALL_LINEAR_BUCKETS>::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<const LANES: usize>(
);
}
} else if all_docs_in_bounds && num_time_buckets <= NUM_LARGE_LINEAR_BUCKETS {
if let Some(resolver) = FusedLinearBucketResolver::<NUM_LARGE_LINEAR_BUCKETS>::new(
if let Some(resolver) = LinearBucketResolver::<NUM_LARGE_LINEAR_BUCKETS>::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<const LANES: usize>(
}
}
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<const LANES: usize>(
)
}
fn build_fused_collector_with_resolver<R: FusedBucketResolver, const LANES: usize>(
fn build_flattened_collector_with_resolver<R: BucketResolver, const LANES: usize>(
agg_data: &mut AggregationsSegmentCtx,
terms_req_data: &TermsAggReqData,
hist_req_data: HistogramAggReqData,
@@ -681,7 +681,7 @@ fn build_fused_collector_with_resolver<R: FusedBucketResolver, const LANES: usiz
.limits
.add_memory_consumed(memory_consumption as u64)?;
Ok(Box::new(FusedTermHistogramCollector::<R, LANES> {
Ok(Box::new(FlattenedTermHistogramCollector::<R, LANES> {
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.
+5 -4
View File
@@ -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,
+16 -1
View File
@@ -266,7 +266,7 @@ impl<const COLUMN_TYPE_ID: u8> SegmentAggregationCollector
return Err(TantivyError::InvalidArgument(format!(
"Unsupported stats type for stats aggregation: {:?}",
self.collecting_for
)))
)));
}
};
@@ -285,6 +285,21 @@ impl<const COLUMN_TYPE_ID: u8> 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::<COLUMN_TYPE_ID>(
&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,