From e91faa9df8885506954053a2cbc9959a7e975ab6 Mon Sep 17 00:00:00 2001 From: "Lei, HUANG" <6406592+v0y4g3r@users.noreply.github.com> Date: Tue, 22 Sep 2026 15:52:54 +0000 Subject: [PATCH] perf(mito2): lazily decode dense primary key columns (#9226) * perf(mito2): lazily decode dense primary key columns Signed-off-by: Lei, HUANG * perf(mito2): bypass lazy decoding for full primary keys Signed-off-by: Lei, HUANG * fix(mito-codec): preserve prefix decoding errors Signed-off-by: Lei, HUANG * refactor(mito-codec): align encoded length helper naming Rename encoded_length to encoded_len and update all callers to match the other length helpers in the module. Signed-off-by: Lei, HUANG * refactor(mito2): clarify conditional dense key decoding Rename decode_dense_pk to ensure_dense_pk_decoded so callers can see that existing decoded values are preserved and only missing caches are populated. Signed-off-by: Lei, HUANG * refactor(mito-codec): share string framing in row converter Move encoded_string_len to the parent module so Dense and Sparse use the same framing helper without depending on each other. Preserve its implementation and visibility. Signed-off-by: Lei, HUANG --------- Signed-off-by: Lei, HUANG --- src/mito-codec/src/row_converter.rs | 30 ++ src/mito-codec/src/row_converter/dense.rs | 481 ++++++++++++++---- .../src/row_converter/sparse/checked.rs | 45 +- src/mito2/benches/bench_dense_index_update.rs | 162 ++++-- src/mito2/benches/bench_pk_tag_column.rs | 117 ++++- src/mito2/src/read.rs | 127 ++++- src/mito2/src/sst/index.rs | 53 ++ src/mito2/src/sst/index/column_test.rs | 90 +++- src/mito2/src/sst/index/indexer/abort.rs | 1 + src/mito2/src/sst/index/indexer/finish.rs | 1 + src/mito2/src/sst/index/indexer/update.rs | 22 + src/mito2/src/sst/parquet/flat_format.rs | 247 ++++++++- src/mito2/src/test_util/bench_util.rs | 20 +- 13 files changed, 1167 insertions(+), 229 deletions(-) diff --git a/src/mito-codec/src/row_converter.rs b/src/mito-codec/src/row_converter.rs index fae2997182c..60bbf1935e0 100644 --- a/src/mito-codec/src/row_converter.rs +++ b/src/mito-codec/src/row_converter.rs @@ -134,6 +134,11 @@ pub trait PrimaryKeyCodec: Send + Sync + Debug { /// Returns the encoding type of the primary key. fn encoding(&self) -> PrimaryKeyEncoding; + /// Returns the dense codec for schema-aware positional access, if applicable. + fn as_dense(&self) -> Option<&DensePrimaryKeyCodec> { + None + } + /// Decodes the primary key from the given bytes. /// /// Returns a [`CompositeValues`] that follows the primary key ordering. @@ -166,3 +171,28 @@ pub fn build_primary_key_codec_with_fields( } } } + +/// Finds the checked boundary of an Option, shared by Dense and Sparse. +/// This validates framing, not UTF-8; consumers decoding strings validate UTF-8. +pub(crate) fn encoded_string_len(bytes: &[u8]) -> memcomparable::Result { + match bytes.first().copied().ok_or(memcomparable::Error::Eof)? { + 0 => return Ok(1), + 1 => {} + marker => return Err(memcomparable::Error::InvalidTagEncoding(marker as usize)), + } + match bytes.get(1).copied().ok_or(memcomparable::Error::Eof)? { + 0 => return Ok(2), + 1 => {} + marker => return Err(memcomparable::Error::InvalidBytesEncoding(marker)), + } + let mut end = 2; + loop { + let chunk = bytes.get(end..end + 9).ok_or(memcomparable::Error::Eof)?; + end += 9; + match chunk[8] { + 1..=8 => return Ok(end), + 9 => {} + marker => return Err(memcomparable::Error::InvalidBytesEncoding(marker)), + } + } +} diff --git a/src/mito-codec/src/row_converter/dense.rs b/src/mito-codec/src/row_converter/dense.rs index 56289a08b11..ad85f4ac8f1 100644 --- a/src/mito-codec/src/row_converter/dense.rs +++ b/src/mito-codec/src/row_converter/dense.rs @@ -27,7 +27,7 @@ use datatypes::value::ValueRef; use memcomparable::{Deserializer, Serializer}; use paste::paste; use serde::{Deserialize, Serialize}; -use snafu::ResultExt; +use snafu::{ResultExt, ensure}; use store_api::codec::PrimaryKeyEncoding; use store_api::metadata::{RegionMetadata, RegionMetadataRef}; use store_api::storage::ColumnId; @@ -38,7 +38,7 @@ use crate::error::{ use crate::key_values::KeyValue; use crate::primary_key_filter::DensePrimaryKeyFilter; use crate::row_converter::{ - CompositeValues, PrimaryKeyCodec, PrimaryKeyCodecExt, PrimaryKeyFilter, + CompositeValues, PrimaryKeyCodec, PrimaryKeyCodecExt, PrimaryKeyFilter, encoded_string_len, }; /// Field to serialize and deserialize value in memcomparable format. @@ -290,84 +290,112 @@ impl SortField { bytes: &[u8], deserializer: &mut Deserializer<&[u8]>, ) -> Result { - let len = self.encoded_field_len(&bytes[deserializer.position()..])?; + let pos = deserializer.position(); + let remaining = bytes + .get(pos..) + .ok_or(memcomparable::Error::Eof) + .context(error::DeserializeFieldSnafu)?; + let len = Self::encoded_len(self.encode_data_type(), remaining)?; + if len > 1 && self.encode_data_type().is_boolean() && remaining[1] > 1 { + return Err(memcomparable::Error::InvalidBoolEncoding(remaining[1])) + .context(error::DeserializeFieldSnafu); + } deserializer.advance(len); Ok(len) } /// Checks field boundaries before any unchecked reads in memcomparable. - fn encoded_field_len(&self, bytes: &[u8]) -> Result { - let invalid = || { - error::InvalidDensePrimaryKeySnafu { - reason: "truncated field or invalid encoding", + fn encoded_len(data_type: &ConcreteDataType, bytes: &[u8]) -> Result { + let marker = bytes + .first() + .copied() + .ok_or(memcomparable::Error::Eof) + .context(error::DeserializeFieldSnafu)?; + match marker { + 0 => return Ok(1), + 1 => {} + value => { + return Err(memcomparable::Error::InvalidTagEncoding(value as usize)) + .context(error::DeserializeFieldSnafu); } - .build() - }; - match bytes.first() { - Some(0) => return Ok(1), - Some(1) => {} - _ => return Err(invalid()), } - let len = match self.encode_data_type() { - ConcreteDataType::Boolean(_) - | ConcreteDataType::Int8(_) - | ConcreteDataType::UInt8(_) => 2, + let to_skip = match data_type { + ConcreteDataType::Boolean(_) => 2, + ConcreteDataType::Int8(_) | ConcreteDataType::UInt8(_) => 2, ConcreteDataType::Int16(_) | ConcreteDataType::UInt16(_) => 3, - ConcreteDataType::Int32(_) - | ConcreteDataType::UInt32(_) - | ConcreteDataType::Float32(_) - | ConcreteDataType::Date(_) => 5, - ConcreteDataType::Int64(_) - | ConcreteDataType::UInt64(_) - | ConcreteDataType::Float64(_) - | ConcreteDataType::Timestamp(_) => 9, - ConcreteDataType::Time(_) | ConcreteDataType::Duration(_) => 10, + ConcreteDataType::Int32(_) | ConcreteDataType::UInt32(_) => 5, + ConcreteDataType::Int64(_) | ConcreteDataType::UInt64(_) => 9, + ConcreteDataType::Float32(_) => 5, + ConcreteDataType::Float64(_) => 9, + ConcreteDataType::Binary(_) + | ConcreteDataType::Json(_) + | ConcreteDataType::Vector(_) => { + // Binary is encoded as a sequence of bytes, not chunked strings. + return encoded_binary_len(bytes).context(error::DeserializeFieldSnafu); + } + ConcreteDataType::String(_) => { + return encoded_string_len(bytes).context(error::DeserializeFieldSnafu); + } + ConcreteDataType::Date(_) => 5, + ConcreteDataType::Timestamp(_) => 9, // We treat timestamp as Option + ConcreteDataType::Time(_) => 10, // i64 and 1 byte time unit + ConcreteDataType::Duration(_) => 10, ConcreteDataType::Interval(IntervalType::YearMonth(_)) => 5, ConcreteDataType::Interval(IntervalType::DayTime(_)) => 9, ConcreteDataType::Interval(IntervalType::MonthDayNano(_)) => 17, ConcreteDataType::Decimal128(_) => 19, - ConcreteDataType::Binary(_) - | ConcreteDataType::Json(_) - | ConcreteDataType::Vector(_) => { - // A serde sequence: each byte is preceded by 1, terminated by 0. - let mut pos = 1; - while bytes.get(pos) == Some(&1) { - pos += 2; - } - if bytes.get(pos) != Some(&0) { - return Err(invalid()); - } - pos + 1 - } - ConcreteDataType::String(_) => { - match bytes.get(1) { - Some(0) => return Ok(2), - Some(1) => {} - _ => return Err(invalid()), - } - let mut pos = 2; - loop { - match bytes.get(pos + 8) { - Some(1..=8) => break pos + 9, - Some(9) => pos += 9, - _ => return Err(invalid()), - } - } - } - data_type => { - return error::NotSupportedFieldSnafu { + ConcreteDataType::Null(_) + | ConcreteDataType::List(_) + | ConcreteDataType::Struct(_) + | ConcreteDataType::Dictionary(_) => { + return NotSupportedFieldSnafu { data_type: data_type.clone(), } .fail(); } }; - snafu::ensure!( - bytes.len() >= len, - error::InvalidDensePrimaryKeySnafu { - reason: "truncated field", - } - ); - Ok(len) + if bytes.len() < to_skip { + return Err(memcomparable::Error::Eof).context(error::DeserializeFieldSnafu); + } + Ok(to_skip) + } + + /// Decodes primitive fields straight into a reference value for column + /// builders, avoiding owned Value construction and conversion. The caller has + /// already checked the encoded boundary. + fn deserialize_primitive_ref( + data_type: &ConcreteDataType, + deserializer: &mut Deserializer, + ) -> Option>> { + macro_rules! decode_primitive { + ($($variant:ident, $native:ty),* $(,)?) => { + match data_type { + $(ConcreteDataType::$variant(_) => Some( + Option::<$native>::deserialize(deserializer) + .map(ValueRef::from) + .context(error::DeserializeFieldSnafu) + ),)* + _ => None, + } + }; + } + decode_primitive!( + Boolean, bool, Int8, i8, Int16, i16, Int32, i32, Int64, i64, UInt8, u8, UInt16, u16, + UInt32, u32, UInt64, u64, Float32, f32, Float64, f64, + ) + } +} + +/// The Option marker has already been checked by SortField. +fn encoded_binary_len(bytes: &[u8]) -> memcomparable::Result { + let mut current = 1; + loop { + match bytes.get(current).copied() { + Some(0) => return Ok(current + 1), + Some(1) if bytes.get(current + 1).is_some() => current += 2, + Some(1) | None => return Err(memcomparable::Error::Eof), + Some(marker) => return Err(memcomparable::Error::InvalidSeqEncoding(marker)), + } } } @@ -444,7 +472,16 @@ impl DensePrimaryKeyCodec { return Ok(index); } let start = deserializer.position(); - let len = field.encoded_field_len(&bytes[start..])?; + // Preserve the prefix error contract used by primary-key range mapping. + let len = SortField::encoded_len(field.encode_data_type(), &bytes[start..]).map_err( + |source| match source { + error::Error::DeserializeField { .. } => error::InvalidDensePrimaryKeySnafu { + reason: "truncated field or invalid encoding", + } + .build(), + source => source, + }, + )?; // Production range mapping only needs boundaries, not allocated field values. #[cfg(any(debug_assertions, test))] field.deserialize(&mut Deserializer::new(&bytes[start..start + len]))?; @@ -470,6 +507,73 @@ impl DensePrimaryKeyCodec { Ok(values) } + /// Iterates over all source PK fields in schema order, checking each field's + /// encoded boundary before deserializing it. Stops after the first error. + /// + /// Unlike positional access, this requires no offsets or value cache. Callers + /// can consume each value immediately instead of retaining the entire key. + pub fn decode_dense_iter<'a>( + &'a self, + bytes: &'a [u8], + ) -> impl Iterator> + 'a { + self.ordered_primary_key_columns.iter().scan( + Some(Deserializer::new(bytes)), + move |state, (id, field)| { + let deserializer = state.as_mut()?; + let remaining = &bytes[deserializer.position()..]; + let decoded = SortField::encoded_len(field.encode_data_type(), remaining) + .and_then(|_| field.deserialize(deserializer).map(|value| (*id, value))); + if decoded.is_err() { + *state = None; + } + Some(decoded) + }, + ) + } + + /// Returns the column ids and sort fields in encoded primary-key order. + pub fn fields(&self) -> &[(ColumnId, SortField)] { + &self.ordered_primary_key_columns + } + + /// Decodes all fields in source order into a consumer, checking boundaries. + /// Strings borrow the reusable buffer for the duration of each callback, so + /// materializing them into column builders needs no per-value owned string. + pub fn decode_dense_with( + &self, + bytes: &[u8], + value_buf: &mut Vec, + mut consume: impl FnMut(usize, ValueRef<'_>), + ) -> Result<()> { + let mut deserializer = Deserializer::new(bytes); + for (pos, (_, field)) in self.ordered_primary_key_columns.iter().enumerate() { + let data_type = field.encode_data_type(); + let remaining = &bytes[deserializer.position()..]; + SortField::encoded_len(data_type, remaining)?; + if data_type.is_string() && remaining[0] != 0 { + deserializer.advance(1); + deserializer + .read_bytes_into(value_buf) + .context(error::DeserializeFieldSnafu)?; + let value = std::str::from_utf8(value_buf).map_err(|err| { + error::InvalidDensePrimaryKeySnafu { + reason: format!("string is not valid UTF-8: {err}"), + } + .build() + })?; + consume(pos, ValueRef::String(value)); + } else if let Some(value) = + SortField::deserialize_primitive_ref(data_type, &mut deserializer) + { + consume(pos, value?); + } else { + let value = field.deserialize(&mut deserializer)?; + consume(pos, value.as_value_ref()); + } + } + Ok(()) + } + /// Returns the field at `pos`. /// /// # Panics @@ -488,61 +592,61 @@ impl DensePrimaryKeyCodec { offsets_buf: &mut Vec, deserializer: &mut Deserializer<&[u8]>, ) -> Result { - if pos < offsets_buf.len() { - // We computed the offset before. - let offset = offsets_buf[pos]; - deserializer.advance(offset); - return Ok(offset); - } - + ensure!( + pos < self.num_fields(), + error::InvalidDensePrimaryKeySnafu { + reason: format!( + "field position {pos} exceeds field count {}", + self.num_fields() + ), + } + ); if offsets_buf.is_empty() { - let mut offset = 0; - // Skip values before `pos`. - for i in 0..pos { - // Offset to skip before reading value i. - offsets_buf.push(offset); - let skip = self.field_at(i).skip_deserialize(bytes, deserializer)?; - offset += skip; - } - // Offset to skip before reading this value. - offsets_buf.push(offset); - Ok(offset) - } else { - // Offsets are not enough. - let value_start = offsets_buf.len() - 1; - // Advances to decode value at `value_start`. - let mut offset = offsets_buf[value_start]; - deserializer.advance(offset); - for i in value_start..pos { - // Skip value i. - let skip = self.field_at(i).skip_deserialize(bytes, deserializer)?; - // Offset for the value at i + 1. - offset += skip; - offsets_buf.push(offset); - } - Ok(offset) + // The schema bounds this cache. Reserve once even when the first + // requested field is near the end of a wide key. + offsets_buf.reserve_exact(self.num_fields()); + offsets_buf.push(0); } + // Start at the requested field if cached, otherwise resume discovery + // from the furthest known boundary. + let value_start = pos.min(offsets_buf.len() - 1); + let mut offset = offsets_buf[value_start]; + ensure!( + offset <= bytes.len(), + error::InvalidDensePrimaryKeySnafu { + reason: "cached offset exceeds key length", + } + ); + deserializer.advance(offset); + for i in value_start..pos { + offset += self.field_at(i).skip_deserialize(bytes, deserializer)?; + offsets_buf.push(offset); + } + Ok(offset) } /// Decode value at `pos` in `bytes`. /// - /// The i-th element in offsets buffer is how many bytes to skip in order to read value at `pos`. + /// The i-th element in the offsets buffer is the start of field i. The buffer + /// must only be reused for the same encoded key and codec. + #[inline] pub fn decode_value_at( &self, bytes: &[u8], pos: usize, offsets_buf: &mut Vec, ) -> Result { - let mut deserializer = Deserializer::new(bytes); - self.advance_to_value_at(bytes, pos, offsets_buf, &mut deserializer)?; - - self.field_at(pos).deserialize(&mut deserializer) + let encoded = self.encoded_value_at(bytes, pos, offsets_buf)?; + self.field_at(pos) + .deserialize(&mut Deserializer::new(encoded)) } /// Returns the encoded bytes at `pos` in `bytes`. /// - /// The i-th element in offsets buffer is how many bytes to skip in order to read value at - /// `pos`. + /// The i-th element in the offsets buffer is the start of field i. The buffer + /// must only be reused for the same encoded key and codec. Validates the + /// framing up to this field, without validating or decoding later fields. + #[inline] pub fn encoded_value_at<'a>( &self, bytes: &'a [u8], @@ -555,6 +659,11 @@ impl DensePrimaryKeyCodec { let len = self .field_at(pos) .skip_deserialize(bytes, &mut deserializer)?; + // We already found the next field's start while validating this one. + // Reuse it instead of scanning this field again on the next lookup. + if offsets_buf.len() == pos + 1 && pos + 1 < self.num_fields() { + offsets_buf.push(offset + len); + } Ok(&bytes[offset..offset + len]) } @@ -600,6 +709,10 @@ impl PrimaryKeyCodec for DensePrimaryKeyCodec { PrimaryKeyEncoding::Dense } + fn as_dense(&self) -> Option<&DensePrimaryKeyCodec> { + Some(self) + } + fn primary_key_filter( &self, metadata: &RegionMetadataRef, @@ -663,6 +776,78 @@ mod tests { } let decoded = encoder.decode(&result).unwrap().into_dense(); assert_eq!(decoded, row); + let mut value_buf = Vec::new(); + let mut borrowed = Vec::new(); + encoder + .decode_dense_with(&result, &mut value_buf, |_, value| { + borrowed.push(Value::from(value)) + }) + .unwrap(); + assert_eq!(borrowed, row); + let sequential = encoder + .decode_dense_iter(&result) + .collect::>>() + .unwrap(); + assert_eq!( + sequential + .iter() + .map(|(_, value)| value) + .collect::>(), + row.iter().collect::>() + ); + for end in 0..result.len() { + assert!( + encoder + .decode_dense_with(&result[..end], &mut value_buf, |_, _| {}) + .is_err() + ); + assert!( + encoder + .decode_dense_iter(&result[..end]) + .collect::>>() + .is_err(), + "truncated at {end}" + ); + } + // Every supported type must have the same boundaries in positional and + // full decoding. Exercise every truncation both at the requested field + // and while skipping it to reach a later field. + let mut offsets = Vec::new(); + let mut end = 0; + for (pos, value) in row.iter().enumerate() { + let encoded = encoder + .encoded_value_at(&result, pos, &mut offsets) + .unwrap(); + assert_eq!( + encoder + .field_at(pos) + .deserialize(&mut Deserializer::new(encoded)) + .unwrap(), + *value + ); + let start = end; + end += encoded.len(); + for len in start..end { + assert!( + encoder + .encoded_value_at(&result[..len], pos, &mut Vec::new()) + .is_err() + ); + assert!( + encoder + .decode_value_at(&result[..len], pos, &mut Vec::new()) + .is_err() + ); + if pos + 1 < row.len() { + assert!( + encoder + .encoded_value_at(&result[..len], pos + 1, &mut Vec::new()) + .is_err() + ); + } + } + } + assert_eq!(end, result.len()); let mut decoded = Vec::new(); let mut offsets = Vec::new(); // Iter two times to test offsets buffer. @@ -871,6 +1056,96 @@ mod tests { } } + #[test] + fn test_positional_decode_invalid_encoding() { + for (ty, bytes) in [ + (ConcreteDataType::int64_datatype(), vec![2]), + (ConcreteDataType::boolean_datatype(), vec![1, 2]), + (ConcreteDataType::binary_datatype(), vec![1, 2]), + (ConcreteDataType::binary_datatype(), vec![1, 1, 42, 2]), + (ConcreteDataType::string_datatype(), vec![1, 2]), + ( + ConcreteDataType::string_datatype(), + vec![1, 1, 0, 0, 0, 0, 0, 0, 0, 0, 0], + ), + ( + ConcreteDataType::string_datatype(), + vec![1, 1, 0, 0, 0, 0, 0, 0, 0, 0, 10], + ), + ] { + let codec = DensePrimaryKeyCodec::with_fields(vec![ + (0, SortField::new(ty)), + (1, SortField::new(ConcreteDataType::int64_datatype())), + ]); + for pos in [0, 1, 2, usize::MAX] { + assert!( + codec + .encoded_value_at(&bytes, pos, &mut Vec::new()) + .is_err() + ); + assert!(codec.decode_value_at(&bytes, pos, &mut Vec::new()).is_err()); + } + assert!( + codec + .decode_dense_iter(&bytes) + .collect::>>() + .is_err() + ); + assert!( + codec + .decode_dense_with(&bytes, &mut Vec::new(), |_, _| {}) + .is_err() + ); + } + let codec = DensePrimaryKeyCodec::with_fields(vec![ + (0, SortField::new(ConcreteDataType::string_datatype())), + (1, SortField::new(ConcreteDataType::int64_datatype())), + ]); + let bytes = codec + .encode([ValueRef::String("abcdefghijk"), ValueRef::Int64(42)].into_iter()) + .unwrap(); + for pos in [2, usize::MAX] { + assert!( + codec + .encoded_value_at(&bytes, pos, &mut Vec::new()) + .is_err() + ); + assert!(codec.decode_value_at(&bytes, pos, &mut Vec::new()).is_err()); + } + for offset in [bytes.len(), bytes.len() + 1, usize::MAX] { + for pos in [0, 1] { + assert!( + codec + .encoded_value_at(&bytes, pos, &mut vec![offset]) + .is_err() + ); + assert!( + codec + .decode_value_at(&bytes, pos, &mut vec![offset]) + .is_err() + ); + } + } + let mut invalid_utf8 = bytes; + invalid_utf8[2] = 0xff; + assert!( + codec + .decode_dense_with(&invalid_utf8, &mut Vec::new(), |_, _| {}) + .is_err() + ); + assert!( + codec + .decode_dense_iter(&invalid_utf8) + .collect::>>() + .is_err() + ); + assert!( + codec + .decode_value_at(&invalid_utf8, 0, &mut Vec::new()) + .is_err() + ); + } + #[test] fn test_memcmp_dictionary() { // Test Dictionary diff --git a/src/mito-codec/src/row_converter/sparse/checked.rs b/src/mito-codec/src/row_converter/sparse/checked.rs index 1e57c13cc7f..0cc9ed8211a 100644 --- a/src/mito-codec/src/row_converter/sparse/checked.rs +++ b/src/mito-codec/src/row_converter/sparse/checked.rs @@ -17,6 +17,7 @@ use snafu::{OptionExt, ensure}; use store_api::storage::ColumnId; use crate::error::{InvalidSparsePrimaryKeySnafu, Result}; +use crate::row_converter::encoded_string_len; use crate::row_converter::sparse::{ COLUMN_ID_ENCODE_SIZE, RESERVED_COLUMN_ID_TABLE_ID, RESERVED_COLUMN_ID_TSID, SparseOffsetsCache, TABLE_ID_VALUE_OFFSET, TAGS_START_OFFSET, TSID_VALUE_OFFSET, @@ -141,45 +142,13 @@ impl<'a, 'b> SparsePrimaryKeyView<'a, 'b> { /// Finds the end of an Option without allocating or reading its payload. /// Unlike memcomparable's skip_bytes, all advances are checked for truncated input. fn encoded_label(bytes: &[u8]) -> Result<&[u8]> { - match bytes.first() { - Some(0) => return Ok(&bytes[..1]), - Some(1) => {} - _ => { - return InvalidSparsePrimaryKeySnafu { - reason: "invalid label null marker", - } - .fail(); + let len = encoded_string_len(bytes).map_err(|error| { + InvalidSparsePrimaryKeySnafu { + reason: error.to_string(), } - } - match bytes.get(1) { - Some(0) => return Ok(&bytes[..2]), - Some(1) => {} - _ => { - return InvalidSparsePrimaryKeySnafu { - reason: "invalid label bytes marker", - } - .fail(); - } - } - let mut end = 2; - loop { - let chunk = bytes - .get(end..end + 9) - .context(InvalidSparsePrimaryKeySnafu { - reason: "truncated label chunk", - })?; - end += 9; - match chunk[8] { - 1..=8 => return Ok(&bytes[..end]), - 9 => {} - _ => { - return InvalidSparsePrimaryKeySnafu { - reason: "invalid label chunk length", - } - .fail(); - } - } - } + .build() + })?; + Ok(&bytes[..len]) } #[cfg(test)] diff --git a/src/mito2/benches/bench_dense_index_update.rs b/src/mito2/benches/bench_dense_index_update.rs index f81d30eb4a6..6f5ccfc080b 100644 --- a/src/mito2/benches/bench_dense_index_update.rs +++ b/src/mito2/benches/bench_dense_index_update.rs @@ -30,8 +30,8 @@ use api::v1::SemanticType; use common_test_util::temp_dir::create_temp_dir; use criterion::{Criterion, Throughput, criterion_group, criterion_main}; use datatypes::arrow::array::{ - ArrayRef, BinaryDictionaryBuilder, StringDictionaryBuilder, TimestampMillisecondArray, - UInt8Array, UInt64Array, + ArrayRef, BinaryArray, BinaryDictionaryBuilder, DictionaryArray, StringDictionaryBuilder, + TimestampMillisecondArray, UInt8Array, UInt64Array, }; use datatypes::arrow::datatypes::UInt32Type; use datatypes::arrow::record_batch::RecordBatch; @@ -39,6 +39,8 @@ use datatypes::data_type::ConcreteDataType; use datatypes::schema::{ColumnSchema, SkippingIndexOptions}; use datatypes::value::Value; use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt}; +use mito2::read::{Batch, BatchBuilder}; +use mito2::sst::index::Indexer; use mito2::sst::index::intermediate::IntermediateManager; use mito2::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema}; use mito2::test_util::bench_util::{BloomFilterIndexer, InvertedIndexer}; @@ -284,55 +286,115 @@ fn bench_dense_index_update(c: &mut Criterion) { } for shape in shapes { let (metadata, inverted_columns, batch) = input(&shape); - let run = |iterations| { - runtime.block_on(async { - let mut elapsed = Duration::ZERO; - for _ in 0..iterations { - let file_id = FileId::random(); - let mut inverted = shape.inverted.then(|| { - InvertedIndexer::new( - file_id, - &metadata, - intermediate.clone(), - None, - NonZeroUsize::new(SEGMENT_ROWS).unwrap(), - inverted_columns.clone(), - ) - }); - let mut bloom = if shape.bloom { - BloomFilterIndexer::new(file_id, &metadata, intermediate.clone(), None) - .unwrap() - } else { - None - }; - if allocations { - ALLOCATIONS.set(Some((0, 0))); + for encoded_only in [false, true] { + let name = if encoded_only { + format!("encoded/{}", shape.name) + } else { + shape.name.clone() + }; + let encoded_batches: Vec = if encoded_only { + let pk = batch + .column_by_name("__primary_key") + .unwrap() + .as_any() + .downcast_ref::>() + .unwrap(); + let values = pk.values().as_any().downcast_ref::().unwrap(); + (0..ROWS) + .step_by(shape.rows_per_key) + .map(|row| { + let mut builder = + BatchBuilder::new(values.value(pk.keys().value(row) as usize).to_vec()); + builder + .timestamps_array( + batch + .column_by_name("ts") + .unwrap() + .slice(row, shape.rows_per_key), + ) + .unwrap(); + builder + .sequences_array( + batch + .column_by_name("__sequence") + .unwrap() + .slice(row, shape.rows_per_key), + ) + .unwrap(); + builder + .op_types_array( + batch + .column_by_name("__op_type") + .unwrap() + .slice(row, shape.rows_per_key), + ) + .unwrap(); + if shape.field { + builder + .push_field_array( + shape.tags, + batch + .column_by_name(&format!("col_{}", shape.tags)) + .unwrap() + .slice(row, shape.rows_per_key), + ) + .unwrap(); + } + builder.build().unwrap() + }) + .collect() + } else { + Vec::new() + }; + let run = |iterations| { + runtime.block_on(async { + let mut elapsed = Duration::ZERO; + for _ in 0..iterations { + let file_id = FileId::random(); + let inverted = shape.inverted.then(|| { + InvertedIndexer::new( + file_id, + &metadata, + intermediate.clone(), + None, + NonZeroUsize::new(SEGMENT_ROWS).unwrap(), + inverted_columns.clone(), + ) + }); + let bloom = if shape.bloom { + BloomFilterIndexer::new(file_id, &metadata, intermediate.clone(), None) + .unwrap() + } else { + None + }; + let mut indexer = Indexer::for_bench(&metadata, inverted, bloom); + // Fresh caches for each iteration; both creators share each Batch. + let mut encoded_batches = encoded_batches.clone(); + if allocations { + ALLOCATIONS.set(Some((0, 0))); + } + let start = Instant::now(); + if encoded_only { + for batch in &mut encoded_batches { + indexer.update(batch).await; + } + } else { + indexer.update_flat(&batch).await; + } + elapsed += start.elapsed(); + if allocations { + report_allocations(&name); + } + indexer.abort().await; } - let start = Instant::now(); - if let Some(indexer) = &mut inverted { - indexer.update_flat(&batch).await.unwrap(); - } - if let Some(indexer) = &mut bloom { - indexer.update_flat(&batch).await.unwrap(); - } - elapsed += start.elapsed(); - if allocations { - report_allocations(&shape.name); - } - if let Some(indexer) = &mut inverted { - indexer.abort().await.unwrap(); - } - if let Some(indexer) = &mut bloom { - indexer.abort().await.unwrap(); - } - } - elapsed - }) - }; - if allocations { - run(1); - } else { - group.bench_function(&shape.name, |b| b.iter_custom(&run)); + elapsed + }) + }; + if allocations { + run(1); + } else { + group.bench_function(&name, |b| b.iter_custom(&run)); + } } } group.finish(); diff --git a/src/mito2/benches/bench_pk_tag_column.rs b/src/mito2/benches/bench_pk_tag_column.rs index d3ded3d0b05..9e191a69511 100644 --- a/src/mito2/benches/bench_pk_tag_column.rs +++ b/src/mito2/benches/bench_pk_tag_column.rs @@ -24,6 +24,7 @@ use std::hint::black_box; use std::sync::Arc; +use std::time::Duration; use api::v1::SemanticType; use criterion::{BenchmarkId, Criterion, criterion_group, criterion_main}; @@ -35,9 +36,12 @@ use datatypes::arrow::datatypes::UInt32Type; use datatypes::arrow::record_batch::RecordBatch; use datatypes::data_type::ConcreteDataType; use datatypes::schema::ColumnSchema; -use mito_codec::row_converter::SparsePrimaryKeyCodec; +use datatypes::value::Value; +use mito_codec::row_converter::{ + DensePrimaryKeyCodec, PrimaryKeyCodecExt, SortField, SparsePrimaryKeyCodec, +}; use mito2::sst::parquet::flat_format::decode_primary_keys; -use mito2::test_util::bench_util::tag_filter_for_bench; +use mito2::test_util::bench_util::{pk_materializer_for_bench, tag_filter_for_bench}; use store_api::codec::PrimaryKeyEncoding; use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder, RegionMetadataRef}; use store_api::storage::RegionId; @@ -237,5 +241,112 @@ fn bench_pk_tag_filters(c: &mut Criterion) { group.finish(); } -criterion_group!(benches, bench_pk_tag_column, bench_pk_tag_filters); +/// Measures encoded-only Dense inputs, including the cost of discovering offsets. +fn bench_dense_pk_tag_column(c: &mut Criterion) { + const TAGS: usize = 40; + let mut group = c.benchmark_group("dense_pk_tag_column"); + group.sample_size(30); + group.warm_up_time(Duration::from_secs(1)); + group.measurement_time(Duration::from_secs(3)); + for numeric in [false, true] { + let ty = if numeric { + ConcreteDataType::uint64_datatype() + } else { + ConcreteDataType::string_datatype() + }; + let codec = DensePrimaryKeyCodec::with_fields( + (0..TAGS) + .map(|id| (id as u32, SortField::new(ty.clone()))) + .collect(), + ); + let mut metadata = RegionMetadataBuilder::new(RegionId::new(1, 1)); + for id in 0..TAGS { + metadata.push_column_metadata(ColumnMetadata { + column_id: id as u32, + column_schema: ColumnSchema::new(format!("tag_{id}"), ty.clone(), true), + semantic_type: SemanticType::Tag, + }); + } + metadata + .push_column_metadata(ColumnMetadata { + column_id: TAGS as u32, + column_schema: ColumnSchema::new( + "ts", + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ), + semantic_type: SemanticType::Timestamp, + }) + .primary_key((0..TAGS as u32).collect()); + let metadata = Arc::new(metadata.build().unwrap()); + for rows_per_key in [1, 32] { + let mut keys = BinaryDictionaryBuilder::::new(); + for series in 0..ROWS / rows_per_key { + let values: Vec<_> = (0..TAGS) + .map(|id| { + if numeric { + Value::UInt64((series + id) as u64) + } else { + Value::from(format!( + "tag-{id:03}-value-{series:010}-abcdefghijklmnopqrstuvwxyz" + )) + } + }) + .collect(); + let pk = codec + .encode(values.iter().map(Value::as_value_ref)) + .unwrap(); + for _ in 0..rows_per_key { + keys.append(&pk).unwrap(); + } + } + let batch = RecordBatch::try_from_iter([ + ( + "ts", + Arc::new(TimestampMillisecondArray::from_iter_values(0..ROWS as i64)) + as ArrayRef, + ), + ("__primary_key", Arc::new(keys.finish()) as ArrayRef), + ( + "__sequence", + Arc::new(UInt64Array::from(vec![1; ROWS])) as ArrayRef, + ), + ( + "__op_type", + Arc::new(UInt8Array::from(vec![1; ROWS])) as ArrayRef, + ), + ]) + .unwrap(); + for (projection, start, count) in [ + ("first", 0, 1), + ("last", 39, 1), + ("4last", 36, 4), + ("all", 0, 40), + ] { + let kind = if numeric { "numeric" } else { "string" }; + let materialize = pk_materializer_for_bench( + metadata.clone(), + (start as u32..(start + count) as u32) + .chain([TAGS as u32]) + .collect(), + batch.schema(), + ); + assert_eq!(materialize(batch.clone()).unwrap().num_columns(), count + 4); + group.bench_function(format!("{kind}_{projection}_{rows_per_key}rpk"), |b| { + b.iter(|| { + black_box(materialize(black_box(batch.clone())).unwrap()); + }); + }); + } + } + } + group.finish(); +} + +criterion_group!( + benches, + bench_pk_tag_column, + bench_pk_tag_filters, + bench_dense_pk_tag_column +); criterion_main!(benches); diff --git a/src/mito2/src/read.rs b/src/mito2/src/read.rs index 72e778851ee..d5abe29f021 100644 --- a/src/mito2/src/read.rs +++ b/src/mito2/src/read.rs @@ -65,7 +65,7 @@ use datatypes::vectors::{ }; use futures::TryStreamExt; use futures::stream::BoxStream; -use mito_codec::row_converter::{CompositeValues, PrimaryKeyCodec}; +use mito_codec::row_converter::{CompositeValues, DensePrimaryKeyCodec, PrimaryKeyCodec}; use snafu::{OptionExt, ResultExt, ensure}; use store_api::storage::{ColumnId, SequenceNumber, SequenceRange}; @@ -119,6 +119,8 @@ pub struct Batch { primary_key: Vec, /// Possibly decoded `primary_key` values. Some places would decode it in advance. pk_values: Option, + /// Lazily decoded Dense columns, shared by consumers such as index creators. + dense_pk_cache: Option>, /// Timestamps of rows, should be sorted and not null. timestamps: VectorRef, /// Sequences of rows @@ -135,6 +137,13 @@ pub struct Batch { fields_idx: Option>, } +/// Only used when no fully decoded primary key was supplied by the reader. +#[derive(Debug, PartialEq, Clone)] +struct DensePkCache { + offsets: Vec, + values: Vec>, +} + impl Batch { /// Creates a new batch. pub fn new( @@ -173,12 +182,14 @@ impl Batch { /// Sets possibly decoded primary-key values. pub fn set_pk_values(&mut self, pk_values: CompositeValues) { self.pk_values = Some(pk_values); + self.dense_pk_cache = None; } /// Removes possibly decoded primary-key values. For testing only. #[cfg(any(test, feature = "test"))] pub fn remove_pk_values(&mut self) { self.pk_values = None; + self.dense_pk_cache = None; } /// Returns fields in the batch. @@ -214,6 +225,7 @@ impl Batch { Self { primary_key: vec![], pk_values: None, + dense_pk_cache: None, timestamps: Arc::new(TimestampMillisecondVectorBuilder::with_capacity(0).finish()), sequences: Arc::new(UInt64VectorBuilder::with_capacity(0).finish()), op_types: Arc::new(UInt8VectorBuilder::with_capacity(0).finish()), @@ -269,6 +281,7 @@ impl Batch { /// Be sure to update that field as well. pub fn set_primary_key(&mut self, primary_key: Vec) { self.primary_key = primary_key; + self.dense_pk_cache = None; } /// Slice the batch, returning a new batch. @@ -290,6 +303,7 @@ impl Batch { // this becomes a bottleneck. primary_key: self.primary_key.clone(), pk_values: self.pk_values.clone(), + dense_pk_cache: self.dense_pk_cache.clone(), timestamps: self.timestamps.slice(offset, length), sequences: Arc::new(self.sequences.get_slice(offset, length)), op_types: Arc::new(self.op_types.get_slice(offset, length)), @@ -772,9 +786,26 @@ impl Batch { )) } + /// Prepares a shared full-key cache when all Dense PK columns are needed. + /// Explicitly supplied decoded/defaulted values take precedence. + pub(crate) fn ensure_dense_pk_decoded(&mut self, codec: &DensePrimaryKeyCodec) -> Result<()> { + if self.pk_values.is_none() { + // A fallible iterator has no nonzero lower size hint. Reserve the + // known field count instead of growing its collected Vec per key. + let mut values = Vec::with_capacity(codec.num_fields()); + for value in codec.decode_dense_iter(&self.primary_key) { + values.push(value.context(DecodeSnafu)?); + } + self.set_pk_values(CompositeValues::Dense(values)); + } + Ok(()) + } + /// Returns the value of the column in the primary key. /// - /// Lazily decodes the primary key and caches the result. + /// Reuses predecoded values when available. Otherwise Dense keys decode only + /// the requested column, sharing offsets and values across callers; Sparse + /// keys cache the full decode. pub fn pk_col_value( &mut self, codec: &dyn PrimaryKeyCodec, @@ -782,6 +813,24 @@ impl Batch { column_id: ColumnId, ) -> Result> { if self.pk_values.is_none() { + if let Some(codec) = codec.as_dense() { + if col_idx_in_pk >= codec.num_fields() { + return Ok(None); + } + let cache = self.dense_pk_cache.get_or_insert_with(|| { + Box::new(DensePkCache { + offsets: Vec::new(), + values: vec![None; codec.num_fields()], + }) + }); + if cache.values[col_idx_in_pk].is_none() { + let value = codec + .decode_value_at(&self.primary_key, col_idx_in_pk, &mut cache.offsets) + .context(DecodeSnafu)?; + cache.values[col_idx_in_pk] = Some(value); + } + return Ok(cache.values[col_idx_in_pk].as_ref()); + } self.pk_values = Some(codec.decode(&self.primary_key).context(DecodeSnafu)?); } @@ -1098,6 +1147,7 @@ impl BatchBuilder { Ok(Batch { primary_key: self.primary_key, pk_values: None, + dense_pk_cache: None, timestamps, sequences, op_types, @@ -1254,6 +1304,79 @@ mod tests { use crate::error::Error; use crate::test_util::new_batch_builder; + #[test] + fn dense_pk_columns_are_lazy_and_reset_with_the_key() { + use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt, SortField}; + + let codec = DensePrimaryKeyCodec::with_fields(vec![ + (7, SortField::new(ConcreteDataType::string_datatype())), + (3, SortField::new(ConcreteDataType::int64_datatype())), + (9, SortField::new(ConcreteDataType::string_datatype())), + ]); + let values = [ + Value::from("中文\0abcdefgh"), + Value::Int64(-42), + Value::Null, + ]; + let key = codec + .encode(values.iter().map(Value::as_value_ref)) + .unwrap(); + let mut batch = new_batch_without_fields(&[1, 2], &[1, 1], &[OpType::Put, OpType::Put]); + batch.set_primary_key(key); + for pos in [2, 0, 1, 2] { + assert_eq!( + batch.pk_col_value(&codec, pos, [7, 3, 9][pos]).unwrap(), + Some(&values[pos]) + ); + } + assert!(batch.pk_col_value(&codec, 3, 99).unwrap().is_none()); + let mut sliced = batch.slice(1, 1); + assert_eq!(sliced.pk_col_value(&codec, 0, 7).unwrap(), Some(&values[0])); + + // A corrupt unrequested suffix must not force whole-key decoding. + batch.set_primary_key(vec![0, 1]); + assert_eq!( + batch.pk_col_value(&codec, 0, 7).unwrap(), + Some(&Value::Null) + ); + assert!(batch.pk_col_value(&codec, 1, 3).is_err()); + batch.set_primary_key( + codec + .encode( + [ + ValueRef::String(""), + ValueRef::Int64(8), + ValueRef::String("new"), + ] + .into_iter(), + ) + .unwrap(), + ); + assert_eq!( + batch.pk_col_value(&codec, 2, 9).unwrap(), + Some(&Value::from("new")) + ); + assert_eq!( + batch.pk_col_value(&codec, 0, 7).unwrap(), + Some(&Value::from("")) + ); + + // Schema compatibility may supply already decoded/defaulted values. + batch.set_pk_values(CompositeValues::Dense(vec![(7, Value::from("default"))])); + batch.ensure_dense_pk_decoded(&codec).unwrap(); + assert_eq!( + batch.pk_col_value(&codec, 0, 7).unwrap(), + Some(&Value::from("default")) + ); + assert!(batch.pk_col_value(&codec, 1, 3).unwrap().is_none()); + batch.remove_pk_values(); + batch.ensure_dense_pk_decoded(&codec).unwrap(); + assert_eq!( + batch.pk_col_value(&codec, 1, 3).unwrap(), + Some(&Value::Int64(8)) + ); + } + fn new_batch( timestamps: &[i64], sequences: &[u64], diff --git a/src/mito2/src/sst/index.rs b/src/mito2/src/sst/index.rs index ed3ec36a473..2dd29e9803e 100644 --- a/src/mito2/src/sst/index.rs +++ b/src/mito2/src/sst/index.rs @@ -37,11 +37,13 @@ use std::sync::Arc; use bloom_filter::creator::BloomFilterIndexer; use common_telemetry::{debug, error, info, warn}; use datatypes::arrow::record_batch::RecordBatch; +use mito_codec::row_converter::DensePrimaryKeyCodec; use object_store::ObjectStore; use puffin_manager::SstPuffinManager; use smallvec::{SmallVec, smallvec}; use snafu::ResultExt; use statistics::{ByteCount, RowCount}; +use store_api::codec::PrimaryKeyEncoding; use store_api::metadata::RegionMetadataRef; use store_api::storage::{ColumnId, FileId, RegionId}; use strum::IntoStaticStr; @@ -272,6 +274,8 @@ pub struct Indexer { last_mem_fulltext_index: usize, bloom_filter_indexer: Option, last_mem_bloom_filter: usize, + /// Present only when the active creators together need every Dense PK field. + dense_pk_decoder: Option, #[cfg(feature = "vector_index")] vector_indexer: Option, #[cfg(feature = "vector_index")] @@ -280,6 +284,54 @@ pub struct Indexer { } impl Indexer { + /// Wraps real creators for update-only benchmarks without a Puffin output. + #[cfg(feature = "testing")] + pub fn for_bench( + metadata: &RegionMetadataRef, + inverted_indexer: Option, + bloom_filter_indexer: Option, + ) -> Self { + let mut indexer = Self { + region_id: metadata.region_id, + inverted_indexer, + bloom_filter_indexer, + ..Default::default() + }; + indexer.prepare_dense_pk_decoder(metadata); + indexer + } + + /// Computes demand once from active creators. Field columns and duplicate + /// requests across creators cannot turn a partial PK projection into a full one. + fn prepare_dense_pk_decoder(&mut self, metadata: &RegionMetadataRef) { + if metadata.primary_key_encoding != PrimaryKeyEncoding::Dense + || metadata.primary_key.is_empty() + { + return; + } + let requested: HashSet<_> = self + .inverted_indexer + .iter() + .flat_map(|indexer| indexer.column_ids()) + .chain( + self.bloom_filter_indexer + .iter() + .flat_map(|indexer| indexer.column_ids()), + ) + .collect(); + if metadata.primary_key.iter().all(|id| requested.contains(id)) { + self.dense_pk_decoder = Some(DensePrimaryKeyCodec::new(metadata)); + } + } + + /// Called before any creator consumes the batch, so the full decode is shared. + fn prepare_primary_key(&self, batch: &mut Batch) -> Result<()> { + if let Some(codec) = &self.dense_pk_decoder { + batch.ensure_dense_pk_decoded(codec)?; + } + Ok(()) + } + /// Updates the index with the given batch. pub async fn update(&mut self, batch: &mut Batch) { self.do_update(batch).await; @@ -398,6 +450,7 @@ impl IndexerBuilder for IndexerBuilderImpl { self.build_inverted_indexer(region_file_id.file_id(), row_group_size); indexer.fulltext_indexer = self.build_fulltext_indexer(region_file_id.file_id()).await; indexer.bloom_filter_indexer = self.build_bloom_filter_indexer(region_file_id.file_id()); + indexer.prepare_dense_pk_decoder(&self.metadata); #[cfg(feature = "vector_index")] { indexer.vector_indexer = self.build_vector_indexer(region_file_id.file_id()); diff --git a/src/mito2/src/sst/index/column_test.rs b/src/mito2/src/sst/index/column_test.rs index 146bed120c0..34ff1c19fde 100644 --- a/src/mito2/src/sst/index/column_test.rs +++ b/src/mito2/src/sst/index/column_test.rs @@ -29,7 +29,9 @@ use index::bitmap::{Bitmap, BitmapType}; use index::bloom_filter::reader::{BloomFilterReader, BloomFilterReaderImpl}; use index::inverted_index::format::reader::{InvertedIndexBlobReader, InvertedIndexReader}; use mito_codec::index::IndexValueCodec; -use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt, SortField}; +use mito_codec::row_converter::{ + DensePrimaryKeyCodec, PrimaryKeyCodec, PrimaryKeyCodecExt, SortField, +}; use object_store::ObjectStore; use object_store::services::Memory; use prost::Message; @@ -201,13 +203,85 @@ async fn materialized_runs_match_legacy_index_bytes_and_bloom_lookups() { vector_index_config: Default::default(), }; let mut expected = None; - // The legacy Batch path is an independent whole-PK decode/owned-encoding oracle. + // Demand is the union of active creators' PK columns, not their counts or + // all columns in the metadata. A malformed unneeded suffix distinguishes + // the lazy path from a premature whole-key decode. + for (inverted, bloom, disable_bloom, all_pk) in [ + (vec![0, 1], vec![], false, true), + (vec![], vec![0, 1], false, true), + (vec![0], vec![1], false, true), + (vec![0, 2], vec![0, 2], false, false), + (vec![0], vec![1], true, false), + (vec![2], vec![2], false, false), + ] { + let mut case_metadata = RegionMetadataBuilder::new(metadata.region_id); + for column in &metadata.column_metadatas { + let mut column = column.clone(); + column + .column_schema + .set_inverted_index(inverted.contains(&column.column_id)); + if !bloom.contains(&column.column_id) { + column.column_schema.unset_skipping_options().unwrap(); + } + case_metadata.push_column_metadata(column); + } + case_metadata.primary_key(metadata.primary_key.clone()); + let mut case_builder = builder.clone(); + case_builder.metadata = Arc::new(case_metadata.build().unwrap()); + if disable_bloom { + case_builder.bloom_filter_index_config.create_on_flush = crate::config::Mode::Disable; + } + let mut indexer = case_builder + .build( + RegionFileId::new(metadata.region_id, FileId::random()), + 0, + None, + ) + .await; + let make_batch = |key| { + let mut batch_builder = BatchBuilder::new(key); + batch_builder + .timestamps_array(batch.column(3).slice(0, 1)) + .unwrap(); + batch_builder + .sequences_array(batch.column(5).slice(0, 1)) + .unwrap(); + batch_builder + .op_types_array(batch.column(6).slice(0, 1)) + .unwrap(); + batch_builder.build().unwrap() + }; + let mut valid = make_batch(encoded[0].clone()); + indexer.prepare_primary_key(&mut valid).unwrap(); + if all_pk { + assert_eq!( + valid.pk_values(), + Some( + &DensePrimaryKeyCodec::new(&metadata) + .decode(&encoded[0]) + .unwrap() + ) + ); + } else { + assert!(valid.pk_values().is_none()); + } + let mut truncated = make_batch(vec![0, 1]); + assert_eq!(indexer.prepare_primary_key(&mut truncated).is_err(), all_pk); + assert!(truncated.pk_values().is_none()); + indexer.abort().await; + // Aborted creators must not retain their previous all-PK demand. + indexer.prepare_primary_key(&mut truncated).unwrap(); + } + // Explicit full decoding is the oracle for lazy encoded-only Batch inputs. // One-row slices disable run merging; a tag-only projection also exercises the no-PK fallback. - for mode in ["legacy", "runs", "single_rows", "no_pk"] { + for mode in ["eager", "legacy", "lazy", "runs", "single_rows", "no_pk"] { let file = RegionFileId::new(metadata.region_id, FileId::random()); let mut indexer = builder.build(file, 0, None).await; + if mode == "lazy" { + indexer.dense_pk_decoder = None; + } match mode { - "legacy" => { + "eager" | "legacy" | "lazy" => { for (row, key) in encoded.iter().enumerate() { let mut old = BatchBuilder::new(key.clone()); old.push_field_array(2, batch.column(2).slice(row, 1)) @@ -215,7 +289,13 @@ async fn materialized_runs_match_legacy_index_bytes_and_bloom_lookups() { old.timestamps_array(batch.column(3).slice(row, 1)).unwrap(); old.sequences_array(batch.column(5).slice(row, 1)).unwrap(); old.op_types_array(batch.column(6).slice(row, 1)).unwrap(); - indexer.update(&mut old.build().unwrap()).await; + let mut old = old.build().unwrap(); + if mode == "eager" { + old.set_pk_values( + DensePrimaryKeyCodec::new(&metadata).decode(key).unwrap(), + ); + } + indexer.update(&mut old).await; } } "single_rows" => { diff --git a/src/mito2/src/sst/index/indexer/abort.rs b/src/mito2/src/sst/index/indexer/abort.rs index a928eb186da..24cbb797cbd 100644 --- a/src/mito2/src/sst/index/indexer/abort.rs +++ b/src/mito2/src/sst/index/indexer/abort.rs @@ -20,6 +20,7 @@ use crate::sst::index::Indexer; impl Indexer { pub(crate) async fn do_abort(&mut self) { + self.dense_pk_decoder = None; self.do_abort_inverted_index().await; self.do_abort_fulltext_index().await; self.do_abort_bloom_filter().await; diff --git a/src/mito2/src/sst/index/indexer/finish.rs b/src/mito2/src/sst/index/indexer/finish.rs index c7966823ddf..823cd5c2b8a 100644 --- a/src/mito2/src/sst/index/indexer/finish.rs +++ b/src/mito2/src/sst/index/indexer/finish.rs @@ -27,6 +27,7 @@ use crate::sst::index::{ impl Indexer { pub(crate) async fn do_finish(&mut self) -> IndexOutput { + self.dense_pk_decoder = None; let mut output = IndexOutput::default(); let Some(mut writer) = self.build_puffin_writer().await else { diff --git a/src/mito2/src/sst/index/indexer/update.rs b/src/mito2/src/sst/index/indexer/update.rs index 095da337f49..0168bee876e 100644 --- a/src/mito2/src/sst/index/indexer/update.rs +++ b/src/mito2/src/sst/index/indexer/update.rs @@ -24,6 +24,11 @@ impl Indexer { return; } + if !self.do_prepare_primary_key(batch) { + self.do_abort().await; + return; + } + if !self.do_update_inverted_index(batch).await { self.do_abort().await; } @@ -39,6 +44,23 @@ impl Indexer { } } + /// Handles decode errors before entering asynchronous cleanup, following the + /// creators' update policy without carrying a decode error across an await. + fn do_prepare_primary_key(&self, batch: &mut Batch) -> bool { + let Err(err) = self.prepare_primary_key(batch) else { + return true; + }; + if cfg!(any(test, feature = "test")) { + panic!( + "Failed to decode primary key for indexes, region_id: {}, file_id: {}, err: {:?}", + self.region_id, self.file_id, err + ); + } else { + warn!(err; "Failed to decode primary key for indexes, region_id: {}, file_id: {}", self.region_id, self.file_id); + } + false + } + /// Returns false if the update failed. async fn do_update_inverted_index(&mut self, batch: &mut Batch) -> bool { let Some(creator) = self.inverted_indexer.as_mut() else { diff --git a/src/mito2/src/sst/parquet/flat_format.rs b/src/mito2/src/sst/parquet/flat_format.rs index a99fc201baf..34fb8cb873c 100644 --- a/src/mito2/src/sst/parquet/flat_format.rs +++ b/src/mito2/src/sst/parquet/flat_format.rs @@ -46,7 +46,7 @@ use mito_codec::row_converter::sparse::{ RESERVED_COLUMN_ID_TABLE_ID, RESERVED_COLUMN_ID_TSID, SparsePrimaryKeyView, }; use mito_codec::row_converter::{ - CompositeValues, PrimaryKeyCodec, SparseOffsetsCache, build_primary_key_codec, + DensePrimaryKeyCodec, PrimaryKeyCodec, SparseOffsetsCache, build_primary_key_codec, }; use parquet::file::metadata::RowGroupMetaData; use snafu::{OptionExt, ResultExt, ensure}; @@ -676,8 +676,9 @@ pub(crate) fn sst_column_id_indices(metadata: &RegionMetadata) -> HashMap { - let mut decoded_pk_values = Vec::with_capacity(distinct_keys.len()); - for key in &distinct_keys { - let pk_bytes = pk_values_array.value(*key as usize); - let decoded_value = codec.decode(pk_bytes).context(DecodeSnafu)?; - decoded_pk_values.push(decoded_value); - } - DecodedKeysInner::Dense(decoded_pk_values) - } + PrimaryKeyEncoding::Dense => DecodedKeysInner::Dense { + codec: codec + .as_dense() + .context(InvalidRecordBatchSnafu { + reason: "expected dense primary key codec", + })? + .clone(), + offsets: Vec::new(), + values: pk_values_array.clone(), + distinct_keys, + }, }; Ok(DecodedPrimaryKeys { @@ -748,10 +751,16 @@ pub fn decode_primary_keys( }) } -/// Eagerly decoded dense primary keys or lazily extracted sparse primary keys. +/// Encoded primary keys and encoding-specific lookup state. enum DecodedKeysInner { - /// Dense primary keys are decoded once and shared by all tag columns. - Dense(Vec), + /// Dense keys share positional offsets across projected tag columns. + Dense { + codec: DensePrimaryKeyCodec, + values: BinaryArray, + distinct_keys: Vec, + /// Initialized by positional access; full extraction needs no offsets. + offsets: Vec>, + }, /// Sparse primary keys stay encoded; each tag column extracts only its own /// values from the raw keys. Sparse { @@ -811,7 +820,7 @@ impl DecodedPrimaryKeys { /// Gets a tag column array by column id and data type. /// /// For sparse encoding, extracts the column lazily from the encoded keys. - /// For dense encoding, uses pk_index to get values from the decoded keys. + /// For dense encoding, uses pk_index to decode only the requested field. pub fn get_tag_column( &mut self, column_id: ColumnId, @@ -827,18 +836,25 @@ impl DecodedPrimaryKeys { // Gets values from the primary key. let values_vector = match inner { - DecodedKeysInner::Dense(decoded_pk_values) => { - let mut builder = column_type.create_mutable_vector(decoded_pk_values.len()); - for decoded in decoded_pk_values { - let CompositeValues::Dense(dense) = decoded else { - return InvalidRecordBatchSnafu { - reason: "expected dense primary key values", - } - .fail(); - }; - let pk_idx = pk_index.expect("pk_index required for dense encoding"); - if pk_idx < dense.len() { - builder.push_value_ref(&dense[pk_idx].1.as_value_ref()); + DecodedKeysInner::Dense { + codec, + values, + distinct_keys, + offsets, + } => { + let pk_idx = pk_index.context(InvalidRecordBatchSnafu { + reason: "pk_index required for dense encoding", + })?; + let mut builder = column_type.create_mutable_vector(distinct_keys.len()); + if offsets.is_empty() { + *offsets = vec![Vec::new(); distinct_keys.len()]; + } + for (&key, offsets) in distinct_keys.iter().zip(offsets) { + if pk_idx < codec.num_fields() { + let value = codec + .decode_value_at(values.value(key as usize), pk_idx, offsets) + .context(DecodeSnafu)?; + builder.push_value_ref(&value.as_value_ref()); } else { builder.push_null(); } @@ -872,6 +888,49 @@ impl DecodedPrimaryKeys { } } + /// Materializes all Dense tags in source-schema order with one sequential + /// traversal per key. Values are consumed immediately by their Arrow builders. + pub fn get_dense_tag_columns(&self) -> Result> { + let DecodedKeysInner::Dense { + codec, + values, + distinct_keys, + .. + } = &self.inner + else { + return InvalidRecordBatchSnafu { + reason: "expected dense primary key values", + } + .fail(); + }; + let mut builders: Vec<_> = codec + .fields() + .iter() + .map(|(_, field)| field.data_type().create_mutable_vector(distinct_keys.len())) + .collect(); + let mut value_buf = Vec::new(); + for &key in distinct_keys { + codec + .decode_dense_with(values.value(key as usize), &mut value_buf, |pos, value| { + builders[pos].push_value_ref(&value); + }) + .context(DecodeSnafu)?; + } + codec + .fields() + .iter() + .zip(builders) + .map(|((_, field), mut builder)| { + let values = builder.to_vector().to_arrow_array(); + if field.data_type().is_string() { + Ok(Arc::new(DictionaryArray::new(self.keys_array.clone(), values)) as ArrayRef) + } else { + take(&values, &self.keys_array, None).context(ComputeArrowSnafu) + } + }) + .collect() + } + /// Gets multiple sparse tag column arrays in one pass over the distinct keys, /// sharing each key's offset discovery between the columns. /// @@ -1004,6 +1063,8 @@ impl FlatConvertFormat { }) .collect(); decoded_columns.extend(decoded_pks.get_sparse_tag_columns(&columns)?); + } else if self.projected_primary_keys.len() == self.metadata.primary_key.len() { + decoded_columns.extend(decoded_pks.get_dense_tag_columns()?); } else { for (column_id, pk_index, column_index) in &self.projected_primary_keys { let column_metadata = &self.metadata.column_metadatas[*column_index]; @@ -1086,6 +1147,138 @@ mod tests { flat_sst_arrow_schema_column_num, override_pk_field_to_binary, to_flat_sst_arrow_schema, }; + #[test] + fn dense_tag_columns_match_eager_decoding() { + use datatypes::arrow::datatypes::UInt32Type; + use datatypes::value::Value; + use datatypes::vectors::Helper; + use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodecExt, SortField}; + + let columns = [ + (17, ConcreteDataType::string_datatype()), + (9, ConcreteDataType::int64_datatype()), + (2, ConcreteDataType::binary_datatype()), + ]; + let codec = DensePrimaryKeyCodec::with_fields( + columns + .iter() + .map(|(id, ty)| (*id, SortField::new(ty.clone()))) + .collect(), + ); + let rows = [ + vec![ + Value::from("中文\0abcdefgh"), + Value::Int64(-42), + Value::Binary(vec![0, 1, 255].into()), + ], + vec![Value::from(""), Value::Null, Value::Binary(vec![].into())], + vec![Value::Null, Value::Int64(i64::MAX), Value::Null], + ]; + let mut encoded: Vec<_> = rows + .iter() + .map(|row| codec.encode(row.iter().map(Value::as_value_ref)).unwrap()) + .collect(); + // Unreferenced dictionary entries must not be inspected. + encoded.push(vec![255]); + let row_keys = [2, 2, 0, 1, 0, 0]; + let pk = DictionaryArray::::new( + UInt32Array::from(row_keys.to_vec()), + Arc::new(BinaryArray::from_iter_values(&encoded)), + ); + let batch = RecordBatch::try_from_iter([ + ( + "ts", + Arc::new(TimestampMillisecondArray::from_iter_values(0..6)) as ArrayRef, + ), + ("__primary_key", Arc::new(pk) as ArrayRef), + ( + "__sequence", + Arc::new(UInt64Array::from(vec![1; 6])) as ArrayRef, + ), + ( + "__op_type", + Arc::new(UInt8Array::from(vec![1; 6])) as ArrayRef, + ), + ]) + .unwrap(); + let mut metadata = RegionMetadataBuilder::new(RegionId::new(1, 1)); + for (id, ty) in columns.iter().rev() { + metadata.push_column_metadata(ColumnMetadata { + column_id: *id, + column_schema: ColumnSchema::new(format!("tag_{id}"), ty.clone(), true), + semantic_type: SemanticType::Tag, + }); + } + metadata + .push_column_metadata(ColumnMetadata { + column_id: 23, + column_schema: ColumnSchema::new( + "ts", + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ), + semantic_type: SemanticType::Timestamp, + }) + .primary_key(vec![17, 9, 2]); + let metadata = Arc::new(metadata.build().unwrap()); + let format = FlatReadFormat::new( + metadata.clone(), + ReadColumns::new([17, 9, 2, 23]), + Some(batch.schema()), + "test", + false, + ) + .unwrap(); + for (offset, len) in [(0, 6), (1, 4), (3, 0)] { + let mut decoded = decode_primary_keys(&codec, &batch.slice(offset, len)).unwrap(); + let all = decoded.get_dense_tag_columns().unwrap(); + let converted = format + .convert_batch(batch.slice(offset, len), None) + .unwrap(); + // Out-of-order and repeated projections share the same key offsets. + for pos in [2, 0, 1, 0] { + let array = decoded + .get_tag_column(columns[pos].0, Some(pos), &columns[pos].1) + .unwrap(); + assert_eq!(array.to_data(), all[pos].to_data()); + assert_eq!(converted.column(pos).to_data(), all[pos].to_data()); + let actual = Helper::try_into_vector(array).unwrap(); + for (row, &key) in row_keys[offset..offset + len].iter().enumerate() { + let eager = codec + .decode_dense_without_column_id(&encoded[key as usize]) + .unwrap(); + assert_eq!(actual.get(row), eager[pos]); + } + } + // A position beyond the source schema remains NULL, never decoded + // with a newer schema's field order or type. + let missing = decoded + .get_tag_column(99, Some(3), &ConcreteDataType::int64_datatype()) + .unwrap(); + assert_eq!(missing.null_count(), len); + assert!(decoded.get_tag_column(17, None, &columns[0].1).is_err()); + } + let mut malformed_columns = batch.columns().to_vec(); + malformed_columns[1] = Arc::new(DictionaryArray::::new( + UInt32Array::from(vec![0; 6]), + Arc::new(BinaryArray::from(vec![&[0, 1][..]])), + )); + let malformed = RecordBatch::try_new(batch.schema(), malformed_columns).unwrap(); + assert!(format.convert_batch(malformed.clone(), None).is_err()); + // A partial projection must still avoid decoding an unneeded broken suffix. + let partial = FlatReadFormat::new( + metadata, + ReadColumns::new([17, 23]), + Some(batch.schema()), + "test", + false, + ) + .unwrap(); + let partial = partial.convert_batch(malformed, None).unwrap(); + let tag = Helper::try_into_vector(partial.column(0).clone()).unwrap(); + assert!((0..6).all(|row| tag.get(row).is_null())); + } + /// Builds a `RegionMetadata` with the given number of tags and fields. fn build_metadata( num_tags: usize, diff --git a/src/mito2/src/test_util/bench_util.rs b/src/mito2/src/test_util/bench_util.rs index ac635d2799a..e0334676e5b 100644 --- a/src/mito2/src/test_util/bench_util.rs +++ b/src/mito2/src/test_util/bench_util.rs @@ -21,6 +21,7 @@ use api::v1::value::ValueData; use api::v1::{Row, Rows, SemanticType}; use datafusion_common::Column; use datafusion_expr::{Expr, lit}; +use datatypes::arrow::datatypes::SchemaRef; use datatypes::arrow::record_batch::RecordBatch; use datatypes::data_type::ConcreteDataType; use datatypes::schema::ColumnSchema; @@ -30,7 +31,7 @@ use rand::seq::IndexedRandom; use store_api::metadata::{ ColumnMetadata, RegionMetadata, RegionMetadataBuilder, RegionMetadataRef, }; -use store_api::storage::RegionId; +use store_api::storage::{ColumnId, RegionId}; use table::predicate::Predicate; use crate::memtable::KeyValues; @@ -42,6 +43,23 @@ use crate::sst::parquet::flat_format::FlatReadFormat; use crate::sst::parquet::reader::SimpleFilterContext; use crate::test_util::memtable_util::region_metadata_to_row_schema; +/// Builds the actual encoded-PK-to-flat conversion used by the SST reader. +pub fn pk_materializer_for_bench( + metadata: RegionMetadataRef, + columns: Vec, + file_schema: SchemaRef, +) -> impl Fn(RecordBatch) -> crate::error::Result { + let format = FlatReadFormat::new( + metadata, + ReadColumns::new(columns), + Some(file_schema), + "bench", + false, + ) + .unwrap(); + move |batch| format.convert_batch(batch, None) +} + /// Builds a precise-filter benchmark with tags left encoded in the primary key. pub fn tag_filter_for_bench( metadata: RegionMetadataRef,