From 46ee982b520da71a1cadce588eca8244140cf394 Mon Sep 17 00:00:00 2001 From: Pascal Seitz Date: Mon, 14 Sep 2026 17:26:01 +0200 Subject: [PATCH 1/9] add tie_breaker option --- columnar/Cargo.toml | 2 +- columnar/src/column/mod.rs | 1 + columnar/src/column/serialize.rs | 17 +++++++ columnar/src/column_values/mod.rs | 1 + .../u64_based/blockwise_linear.rs | 2 +- columnar/src/column_values/u64_based/mod.rs | 2 +- columnar/src/columnar/merge/mod.rs | 28 +++++++++++ columnar/src/columnar/mod.rs | 2 +- columnar/src/columnar/reader/mod.rs | 19 ++++++++ columnar/src/columnar/writer/mod.rs | 48 ++++++++++++------- columnar/src/lib.rs | 2 +- src/fastfield/mod.rs | 22 +++++++++ src/fastfield/plugin.rs | 18 ++++++- src/fastfield/writer.rs | 7 +++ src/schema/numeric_options.rs | 25 ++++++++++ src/schema/schema.rs | 12 +++++ 16 files changed, 185 insertions(+), 23 deletions(-) diff --git a/columnar/Cargo.toml b/columnar/Cargo.toml index 10b49a24d..6e07375c5 100644 --- a/columnar/Cargo.toml +++ b/columnar/Cargo.toml @@ -18,11 +18,11 @@ common = { version= "0.11", path = "../common", package = "tantivy-common" } tantivy-bitpacker = { version= "0.10", path = "../bitpacker/" } serde = "1.0.152" downcast-rs = "2.0.1" +rand = "0.9" [dev-dependencies] proptest = "1" more-asserts = "0.3.1" -rand = "0.9" binggan = "0.17.0" [[bench]] diff --git a/columnar/src/column/mod.rs b/columnar/src/column/mod.rs index f6a50b45f..3bc61cba0 100644 --- a/columnar/src/column/mod.rs +++ b/columnar/src/column/mod.rs @@ -8,6 +8,7 @@ use std::sync::Arc; use common::BinarySerializable; pub use dictionary_encoded::{BytesColumn, StrColumn}; +pub(crate) use serialize::serialize_generated_tie_breaker_column; pub use serialize::{ open_column_bytes, open_column_str, open_column_u64, open_column_u128, open_column_u128_as_compact_u64, serialize_column_mappable_to_u64, diff --git a/columnar/src/column/serialize.rs b/columnar/src/column/serialize.rs index e2127933b..75fa52f9a 100644 --- a/columnar/src/column/serialize.rs +++ b/columnar/src/column/serialize.rs @@ -3,6 +3,7 @@ use std::io::Write; use std::sync::Arc; use common::OwnedBytes; +use rand::Rng; use sstable::Dictionary; use crate::column::{BytesColumn, Column}; @@ -25,6 +26,22 @@ pub fn serialize_column_mappable_to_u128( Ok(()) } +pub(crate) fn serialize_generated_tie_breaker_column( + num_docs: u32, + output: &mut impl Write, +) -> io::Result<()> { + let block_size = crate::column_values::BLOCK_SIZE; + let max_start = u32::MAX - (block_size - 1); + let mut rng = rand::rng(); + let mut values = Vec::with_capacity(num_docs as usize); + for block_start_doc in (0..num_docs).step_by(block_size as usize) { + let start = rng.random_range(0..=max_start); + let block_len = (num_docs - block_start_doc).min(block_size); + values.extend((0..block_len).map(|offset| (start + offset) as u64)); + } + serialize_column_mappable_to_u64(SerializableColumnIndex::Full, &&values[..], output) +} + pub fn serialize_column_mappable_to_u64( column_index: SerializableColumnIndex<'_>, column_values: &impl Iterable, diff --git a/columnar/src/column_values/mod.rs b/columnar/src/column_values/mod.rs index 64bc69b25..0cb4ac1db 100644 --- a/columnar/src/column_values/mod.rs +++ b/columnar/src/column_values/mod.rs @@ -26,6 +26,7 @@ mod monotonic_column; pub(crate) use merge::MergedColumnValues; pub use stats::ColumnStats; +pub(crate) use u64_based::blockwise_linear::BLOCK_SIZE; pub use u64_based::{ ALL_U64_CODEC_TYPES, CodecType, load_u64_based_column_values, serialize_and_load_u64_based_column_values, serialize_u64_based_column_values, diff --git a/columnar/src/column_values/u64_based/blockwise_linear.rs b/columnar/src/column_values/u64_based/blockwise_linear.rs index b60bf5bad..b8d03794f 100644 --- a/columnar/src/column_values/u64_based/blockwise_linear.rs +++ b/columnar/src/column_values/u64_based/blockwise_linear.rs @@ -11,7 +11,7 @@ use crate::column_values::u64_based::line::Line; use crate::column_values::u64_based::{ColumnCodec, ColumnCodecEstimator, ColumnStats}; use crate::column_values::{ColumnValues, VecColumn}; -const BLOCK_SIZE: u32 = 512u32; +pub(crate) const BLOCK_SIZE: u32 = 512u32; #[derive(Debug, Default)] struct Block { diff --git a/columnar/src/column_values/u64_based/mod.rs b/columnar/src/column_values/u64_based/mod.rs index aa2d9818b..42ef335ad 100644 --- a/columnar/src/column_values/u64_based/mod.rs +++ b/columnar/src/column_values/u64_based/mod.rs @@ -1,5 +1,5 @@ mod bitpacked; -mod blockwise_linear; +pub(crate) mod blockwise_linear; mod line; mod linear; mod stats_collector; diff --git a/columnar/src/columnar/merge/mod.rs b/columnar/src/columnar/merge/mod.rs index 4f7739f4d..c1a73a97a 100644 --- a/columnar/src/columnar/merge/mod.rs +++ b/columnar/src/columnar/merge/mod.rs @@ -79,6 +79,23 @@ pub fn merge_columnar( required_columns: &[(String, ColumnType)], merge_row_order: MergeRowOrder, output: &mut impl io::Write, +) -> io::Result<()> { + merge_columnar_with_tie_breakers( + columnar_readers, + required_columns, + &[], + merge_row_order, + output, + ) +} + +/// Merges columnars and regenerates the named tie-breaker columns in the resulting row order. +pub fn merge_columnar_with_tie_breakers( + columnar_readers: &[&ColumnarReader], + required_columns: &[(String, ColumnType)], + tie_breaker_columns: &[String], + merge_row_order: MergeRowOrder, + output: &mut impl io::Write, ) -> io::Result<()> { let mut serializer = ColumnarSerializer::new(output); let num_docs_per_columnar = columnar_readers @@ -89,6 +106,17 @@ pub fn merge_columnar( let columns_to_merge = group_columns_for_merge(columnar_readers, required_columns)?; for res in columns_to_merge { let ((column_name, _column_type_category), grouped_columns) = res; + if tie_breaker_columns.iter().any(|name| name == &column_name) { + let mut column_serializer = + serializer.start_serialize_column(column_name.as_bytes(), ColumnType::U64); + crate::column::serialize_generated_tie_breaker_column( + merge_row_order.num_rows(), + &mut column_serializer, + )?; + column_serializer.finalize()?; + continue; + } + let grouped_columns = grouped_columns.open(&merge_row_order)?; if grouped_columns.is_empty() { continue; diff --git a/columnar/src/columnar/mod.rs b/columnar/src/columnar/mod.rs index 0c23f9e3b..d98fa70c3 100644 --- a/columnar/src/columnar/mod.rs +++ b/columnar/src/columnar/mod.rs @@ -10,7 +10,7 @@ pub use format_version::{CURRENT_VERSION, Version}; pub(crate) use merge::ColumnTypeCategory; pub use merge::{ MergeRowOrder, ShuffleMergeOrder, StackMergeOrder, compute_merged_term_ord_mapping, - merge_columnar, + merge_columnar, merge_columnar_with_tie_breakers, }; pub use reader::ColumnarReader; pub use writer::ColumnarWriter; diff --git a/columnar/src/columnar/reader/mod.rs b/columnar/src/columnar/reader/mod.rs index 8592b3a22..30734ea79 100644 --- a/columnar/src/columnar/reader/mod.rs +++ b/columnar/src/columnar/reader/mod.rs @@ -236,6 +236,25 @@ mod tests { assert_eq!(columns[1].1.column_type(), ColumnType::U64); } + #[test] + fn test_generated_tie_breaker_column() { + let mut columnar_writer = ColumnarWriter::default(); + columnar_writer.record_tie_breaker_column("tie"); + let mut buffer = Vec::new(); + columnar_writer.serialize(1_025, None, &mut buffer).unwrap(); + + let columnar = ColumnarReader::open(buffer).unwrap(); + let handles = columnar.read_columns("tie").unwrap(); + let column = handles[0].open_u64_lenient().unwrap().unwrap(); + assert_eq!(column.index.get_cardinality(), crate::Cardinality::Full); + for block_start in [0, 512, 1_024] { + let block_end = (block_start + 512).min(1_025); + for doc in block_start + 1..block_end { + assert_eq!(column.first(doc), Some(column.first(doc - 1).unwrap() + 1)); + } + } + } + #[test] fn test_list_columns_strict_typing_prevents_coercion() { let mut columnar_writer = ColumnarWriter::default(); diff --git a/columnar/src/columnar/writer/mod.rs b/columnar/src/columnar/writer/mod.rs index 999ccd058..049a8a6f4 100644 --- a/columnar/src/columnar/writer/mod.rs +++ b/columnar/src/columnar/writer/mod.rs @@ -3,6 +3,7 @@ mod column_writers; mod serializer; mod value_index; +use std::collections::HashSet; use std::io; use std::net::Ipv6Addr; @@ -54,6 +55,7 @@ pub struct ColumnarWriter { ip_addr_field_hash_map: ArenaHashMap, bytes_field_hash_map: ArenaHashMap, str_field_hash_map: ArenaHashMap, + generated_tie_breaker_columns: HashSet>, arena: MemoryArena, // Dictionaries used to store dictionary-encoded values. dictionaries: Vec, @@ -217,6 +219,13 @@ impl ColumnarWriter { } } + /// Registers a full `u64` column whose values are generated when this writer is serialized. + pub fn record_tie_breaker_column(&mut self, column_name: &str) { + self.record_column_type(column_name, ColumnType::U64, false); + self.generated_tie_breaker_columns + .insert(column_name.as_bytes().to_vec()); + } + pub fn record_numerical + Copy>( &mut self, doc: RowId, @@ -436,24 +445,31 @@ impl ColumnarWriter { column_serializer.finalize()?; } ColumnType::F64 | ColumnType::I64 | ColumnType::U64 => { - let numerical_column_writer: NumericalColumnWriter = - self.numerical_field_hash_map.read(addr); - let cardinality = numerical_column_writer.cardinality(num_docs); let mut column_serializer = serializer.start_serialize_column(column_name, column_type); - let numerical_type = column_type.numerical_type().unwrap(); - serialize_numerical_column( - cardinality, - num_docs, - numerical_type, - numerical_column_writer.operation_iterator( - arena, - old_to_new_row_ids, - &mut symbol_byte_buffer, - ), - buffers, - &mut column_serializer, - )?; + if self.generated_tie_breaker_columns.contains(column_name) { + crate::column::serialize_generated_tie_breaker_column( + num_docs, + &mut column_serializer, + )?; + } else { + let numerical_column_writer: NumericalColumnWriter = + self.numerical_field_hash_map.read(addr); + let cardinality = numerical_column_writer.cardinality(num_docs); + let numerical_type = column_type.numerical_type().unwrap(); + serialize_numerical_column( + cardinality, + num_docs, + numerical_type, + numerical_column_writer.operation_iterator( + arena, + old_to_new_row_ids, + &mut symbol_byte_buffer, + ), + buffers, + &mut column_serializer, + )?; + } column_serializer.finalize()?; } ColumnType::DateTime => { diff --git a/columnar/src/lib.rs b/columnar/src/lib.rs index 1da8d9604..60e0e54ea 100644 --- a/columnar/src/lib.rs +++ b/columnar/src/lib.rs @@ -42,7 +42,7 @@ pub use column_values::{ pub use columnar::{ CURRENT_VERSION, ColumnType, ColumnarReader, ColumnarWriter, HasAssociatedColumnType, MergeRowOrder, ShuffleMergeOrder, StackMergeOrder, Version, compute_merged_term_ord_mapping, - merge_columnar, + merge_columnar, merge_columnar_with_tie_breakers, }; use sstable::VoidSSTable; pub use value::{NumericalType, NumericalValue}; diff --git a/src/fastfield/mod.rs b/src/fastfield/mod.rs index d56dc27a8..68bb0e758 100644 --- a/src/fastfield/mod.rs +++ b/src/fastfield/mod.rs @@ -148,6 +148,28 @@ mod tests { Ok(()) } + #[test] + fn test_generated_tie_breaker_fast_field() { + let mut schema_builder = Schema::builder(); + schema_builder.add_tie_breaker_field("tie"); + let schema = schema_builder.build(); + let mut writer = FastFieldsWriter::from_schema(&schema).unwrap(); + for _ in 0..1_025 { + writer.add_document(&TantivyDocument::default()).unwrap(); + } + let mut bytes = Vec::new(); + writer.serialize(&mut bytes, None).unwrap(); + + let readers = FastFieldReaders::open(bytes.into(), schema).unwrap(); + let values = readers.u64("tie").unwrap().first_or_default_col(0); + for block_start in [0, 512, 1_024] { + let block_end = (block_start + 512).min(1_025); + for doc in block_start + 1..block_end { + assert_eq!(values.get_val(doc), values.get_val(doc - 1) + 1); + } + } + } + #[test] fn test_intfastfield_large() { let path = Path::new("test"); diff --git a/src/fastfield/plugin.rs b/src/fastfield/plugin.rs index eb33b6816..37eabc4ae 100644 --- a/src/fastfield/plugin.rs +++ b/src/fastfield/plugin.rs @@ -18,7 +18,7 @@ use crate::index::{SegmentComponent, SegmentReader}; use crate::indexer::doc_id_mapping::{DocIdMapping, MappingType, SegmentDocIdMapping}; use crate::plugin::{PluginMergeContext, PluginWriter, PluginWriterContext, SegmentPlugin}; use crate::schema::document::Document; -use crate::schema::{value_type_to_column_type, Schema}; +use crate::schema::{value_type_to_column_type, FieldType, Schema}; use crate::space_usage::{ComponentSpaceUsage, FAST_FIELDS}; use crate::Segment; @@ -43,6 +43,7 @@ impl SegmentPlugin for FastFieldsPlugin { ctx.target_segment.index().directory().open_write(&path)?; let required_columns = extract_fast_field_required_columns(ctx.schema); + let tie_breaker_columns = extract_tie_breaker_columns(ctx.schema); let columnars: Vec<&ColumnarReader> = ctx .readers .iter() @@ -53,9 +54,10 @@ impl SegmentPlugin for FastFieldsPlugin { let doc_id_mapping = ctx.doc_id_mapping.clone(); let merge_row_order = convert_to_merge_order(&columnars[..], doc_id_mapping); - columnar::merge_columnar( + columnar::merge_columnar_with_tie_breakers( &columnars[..], &required_columns, + &tie_breaker_columns, merge_row_order, &mut fast_field_wrt, )?; @@ -173,6 +175,18 @@ fn convert_to_merge_order( } } +fn extract_tie_breaker_columns(schema: &Schema) -> Vec { + schema + .fields() + .filter_map(|(_, field_entry)| match field_entry.field_type() { + FieldType::U64(options) if options.is_tie_breaker() => { + Some(field_entry.name().to_string()) + } + _ => None, + }) + .collect() +} + fn extract_fast_field_required_columns(schema: &Schema) -> Vec<(String, ColumnType)> { schema .fields() diff --git a/src/fastfield/writer.rs b/src/fastfield/writer.rs index 9bca41357..d74b68ae4 100644 --- a/src/fastfield/writer.rs +++ b/src/fastfield/writer.rs @@ -52,6 +52,13 @@ impl FastFieldsWriter { if !field_entry.field_type().is_fast() { continue; } + if matches!( + field_entry.field_type(), + FieldType::U64(options) if options.is_tie_breaker() + ) { + columnar_writer.record_tie_breaker_column(field_entry.name()); + continue; + } fast_field_names[field_id.field_id() as usize] = Some(field_entry.name().to_string()); let value_type = field_entry.field_type().value_type(); if let FieldType::Date(date_options) = field_entry.field_type() { diff --git a/src/schema/numeric_options.rs b/src/schema/numeric_options.rs index db36b523e..39f6fffd8 100644 --- a/src/schema/numeric_options.rs +++ b/src/schema/numeric_options.rs @@ -16,6 +16,8 @@ pub struct NumericOptions { stored: bool, #[serde(skip_serializing_if = "is_false")] coerce: bool, + #[serde(skip_serializing_if = "is_false")] + tie_breaker: bool, } fn is_false(val: &bool) -> bool { @@ -37,6 +39,8 @@ struct NumericOptionsDeser { stored: bool, #[serde(default)] coerce: bool, + #[serde(default)] + tie_breaker: bool, } impl From for NumericOptions { @@ -47,6 +51,7 @@ impl From for NumericOptions { fast: deser.fast, stored: deser.stored, coerce: deser.coerce, + tie_breaker: deser.tie_breaker, } } } @@ -82,6 +87,16 @@ impl NumericOptions { self.coerce } + pub(crate) fn is_tie_breaker(&self) -> bool { + self.tie_breaker + } + + pub(crate) fn set_tie_breaker(mut self) -> Self { + self.fast = true; + self.tie_breaker = true; + self + } + /// Try to coerce values if they are not a number. Defaults to false. #[must_use] pub fn set_coerce(mut self) -> Self { @@ -145,6 +160,7 @@ impl From for NumericOptions { stored: false, fast: false, coerce: true, + tie_breaker: false, } } } @@ -157,6 +173,7 @@ impl From for NumericOptions { stored: false, fast: true, coerce: false, + tie_breaker: false, } } } @@ -169,6 +186,7 @@ impl From for NumericOptions { stored: true, fast: false, coerce: false, + tie_breaker: false, } } } @@ -181,6 +199,7 @@ impl From for NumericOptions { stored: false, fast: false, coerce: false, + tie_breaker: false, } } } @@ -196,6 +215,7 @@ impl> BitOr for NumericOptions { stored: self.stored | other.stored, fast: self.fast | other.fast, coerce: self.coerce | other.coerce, + tie_breaker: self.tie_breaker | other.tie_breaker, } } } @@ -230,6 +250,7 @@ mod tests { fast: false, stored: false, coerce: false, + tie_breaker: false, } ); } @@ -249,6 +270,7 @@ mod tests { fast: false, stored: false, coerce: false, + tie_breaker: false, } ); } @@ -269,6 +291,7 @@ mod tests { fast: false, stored: false, coerce: false, + tie_breaker: false, } ); } @@ -290,6 +313,7 @@ mod tests { fast: false, stored: false, coerce: false, + tie_breaker: false, } ); } @@ -312,6 +336,7 @@ mod tests { fast: false, stored: false, coerce: true, + tie_breaker: false, } ); } diff --git a/src/schema/schema.rs b/src/schema/schema.rs index 79414473e..833174847 100644 --- a/src/schema/schema.rs +++ b/src/schema/schema.rs @@ -57,6 +57,18 @@ impl SchemaBuilder { self.add_field(field_entry) } + /// Adds a generated tie-breaker fast field. + /// + /// The field is exposed as a `u64` fast field, but its generated values fit in a `u32`. + /// Values supplied by documents for this field are ignored. + /// + /// # Panics + /// + /// Panics when field already exists. + pub fn add_tie_breaker_field(&mut self, field_name_str: &str) -> Field { + self.add_u64_field(field_name_str, NumericOptions::default().set_tie_breaker()) + } + /// Adds a new i64 field. /// Returns the associated field handle /// From e7ba845037316e8bef2ba497dbcad3a70809aa23 Mon Sep 17 00:00:00 2001 From: Pascal Seitz Date: Mon, 14 Sep 2026 18:35:58 +0200 Subject: [PATCH 2/9] add FieldType::TieBreaker --- columnar/benches/bench_merge.rs | 9 +++++++- columnar/src/column/serialize.rs | 9 +++++--- columnar/src/column_values/mod.rs | 3 +-- columnar/src/columnar/merge/mod.rs | 17 +-------------- columnar/src/columnar/merge/tests.rs | 8 ++++++++ columnar/src/columnar/mod.rs | 2 +- columnar/src/compat_tests.rs | 9 +++++++- columnar/src/lib.rs | 2 +- columnar/src/tests.rs | 8 +++++++- src/fastfield/mod.rs | 18 +++++++++++++--- src/fastfield/plugin.rs | 10 ++++----- src/fastfield/writer.rs | 5 +---- src/index/inverted_index_plugin.rs | 6 +++--- src/postings/per_field_postings_writer.rs | 6 ++++-- src/query/query_parser/query_parser.rs | 8 ++++---- src/schema/field_entry.rs | 7 ++++++- src/schema/field_type.rs | 19 +++++++++++++---- src/schema/numeric_options.rs | 25 ----------------------- src/schema/schema.rs | 2 +- 19 files changed, 94 insertions(+), 79 deletions(-) diff --git a/columnar/benches/bench_merge.rs b/columnar/benches/bench_merge.rs index a4b6c3b3f..4bd2b0cbf 100644 --- a/columnar/benches/bench_merge.rs +++ b/columnar/benches/bench_merge.rs @@ -40,7 +40,14 @@ fn main() { let columnar_readers = columnar_readers.iter().collect::>(); let merge_row_order = StackMergeOrder::stack(&columnar_readers[..]); - merge_columnar(&columnar_readers, &[], merge_row_order.into(), &mut out).unwrap(); + merge_columnar( + &columnar_readers, + &[], + &[], + merge_row_order.into(), + &mut out, + ) + .unwrap(); Some(out.len() as u64) }, ); diff --git a/columnar/src/column/serialize.rs b/columnar/src/column/serialize.rs index 75fa52f9a..b13d5bcb3 100644 --- a/columnar/src/column/serialize.rs +++ b/columnar/src/column/serialize.rs @@ -30,12 +30,15 @@ pub(crate) fn serialize_generated_tie_breaker_column( num_docs: u32, output: &mut impl Write, ) -> io::Result<()> { - let block_size = crate::column_values::BLOCK_SIZE; + let block_size = crate::column_values::u64_based::blockwise_linear::BLOCK_SIZE; let max_start = u32::MAX - (block_size - 1); - let mut rng = rand::rng(); + let mut pseudo_random = rand::rng().random::(); let mut values = Vec::with_capacity(num_docs as usize); for block_start_doc in (0..num_docs).step_by(block_size as usize) { - let start = rng.random_range(0..=max_start); + let start = ((pseudo_random as u64 * (u64::from(max_start) + 1)) >> 32) as u32; + pseudo_random = pseudo_random + .wrapping_mul(1_664_525) + .wrapping_add(1_013_904_223); let block_len = (num_docs - block_start_doc).min(block_size); values.extend((0..block_len).map(|offset| (start + offset) as u64)); } diff --git a/columnar/src/column_values/mod.rs b/columnar/src/column_values/mod.rs index 0cb4ac1db..4911012af 100644 --- a/columnar/src/column_values/mod.rs +++ b/columnar/src/column_values/mod.rs @@ -19,14 +19,13 @@ pub(crate) mod monotonic_mapping; pub(crate) mod monotonic_mapping_u128; mod stats; mod u128_based; -mod u64_based; +pub(crate) mod u64_based; mod vec_column; mod monotonic_column; pub(crate) use merge::MergedColumnValues; pub use stats::ColumnStats; -pub(crate) use u64_based::blockwise_linear::BLOCK_SIZE; pub use u64_based::{ ALL_U64_CODEC_TYPES, CodecType, load_u64_based_column_values, serialize_and_load_u64_based_column_values, serialize_u64_based_column_values, diff --git a/columnar/src/columnar/merge/mod.rs b/columnar/src/columnar/merge/mod.rs index c1a73a97a..24df90f57 100644 --- a/columnar/src/columnar/merge/mod.rs +++ b/columnar/src/columnar/merge/mod.rs @@ -74,23 +74,8 @@ impl From for ColumnTypeCategory { /// /// Reminder: a string and a numerical column may bare the same column name. This is not /// considered a conflict. -pub fn merge_columnar( - columnar_readers: &[&ColumnarReader], - required_columns: &[(String, ColumnType)], - merge_row_order: MergeRowOrder, - output: &mut impl io::Write, -) -> io::Result<()> { - merge_columnar_with_tie_breakers( - columnar_readers, - required_columns, - &[], - merge_row_order, - output, - ) -} - /// Merges columnars and regenerates the named tie-breaker columns in the resulting row order. -pub fn merge_columnar_with_tie_breakers( +pub fn merge_columnar( columnar_readers: &[&ColumnarReader], required_columns: &[(String, ColumnType)], tie_breaker_columns: &[String], diff --git a/columnar/src/columnar/merge/tests.rs b/columnar/src/columnar/merge/tests.rs index 8c812047a..56fe950c2 100644 --- a/columnar/src/columnar/merge/tests.rs +++ b/columnar/src/columnar/merge/tests.rs @@ -209,6 +209,7 @@ fn test_merge_columnar_numbers() { crate::columnar::merge_columnar( columnars, &[], + &[], MergeRowOrder::Stack(stack_merge_order), &mut buffer, ) @@ -237,6 +238,7 @@ fn test_merge_columnar_texts() { crate::columnar::merge_columnar( columnars, &[], + &[], MergeRowOrder::Stack(stack_merge_order), &mut buffer, ) @@ -286,6 +288,7 @@ fn test_merge_columnar_byte() { crate::columnar::merge_columnar( columnars, &[], + &[], MergeRowOrder::Stack(stack_merge_order), &mut buffer, ) @@ -342,6 +345,7 @@ fn test_merge_columnar_byte_with_missing() { crate::columnar::merge_columnar( columnars, &[], + &[], MergeRowOrder::Stack(stack_merge_order), &mut buffer, ) @@ -394,6 +398,7 @@ fn test_merge_columnar_different_types() { crate::columnar::merge_columnar( columnars, &[], + &[], MergeRowOrder::Stack(stack_merge_order), &mut buffer, ) @@ -459,6 +464,7 @@ fn test_merge_columnar_different_empty_cardinality() { crate::columnar::merge_columnar( columnars, &[], + &[], MergeRowOrder::Stack(stack_merge_order), &mut buffer, ) @@ -569,6 +575,7 @@ proptest! { merge_columnar( &columnar_refs, &[], + &[], MergeRowOrder::Stack(stack_merge_order), &mut out, ).unwrap(); @@ -586,6 +593,7 @@ proptest! { merge_columnar( &columnar_refs, &[], + &[], MergeRowOrder::Stack(stack_merge_order), &mut out, ).unwrap(); diff --git a/columnar/src/columnar/mod.rs b/columnar/src/columnar/mod.rs index d98fa70c3..0c23f9e3b 100644 --- a/columnar/src/columnar/mod.rs +++ b/columnar/src/columnar/mod.rs @@ -10,7 +10,7 @@ pub use format_version::{CURRENT_VERSION, Version}; pub(crate) use merge::ColumnTypeCategory; pub use merge::{ MergeRowOrder, ShuffleMergeOrder, StackMergeOrder, compute_merged_term_ord_mapping, - merge_columnar, merge_columnar_with_tie_breakers, + merge_columnar, }; pub use reader::ColumnarReader; pub use writer::ColumnarWriter; diff --git a/columnar/src/compat_tests.rs b/columnar/src/compat_tests.rs index 64e615973..4f8916927 100644 --- a/columnar/src/compat_tests.rs +++ b/columnar/src/compat_tests.rs @@ -74,7 +74,14 @@ fn test_format(path: &str) { let columnar_readers = vec![&reader, &reader2]; let merge_row_order = StackMergeOrder::stack(&columnar_readers[..]); let mut out = Vec::new(); - merge_columnar(&columnar_readers, &[], merge_row_order.into(), &mut out).unwrap(); + merge_columnar( + &columnar_readers, + &[], + &[], + merge_row_order.into(), + &mut out, + ) + .unwrap(); let reader = ColumnarReader::open(out).unwrap(); check_columns(&reader); } diff --git a/columnar/src/lib.rs b/columnar/src/lib.rs index 60e0e54ea..1da8d9604 100644 --- a/columnar/src/lib.rs +++ b/columnar/src/lib.rs @@ -42,7 +42,7 @@ pub use column_values::{ pub use columnar::{ CURRENT_VERSION, ColumnType, ColumnarReader, ColumnarWriter, HasAssociatedColumnType, MergeRowOrder, ShuffleMergeOrder, StackMergeOrder, Version, compute_merged_term_ord_mapping, - merge_columnar, merge_columnar_with_tie_breakers, + merge_columnar, }; use sstable::VoidSSTable; pub use value::{NumericalType, NumericalValue}; diff --git a/columnar/src/tests.rs b/columnar/src/tests.rs index 42caa6fe1..3990290b9 100644 --- a/columnar/src/tests.rs +++ b/columnar/src/tests.rs @@ -833,7 +833,7 @@ proptest! { let columnar_readers_arr: Vec<&ColumnarReader> = columnar_readers.iter().collect(); let mut output: Vec = Vec::new(); let stack_merge_order = StackMergeOrder::stack(&columnar_readers_arr[..]).into(); - crate::merge_columnar(&columnar_readers_arr[..], &[], stack_merge_order, &mut output).unwrap(); + crate::merge_columnar(&columnar_readers_arr[..], &[], &[], stack_merge_order, &mut output).unwrap(); let merged_columnar = ColumnarReader::open(output).unwrap(); let concat_rows: Vec> = columnar_docs.iter().flatten().cloned().collect(); let expected_merged_columnar = build_columnar(&concat_rows[..]); @@ -855,6 +855,7 @@ fn test_columnar_merging_empty_columnar() { crate::merge_columnar( &columnar_readers_arr[..], &[], + &[], crate::MergeRowOrder::Stack(stack_merge_order), &mut output, ) @@ -892,6 +893,7 @@ fn test_columnar_merging_number_columns() { crate::merge_columnar( &columnar_readers_arr[..], &[], + &[], crate::MergeRowOrder::Stack(stack_merge_order), &mut output, ) @@ -965,6 +967,7 @@ fn test_columnar_merge_and_remap( crate::merge_columnar( &columnar_readers_ref[..], &[], + &[], shuffle_merge_order.into(), &mut output, ) @@ -1007,6 +1010,7 @@ fn test_columnar_merge_empty() { crate::merge_columnar( &[&columnar_reader_1, &columnar_reader_2], &[], + &[], shuffle_merge_order.into(), &mut output, ) @@ -1033,6 +1037,7 @@ fn test_columnar_merge_single_str_column() { crate::merge_columnar( &[&columnar_reader_1, &columnar_reader_2], &[], + &[], shuffle_merge_order.into(), &mut output, ) @@ -1065,6 +1070,7 @@ fn test_delete_decrease_cardinality() { crate::merge_columnar( &[&columnar_reader_1, &columnar_reader_2], &[], + &[], shuffle_merge_order.into(), &mut output, ) diff --git a/src/fastfield/mod.rs b/src/fastfield/mod.rs index 68bb0e758..df3cf270a 100644 --- a/src/fastfield/mod.rs +++ b/src/fastfield/mod.rs @@ -93,8 +93,8 @@ mod tests { use crate::index::SegmentId; use crate::merge_policy::NoMergePolicy; use crate::schema::{ - DateOptions, Facet, FacetOptions, Field, JsonObjectOptions, Schema, SchemaBuilder, - TantivyDocument, TextOptions, FAST, INDEXED, STORED, STRING, TEXT, + DateOptions, Facet, FacetOptions, Field, FieldType, JsonObjectOptions, Schema, + SchemaBuilder, TantivyDocument, TextOptions, FAST, INDEXED, STORED, STRING, TEXT, }; use crate::time::OffsetDateTime; use crate::tokenizer::{ @@ -151,8 +151,20 @@ mod tests { #[test] fn test_generated_tie_breaker_fast_field() { let mut schema_builder = Schema::builder(); - schema_builder.add_tie_breaker_field("tie"); + let tie = schema_builder.add_tie_breaker_field("tie"); let schema = schema_builder.build(); + let entry = schema.get_field_entry(tie); + assert!(matches!(entry.field_type(), FieldType::TieBreaker)); + assert!(entry.is_fast()); + assert!(!entry.is_indexed()); + assert!(!entry.is_stored()); + let schema_json = serde_json::to_string(&schema).unwrap(); + assert!(schema_json.contains(r#""type":"tie_breaker""#)); + assert_eq!( + serde_json::from_str::(&schema_json).unwrap(), + schema + ); + let mut writer = FastFieldsWriter::from_schema(&schema).unwrap(); for _ in 0..1_025 { writer.add_document(&TantivyDocument::default()).unwrap(); diff --git a/src/fastfield/plugin.rs b/src/fastfield/plugin.rs index 37eabc4ae..bdd8dc19a 100644 --- a/src/fastfield/plugin.rs +++ b/src/fastfield/plugin.rs @@ -54,7 +54,7 @@ impl SegmentPlugin for FastFieldsPlugin { let doc_id_mapping = ctx.doc_id_mapping.clone(); let merge_row_order = convert_to_merge_order(&columnars[..], doc_id_mapping); - columnar::merge_columnar_with_tie_breakers( + columnar::merge_columnar( &columnars[..], &required_columns, &tie_breaker_columns, @@ -178,11 +178,9 @@ fn convert_to_merge_order( fn extract_tie_breaker_columns(schema: &Schema) -> Vec { schema .fields() - .filter_map(|(_, field_entry)| match field_entry.field_type() { - FieldType::U64(options) if options.is_tie_breaker() => { - Some(field_entry.name().to_string()) - } - _ => None, + .filter_map(|(_, field_entry)| { + matches!(field_entry.field_type(), FieldType::TieBreaker) + .then(|| field_entry.name().to_string()) }) .collect() } diff --git a/src/fastfield/writer.rs b/src/fastfield/writer.rs index d74b68ae4..e387d9541 100644 --- a/src/fastfield/writer.rs +++ b/src/fastfield/writer.rs @@ -52,10 +52,7 @@ impl FastFieldsWriter { if !field_entry.field_type().is_fast() { continue; } - if matches!( - field_entry.field_type(), - FieldType::U64(options) if options.is_tie_breaker() - ) { + if matches!(field_entry.field_type(), FieldType::TieBreaker) { columnar_writer.record_tie_breaker_column(field_entry.name()); continue; } diff --git a/src/index/inverted_index_plugin.rs b/src/index/inverted_index_plugin.rs index add130c75..68ddf3e48 100644 --- a/src/index/inverted_index_plugin.rs +++ b/src/index/inverted_index_plugin.rs @@ -390,9 +390,9 @@ impl InvertedIndexPluginWriter { self.fieldnorms_writer.record(doc_id, field, num_vals); } } - // Custom fields are not indexed; the `is_indexed()` guard above skips them. - FieldType::Custom(_) => { - unreachable!("the inverted index does not support custom field types") + // These fields are not indexed; the `is_indexed()` guard above skips them. + FieldType::TieBreaker | FieldType::Custom(_) => { + unreachable!("the inverted index does not support this field type") } } } diff --git a/src/postings/per_field_postings_writer.rs b/src/postings/per_field_postings_writer.rs index 5ec4e6a49..7147c356e 100644 --- a/src/postings/per_field_postings_writer.rs +++ b/src/postings/per_field_postings_writer.rs @@ -69,8 +69,10 @@ fn posting_writer_from_field_entry(field_entry: &FieldEntry) -> Box::default().into() } } - // Custom fields are never indexed, so this writer is never fed terms. It only needs to + // These fields are never indexed, so this writer is never fed terms. It only needs to // occupy the per-field slot (the vector is indexed by field id). - FieldType::Custom(_) => Box::>::default(), + FieldType::TieBreaker | FieldType::Custom(_) => { + Box::>::default() + } } } diff --git a/src/query/query_parser/query_parser.rs b/src/query/query_parser/query_parser.rs index bc97f1456..73416ee2d 100644 --- a/src/query/query_parser/query_parser.rs +++ b/src/query/query_parser/query_parser.rs @@ -453,7 +453,7 @@ impl QueryParser { ))); } match *field_type { - FieldType::U64(_) => { + FieldType::U64(_) | FieldType::TieBreaker => { let val: u64 = u64::from_str(phrase)?; Ok(Term::from_field_u64(field, val)) } @@ -635,10 +635,10 @@ impl QueryParser { let term = Term::from_field_ip_addr(field, ip_v6); Ok(vec![LogicalLiteral::Term(term)]) } - // Custom fields are not indexed, so the `is_indexed()` guard above returns + // These fields are not indexed, so the `is_indexed()` guard above returns // `FieldNotIndexed` before this match. - FieldType::Custom(_) => { - unreachable!("the query parser does not support custom field types") + FieldType::TieBreaker | FieldType::Custom(_) => { + unreachable!("the query parser does not support this field type") } } } diff --git a/src/schema/field_entry.rs b/src/schema/field_entry.rs index 5dddb5799..f7816b218 100644 --- a/src/schema/field_entry.rs +++ b/src/schema/field_entry.rs @@ -41,6 +41,11 @@ impl FieldEntry { Self::new(field_name, FieldType::U64(int_options)) } + /// Creates a generated tie-breaker field entry. + pub fn new_tie_breaker(field_name: String) -> FieldEntry { + Self::new(field_name, FieldType::TieBreaker) + } + /// Creates a new i64 field entry. pub fn new_i64(field_name: String, int_options: NumericOptions) -> FieldEntry { Self::new(field_name, FieldType::I64(int_options)) @@ -135,7 +140,7 @@ impl FieldEntry { FieldType::Bytes(ref options) => options.is_stored(), FieldType::JsonObject(ref options) => options.is_stored(), FieldType::IpAddr(ref options) => options.is_stored(), - FieldType::Custom(_) => false, + FieldType::TieBreaker | FieldType::Custom(_) => false, } } } diff --git a/src/schema/field_type.rs b/src/schema/field_type.rs index bdf8cdd60..30d1e6af1 100644 --- a/src/schema/field_type.rs +++ b/src/schema/field_type.rs @@ -190,6 +190,8 @@ pub enum FieldType { Str(TextOptions), /// Unsigned 64-bits integers field type configuration U64(NumericOptions), + /// Generated tie-breaker field, exposed as a `u64` fast field. + TieBreaker, /// Signed 64-bits integers 64 field type configuration I64(NumericOptions), /// 64-bits float 64 field type configuration @@ -217,7 +219,7 @@ impl FieldType { pub fn value_type(&self) -> Type { match *self { FieldType::Str(_) => Type::Str, - FieldType::U64(_) => Type::U64, + FieldType::U64(_) | FieldType::TieBreaker => Type::U64, FieldType::I64(_) => Type::I64, FieldType::F64(_) => Type::F64, FieldType::Bool(_) => Type::Bool, @@ -273,7 +275,7 @@ impl FieldType { FieldType::Bytes(ref bytes_options) => bytes_options.is_indexed(), FieldType::JsonObject(ref json_object_options) => json_object_options.is_indexed(), FieldType::IpAddr(ref ip_addr_options) => ip_addr_options.is_indexed(), - FieldType::Custom(_) => false, + FieldType::TieBreaker | FieldType::Custom(_) => false, } } @@ -311,6 +313,7 @@ impl FieldType { FieldType::IpAddr(ref ip_addr_options) => ip_addr_options.is_fast(), FieldType::Facet(_) => true, FieldType::JsonObject(ref json_object_options) => json_object_options.is_fast(), + FieldType::TieBreaker => true, FieldType::Custom(_) => false, } } @@ -331,7 +334,7 @@ impl FieldType { FieldType::Bytes(ref bytes_options) => bytes_options.fieldnorms(), FieldType::JsonObject(ref _json_object_options) => false, FieldType::IpAddr(ref ip_addr_options) => ip_addr_options.fieldnorms(), - FieldType::Custom(_) => false, + FieldType::TieBreaker | FieldType::Custom(_) => false, } } @@ -383,7 +386,7 @@ impl FieldType { None } } - FieldType::Custom(_) => None, + FieldType::TieBreaker | FieldType::Custom(_) => None, } } @@ -486,6 +489,10 @@ impl FieldType { Ok(OwnedValue::IpAddr(ip_addr.into_ipv6_addr())) } + FieldType::TieBreaker => Err(ValueParsingError::TypeError { + expected: "a generated tie-breaker field", + json: JsonValue::String(field_text), + }), FieldType::Custom(_) => Err(custom_not_json_error(JsonValue::String(field_text))), }, JsonValue::Number(field_val_num) => match self { @@ -545,6 +552,10 @@ impl FieldType { expected: "a string with an ip addr", json: JsonValue::Number(field_val_num), }), + FieldType::TieBreaker => Err(ValueParsingError::TypeError { + expected: "a generated tie-breaker field", + json: JsonValue::Number(field_val_num), + }), FieldType::Custom(_) => { Err(custom_not_json_error(JsonValue::Number(field_val_num))) } diff --git a/src/schema/numeric_options.rs b/src/schema/numeric_options.rs index 39f6fffd8..db36b523e 100644 --- a/src/schema/numeric_options.rs +++ b/src/schema/numeric_options.rs @@ -16,8 +16,6 @@ pub struct NumericOptions { stored: bool, #[serde(skip_serializing_if = "is_false")] coerce: bool, - #[serde(skip_serializing_if = "is_false")] - tie_breaker: bool, } fn is_false(val: &bool) -> bool { @@ -39,8 +37,6 @@ struct NumericOptionsDeser { stored: bool, #[serde(default)] coerce: bool, - #[serde(default)] - tie_breaker: bool, } impl From for NumericOptions { @@ -51,7 +47,6 @@ impl From for NumericOptions { fast: deser.fast, stored: deser.stored, coerce: deser.coerce, - tie_breaker: deser.tie_breaker, } } } @@ -87,16 +82,6 @@ impl NumericOptions { self.coerce } - pub(crate) fn is_tie_breaker(&self) -> bool { - self.tie_breaker - } - - pub(crate) fn set_tie_breaker(mut self) -> Self { - self.fast = true; - self.tie_breaker = true; - self - } - /// Try to coerce values if they are not a number. Defaults to false. #[must_use] pub fn set_coerce(mut self) -> Self { @@ -160,7 +145,6 @@ impl From for NumericOptions { stored: false, fast: false, coerce: true, - tie_breaker: false, } } } @@ -173,7 +157,6 @@ impl From for NumericOptions { stored: false, fast: true, coerce: false, - tie_breaker: false, } } } @@ -186,7 +169,6 @@ impl From for NumericOptions { stored: true, fast: false, coerce: false, - tie_breaker: false, } } } @@ -199,7 +181,6 @@ impl From for NumericOptions { stored: false, fast: false, coerce: false, - tie_breaker: false, } } } @@ -215,7 +196,6 @@ impl> BitOr for NumericOptions { stored: self.stored | other.stored, fast: self.fast | other.fast, coerce: self.coerce | other.coerce, - tie_breaker: self.tie_breaker | other.tie_breaker, } } } @@ -250,7 +230,6 @@ mod tests { fast: false, stored: false, coerce: false, - tie_breaker: false, } ); } @@ -270,7 +249,6 @@ mod tests { fast: false, stored: false, coerce: false, - tie_breaker: false, } ); } @@ -291,7 +269,6 @@ mod tests { fast: false, stored: false, coerce: false, - tie_breaker: false, } ); } @@ -313,7 +290,6 @@ mod tests { fast: false, stored: false, coerce: false, - tie_breaker: false, } ); } @@ -336,7 +312,6 @@ mod tests { fast: false, stored: false, coerce: true, - tie_breaker: false, } ); } diff --git a/src/schema/schema.rs b/src/schema/schema.rs index 833174847..ce7bb5804 100644 --- a/src/schema/schema.rs +++ b/src/schema/schema.rs @@ -66,7 +66,7 @@ impl SchemaBuilder { /// /// Panics when field already exists. pub fn add_tie_breaker_field(&mut self, field_name_str: &str) -> Field { - self.add_u64_field(field_name_str, NumericOptions::default().set_tie_breaker()) + self.add_field(FieldEntry::new_tie_breaker(field_name_str.to_string())) } /// Adds a new i64 field. From a4cde3a22c3f76a876d240ebdbf1b542ea7604ce Mon Sep 17 00:00:00 2001 From: Pascal Seitz Date: Tue, 15 Sep 2026 09:45:58 +0200 Subject: [PATCH 3/9] Classify tie breakers during columnar merge Represent generated tie breakers as a merge column category so existing values are skipped and regenerated without being loaded. Test grouping, block linearity, and compact merged output. --- columnar/src/columnar/merge/mod.rs | 28 +++++++++-- columnar/src/columnar/merge/tests.rs | 75 +++++++++++++++++++++++++--- 2 files changed, 90 insertions(+), 13 deletions(-) diff --git a/columnar/src/columnar/merge/mod.rs b/columnar/src/columnar/merge/mod.rs index 24df90f57..f9f043bf6 100644 --- a/columnar/src/columnar/merge/mod.rs +++ b/columnar/src/columnar/merge/mod.rs @@ -30,9 +30,10 @@ use crate::{ /// /// See also [README.md]. /// -/// The ordering has to match the ordering of the variants in [ColumnType]. +/// Except for generated tie-breakers, the ordering matches the variants in [ColumnType]. #[derive(Copy, Clone, Eq, PartialOrd, Ord, PartialEq, Hash, Debug)] pub(crate) enum ColumnTypeCategory { + TieBreaker, Numerical, Bytes, Str, @@ -88,10 +89,11 @@ pub fn merge_columnar( .map(|reader| reader.num_docs()) .collect::>(); - let columns_to_merge = group_columns_for_merge(columnar_readers, required_columns)?; + let columns_to_merge = + group_columns_for_merge(columnar_readers, required_columns, tie_breaker_columns)?; for res in columns_to_merge { - let ((column_name, _column_type_category), grouped_columns) = res; - if tie_breaker_columns.iter().any(|name| name == &column_name) { + let ((column_name, column_type_category), grouped_columns) = res; + if column_type_category == ColumnTypeCategory::TieBreaker { let mut column_serializer = serializer.start_serialize_column(column_name.as_bytes(), ColumnType::U64); crate::column::serialize_generated_tie_breaker_column( @@ -411,21 +413,37 @@ fn is_empty_after_merge( fn group_columns_for_merge<'a>( columnar_readers: &'a [&'a ColumnarReader], required_columns: &'a [(String, ColumnType)], + tie_breaker_columns: &[String], ) -> io::Result> { let mut columns: BTreeMap<(String, ColumnTypeCategory), GroupedColumnsHandle> = BTreeMap::new(); + let tie_breaker_columns: HashSet<&str> = + tie_breaker_columns.iter().map(String::as_str).collect(); for &(ref column_name, column_type) in required_columns { + if tie_breaker_columns.contains(column_name.as_str()) { + continue; + } columns .entry((column_name.clone(), column_type.into())) .or_insert_with(|| GroupedColumnsHandle::new(columnar_readers.len())) .require_type(column_type)?; } + for column_name in &tie_breaker_columns { + columns + .entry(((*column_name).to_string(), ColumnTypeCategory::TieBreaker)) + .or_insert_with(|| GroupedColumnsHandle::new(columnar_readers.len())) + .require_type(ColumnType::U64)?; + } + for (columnar_id, columnar_reader) in columnar_readers.iter().enumerate() { let column_name_and_handle = columnar_reader.iter_columns()?; for (column_name, handle) in column_name_and_handle { - let column_category: ColumnTypeCategory = handle.column_type().into(); + if tie_breaker_columns.contains(column_name.as_str()) { + continue; + } + let column_category = handle.column_type().into(); columns .entry((column_name, column_category)) .or_insert_with(|| GroupedColumnsHandle::new(columnar_readers.len())) diff --git a/columnar/src/columnar/merge/tests.rs b/columnar/src/columnar/merge/tests.rs index 56fe950c2..ebd7b0e8d 100644 --- a/columnar/src/columnar/merge/tests.rs +++ b/columnar/src/columnar/merge/tests.rs @@ -30,7 +30,7 @@ fn test_column_coercion_to_u64() { let columnar2 = make_columnar("numbers", &[u64::MAX]); let columnars = &[&columnar1, &columnar2]; let column_map: BTreeMap<(String, ColumnTypeCategory), GroupedColumnsHandle> = - group_columns_for_merge(columnars, &[]).unwrap(); + group_columns_for_merge(columnars, &[], &[]).unwrap(); assert_eq!(column_map.len(), 1); assert!(column_map.contains_key(&("numbers".to_string(), ColumnTypeCategory::Numerical))); } @@ -41,7 +41,7 @@ fn test_column_coercion_to_i64() { let columnar2 = make_columnar("numbers", &[2u64]); let columnars = &[&columnar1, &columnar2]; let column_map: BTreeMap<(String, ColumnTypeCategory), GroupedColumnsHandle> = - group_columns_for_merge(columnars, &[]).unwrap(); + group_columns_for_merge(columnars, &[], &[]).unwrap(); assert_eq!(column_map.len(), 1); assert!(column_map.contains_key(&("numbers".to_string(), ColumnTypeCategory::Numerical))); } @@ -65,19 +65,77 @@ fn test_group_columns_with_required_column() { let columnar2 = make_columnar("numbers", &[2u64]); let columnars = &[&columnar1, &columnar2]; let column_map: BTreeMap<(String, ColumnTypeCategory), GroupedColumnsHandle> = - group_columns_for_merge(columnars, &[("numbers".to_string(), ColumnType::U64)]).unwrap(); + group_columns_for_merge(columnars, &[("numbers".to_string(), ColumnType::U64)], &[]) + .unwrap(); assert_eq!(column_map.len(), 1); assert!(column_map.contains_key(&("numbers".to_string(), ColumnTypeCategory::Numerical))); } +#[test] +fn test_group_tie_breaker_column() { + let columnar = make_columnar("tie", &[1u64]); + let column_map = group_columns_for_merge(&[&columnar], &[], &["tie".to_string()]).unwrap(); + assert_eq!(column_map.len(), 1); + let grouped = column_map + .get(&("tie".to_string(), ColumnTypeCategory::TieBreaker)) + .unwrap(); + assert!(grouped.columns[0].is_none()); +} + +#[test] +fn test_merge_generated_tie_breaker_columns_stays_small() { + fn make_tie_breaker_columnar(num_docs: u32) -> ColumnarReader { + let mut writer = ColumnarWriter::default(); + writer.record_tie_breaker_column("tie"); + let mut buffer = Vec::new(); + writer.serialize(num_docs, None, &mut buffer).unwrap(); + ColumnarReader::open(buffer).unwrap() + } + + let columnar1 = make_tie_breaker_columnar(10_000); + let columnar2 = make_tie_breaker_columnar(10_000); + let columnars = &[&columnar1, &columnar2]; + let mut buffer = Vec::new(); + merge_columnar( + columnars, + &[("tie".to_string(), ColumnType::U64)], + &["tie".to_string()], + StackMergeOrder::stack(columnars).into(), + &mut buffer, + ) + .unwrap(); + + assert!( + buffer.len() <= 600, + "merged columnar is {} bytes", + buffer.len() + ); + let merged = ColumnarReader::open(buffer).unwrap(); + let column = merged.read_columns("tie").unwrap()[0] + .open_u64_lenient() + .unwrap() + .unwrap(); + assert_eq!(merged.num_docs(), 20_000); + assert_eq!(column.get_cardinality(), Cardinality::Full); + for block_start in (0..20_000).step_by(512) { + let block_end = (block_start + 512).min(20_000); + for doc in block_start + 1..block_end { + assert_eq!(column.first(doc), Some(column.first(doc - 1).unwrap() + 1)); + } + } +} + #[test] fn test_group_columns_required_column_with_no_existing_columns() { let columnar1 = make_columnar("numbers", &[2u64]); let columnar2 = make_columnar("numbers", &[2u64]); let columnars = &[&columnar1, &columnar2]; - let column_map: BTreeMap<_, _> = - group_columns_for_merge(columnars, &[("required_col".to_string(), ColumnType::Str)]) - .unwrap(); + let column_map: BTreeMap<_, _> = group_columns_for_merge( + columnars, + &[("required_col".to_string(), ColumnType::Str)], + &[], + ) + .unwrap(); assert_eq!(column_map.len(), 2); let columns = &column_map .get(&("required_col".to_string(), ColumnTypeCategory::Str)) @@ -94,7 +152,8 @@ fn test_group_columns_required_column_is_above_all_columns_have_the_same_type_ru let columnar2 = make_columnar("numbers", &[2i64]); let columnars = &[&columnar1, &columnar2]; let column_map: BTreeMap<(String, ColumnTypeCategory), GroupedColumnsHandle> = - group_columns_for_merge(columnars, &[("numbers".to_string(), ColumnType::U64)]).unwrap(); + group_columns_for_merge(columnars, &[("numbers".to_string(), ColumnType::U64)], &[]) + .unwrap(); assert_eq!(column_map.len(), 1); assert!(column_map.contains_key(&("numbers".to_string(), ColumnTypeCategory::Numerical))); } @@ -105,7 +164,7 @@ fn test_missing_column() { let columnar2 = make_columnar("numbers2", &[2u64]); let columnars = &[&columnar1, &columnar2]; let column_map: BTreeMap<(String, ColumnTypeCategory), GroupedColumnsHandle> = - group_columns_for_merge(columnars, &[]).unwrap(); + group_columns_for_merge(columnars, &[], &[]).unwrap(); assert_eq!(column_map.len(), 2); assert!(column_map.contains_key(&("numbers".to_string(), ColumnTypeCategory::Numerical))); { From 4dddea5a5f3050394f711c06075148965a752c1f Mon Sep 17 00:00:00 2001 From: Pascal Seitz Date: Tue, 15 Sep 2026 17:55:32 +0200 Subject: [PATCH 4/9] Preserve tie breakers during columnar merges Generate one random-offset linear tie-breaker sequence per segment and preserve its values through normal columnar merging. Enable the linear codec to keep initial segments compact. --- columnar/benches/bench_merge.rs | 9 +--- columnar/src/column/serialize.rs | 31 +++++++------ columnar/src/columnar/merge/mod.rs | 39 ++--------------- columnar/src/columnar/merge/tests.rs | 65 ++++++++++------------------ columnar/src/columnar/reader/mod.rs | 7 +-- columnar/src/compat_tests.rs | 9 +--- columnar/src/tests.rs | 8 +--- src/fastfield/mod.rs | 7 +-- src/fastfield/plugin.rs | 14 +----- 9 files changed, 52 insertions(+), 137 deletions(-) diff --git a/columnar/benches/bench_merge.rs b/columnar/benches/bench_merge.rs index 4bd2b0cbf..a4b6c3b3f 100644 --- a/columnar/benches/bench_merge.rs +++ b/columnar/benches/bench_merge.rs @@ -40,14 +40,7 @@ fn main() { let columnar_readers = columnar_readers.iter().collect::>(); let merge_row_order = StackMergeOrder::stack(&columnar_readers[..]); - merge_columnar( - &columnar_readers, - &[], - &[], - merge_row_order.into(), - &mut out, - ) - .unwrap(); + merge_columnar(&columnar_readers, &[], merge_row_order.into(), &mut out).unwrap(); Some(out.len() as u64) }, ); diff --git a/columnar/src/column/serialize.rs b/columnar/src/column/serialize.rs index b13d5bcb3..f18d18199 100644 --- a/columnar/src/column/serialize.rs +++ b/columnar/src/column/serialize.rs @@ -30,19 +30,24 @@ pub(crate) fn serialize_generated_tie_breaker_column( num_docs: u32, output: &mut impl Write, ) -> io::Result<()> { - let block_size = crate::column_values::u64_based::blockwise_linear::BLOCK_SIZE; - let max_start = u32::MAX - (block_size - 1); - let mut pseudo_random = rand::rng().random::(); - let mut values = Vec::with_capacity(num_docs as usize); - for block_start_doc in (0..num_docs).step_by(block_size as usize) { - let start = ((pseudo_random as u64 * (u64::from(max_start) + 1)) >> 32) as u32; - pseudo_random = pseudo_random - .wrapping_mul(1_664_525) - .wrapping_add(1_013_904_223); - let block_len = (num_docs - block_start_doc).min(block_size); - values.extend((0..block_len).map(|offset| (start + offset) as u64)); - } - serialize_column_mappable_to_u64(SerializableColumnIndex::Full, &&values[..], output) + let max_start = u32::MAX - num_docs.saturating_sub(1); + let random = rand::rng().random::(); + let start = ((random as u64 * (u64::from(max_start) + 1)) >> 32) as u32; + let values: Vec = (0..num_docs) + .map(|offset| (start + offset) as u64) + .collect(); + let column_index_num_bytes = serialize_column_index(SerializableColumnIndex::Full, output)?; + serialize_u64_based_column_values( + &&values[..], + &[ + CodecType::Bitpacked, + CodecType::Linear, + CodecType::BlockwiseLinear, + ], + output, + )?; + output.write_all(&column_index_num_bytes.to_le_bytes())?; + Ok(()) } pub fn serialize_column_mappable_to_u64( diff --git a/columnar/src/columnar/merge/mod.rs b/columnar/src/columnar/merge/mod.rs index f9f043bf6..4f7739f4d 100644 --- a/columnar/src/columnar/merge/mod.rs +++ b/columnar/src/columnar/merge/mod.rs @@ -30,10 +30,9 @@ use crate::{ /// /// See also [README.md]. /// -/// Except for generated tie-breakers, the ordering matches the variants in [ColumnType]. +/// The ordering has to match the ordering of the variants in [ColumnType]. #[derive(Copy, Clone, Eq, PartialOrd, Ord, PartialEq, Hash, Debug)] pub(crate) enum ColumnTypeCategory { - TieBreaker, Numerical, Bytes, Str, @@ -75,11 +74,9 @@ impl From for ColumnTypeCategory { /// /// Reminder: a string and a numerical column may bare the same column name. This is not /// considered a conflict. -/// Merges columnars and regenerates the named tie-breaker columns in the resulting row order. pub fn merge_columnar( columnar_readers: &[&ColumnarReader], required_columns: &[(String, ColumnType)], - tie_breaker_columns: &[String], merge_row_order: MergeRowOrder, output: &mut impl io::Write, ) -> io::Result<()> { @@ -89,21 +86,9 @@ pub fn merge_columnar( .map(|reader| reader.num_docs()) .collect::>(); - let columns_to_merge = - group_columns_for_merge(columnar_readers, required_columns, tie_breaker_columns)?; + let columns_to_merge = group_columns_for_merge(columnar_readers, required_columns)?; for res in columns_to_merge { - let ((column_name, column_type_category), grouped_columns) = res; - if column_type_category == ColumnTypeCategory::TieBreaker { - let mut column_serializer = - serializer.start_serialize_column(column_name.as_bytes(), ColumnType::U64); - crate::column::serialize_generated_tie_breaker_column( - merge_row_order.num_rows(), - &mut column_serializer, - )?; - column_serializer.finalize()?; - continue; - } - + let ((column_name, _column_type_category), grouped_columns) = res; let grouped_columns = grouped_columns.open(&merge_row_order)?; if grouped_columns.is_empty() { continue; @@ -413,37 +398,21 @@ fn is_empty_after_merge( fn group_columns_for_merge<'a>( columnar_readers: &'a [&'a ColumnarReader], required_columns: &'a [(String, ColumnType)], - tie_breaker_columns: &[String], ) -> io::Result> { let mut columns: BTreeMap<(String, ColumnTypeCategory), GroupedColumnsHandle> = BTreeMap::new(); - let tie_breaker_columns: HashSet<&str> = - tie_breaker_columns.iter().map(String::as_str).collect(); for &(ref column_name, column_type) in required_columns { - if tie_breaker_columns.contains(column_name.as_str()) { - continue; - } columns .entry((column_name.clone(), column_type.into())) .or_insert_with(|| GroupedColumnsHandle::new(columnar_readers.len())) .require_type(column_type)?; } - for column_name in &tie_breaker_columns { - columns - .entry(((*column_name).to_string(), ColumnTypeCategory::TieBreaker)) - .or_insert_with(|| GroupedColumnsHandle::new(columnar_readers.len())) - .require_type(ColumnType::U64)?; - } - for (columnar_id, columnar_reader) in columnar_readers.iter().enumerate() { let column_name_and_handle = columnar_reader.iter_columns()?; for (column_name, handle) in column_name_and_handle { - if tie_breaker_columns.contains(column_name.as_str()) { - continue; - } - let column_category = handle.column_type().into(); + let column_category: ColumnTypeCategory = handle.column_type().into(); columns .entry((column_name, column_category)) .or_insert_with(|| GroupedColumnsHandle::new(columnar_readers.len())) diff --git a/columnar/src/columnar/merge/tests.rs b/columnar/src/columnar/merge/tests.rs index ebd7b0e8d..78c2bd804 100644 --- a/columnar/src/columnar/merge/tests.rs +++ b/columnar/src/columnar/merge/tests.rs @@ -30,7 +30,7 @@ fn test_column_coercion_to_u64() { let columnar2 = make_columnar("numbers", &[u64::MAX]); let columnars = &[&columnar1, &columnar2]; let column_map: BTreeMap<(String, ColumnTypeCategory), GroupedColumnsHandle> = - group_columns_for_merge(columnars, &[], &[]).unwrap(); + group_columns_for_merge(columnars, &[]).unwrap(); assert_eq!(column_map.len(), 1); assert!(column_map.contains_key(&("numbers".to_string(), ColumnTypeCategory::Numerical))); } @@ -41,7 +41,7 @@ fn test_column_coercion_to_i64() { let columnar2 = make_columnar("numbers", &[2u64]); let columnars = &[&columnar1, &columnar2]; let column_map: BTreeMap<(String, ColumnTypeCategory), GroupedColumnsHandle> = - group_columns_for_merge(columnars, &[], &[]).unwrap(); + group_columns_for_merge(columnars, &[]).unwrap(); assert_eq!(column_map.len(), 1); assert!(column_map.contains_key(&("numbers".to_string(), ColumnTypeCategory::Numerical))); } @@ -65,23 +65,11 @@ fn test_group_columns_with_required_column() { let columnar2 = make_columnar("numbers", &[2u64]); let columnars = &[&columnar1, &columnar2]; let column_map: BTreeMap<(String, ColumnTypeCategory), GroupedColumnsHandle> = - group_columns_for_merge(columnars, &[("numbers".to_string(), ColumnType::U64)], &[]) - .unwrap(); + group_columns_for_merge(columnars, &[("numbers".to_string(), ColumnType::U64)]).unwrap(); assert_eq!(column_map.len(), 1); assert!(column_map.contains_key(&("numbers".to_string(), ColumnTypeCategory::Numerical))); } -#[test] -fn test_group_tie_breaker_column() { - let columnar = make_columnar("tie", &[1u64]); - let column_map = group_columns_for_merge(&[&columnar], &[], &["tie".to_string()]).unwrap(); - assert_eq!(column_map.len(), 1); - let grouped = column_map - .get(&("tie".to_string(), ColumnTypeCategory::TieBreaker)) - .unwrap(); - assert!(grouped.columns[0].is_none()); -} - #[test] fn test_merge_generated_tie_breaker_columns_stays_small() { fn make_tie_breaker_columnar(num_docs: u32) -> ColumnarReader { @@ -89,24 +77,28 @@ fn test_merge_generated_tie_breaker_columns_stays_small() { writer.record_tie_breaker_column("tie"); let mut buffer = Vec::new(); writer.serialize(num_docs, None, &mut buffer).unwrap(); + assert!( + buffer.len() <= 100, + "input columnar is {} bytes", + buffer.len() + ); ColumnarReader::open(buffer).unwrap() } - let columnar1 = make_tie_breaker_columnar(10_000); - let columnar2 = make_tie_breaker_columnar(10_000); - let columnars = &[&columnar1, &columnar2]; + let columnars: Vec = + (0..4).map(|_| make_tie_breaker_columnar(250_000)).collect(); + let columnar_refs: Vec<&ColumnarReader> = columnars.iter().collect(); let mut buffer = Vec::new(); merge_columnar( - columnars, + &columnar_refs, &[("tie".to_string(), ColumnType::U64)], - &["tie".to_string()], - StackMergeOrder::stack(columnars).into(), + StackMergeOrder::stack(&columnar_refs).into(), &mut buffer, ) .unwrap(); assert!( - buffer.len() <= 600, + buffer.len() <= 35_000, "merged columnar is {} bytes", buffer.len() ); @@ -115,11 +107,10 @@ fn test_merge_generated_tie_breaker_columns_stays_small() { .open_u64_lenient() .unwrap() .unwrap(); - assert_eq!(merged.num_docs(), 20_000); + assert_eq!(merged.num_docs(), 1_000_000); assert_eq!(column.get_cardinality(), Cardinality::Full); - for block_start in (0..20_000).step_by(512) { - let block_end = (block_start + 512).min(20_000); - for doc in block_start + 1..block_end { + for split_start in (0..1_000_000).step_by(250_000) { + for doc in split_start + 1..split_start + 250_000 { assert_eq!(column.first(doc), Some(column.first(doc - 1).unwrap() + 1)); } } @@ -130,12 +121,9 @@ fn test_group_columns_required_column_with_no_existing_columns() { let columnar1 = make_columnar("numbers", &[2u64]); let columnar2 = make_columnar("numbers", &[2u64]); let columnars = &[&columnar1, &columnar2]; - let column_map: BTreeMap<_, _> = group_columns_for_merge( - columnars, - &[("required_col".to_string(), ColumnType::Str)], - &[], - ) - .unwrap(); + let column_map: BTreeMap<_, _> = + group_columns_for_merge(columnars, &[("required_col".to_string(), ColumnType::Str)]) + .unwrap(); assert_eq!(column_map.len(), 2); let columns = &column_map .get(&("required_col".to_string(), ColumnTypeCategory::Str)) @@ -152,8 +140,7 @@ fn test_group_columns_required_column_is_above_all_columns_have_the_same_type_ru let columnar2 = make_columnar("numbers", &[2i64]); let columnars = &[&columnar1, &columnar2]; let column_map: BTreeMap<(String, ColumnTypeCategory), GroupedColumnsHandle> = - group_columns_for_merge(columnars, &[("numbers".to_string(), ColumnType::U64)], &[]) - .unwrap(); + group_columns_for_merge(columnars, &[("numbers".to_string(), ColumnType::U64)]).unwrap(); assert_eq!(column_map.len(), 1); assert!(column_map.contains_key(&("numbers".to_string(), ColumnTypeCategory::Numerical))); } @@ -164,7 +151,7 @@ fn test_missing_column() { let columnar2 = make_columnar("numbers2", &[2u64]); let columnars = &[&columnar1, &columnar2]; let column_map: BTreeMap<(String, ColumnTypeCategory), GroupedColumnsHandle> = - group_columns_for_merge(columnars, &[], &[]).unwrap(); + group_columns_for_merge(columnars, &[]).unwrap(); assert_eq!(column_map.len(), 2); assert!(column_map.contains_key(&("numbers".to_string(), ColumnTypeCategory::Numerical))); { @@ -268,7 +255,6 @@ fn test_merge_columnar_numbers() { crate::columnar::merge_columnar( columnars, &[], - &[], MergeRowOrder::Stack(stack_merge_order), &mut buffer, ) @@ -297,7 +283,6 @@ fn test_merge_columnar_texts() { crate::columnar::merge_columnar( columnars, &[], - &[], MergeRowOrder::Stack(stack_merge_order), &mut buffer, ) @@ -347,7 +332,6 @@ fn test_merge_columnar_byte() { crate::columnar::merge_columnar( columnars, &[], - &[], MergeRowOrder::Stack(stack_merge_order), &mut buffer, ) @@ -404,7 +388,6 @@ fn test_merge_columnar_byte_with_missing() { crate::columnar::merge_columnar( columnars, &[], - &[], MergeRowOrder::Stack(stack_merge_order), &mut buffer, ) @@ -457,7 +440,6 @@ fn test_merge_columnar_different_types() { crate::columnar::merge_columnar( columnars, &[], - &[], MergeRowOrder::Stack(stack_merge_order), &mut buffer, ) @@ -523,7 +505,6 @@ fn test_merge_columnar_different_empty_cardinality() { crate::columnar::merge_columnar( columnars, &[], - &[], MergeRowOrder::Stack(stack_merge_order), &mut buffer, ) @@ -634,7 +615,6 @@ proptest! { merge_columnar( &columnar_refs, &[], - &[], MergeRowOrder::Stack(stack_merge_order), &mut out, ).unwrap(); @@ -652,7 +632,6 @@ proptest! { merge_columnar( &columnar_refs, &[], - &[], MergeRowOrder::Stack(stack_merge_order), &mut out, ).unwrap(); diff --git a/columnar/src/columnar/reader/mod.rs b/columnar/src/columnar/reader/mod.rs index 30734ea79..198f659de 100644 --- a/columnar/src/columnar/reader/mod.rs +++ b/columnar/src/columnar/reader/mod.rs @@ -247,11 +247,8 @@ mod tests { let handles = columnar.read_columns("tie").unwrap(); let column = handles[0].open_u64_lenient().unwrap().unwrap(); assert_eq!(column.index.get_cardinality(), crate::Cardinality::Full); - for block_start in [0, 512, 1_024] { - let block_end = (block_start + 512).min(1_025); - for doc in block_start + 1..block_end { - assert_eq!(column.first(doc), Some(column.first(doc - 1).unwrap() + 1)); - } + for doc in 1..1_025 { + assert_eq!(column.first(doc), Some(column.first(doc - 1).unwrap() + 1)); } } diff --git a/columnar/src/compat_tests.rs b/columnar/src/compat_tests.rs index 4f8916927..64e615973 100644 --- a/columnar/src/compat_tests.rs +++ b/columnar/src/compat_tests.rs @@ -74,14 +74,7 @@ fn test_format(path: &str) { let columnar_readers = vec![&reader, &reader2]; let merge_row_order = StackMergeOrder::stack(&columnar_readers[..]); let mut out = Vec::new(); - merge_columnar( - &columnar_readers, - &[], - &[], - merge_row_order.into(), - &mut out, - ) - .unwrap(); + merge_columnar(&columnar_readers, &[], merge_row_order.into(), &mut out).unwrap(); let reader = ColumnarReader::open(out).unwrap(); check_columns(&reader); } diff --git a/columnar/src/tests.rs b/columnar/src/tests.rs index 3990290b9..42caa6fe1 100644 --- a/columnar/src/tests.rs +++ b/columnar/src/tests.rs @@ -833,7 +833,7 @@ proptest! { let columnar_readers_arr: Vec<&ColumnarReader> = columnar_readers.iter().collect(); let mut output: Vec = Vec::new(); let stack_merge_order = StackMergeOrder::stack(&columnar_readers_arr[..]).into(); - crate::merge_columnar(&columnar_readers_arr[..], &[], &[], stack_merge_order, &mut output).unwrap(); + crate::merge_columnar(&columnar_readers_arr[..], &[], stack_merge_order, &mut output).unwrap(); let merged_columnar = ColumnarReader::open(output).unwrap(); let concat_rows: Vec> = columnar_docs.iter().flatten().cloned().collect(); let expected_merged_columnar = build_columnar(&concat_rows[..]); @@ -855,7 +855,6 @@ fn test_columnar_merging_empty_columnar() { crate::merge_columnar( &columnar_readers_arr[..], &[], - &[], crate::MergeRowOrder::Stack(stack_merge_order), &mut output, ) @@ -893,7 +892,6 @@ fn test_columnar_merging_number_columns() { crate::merge_columnar( &columnar_readers_arr[..], &[], - &[], crate::MergeRowOrder::Stack(stack_merge_order), &mut output, ) @@ -967,7 +965,6 @@ fn test_columnar_merge_and_remap( crate::merge_columnar( &columnar_readers_ref[..], &[], - &[], shuffle_merge_order.into(), &mut output, ) @@ -1010,7 +1007,6 @@ fn test_columnar_merge_empty() { crate::merge_columnar( &[&columnar_reader_1, &columnar_reader_2], &[], - &[], shuffle_merge_order.into(), &mut output, ) @@ -1037,7 +1033,6 @@ fn test_columnar_merge_single_str_column() { crate::merge_columnar( &[&columnar_reader_1, &columnar_reader_2], &[], - &[], shuffle_merge_order.into(), &mut output, ) @@ -1070,7 +1065,6 @@ fn test_delete_decrease_cardinality() { crate::merge_columnar( &[&columnar_reader_1, &columnar_reader_2], &[], - &[], shuffle_merge_order.into(), &mut output, ) diff --git a/src/fastfield/mod.rs b/src/fastfield/mod.rs index df3cf270a..243f3571a 100644 --- a/src/fastfield/mod.rs +++ b/src/fastfield/mod.rs @@ -174,11 +174,8 @@ mod tests { let readers = FastFieldReaders::open(bytes.into(), schema).unwrap(); let values = readers.u64("tie").unwrap().first_or_default_col(0); - for block_start in [0, 512, 1_024] { - let block_end = (block_start + 512).min(1_025); - for doc in block_start + 1..block_end { - assert_eq!(values.get_val(doc), values.get_val(doc - 1) + 1); - } + for doc in 1..1_025 { + assert_eq!(values.get_val(doc), values.get_val(doc - 1) + 1); } } diff --git a/src/fastfield/plugin.rs b/src/fastfield/plugin.rs index bdd8dc19a..eb33b6816 100644 --- a/src/fastfield/plugin.rs +++ b/src/fastfield/plugin.rs @@ -18,7 +18,7 @@ use crate::index::{SegmentComponent, SegmentReader}; use crate::indexer::doc_id_mapping::{DocIdMapping, MappingType, SegmentDocIdMapping}; use crate::plugin::{PluginMergeContext, PluginWriter, PluginWriterContext, SegmentPlugin}; use crate::schema::document::Document; -use crate::schema::{value_type_to_column_type, FieldType, Schema}; +use crate::schema::{value_type_to_column_type, Schema}; use crate::space_usage::{ComponentSpaceUsage, FAST_FIELDS}; use crate::Segment; @@ -43,7 +43,6 @@ impl SegmentPlugin for FastFieldsPlugin { ctx.target_segment.index().directory().open_write(&path)?; let required_columns = extract_fast_field_required_columns(ctx.schema); - let tie_breaker_columns = extract_tie_breaker_columns(ctx.schema); let columnars: Vec<&ColumnarReader> = ctx .readers .iter() @@ -57,7 +56,6 @@ impl SegmentPlugin for FastFieldsPlugin { columnar::merge_columnar( &columnars[..], &required_columns, - &tie_breaker_columns, merge_row_order, &mut fast_field_wrt, )?; @@ -175,16 +173,6 @@ fn convert_to_merge_order( } } -fn extract_tie_breaker_columns(schema: &Schema) -> Vec { - schema - .fields() - .filter_map(|(_, field_entry)| { - matches!(field_entry.field_type(), FieldType::TieBreaker) - .then(|| field_entry.name().to_string()) - }) - .collect() -} - fn extract_fast_field_required_columns(schema: &Schema) -> Vec<(String, ColumnType)> { schema .fields() From 6a6378fa1fd8c66899948970d6e00c4f0a46bc46 Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Thu, 1 Oct 2026 14:24:04 +0200 Subject: [PATCH 5/9] feat: add generated tie-breaker fast fields --- CHANGELOG.md | 4 ++++ columnar/src/column/serialize.rs | 12 +++++------- columnar/src/column_values/mod.rs | 2 +- .../u64_based/blockwise_linear.rs | 2 +- columnar/src/column_values/u64_based/mod.rs | 2 +- src/schema/document/default_document.rs | 18 +++++++++++++++++- src/schema/field_type.rs | 2 ++ src/schema/schema.rs | 4 ++++ 8 files changed, 35 insertions(+), 11 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 395dbc6b5..1056470b4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -3,6 +3,10 @@ Tantivy 0.27.0 ## Breaking change - `.set_fast(..)` now takes a &str. The same behavior as `.set_fast(None)` can be obtained with .set_fast(tantivy::tokenizer::RAW_TOKENIZER_NAME). +- Added the `FieldType::TieBreaker` variant; exhaustive matches on `FieldType` need a new arm. + +## Features/Improvements +- Add generated tie-breaker fast fields via `SchemaBuilder::add_tie_breaker_field`. Values are consecutive within a segment, start at a random offset, and are preserved across merges. Tantivy 0.26.2 diff --git a/columnar/src/column/serialize.rs b/columnar/src/column/serialize.rs index f18d18199..890832c7b 100644 --- a/columnar/src/column/serialize.rs +++ b/columnar/src/column/serialize.rs @@ -30,15 +30,13 @@ pub(crate) fn serialize_generated_tie_breaker_column( num_docs: u32, output: &mut impl Write, ) -> io::Result<()> { - let max_start = u32::MAX - num_docs.saturating_sub(1); - let random = rand::rng().random::(); - let start = ((random as u64 * (u64::from(max_start) + 1)) >> 32) as u32; - let values: Vec = (0..num_docs) - .map(|offset| (start + offset) as u64) - .collect(); + let max_start = (u32::MAX - num_docs.saturating_sub(1)) as u64; + let start: u64 = rand::rng().random_range(0..=max_start); + let end = start + num_docs as u64; + let values = start..end; let column_index_num_bytes = serialize_column_index(SerializableColumnIndex::Full, output)?; serialize_u64_based_column_values( - &&values[..], + &values, &[ CodecType::Bitpacked, CodecType::Linear, diff --git a/columnar/src/column_values/mod.rs b/columnar/src/column_values/mod.rs index 87d8cca32..486f2e385 100644 --- a/columnar/src/column_values/mod.rs +++ b/columnar/src/column_values/mod.rs @@ -19,7 +19,7 @@ pub(crate) mod monotonic_mapping; pub(crate) mod monotonic_mapping_u128; mod stats; mod u128_based; -pub(crate) mod u64_based; +mod u64_based; mod vec_column; mod monotonic_column; diff --git a/columnar/src/column_values/u64_based/blockwise_linear.rs b/columnar/src/column_values/u64_based/blockwise_linear.rs index b8d03794f..b60bf5bad 100644 --- a/columnar/src/column_values/u64_based/blockwise_linear.rs +++ b/columnar/src/column_values/u64_based/blockwise_linear.rs @@ -11,7 +11,7 @@ use crate::column_values::u64_based::line::Line; use crate::column_values::u64_based::{ColumnCodec, ColumnCodecEstimator, ColumnStats}; use crate::column_values::{ColumnValues, VecColumn}; -pub(crate) const BLOCK_SIZE: u32 = 512u32; +const BLOCK_SIZE: u32 = 512u32; #[derive(Debug, Default)] struct Block { diff --git a/columnar/src/column_values/u64_based/mod.rs b/columnar/src/column_values/u64_based/mod.rs index 59ae64e09..313cd2251 100644 --- a/columnar/src/column_values/u64_based/mod.rs +++ b/columnar/src/column_values/u64_based/mod.rs @@ -1,5 +1,5 @@ mod bitpacked; -pub(crate) mod blockwise_linear; +mod blockwise_linear; mod line; mod linear; mod stats_collector; diff --git a/src/schema/document/default_document.rs b/src/schema/document/default_document.rs index be66956e3..de4e0e824 100644 --- a/src/schema/document/default_document.rs +++ b/src/schema/document/default_document.rs @@ -12,7 +12,7 @@ use crate::schema::document::{ DeserializeError, Document, DocumentDeserialize, DocumentDeserializer, }; use crate::schema::field_type::ValueParsingError; -use crate::schema::{Facet, Field, NamedFieldDocument, OwnedValue, Schema}; +use crate::schema::{Facet, Field, FieldType, NamedFieldDocument, OwnedValue, Schema}; use crate::tokenizer::PreTokenizedString; #[repr(C, packed)] @@ -219,6 +219,9 @@ impl CompactDoc { if let Ok(field) = schema.get_field(&field_name) { let field_entry = schema.get_field_entry(field); let field_type = field_entry.field_type(); + if matches!(field_type, FieldType::TieBreaker) { + continue; + } match json_value { serde_json::Value::Array(json_items) => { for json_item in json_items { @@ -752,6 +755,19 @@ mod tests { let _json = doc.to_named_doc(&schema); } + #[test] + fn test_parse_json_ignores_tie_breaker_values() { + let mut schema_builder = Schema::builder(); + let title = schema_builder.add_text_field("title", TEXT); + let tie = schema_builder.add_tie_breaker_field("tie"); + let schema = schema_builder.build(); + let doc = + TantivyDocument::parse_json(&schema, r#"{"title": "hello", "tie": [1, "two", null]}"#) + .unwrap(); + assert!(doc.get_first(title).is_some()); + assert!(doc.get_first(tie).is_none()); + } + #[test] fn test_json_value() { let json_str = r#"{ diff --git a/src/schema/field_type.rs b/src/schema/field_type.rs index 30d1e6af1..f1297a3c2 100644 --- a/src/schema/field_type.rs +++ b/src/schema/field_type.rs @@ -191,6 +191,8 @@ pub enum FieldType { /// Unsigned 64-bits integers field type configuration U64(NumericOptions), /// Generated tie-breaker field, exposed as a `u64` fast field. + /// + /// Values are almost always distinct, but not guaranteed to be unique across segments. TieBreaker, /// Signed 64-bits integers 64 field type configuration I64(NumericOptions), diff --git a/src/schema/schema.rs b/src/schema/schema.rs index ce7bb5804..3ce65ae97 100644 --- a/src/schema/schema.rs +++ b/src/schema/schema.rs @@ -62,6 +62,10 @@ impl SchemaBuilder { /// The field is exposed as a `u64` fast field, but its generated values fit in a `u32`. /// Values supplied by documents for this field are ignored. /// + /// Values are consecutive within a segment, starting at a random offset. Ranges of + /// different segments may overlap, so values are almost always distinct but not + /// guaranteed to be unique. + /// /// # Panics /// /// Panics when field already exists. From af3ad7625115589fa77c4d0fcc77747f82597763 Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Thu, 1 Oct 2026 14:30:35 +0200 Subject: [PATCH 6/9] Reject tie-breaker fields as index sort fields --- src/index/index.rs | 6 ++++++ src/indexer/doc_id_mapping.rs | 23 +++++++++++++++++++++++ 2 files changed, 29 insertions(+) diff --git a/src/index/index.rs b/src/index/index.rs index 0fc6993fb..250674775 100644 --- a/src/index/index.rs +++ b/src/index/index.rs @@ -320,6 +320,12 @@ impl IndexBuilder { )) })?; let entry = schema.get_field_entry(schema_field); + if matches!(entry.field_type(), FieldType::TieBreaker) { + return Err(TantivyError::InvalidArgument(format!( + "Field {} is a tie-breaker field and cannot be used to sort an index", + sort_by_field.field + ))); + } if !entry.is_fast() { return Err(TantivyError::InvalidArgument(format!( "Field {} is no fast field. Field needs to be a single value fast field \ diff --git a/src/indexer/doc_id_mapping.rs b/src/indexer/doc_id_mapping.rs index 4819330da..dd5739785 100644 --- a/src/indexer/doc_id_mapping.rs +++ b/src/indexer/doc_id_mapping.rs @@ -721,6 +721,29 @@ mod tests_indexsorting { assert!(matches!(error, TantivyError::InvalidArgument(_))); } + #[test] + fn test_index_builder_rejects_tie_breaker_sort_by_field() { + for order in [Order::Asc, Order::Desc] { + let mut schema_builder = Schema::builder(); + schema_builder.add_tie_breaker_field("tie"); + let schema = schema_builder.build(); + let settings = IndexSettings { + sort_by_field: Some(IndexSortByField { + field: "tie".to_string(), + order, + }), + ..Default::default() + }; + + let error = Index::builder() + .schema(schema) + .settings(settings) + .create_in_ram() + .unwrap_err(); + assert!(matches!(error, TantivyError::InvalidArgument(_))); + } + } + #[test] fn test_doc_mapping() { let doc_mapping = DocIdMapping::from_new_id_to_old_id(vec![3, 2, 5]); From a68b80ea9468ce7b47b7f97b28738d14c463b27d Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Thu, 1 Oct 2026 16:25:38 +0200 Subject: [PATCH 7/9] Draw tie-breaker starts from the full u64 range --- columnar/src/column/serialize.rs | 2 +- columnar/src/columnar/merge/tests.rs | 4 ++-- src/schema/schema.rs | 4 ++-- 3 files changed, 5 insertions(+), 5 deletions(-) diff --git a/columnar/src/column/serialize.rs b/columnar/src/column/serialize.rs index 890832c7b..dbe787a4d 100644 --- a/columnar/src/column/serialize.rs +++ b/columnar/src/column/serialize.rs @@ -30,7 +30,7 @@ pub(crate) fn serialize_generated_tie_breaker_column( num_docs: u32, output: &mut impl Write, ) -> io::Result<()> { - let max_start = (u32::MAX - num_docs.saturating_sub(1)) as u64; + let max_start = u64::MAX - u64::from(num_docs); let start: u64 = rand::rng().random_range(0..=max_start); let end = start + num_docs as u64; let values = start..end; diff --git a/columnar/src/columnar/merge/tests.rs b/columnar/src/columnar/merge/tests.rs index 78c2bd804..04c54fa3f 100644 --- a/columnar/src/columnar/merge/tests.rs +++ b/columnar/src/columnar/merge/tests.rs @@ -78,7 +78,7 @@ fn test_merge_generated_tie_breaker_columns_stays_small() { let mut buffer = Vec::new(); writer.serialize(num_docs, None, &mut buffer).unwrap(); assert!( - buffer.len() <= 100, + buffer.len() <= 128, "input columnar is {} bytes", buffer.len() ); @@ -98,7 +98,7 @@ fn test_merge_generated_tie_breaker_columns_stays_small() { .unwrap(); assert!( - buffer.len() <= 35_000, + buffer.len() <= 45_000, "merged columnar is {} bytes", buffer.len() ); diff --git a/src/schema/schema.rs b/src/schema/schema.rs index 3ce65ae97..74f30dcb7 100644 --- a/src/schema/schema.rs +++ b/src/schema/schema.rs @@ -59,8 +59,8 @@ impl SchemaBuilder { /// Adds a generated tie-breaker fast field. /// - /// The field is exposed as a `u64` fast field, but its generated values fit in a `u32`. - /// Values supplied by documents for this field are ignored. + /// The field is exposed as a `u64` fast field and its generated values span the full `u64` + /// range. Values supplied by documents for this field are ignored. /// /// Values are consecutive within a segment, starting at a random offset. Ranges of /// different segments may overlap, so values are almost always distinct but not From 547524adbe0426a71a7a35d835e6be753bcabc11 Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Fri, 2 Oct 2026 09:08:06 +0200 Subject: [PATCH 8/9] fix: keep tiebreaker values u32-compatible for now --- columnar/src/column/serialize.rs | 3 ++- src/schema/schema.rs | 5 +++-- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/columnar/src/column/serialize.rs b/columnar/src/column/serialize.rs index dbe787a4d..9747efb8e 100644 --- a/columnar/src/column/serialize.rs +++ b/columnar/src/column/serialize.rs @@ -30,7 +30,8 @@ pub(crate) fn serialize_generated_tie_breaker_column( num_docs: u32, output: &mut impl Write, ) -> io::Result<()> { - let max_start = u64::MAX - u64::from(num_docs); + // TODO: Lift this temporary u32 limit once downstream consumers support the full u64 range. + let max_start = u64::from(u32::MAX - num_docs.saturating_sub(1)); let start: u64 = rand::rng().random_range(0..=max_start); let end = start + num_docs as u64; let values = start..end; diff --git a/src/schema/schema.rs b/src/schema/schema.rs index 74f30dcb7..48af9a750 100644 --- a/src/schema/schema.rs +++ b/src/schema/schema.rs @@ -59,8 +59,9 @@ impl SchemaBuilder { /// Adds a generated tie-breaker fast field. /// - /// The field is exposed as a `u64` fast field and its generated values span the full `u64` - /// range. Values supplied by documents for this field are ignored. + /// The field is exposed as a `u64` fast field. Its generated values are temporarily limited + /// to the `u32` range for downstream compatibility; this limit will eventually be lifted. + /// Values supplied by documents for this field are ignored. /// /// Values are consecutive within a segment, starting at a random offset. Ranges of /// different segments may overlap, so values are almost always distinct but not From ef0bed7ea33bdf98dc0783ba22fc344a54519153 Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Fri, 2 Oct 2026 09:18:57 +0200 Subject: [PATCH 9/9] Simplify tie-breaker start bound and doc comment --- columnar/src/column/serialize.rs | 2 +- src/schema/schema.rs | 10 ++++------ 2 files changed, 5 insertions(+), 7 deletions(-) diff --git a/columnar/src/column/serialize.rs b/columnar/src/column/serialize.rs index 9747efb8e..f50605564 100644 --- a/columnar/src/column/serialize.rs +++ b/columnar/src/column/serialize.rs @@ -31,7 +31,7 @@ pub(crate) fn serialize_generated_tie_breaker_column( output: &mut impl Write, ) -> io::Result<()> { // TODO: Lift this temporary u32 limit once downstream consumers support the full u64 range. - let max_start = u64::from(u32::MAX - num_docs.saturating_sub(1)); + let max_start = (u32::MAX - num_docs) as u64; let start: u64 = rand::rng().random_range(0..=max_start); let end = start + num_docs as u64; let values = start..end; diff --git a/src/schema/schema.rs b/src/schema/schema.rs index 48af9a750..745a1ec30 100644 --- a/src/schema/schema.rs +++ b/src/schema/schema.rs @@ -59,13 +59,11 @@ impl SchemaBuilder { /// Adds a generated tie-breaker fast field. /// - /// The field is exposed as a `u64` fast field. Its generated values are temporarily limited - /// to the `u32` range for downstream compatibility; this limit will eventually be lifted. - /// Values supplied by documents for this field are ignored. + /// The field is exposed as a `u64` fast field. Values are consecutive within a segment, + /// starting at a random offset. Ranges of different segments may overlap, so values are + /// almost always distinct but not guaranteed to be unique. /// - /// Values are consecutive within a segment, starting at a random offset. Ranges of - /// different segments may overlap, so values are almost always distinct but not - /// guaranteed to be unique. + /// Values supplied by documents for this field are ignored. /// /// # Panics ///