mirror of
https://github.com/quickwit-oss/tantivy.git
synced 2026-10-07 04:12:42 +00:00
Merge incoming aggregation buckets into the accumulator
Probe the accumulated map only for incoming keys rather than rehashing every accumulated key on each merge. This removes quadratic work when folding many disjoint multi-terms results. The same helper serves range and composite buckets; existing bucket values still merge left-to-right and the wire representation and pruning rules are unchanged. Cover overlapping/disjoint tuple keys, empty inputs, recursive range subaggregations, postcard round trips, error/count bookkeeping, and merging after pruning. Same benchmark and configuration as 0062d2f0d (Apple M4 Max, rustc 1.98.0): Median milliseconds for 10 / 100 / 1000 inputs: shared 8: 0.006365 / 0.0656 / 0.6449 disjoint 16: 0.0126 / 0.1486 / 1.6882 disjoint 160: 0.1336 / 1.5558 / 20.4019 At 1000 inputs this is 52.4x and 61.6x faster for the disjoint cases, with unchanged measured peak allocation. Shared-key control remains within approximately 2% of baseline. Validation: 308 aggregation tests pass with default features and 308 with quickwit; changed-file rustfmt and git diff checks pass. Clippy --lib --bench agg_bench passes with the pre-existing clippy::drop_non_drop warning allowed (unchanged drop(add_document) in the benchmark).
This commit is contained in:
@@ -1213,19 +1213,20 @@ trait MergeFruits {
|
||||
fn merge_fruits(&mut self, other: Self) -> crate::Result<()>;
|
||||
}
|
||||
|
||||
fn merge_maps<V: MergeFruits + Clone, T: Eq + PartialEq + Hash>(
|
||||
fn merge_maps<V: MergeFruits, T: Eq + Hash>(
|
||||
entries_left: &mut FxHashMap<T, V>,
|
||||
mut entries_right: FxHashMap<T, V>,
|
||||
entries_right: FxHashMap<T, V>,
|
||||
) -> crate::Result<()> {
|
||||
for (name, entry_left) in entries_left.iter_mut() {
|
||||
if let Some(entry_right) = entries_right.remove(name) {
|
||||
entry_left.merge_fruits(entry_right)?;
|
||||
// Visit incoming entries, not the growing accumulator, so folding many results
|
||||
// does not repeatedly hash all previously merged keys.
|
||||
for (key, entry_right) in entries_right {
|
||||
match entries_left.entry(key) {
|
||||
Entry::Occupied(mut entry) => entry.get_mut().merge_fruits(entry_right)?,
|
||||
Entry::Vacant(entry) => {
|
||||
entry.insert(entry_right);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for (key, res) in entries_right.into_iter() {
|
||||
entries_left.entry(key).or_insert(res);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -1698,6 +1699,106 @@ mod tests {
|
||||
assert_range_trees_eq(&tree_left, &tree_expected);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_multi_terms_repeated_merge_and_prune() {
|
||||
let key = |id: u64| {
|
||||
vec![
|
||||
IntermediateKey::Str(format!("host-{}", id % 2)),
|
||||
IntermediateKey::U64(id / 2),
|
||||
]
|
||||
};
|
||||
let make_result =
|
||||
|data: &[(u64, u64)], source: &str| IntermediateBucketResult::MultiTerms {
|
||||
buckets: IntermediateMultiTermsBucketResult {
|
||||
entries: data
|
||||
.iter()
|
||||
.map(|&(id, count)| {
|
||||
(
|
||||
key(id),
|
||||
IntermediateTermBucketEntry {
|
||||
doc_count: count,
|
||||
sub_aggregation: get_sub_test_tree(&[
|
||||
("shared".to_string(), count),
|
||||
(source.to_string(), count),
|
||||
]),
|
||||
},
|
||||
)
|
||||
})
|
||||
.collect(),
|
||||
sum_other_doc_count: 2,
|
||||
doc_count_error_upper_bound: 1,
|
||||
},
|
||||
};
|
||||
let mut merged = IntermediateBucketResult::MultiTerms {
|
||||
buckets: Default::default(),
|
||||
};
|
||||
merged
|
||||
.merge_fruits(make_result(&[(0, 3), (1, 5)], "first"))
|
||||
.unwrap();
|
||||
merged
|
||||
.merge_fruits(make_result(&[(1, 7), (2, 11)], "second"))
|
||||
.unwrap();
|
||||
// A distributed fold may resume after serializing an intermediate result.
|
||||
merged = postcard::from_bytes(&postcard::to_allocvec(&merged).unwrap()).unwrap();
|
||||
merged.merge_fruits(make_result(&[], "empty")).unwrap();
|
||||
merged
|
||||
.merge_fruits(make_result(&[(0, 13), (3, 17)], "third"))
|
||||
.unwrap();
|
||||
|
||||
let IntermediateBucketResult::MultiTerms { buckets } = &mut merged else {
|
||||
panic!("expected multi_terms");
|
||||
};
|
||||
assert_eq!(buckets.entries.len(), 4);
|
||||
assert_eq!(buckets.sum_other_doc_count, 8);
|
||||
assert_eq!(buckets.doc_count_error_upper_bound, 4);
|
||||
for (id, count, sources) in [
|
||||
(0, 16, vec![("first", 3), ("third", 13)]),
|
||||
(1, 12, vec![("first", 5), ("second", 7)]),
|
||||
(2, 11, vec![("second", 11)]),
|
||||
(3, 17, vec![("third", 17)]),
|
||||
] {
|
||||
let entry = &buckets.entries[&key(id)];
|
||||
assert_eq!(entry.doc_count, count);
|
||||
let mut expected = vec![("shared".to_string(), count)];
|
||||
expected.extend(
|
||||
sources
|
||||
.into_iter()
|
||||
.map(|(name, count)| (name.to_string(), count)),
|
||||
);
|
||||
assert_range_trees_eq(&entry.sub_aggregation, &get_sub_test_tree(&expected));
|
||||
}
|
||||
|
||||
let req = serde_json::from_value(serde_json::json!({
|
||||
"terms": [{"field": "host"}, {"field": "path"}],
|
||||
"size": 1,
|
||||
"segment_size": 2
|
||||
}))
|
||||
.unwrap();
|
||||
buckets
|
||||
.prune_intermediate_results(&req, &Default::default(), PruneMode::Intermediate)
|
||||
.unwrap();
|
||||
assert_eq!(buckets.entries.len(), 2);
|
||||
assert_eq!(buckets.sum_other_doc_count, 31); // 8 + 12 + 11
|
||||
assert_eq!(buckets.doc_count_error_upper_bound, 16); // 4 + cutoff 12
|
||||
|
||||
// A previously pruned key can be inserted again in a subsequent fold.
|
||||
merged
|
||||
.merge_fruits(make_result(&[(2, 19)], "fourth"))
|
||||
.unwrap();
|
||||
let IntermediateBucketResult::MultiTerms { buckets } = &mut merged else {
|
||||
unreachable!();
|
||||
};
|
||||
assert_eq!(buckets.entries.len(), 3);
|
||||
assert_eq!(buckets.entries[&key(2)].doc_count, 19);
|
||||
buckets
|
||||
.prune_intermediate_results(&req, &Default::default(), PruneMode::Final)
|
||||
.unwrap();
|
||||
assert_eq!(buckets.entries.len(), 1);
|
||||
assert_eq!(buckets.entries[&key(2)].doc_count, 19);
|
||||
assert_eq!(buckets.sum_other_doc_count, 66); // 31 + 2 + 16 + 17
|
||||
assert_eq!(buckets.doc_count_error_upper_bound, 17); // 16 + 1, no final cutoff
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_prune_intermediate_results_finalizer_size() {
|
||||
use crate::aggregation::bucket::TermsAggregation;
|
||||
|
||||
Reference in New Issue
Block a user