diff --git a/src/aggregation/bucket/multi_terms/mod.rs b/src/aggregation/bucket/multi_terms/mod.rs index 9ff6ebfcf..3da8c58df 100644 --- a/src/aggregation/bucket/multi_terms/mod.rs +++ b/src/aggregation/bucket/multi_terms/mod.rs @@ -4,8 +4,8 @@ use std::sync::Arc; use columnar::column_values::CompactSpaceU64Accessor; use columnar::{ - Column, ColumnType, Dictionary, MonotonicallyMappableToU128, MonotonicallyMappableToU64, - NumericalValue, StrColumn, + Column, ColumnBlockAccessor, ColumnType, Dictionary, MonotonicallyMappableToU128, + MonotonicallyMappableToU64, NumericalValue, StrColumn, }; use rustc_hash::FxHashMap; use serde::{Deserialize, Serialize}; @@ -245,6 +245,26 @@ impl AggregationMapKey for MultiTermsKey { /// One field's candidate `KeyElem`s for the current document. type FieldValues = SmallVec<[KeyElem; 2]>; +/// Fetches one field into the shared block accessor and returns whether its values are dense over +/// `docs`. +#[inline(always)] +fn fetch_field_block( + docs: &[crate::DocId], + field: &MultiTermsFieldAccessor, + block_accessor: &mut ColumnBlockAccessor, +) -> bool { + let cardinality = field.column.get_cardinality(); + let is_full = cardinality.is_full(); + + if cardinality.is_multivalue() { + block_accessor.fetch_block_with_missing_unique_per_doc(docs, &field.column, None); + } else { + block_accessor.fetch_block_with_is_full(docs, &field.column, is_full); + } + + is_full +} + /// Key operations used by the unified collector. /// /// Compact layouts use a packed `u64`; layouts wider than 64 bits use [`MultiTermsKey`]. @@ -257,12 +277,10 @@ trait MultiTermsKeyCodec: Clone + Debug + 'static { fn push(&self, key: &mut Self::Key, field_idx: usize, elem: KeyElem); fn pop(&self, key: &mut Self::Key, field_idx: usize); + /// Pushes one full field's values into every key buffer. Entries belonging to documents + /// already marked invalid may have an incomplete prefix; their contents are discarded later. fn push_full_values(&self, keys: &mut [Self::Key], field_idx: usize, values: I) - where I: IntoIterator { - for (key, val) in keys.iter_mut().zip(values) { - self.push(key, field_idx, KeyElem::new(val)); - } - } + where I: IntoIterator; fn resolve_key( &self, @@ -297,6 +315,14 @@ impl MultiTermsKeyCodec for UnpackedKeyCodec { key.pop(); } + fn push_full_values(&self, keys: &mut [Self::Key], field_idx: usize, values: I) + where I: IntoIterator { + for (key, val) in keys.iter_mut().zip(values) { + debug_assert!(key.len() <= field_idx); + key.push(KeyElem::new(val)); + } + } + fn resolve_key( &self, key: &Self::Key, @@ -348,7 +374,7 @@ where /// materializing per-field value vectors. single_value_keys_buf: Vec, /// Whether each document in the current single-valued block can produce a complete key. - /// `Vec` keeps this as a reusable bitset indexed like `single_value_keys_buf`. + /// This reusable mask is indexed like `single_value_keys_buf`. valid_docs_buf: Vec, /// Reused only when at least one column is multivalued, laid out as `[doc][field]`. field_values_buf: Vec, @@ -639,8 +665,8 @@ where self.single_value_keys_buf .resize_with(docs.len(), || key_codec.new_key()); } - let keys = &mut self.single_value_keys_buf[..docs.len()]; - for key in keys.iter_mut() { + let keys_buf = &mut self.single_value_keys_buf[..docs.len()]; + for key in keys_buf.iter_mut() { key_codec.clear_key(key); } @@ -649,62 +675,49 @@ where // Bulk-decode each full column and build one key per document directly, without // per-document column dispatch or Cartesian recursion. for (field_idx, field) in self.req_data.fields.iter().enumerate() { - block_accessor.fetch_block_with_is_full(docs, &field.column, true); - key_codec.push_full_values(keys, field_idx, block_accessor.iter_vals()); + fetch_field_block(docs, field, block_accessor); + key_codec.push_full_values(keys_buf, field_idx, block_accessor.iter_vals()); } if let Some(sub_agg) = sub_agg { - for (&doc_id, key) in docs.iter().zip(keys.iter()) { + for (&doc_id, key) in docs.iter().zip(keys_buf.iter()) { let bucket_id = buckets.term_entry(key.clone(), bucket_id_provider); sub_agg.push(bucket_id.to_bucket_id(), doc_id); } } else { - for key in keys.iter() { + for key in keys_buf.iter() { buckets.term_entry(key.clone(), bucket_id_provider); } } } else { // Optional columns still have at most one value per document. Decode every field // against the original block and track which document positions can produce a - // complete key. Invalid documents are skipped for later key writes and when the - // bucket map is populated. + // complete key. Full columns may bulk-write throwaway values into an already + // invalid key slot; only valid positions are inserted into the bucket map. self.valid_docs_buf.resize(docs.len(), true); self.valid_docs_buf.fill(true); let valid_docs = &mut self.valid_docs_buf; - let mut num_valid_docs = docs.len(); for (field_idx, field) in self.req_data.fields.iter().enumerate() { - if num_valid_docs == 0 { + if !valid_docs.contains(&true) { break; } - if field.column.get_cardinality().is_full() { - block_accessor.fetch_block_with_is_full(docs, &field.column, true); - if num_valid_docs == docs.len() { - key_codec.push_full_values(keys, field_idx, block_accessor.iter_vals()); - } else { - for (doc_idx, val) in block_accessor.iter_vals().enumerate() { - if valid_docs[doc_idx] { - key_codec.push( - &mut keys[doc_idx], - field_idx, - KeyElem::new(val), - ); - } - } - } + let is_full = fetch_field_block(docs, field, block_accessor); + if is_full { + key_codec.push_full_values(keys_buf, field_idx, block_accessor.iter_vals()); continue; } - debug_assert!(field.column.get_cardinality().is_optional()); - block_accessor.fetch_block_with_is_full(docs, &field.column, false); let missing = self.req_data.missing_accessors[field_idx].as_ref(); - let mut hits = block_accessor + let mut docids_and_vals = block_accessor .iter_docid_vals(docs, &field.column) .peekable(); for (doc_idx, &doc_id) in docs.iter().enumerate() { - debug_assert!(hits.peek().is_none_or(|(hit_doc, _)| *hit_doc >= doc_id)); - let column_elem = hits + debug_assert!(docids_and_vals + .peek() + .is_none_or(|(hit_doc, _)| *hit_doc >= doc_id)); + let column_elem = docids_and_vals .next_if(|(hit_doc, _)| *hit_doc == doc_id) .map(|(_, val)| KeyElem::new(val)); @@ -715,24 +728,23 @@ where if let Some(elem) = column_elem.or_else(|| missing_key_elem_for_doc(missing, doc_id)) { - key_codec.push(&mut keys[doc_idx], field_idx, elem); + key_codec.push(&mut keys_buf[doc_idx], field_idx, elem); } else { valid_docs[doc_idx] = false; - num_valid_docs -= 1; } } - debug_assert!(hits.next().is_none()); + debug_assert!(docids_and_vals.next().is_none()); } if let Some(sub_agg) = sub_agg { - for (doc_idx, (&doc_id, key)) in docs.iter().zip(keys.iter()).enumerate() { + for (doc_idx, (&doc_id, key)) in docs.iter().zip(keys_buf.iter()).enumerate() { if valid_docs[doc_idx] { let bucket_id = buckets.term_entry(key.clone(), bucket_id_provider); sub_agg.push(bucket_id.to_bucket_id(), doc_id); } } } else { - for (doc_idx, key) in keys.iter().enumerate() { + for (doc_idx, key) in keys_buf.iter().enumerate() { if valid_docs[doc_idx] { buckets.term_entry(key.clone(), bucket_id_provider); } @@ -759,17 +771,12 @@ where // `(doc_id, value)` pairs and deduplicate repeated values within one document. let block_accessor = &mut agg_data.column_block_accessor; for (field_idx, field) in self.req_data.fields.iter().enumerate() { - if field.column.get_cardinality().is_full() { - block_accessor.fetch_block_with_is_full(docs, &field.column, true); + let is_full = fetch_field_block(docs, field, block_accessor); + if is_full { for (doc_idx, val) in block_accessor.iter_vals().enumerate() { field_values[doc_idx * num_fields + field_idx].push(KeyElem::new(val)); } } else { - block_accessor.fetch_block_with_missing_unique_per_doc( - docs, - &field.column, - None, - ); let mut doc_idx = 0usize; for (doc_id, val) in block_accessor.iter_docid_vals(docs, &field.column) { while docs