diff --git a/columnar/src/columnar/writer/column_writers.rs b/columnar/src/columnar/writer/column_writers.rs index ff9308a7c..8d24c2562 100644 --- a/columnar/src/columnar/writer/column_writers.rs +++ b/columnar/src/columnar/writer/column_writers.rs @@ -1,11 +1,9 @@ use std::cmp::Ordering; -use std::mem::size_of; use stacker::{ExpUnrolledLinkedList, MemoryArena}; use crate::columnar::writer::column_operation::{ColumnOperation, SymbolValue}; -use crate::dictionary::DictionaryBuilder; -use crate::{Cardinality, NumericalType, NumericalValue, PayloadEncoding, RowId}; +use crate::{Cardinality, NumericalType, NumericalValue, RowId}; #[derive(Copy, Clone, Debug, Eq, PartialEq)] #[repr(u8)] @@ -241,138 +239,6 @@ impl NumericalColumnWriter { } } -#[derive(Copy, Clone)] -pub(crate) struct StrOrBytesColumnWriter { - payload: PayloadColumnWriter, - // If true, when facing a multivalued cardinality, - // values associated to a given document will be sorted. - // - // This is useful for facets. - // - // If false, the order of appearance in the document will be - // observed. - pub(crate) sort_values_within_row: bool, -} - -#[derive(Copy, Clone)] -pub(crate) enum PayloadColumnWriter { - Dictionary(DictionaryEncodedColumnWriter), - Plain(PlainColumnWriter), -} - -#[derive(Copy, Clone)] -pub(crate) struct DictionaryEncodedColumnWriter { - pub(crate) dictionary_id: u32, - column_writer: ColumnWriter, -} - -#[derive(Copy, Clone, Default)] -pub(crate) struct PlainColumnWriter { - column_writer: ColumnWriter, - pub(crate) value_store_id: u32, -} - -#[derive(Default)] -pub(crate) struct PlainValueStore { - concatenated_payloads: Vec, - end_offsets: Vec, -} - -impl PlainValueStore { - fn push(&mut self, value: &[u8]) -> u32 { - let value_id = self.end_offsets.len() as u32; - self.concatenated_payloads.extend_from_slice(value); - self.end_offsets.push(self.concatenated_payloads.len()); - value_id - } - - pub(crate) fn get(&self, value_id: u32) -> &[u8] { - let value_id = value_id as usize; - let end = self.end_offsets[value_id]; - let start = if value_id == 0 { - 0 - } else { - self.end_offsets[value_id - 1] - }; - &self.concatenated_payloads[start..end] - } - - pub(crate) fn mem_usage(&self) -> usize { - self.concatenated_payloads.capacity() + self.end_offsets.capacity() * size_of::() - } -} - -impl StrOrBytesColumnWriter { - pub(crate) fn dictionary(dictionary_id: u32) -> StrOrBytesColumnWriter { - StrOrBytesColumnWriter { - payload: PayloadColumnWriter::Dictionary(DictionaryEncodedColumnWriter { - dictionary_id, - column_writer: ColumnWriter::default(), - }), - sort_values_within_row: false, - } - } - - pub(crate) fn plain(value_store_id: u32) -> StrOrBytesColumnWriter { - StrOrBytesColumnWriter { - payload: PayloadColumnWriter::Plain(PlainColumnWriter { - value_store_id, - ..PlainColumnWriter::default() - }), - sort_values_within_row: false, - } - } - - pub(crate) fn encoding(self) -> PayloadEncoding { - match self.payload { - PayloadColumnWriter::Dictionary(_) => PayloadEncoding::Dictionary, - PayloadColumnWriter::Plain(_) => PayloadEncoding::Plain, - } - } - - pub(crate) fn payload(self) -> PayloadColumnWriter { - self.payload - } - - pub(crate) fn column_writer(self) -> ColumnWriter { - match self.payload { - PayloadColumnWriter::Dictionary(writer) => writer.column_writer, - PayloadColumnWriter::Plain(writer) => writer.column_writer, - } - } - - pub(crate) fn record_bytes( - &mut self, - doc: RowId, - bytes: &[u8], - dictionaries: &mut [DictionaryBuilder], - plain_value_stores: &mut [PlainValueStore], - arena: &mut MemoryArena, - ) { - match &mut self.payload { - PayloadColumnWriter::Dictionary(writer) => { - let unordered_id = - dictionaries[writer.dictionary_id as usize].get_or_allocate_id(bytes, arena); - writer.column_writer.record(doc, unordered_id.0, arena); - } - PayloadColumnWriter::Plain(writer) => { - let value_id = plain_value_stores[writer.value_store_id as usize].push(bytes); - writer.column_writer.record(doc, value_id, arena); - } - } - } - - pub(super) fn operation_iterator<'a>( - &self, - arena: &MemoryArena, - old_to_new_ids: Option<&[RowId]>, - byte_buffer: &'a mut Vec, - ) -> impl Iterator> + 'a + use<'a> { - self.column_writer() - .operation_iterator(arena, old_to_new_ids, byte_buffer) - } -} - #[cfg(test)] mod tests { use super::*; @@ -389,20 +255,6 @@ mod tests { assert_eq!(delta_with_last_doc(Some(1u32), 4u32), DocumentStep::Skipped); } - #[test] - fn test_plain_value_store() { - let mut store = PlainValueStore::default(); - let first = store.push(b"same"); - let empty = store.push(b""); - let duplicate = store.push(b"same"); - - assert_eq!((first, empty, duplicate), (0, 1, 2)); - assert_eq!(store.get(first), b"same"); - assert_eq!(store.get(empty), b""); - assert_eq!(store.get(duplicate), b"same"); - assert!(store.mem_usage() >= 8); - } - #[track_caller] fn test_column_writer_coercion_iter_aux( values: impl Iterator, diff --git a/columnar/src/columnar/writer/dictionary_encoding.rs b/columnar/src/columnar/writer/dictionary_encoding.rs new file mode 100644 index 000000000..64e2ee27d --- /dev/null +++ b/columnar/src/columnar/writer/dictionary_encoding.rs @@ -0,0 +1,65 @@ +use stacker::{ArenaHashMap, MemoryArena}; + +use super::column_operation::ColumnOperation; +use super::column_writers::ColumnWriter; +use crate::RowId; +use crate::dictionary::DictionaryBuilder; + +#[derive(Default)] +pub(super) struct DictionaryEncodedColumnsWriter { + pub(super) bytes_columns: ArenaHashMap, + pub(super) str_columns: ArenaHashMap, + pub(super) dictionaries: Vec, +} + +impl DictionaryEncodedColumnsWriter { + pub(super) fn mem_usage(&self) -> usize { + self.bytes_columns.mem_usage() + + self.str_columns.mem_usage() + + self + .dictionaries + .iter() + .map(DictionaryBuilder::mem_usage) + .sum::() + } +} + +#[derive(Copy, Clone)] +pub(super) struct DictionaryEncodedColumnWriter { + pub(super) dictionary_id: u32, + pub(super) column_writer: ColumnWriter, + // This is used in facets aggregation + pub(super) sort_values_within_row: bool, +} + +impl DictionaryEncodedColumnWriter { + pub(super) fn new(dictionary_id: u32) -> Self { + Self { + dictionary_id, + column_writer: ColumnWriter::default(), + sort_values_within_row: false, + } + } + + pub(super) fn record_bytes( + &mut self, + doc: RowId, + bytes: &[u8], + dictionaries: &mut [DictionaryBuilder], + arena: &mut MemoryArena, + ) { + let unordered_id = + dictionaries[self.dictionary_id as usize].get_or_allocate_id(bytes, arena); + self.column_writer.record(doc, unordered_id.0, arena); + } + + pub(super) fn operation_iterator<'a>( + self, + arena: &MemoryArena, + old_to_new_ids: Option<&[RowId]>, + byte_buffer: &'a mut Vec, + ) -> impl Iterator> + 'a + use<'a> { + self.column_writer + .operation_iterator(arena, old_to_new_ids, byte_buffer) + } +} diff --git a/columnar/src/columnar/writer/mod.rs b/columnar/src/columnar/writer/mod.rs index af32e980a..20d04a6b1 100644 --- a/columnar/src/columnar/writer/mod.rs +++ b/columnar/src/columnar/writer/mod.rs @@ -1,5 +1,7 @@ mod column_operation; mod column_writers; +mod dictionary_encoding; +mod plain; mod serializer; mod value_index; @@ -19,10 +21,11 @@ use crate::column_index::{ }; use crate::column_values::{MonotonicallyMappableToU64, MonotonicallyMappableToU128}; use crate::columnar::column_type::ColumnType; -use crate::columnar::writer::column_writers::{ - ColumnWriter, NumericalColumnWriter, PayloadColumnWriter, PlainValueStore, - StrOrBytesColumnWriter, +use crate::columnar::writer::column_writers::{ColumnWriter, NumericalColumnWriter}; +use crate::columnar::writer::dictionary_encoding::{ + DictionaryEncodedColumnWriter, DictionaryEncodedColumnsWriter, }; +use crate::columnar::writer::plain::{PlainColumnWriter, PlainColumnsWriter, PlainValueStore}; use crate::columnar::writer::value_index::{IndexBuilder, PreallocatedIndexBuilders}; use crate::dictionary::{DictionaryBuilder, TermIdMapping, UnorderedId}; use crate::value::{Coerce, NumericalType, NumericalValue}; @@ -74,13 +77,9 @@ pub struct ColumnarWriter { datetime_field_hash_map: ArenaHashMap, bool_field_hash_map: ArenaHashMap, ip_addr_field_hash_map: ArenaHashMap, - bytes_field_hash_map: ArenaHashMap, - str_field_hash_map: ArenaHashMap, + dictionary_encoded_columns: DictionaryEncodedColumnsWriter, + plain_columns: PlainColumnsWriter, arena: MemoryArena, - // Dictionaries used to store dictionary-encoded values. - dictionaries: Vec, - // Raw stores used by plain string and byte columns. - plain_value_stores: Vec, buffers: SpareBuffers, } @@ -89,20 +88,10 @@ impl ColumnarWriter { self.arena.mem_usage() + self.numerical_field_hash_map.mem_usage() + self.bool_field_hash_map.mem_usage() - + self.bytes_field_hash_map.mem_usage() - + self.str_field_hash_map.mem_usage() + self.ip_addr_field_hash_map.mem_usage() + self.datetime_field_hash_map.mem_usage() - + self - .dictionaries - .iter() - .map(|dict| dict.mem_usage()) - .sum::() - + self - .plain_value_stores - .iter() - .map(PlainValueStore::mem_usage) - .sum::() + + self.dictionary_encoded_columns.mem_usage() + + self.plain_columns.mem_usage() + self.buffers.mem_usage() } @@ -123,51 +112,66 @@ impl ColumnarWriter { .get::(sort_field.as_bytes()) }) else { - let str_or_bytes_column_opt = self - .str_field_hash_map - .get::(sort_field.as_bytes()) - .or_else(|| { - self.bytes_field_hash_map - .get::(sort_field.as_bytes()) - }); - let Some(str_or_bytes_column) = str_or_bytes_column_opt else { - return Vec::new(); - }; - - let mut symbols_buffer = Vec::new(); - return match str_or_bytes_column.payload() { - PayloadColumnWriter::Dictionary(writer) => { - let dictionary_builder = &self.dictionaries[writer.dictionary_id as usize]; - let term_id_mapping = dictionary_builder.build_term_id_mapping(&self.arena); - collect_sort_order_from_ops( - str_or_bytes_column.operation_iterator( - &self.arena, - None, - &mut symbols_buffer, - ), - num_docs, - reversed, - |unordered_id| Some(term_id_mapping.to_ord(UnorderedId(unordered_id)).0), - None, - |a, b| a.cmp(b), - ) - } - PayloadColumnWriter::Plain(writer) => { - let value_store = &self.plain_value_stores[writer.value_store_id as usize]; - collect_sort_order_from_ops( - str_or_bytes_column.operation_iterator( - &self.arena, - None, - &mut symbols_buffer, - ), - num_docs, - reversed, - Some, - None, - |left, right| compare_optional_plain_values(*left, *right, value_store), - ) - } - }; + let column_name = sort_field.as_bytes(); + if let Some(writer) = self + .dictionary_encoded_columns + .str_columns + .get::(column_name) + { + let dictionary_builder = + &self.dictionary_encoded_columns.dictionaries[writer.dictionary_id as usize]; + return collect_dictionary_sort_order( + writer, + dictionary_builder, + &self.arena, + num_docs, + reversed, + ); + } + if let Some(writer) = self + .plain_columns + .str_columns + .get::(column_name) + { + let value_store = &self.plain_columns.value_stores[writer.value_store_id as usize]; + return collect_plain_sort_order( + writer, + value_store, + &self.arena, + num_docs, + reversed, + ); + } + if let Some(writer) = self + .dictionary_encoded_columns + .bytes_columns + .get::(column_name) + { + let dictionary_builder = + &self.dictionary_encoded_columns.dictionaries[writer.dictionary_id as usize]; + return collect_dictionary_sort_order( + writer, + dictionary_builder, + &self.arena, + num_docs, + reversed, + ); + } + if let Some(writer) = self + .plain_columns + .bytes_columns + .get::(column_name) + { + let value_store = &self.plain_columns.value_stores[writer.value_store_id as usize]; + return collect_plain_sort_order( + writer, + value_store, + &self.arena, + num_docs, + reversed, + ); + } + return Vec::new(); }; let mut symbols_buffer = Vec::new(); collect_sort_order_from_ops( @@ -220,7 +224,8 @@ impl ColumnarWriter { /// Records a column type and the payload encoding for a string or byte column. /// /// Dictionary encoding is the only valid encoding for other column types. Re-registering a - /// string or byte column with a different encoding returns an error. + /// string or byte column with a different encoding returns an error. Plain payloads preserve + /// the recorded value order and ignore `sort_values_within_row`. pub fn record_column_type_with_encoding( &mut self, column_name: &str, @@ -240,48 +245,66 @@ impl ColumnarWriter { } match column_type { ColumnType::Str | ColumnType::Bytes => { - let (hash_map, dictionaries, plain_value_stores) = ( - if column_type == ColumnType::Str { - &mut self.str_field_hash_map - } else { - &mut self.bytes_field_hash_map - }, - &mut self.dictionaries, - &mut self.plain_value_stores, - ); - let existing_column = - hash_map.get::(column_name.as_bytes()); - match existing_column { - Some(column_writer) if column_writer.encoding() != encoding => { - return Err(invalid_input( - "column was already registered with a different payload encoding", - )); - } - _ => {} - } - hash_map.mutate_or_create( - column_name.as_bytes(), - |column_opt: Option| { - let mut column_writer = if let Some(column_writer) = column_opt { - column_writer - } else { - match encoding { - PayloadEncoding::Dictionary => { + let DictionaryEncodedColumnsWriter { + bytes_columns: dictionary_bytes_columns, + str_columns: dictionary_str_columns, + dictionaries, + } = &mut self.dictionary_encoded_columns; + let PlainColumnsWriter { + bytes_columns: plain_bytes_columns, + str_columns: plain_str_columns, + value_stores, + } = &mut self.plain_columns; + let (dictionary_columns, plain_columns) = if column_type == ColumnType::Str { + (dictionary_str_columns, plain_str_columns) + } else { + (dictionary_bytes_columns, plain_bytes_columns) + }; + + match encoding { + PayloadEncoding::Dictionary => { + if plain_columns + .get::(column_name.as_bytes()) + .is_some() + { + return Err(invalid_input( + "column was already registered with a different payload encoding", + )); + } + dictionary_columns.mutate_or_create( + column_name.as_bytes(), + |column_opt: Option| { + let mut column_writer = column_opt.unwrap_or_else(|| { let dictionary_id = dictionaries.len() as u32; dictionaries.push(DictionaryBuilder::default()); - StrOrBytesColumnWriter::dictionary(dictionary_id) - } - PayloadEncoding::Plain => { - let value_store_id = plain_value_stores.len() as u32; - plain_value_stores.push(PlainValueStore::default()); - StrOrBytesColumnWriter::plain(value_store_id) - } - } - }; - column_writer.sort_values_within_row = sort_values_within_row; - column_writer - }, - ); + DictionaryEncodedColumnWriter::new(dictionary_id) + }); + column_writer.sort_values_within_row = sort_values_within_row; + column_writer + }, + ); + } + PayloadEncoding::Plain => { + if dictionary_columns + .get::(column_name.as_bytes()) + .is_some() + { + return Err(invalid_input( + "column was already registered with a different payload encoding", + )); + } + plain_columns.mutate_or_create( + column_name.as_bytes(), + |column_opt: Option| { + column_opt.unwrap_or_else(|| { + let value_store_id = value_stores.len() as u32; + value_stores.push(PlainValueStore::default()); + PlainColumnWriter::new(value_store_id) + }) + }, + ); + } + } } ColumnType::Bool => { require_dictionary_encoding(encoding)?; @@ -378,54 +401,64 @@ impl ColumnarWriter { } pub fn record_str(&mut self, doc: RowId, column_name: &str, value: &str) { - let (hash_map, arena, dictionaries, plain_value_stores) = ( - &mut self.str_field_hash_map, - &mut self.arena, - &mut self.dictionaries, - &mut self.plain_value_stores, - ); - hash_map.mutate_or_create( + self.record_bytes_or_str(doc, column_name, value.as_bytes(), ColumnType::Str); + } + + pub fn record_bytes(&mut self, doc: RowId, column_name: &str, value: &[u8]) { + self.record_bytes_or_str(doc, column_name, value, ColumnType::Bytes); + } + + fn record_bytes_or_str( + &mut self, + doc: RowId, + column_name: &str, + value: &[u8], + column_type: ColumnType, + ) { + let DictionaryEncodedColumnsWriter { + bytes_columns: dictionary_bytes_columns, + str_columns: dictionary_str_columns, + dictionaries, + } = &mut self.dictionary_encoded_columns; + let PlainColumnsWriter { + bytes_columns: plain_bytes_columns, + str_columns: plain_str_columns, + value_stores, + } = &mut self.plain_columns; + let (dictionary_columns, plain_columns) = if column_type == ColumnType::Str { + (dictionary_str_columns, plain_str_columns) + } else { + (dictionary_bytes_columns, plain_bytes_columns) + }; + if plain_columns + .get::(column_name.as_bytes()) + .is_some() + { + plain_columns.mutate_or_create( + column_name.as_bytes(), + |column_opt: Option| { + let mut column = column_opt.unwrap(); + column.record_bytes(doc, value, value_stores, &mut self.arena); + column + }, + ); + return; + } + + dictionary_columns.mutate_or_create( column_name.as_bytes(), - |column_opt: Option| { - let mut column: StrOrBytesColumnWriter = column_opt.unwrap_or_else(|| { - // Each column has its own dictionary + |column_opt: Option| { + let mut column = column_opt.unwrap_or_else(|| { let dictionary_id = dictionaries.len() as u32; dictionaries.push(DictionaryBuilder::default()); - StrOrBytesColumnWriter::dictionary(dictionary_id) + DictionaryEncodedColumnWriter::new(dictionary_id) }); - column.record_bytes( - doc, - value.as_bytes(), - dictionaries, - plain_value_stores, - arena, - ); + column.record_bytes(doc, value, dictionaries, &mut self.arena); column }, ); } - pub fn record_bytes(&mut self, doc: RowId, column_name: &str, value: &[u8]) { - let (hash_map, arena, dictionaries, plain_value_stores) = ( - &mut self.bytes_field_hash_map, - &mut self.arena, - &mut self.dictionaries, - &mut self.plain_value_stores, - ); - hash_map.mutate_or_create( - column_name.as_bytes(), - |column_opt: Option| { - let mut column: StrOrBytesColumnWriter = column_opt.unwrap_or_else(|| { - // Each column has its own dictionary - let dictionary_id = dictionaries.len() as u32; - dictionaries.push(DictionaryBuilder::default()); - StrOrBytesColumnWriter::dictionary(dictionary_id) - }); - column.record_bytes(doc, value, dictionaries, plain_value_stores, arena); - column - }, - ); - } pub fn serialize( &mut self, num_docs: RowId, @@ -434,50 +467,86 @@ impl ColumnarWriter { ) -> io::Result<()> { let mut serializer = ColumnarSerializer::new(wrt); - let mut columns: Vec<(&[u8], ColumnType, Addr)> = self + let mut columns: Vec<(&[u8], ColumnType, Option, Addr)> = self .numerical_field_hash_map .iter() .map(|(column_name, addr)| { let numerical_column_writer: NumericalColumnWriter = self.numerical_field_hash_map.read(addr); let column_type = numerical_column_writer.numerical_type().into(); - (column_name, column_type, addr) + (column_name, column_type, None, addr) }) .collect(); + columns.extend(self.dictionary_encoded_columns.bytes_columns.iter().map( + |(column_name, addr)| { + ( + column_name, + ColumnType::Bytes, + Some(PayloadEncoding::Dictionary), + addr, + ) + }, + )); + columns.extend(self.dictionary_encoded_columns.str_columns.iter().map( + |(column_name, addr)| { + ( + column_name, + ColumnType::Str, + Some(PayloadEncoding::Dictionary), + addr, + ) + }, + )); columns.extend( - self.bytes_field_hash_map + self.plain_columns + .bytes_columns .iter() - .map(|(column_name, addr)| (column_name, ColumnType::Bytes, addr)), + .map(|(column_name, addr)| { + ( + column_name, + ColumnType::Bytes, + Some(PayloadEncoding::Plain), + addr, + ) + }), ); columns.extend( - self.str_field_hash_map + self.plain_columns + .str_columns .iter() - .map(|(column_name, addr)| (column_name, ColumnType::Str, addr)), + .map(|(column_name, addr)| { + ( + column_name, + ColumnType::Str, + Some(PayloadEncoding::Plain), + addr, + ) + }), ); columns.extend( self.bool_field_hash_map .iter() - .map(|(column_name, addr)| (column_name, ColumnType::Bool, addr)), + .map(|(column_name, addr)| (column_name, ColumnType::Bool, None, addr)), ); columns.extend( self.ip_addr_field_hash_map .iter() - .map(|(column_name, addr)| (column_name, ColumnType::IpAddr, addr)), + .map(|(column_name, addr)| (column_name, ColumnType::IpAddr, None, addr)), ); columns.extend( self.datetime_field_hash_map .iter() - .map(|(column_name, addr)| (column_name, ColumnType::DateTime, addr)), + .map(|(column_name, addr)| (column_name, ColumnType::DateTime, None, addr)), ); - columns.sort_unstable_by_key(|(column_name, col_type, _)| (*column_name, *col_type)); + columns.sort_unstable_by_key(|(column_name, col_type, _, _)| (*column_name, *col_type)); let (arena, buffers, dictionaries, plain_value_stores) = ( &self.arena, &mut self.buffers, - &self.dictionaries, - &self.plain_value_stores, + &self.dictionary_encoded_columns.dictionaries, + &self.plain_columns.value_stores, ); let mut symbol_byte_buffer: Vec = Vec::new(); - for (column_name, column_type, addr) in columns { + for (column_name, column_type, payload_encoding, addr) in columns { if column_name.contains(&JSON_END_OF_PATH) { // Tantivy uses b'0' as a separator for nested fields in JSON. // Column names with a b'0' are not simply ignored by the columnar (and the inverted @@ -522,26 +591,24 @@ impl ColumnarWriter { column_serializer.finalize()?; } ColumnType::Bytes | ColumnType::Str => { - let str_or_bytes_column_writer: StrOrBytesColumnWriter = - if column_type == ColumnType::Bytes { - self.bytes_field_hash_map.read(addr) - } else { - self.str_field_hash_map.read(addr) - }; - let cardinality = str_or_bytes_column_writer - .column_writer() - .get_cardinality(num_docs); let mut column_serializer = serializer.start_serialize_column(column_name, column_type); - match str_or_bytes_column_writer.payload() { - PayloadColumnWriter::Dictionary(writer) => { + match payload_encoding.unwrap() { + PayloadEncoding::Dictionary => { + let writer: DictionaryEncodedColumnWriter = + if column_type == ColumnType::Bytes { + self.dictionary_encoded_columns.bytes_columns.read(addr) + } else { + self.dictionary_encoded_columns.str_columns.read(addr) + }; + let cardinality = writer.column_writer.get_cardinality(num_docs); let dictionary_builder = &dictionaries[writer.dictionary_id as usize]; serialize_dictionary_bytes_or_str_column( cardinality, num_docs, - str_or_bytes_column_writer.sort_values_within_row, + writer.sort_values_within_row, dictionary_builder, - str_or_bytes_column_writer.operation_iterator( + writer.operation_iterator( arena, old_to_new_row_ids, &mut symbol_byte_buffer, @@ -551,14 +618,19 @@ impl ColumnarWriter { &mut column_serializer, )?; } - PayloadColumnWriter::Plain(writer) => { + PayloadEncoding::Plain => { + let writer: PlainColumnWriter = if column_type == ColumnType::Bytes { + self.plain_columns.bytes_columns.read(addr) + } else { + self.plain_columns.str_columns.read(addr) + }; + let cardinality = writer.column_writer.get_cardinality(num_docs); let value_store = &plain_value_stores[writer.value_store_id as usize]; serialize_plain_bytes_or_str_column( cardinality, num_docs, - str_or_bytes_column_writer.sort_values_within_row, value_store, - str_or_bytes_column_writer.operation_iterator( + writer.operation_iterator( arena, old_to_new_row_ids, &mut symbol_byte_buffer, @@ -617,6 +689,43 @@ impl ColumnarWriter { } } +fn collect_dictionary_sort_order( + writer: DictionaryEncodedColumnWriter, + dictionary_builder: &DictionaryBuilder, + arena: &MemoryArena, + num_docs: RowId, + reversed: bool, +) -> Vec { + let term_id_mapping = dictionary_builder.build_term_id_mapping(arena); + let mut symbols_buffer = Vec::new(); + collect_sort_order_from_ops( + writer.operation_iterator(arena, None, &mut symbols_buffer), + num_docs, + reversed, + |unordered_id| Some(term_id_mapping.to_ord(UnorderedId(unordered_id)).0), + None, + |a, b| a.cmp(b), + ) +} + +fn collect_plain_sort_order( + writer: PlainColumnWriter, + value_store: &PlainValueStore, + arena: &MemoryArena, + num_docs: RowId, + reversed: bool, +) -> Vec { + let mut symbols_buffer = Vec::new(); + collect_sort_order_from_ops( + writer.operation_iterator(arena, None, &mut symbols_buffer), + num_docs, + reversed, + Some, + None, + |left, right| compare_optional_plain_values(*left, *right, value_store), + ) +} + /// Shared sorting pattern for both numeric and Str/Bytes sort fields. /// /// Iterates column operations, fills gaps for missing docs with `default_key`, converts each value @@ -718,7 +827,6 @@ fn serialize_dictionary_bytes_or_str_column( fn serialize_plain_bytes_or_str_column( cardinality: Cardinality, num_docs: RowId, - sort_values_within_row: bool, value_store: &PlainValueStore, operation_it: impl Iterator>, buffers: &mut SpareBuffers, @@ -758,13 +866,6 @@ fn serialize_plain_bytes_or_str_column( let multivalued_index_builder = value_index_builders.borrow_multivalued_index_builder(); consume_operation_iterator(operation_it, multivalued_index_builder, plain_value_ids); let serializable_multivalued_index = multivalued_index_builder.finish(num_docs); - if sort_values_within_row { - sort_plain_values_within_row( - serializable_multivalued_index.start_offsets.boxed_iter(), - plain_value_ids, - value_store, - ); - } SerializableColumnIndex::Multivalued(serializable_multivalued_index) } }; @@ -874,20 +975,6 @@ fn flush_plain_block( Ok(()) } -fn sort_plain_values_within_row( - multivalued_index: impl Iterator, - values: &mut [u32], - value_store: &PlainValueStore, -) { - let mut start_index = 0usize; - for end_index in multivalued_index { - let end_index = end_index as usize; - values[start_index..end_index] - .sort_unstable_by(|left, right| value_store.get(*left).cmp(value_store.get(*right))); - start_index = end_index; - } -} - fn compare_optional_plain_values( left: Option, right: Option, @@ -1156,7 +1243,26 @@ mod tests { use stacker::MemoryArena; use crate::columnar::writer::column_operation::ColumnOperation; - use crate::{Cardinality, NumericalValue}; + use crate::{Cardinality, ColumnType, NumericalValue, PayloadEncoding}; + + #[test] + fn test_sort_values_within_row_is_stored_for_dictionary_encoding() { + let mut writer = super::ColumnarWriter::default(); + writer + .record_column_type_with_encoding( + "dictionary", + ColumnType::Str, + true, + PayloadEncoding::Dictionary, + ) + .unwrap(); + let dictionary_writer: super::DictionaryEncodedColumnWriter = writer + .dictionary_encoded_columns + .str_columns + .get(b"dictionary") + .unwrap(); + assert!(dictionary_writer.sort_values_within_row); + } #[test] fn test_column_writer_required_simple() { diff --git a/columnar/src/columnar/writer/plain.rs b/columnar/src/columnar/writer/plain.rs new file mode 100644 index 000000000..356ca411b --- /dev/null +++ b/columnar/src/columnar/writer/plain.rs @@ -0,0 +1,111 @@ +use std::mem::size_of; + +use stacker::{ArenaHashMap, MemoryArena}; + +use super::column_operation::ColumnOperation; +use super::column_writers::ColumnWriter; +use crate::RowId; + +#[derive(Default)] +pub(super) struct PlainColumnsWriter { + pub(super) bytes_columns: ArenaHashMap, + pub(super) str_columns: ArenaHashMap, + pub(super) value_stores: Vec, +} + +impl PlainColumnsWriter { + pub(super) fn mem_usage(&self) -> usize { + self.bytes_columns.mem_usage() + + self.str_columns.mem_usage() + + self + .value_stores + .iter() + .map(PlainValueStore::mem_usage) + .sum::() + } +} + +#[derive(Copy, Clone)] +pub(super) struct PlainColumnWriter { + pub(super) column_writer: ColumnWriter, + pub(super) value_store_id: u32, +} + +impl PlainColumnWriter { + pub(super) fn new(value_store_id: u32) -> Self { + Self { + column_writer: ColumnWriter::default(), + value_store_id, + } + } + + pub(super) fn record_bytes( + &mut self, + doc: RowId, + bytes: &[u8], + value_stores: &mut [PlainValueStore], + arena: &mut MemoryArena, + ) { + let value_id = value_stores[self.value_store_id as usize].push(bytes); + self.column_writer.record(doc, value_id, arena); + } + + pub(super) fn operation_iterator<'a>( + self, + arena: &MemoryArena, + old_to_new_ids: Option<&[RowId]>, + byte_buffer: &'a mut Vec, + ) -> impl Iterator> + 'a + use<'a> { + self.column_writer + .operation_iterator(arena, old_to_new_ids, byte_buffer) + } +} + +#[derive(Default)] +pub(super) struct PlainValueStore { + concatenated_payloads: Vec, + end_offsets: Vec, +} + +impl PlainValueStore { + fn push(&mut self, value: &[u8]) -> u32 { + let value_id = self.end_offsets.len() as u32; + self.concatenated_payloads.extend_from_slice(value); + self.end_offsets.push(self.concatenated_payloads.len()); + value_id + } + + pub(super) fn get(&self, value_id: u32) -> &[u8] { + let value_id = value_id as usize; + let end = self.end_offsets[value_id]; + let start = if value_id == 0 { + 0 + } else { + self.end_offsets[value_id - 1] + }; + &self.concatenated_payloads[start..end] + } + + fn mem_usage(&self) -> usize { + self.concatenated_payloads.capacity() + self.end_offsets.capacity() * size_of::() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_plain_value_store() { + let mut store = PlainValueStore::default(); + let first = store.push(b"same"); + let empty = store.push(b""); + let duplicate = store.push(b"same"); + + assert_eq!((first, empty, duplicate), (0, 1, 2)); + assert_eq!(store.get(first), b"same"); + assert_eq!(store.get(empty), b""); + assert_eq!(store.get(duplicate), b"same"); + assert!(store.mem_usage() >= 8); + } +} diff --git a/columnar/src/tests.rs b/columnar/src/tests.rs index 2a34724a4..cc0ff7853 100644 --- a/columnar/src/tests.rs +++ b/columnar/src/tests.rs @@ -348,7 +348,7 @@ fn test_plain_bytes_optional_and_non_utf8_roundtrip() { } #[test] -fn test_plain_bytes_multivalued_values_are_sorted_within_row() { +fn test_plain_bytes_multivalued_values_ignore_sort_flag() { let mut buffer = Vec::new(); let mut columnar_writer = ColumnarWriter::default(); columnar_writer @@ -377,7 +377,7 @@ fn test_plain_bytes_multivalued_values_are_sorted_within_row() { accessor .for_each_value(0, |value| values.push(value.to_vec())) .unwrap(); - assert_eq!(values, [b"".to_vec(), b"a".to_vec(), b"z".to_vec()]); + assert_eq!(values, [b"z".to_vec(), b"".to_vec(), b"a".to_vec()]); values.clear(); accessor .for_each_value(1, |value| values.push(value.to_vec()))