diff --git a/src/mito-codec/src/index.rs b/src/mito-codec/src/index.rs index 53ff5cd3c04..968173b878a 100644 --- a/src/mito-codec/src/index.rs +++ b/src/mito-codec/src/index.rs @@ -24,9 +24,7 @@ use store_api::codec::PrimaryKeyEncoding; use store_api::metadata::ColumnMetadata; use store_api::storage::ColumnId; -use crate::error::{ - FieldTypeMismatchSnafu, IndexEncodeNullSnafu, InvalidSparsePrimaryKeySnafu, Result, -}; +use crate::error::{FieldTypeMismatchSnafu, IndexEncodeNullSnafu, Result}; use crate::row_converter::sparse::{ RESERVED_COLUMN_ID_TABLE_ID, RESERVED_COLUMN_ID_TSID, SparsePrimaryKeyView, }; @@ -66,34 +64,14 @@ impl IndexValueCodec { column_id: ColumnId, buffer: &'a mut Vec, ) -> Result> { - let Some(encoded) = pk.encoded_value(column_id)? else { - return Ok(None); - }; - if encoded[0] == 0 { - return Ok(None); - } if matches!( column_id, RESERVED_COLUMN_ID_TABLE_ID | RESERVED_COLUMN_ID_TSID ) { - return Ok(Some(encoded)); + return pk.encoded_value(column_id); } - - buffer.clear(); - // The view has checked the Option marker, bytes marker and every chunk length. - // Reserving once avoids growing the buffer for each 8-byte chunk. - buffer.reserve(encoded.len() - 2); - for chunk in encoded[2..].chunks_exact(9) { - let len = usize::from(chunk[8]).min(8); - buffer.extend_from_slice(&chunk[..len]); - } - std::str::from_utf8(buffer).map_err(|_| { - InvalidSparsePrimaryKeySnafu { - reason: "label is not valid UTF-8", - } - .build() - })?; - Ok(Some(buffer.as_slice())) + pk.label(column_id, buffer) + .map(|value| value.map(str::as_bytes)) } /// Serializes a non-null `ValueRef` using the data type defined in `SortField` and writes diff --git a/src/mito-codec/src/row_converter/sparse.rs b/src/mito-codec/src/row_converter/sparse.rs index b2b3e0e3778..1be2fbacb2b 100644 --- a/src/mito-codec/src/row_converter/sparse.rs +++ b/src/mito-codec/src/row_converter/sparse.rs @@ -310,7 +310,7 @@ impl SparsePrimaryKeyCodec { { for (tag_column_id, tag_value) in row { let value_len = tag_value.len(); - buffer.reserve(6 + value_len / 8 * 9); + buffer.reserve(6 + value_len.div_ceil(8) * 9); buffer.put_u32(tag_column_id); buffer.put_u8(1); buffer.put_u8(!tag_value.is_empty() as u8); @@ -814,6 +814,26 @@ mod tests { assert_eq!(buffer, buffer_by_raw_encoding); } + #[test] + fn raw_encoding_matches_serde_at_chunk_boundaries() { + let codec = SparsePrimaryKeyCodec::schemaless(); + let mut actual = Vec::new(); + for len in [0, 7, 8, 9, 15, 16, 17, 1024] { + let label = "x".repeat(len); + let mut expected = Vec::new(); + codec.encode_internal(42, 7, &mut expected).unwrap(); + codec + .encode_to_vec([(1, ValueRef::String(&label))].into_iter(), &mut expected) + .unwrap(); + actual.clear(); + codec.encode_internal(42, 7, &mut actual).unwrap(); + codec + .encode_raw_tag_value([(1, label.as_bytes())].into_iter(), &mut actual) + .unwrap(); + assert_eq!(actual, expected, "label length {len}"); + } + } + #[test] fn test_encode_to_vec() { let region_metadata = test_region_metadata(); diff --git a/src/mito-codec/src/row_converter/sparse/checked.rs b/src/mito-codec/src/row_converter/sparse/checked.rs index 667f16fe804..1e57c13cc7f 100644 --- a/src/mito-codec/src/row_converter/sparse/checked.rs +++ b/src/mito-codec/src/row_converter/sparse/checked.rs @@ -86,6 +86,47 @@ impl<'a, 'b> SparsePrimaryKeyView<'a, 'b> { Ok(None) } + /// Extracts a UTF-8 label into a reusable buffer, validating it once. + /// Missing and null labels return `None` without modifying the buffer; + /// an empty label returns `Some("")`. Reserved numeric columns are rejected. + /// The returned string borrows `buffer`, whose contents are replaced for each label. + pub fn label<'buf>( + &mut self, + column_id: ColumnId, + buffer: &'buf mut Vec, + ) -> Result> { + ensure!( + !matches!( + column_id, + RESERVED_COLUMN_ID_TABLE_ID | RESERVED_COLUMN_ID_TSID + ), + InvalidSparsePrimaryKeySnafu { + reason: "reserved column is not a string label" + } + ); + let Some(encoded) = self.encoded_value(column_id)? else { + return Ok(None); + }; + if encoded[0] == 0 { + return Ok(None); + } + + buffer.clear(); + // encoded_value checked the markers and every chunk length. Reserve once + // before unchunking so the buffer does not grow for each 8-byte chunk. + buffer.reserve(encoded.len() - 2); + for chunk in encoded[2..].chunks_exact(9) { + let len = usize::from(chunk[8]).min(8); + buffer.extend_from_slice(&chunk[..len]); + } + std::str::from_utf8(buffer).map(Some).map_err(|_| { + InvalidSparsePrimaryKeySnafu { + reason: "label is not valid UTF-8", + } + .build() + }) + } + /// Returns the table id from the non-null prefix validated by [`Self::new`]. pub fn table_id(&self) -> u32 { (&self.pk[TABLE_ID_VALUE_OFFSET + 1..]).get_u32() @@ -265,4 +306,42 @@ mod tests { let mut view = SparsePrimaryKeyView::new(&invalid_utf8, &mut cache).unwrap(); assert!(IndexValueCodec::encode_sparse_value(&mut view, 1, &mut Vec::new()).is_err()); } + + #[test] + fn label_reuses_buffer_across_null_empty_invalid_and_valid_values() { + let codec = SparsePrimaryKeyCodec::schemaless(); + let mut cache = SparseOffsetsCache::new(); + let mut buffer = Vec::with_capacity(256); + let ptr = buffer.as_ptr(); + for label in [ + b"1234567\xe4\xb8\xad\0\xe6\x96\x87".as_slice(), + b"", + b"\xff", + b"after-error", + ] { + let mut pk = Vec::new(); + codec.encode_internal(42, 7, &mut pk).unwrap(); + codec + .encode_raw_tag_value([(1, label)].into_iter(), &mut pk) + .unwrap(); + pk.extend_from_slice(&2_u32.to_be_bytes()); + pk.push(0); + let mut view = SparsePrimaryKeyView::new(&pk, &mut cache).unwrap(); + let actual = view.label(1, &mut buffer); + match std::str::from_utf8(label) { + Ok(expected) => assert_eq!(actual.unwrap(), Some(expected)), + Err(_) => assert!(actual.is_err()), + } + assert_eq!(view.label(2, &mut buffer).unwrap(), None); + assert_eq!(view.label(3, &mut buffer).unwrap(), None); + // Null/missing values do not overwrite the scratch buffer, matching index extraction. + assert_eq!(buffer, label); + assert_eq!(buffer.as_ptr(), ptr); + assert!( + view.label(RESERVED_COLUMN_ID_TABLE_ID, &mut buffer) + .is_err() + ); + assert!(view.label(RESERVED_COLUMN_ID_TSID, &mut buffer).is_err()); + } + } } diff --git a/src/mito2/benches/bench_pk_tag_column.rs b/src/mito2/benches/bench_pk_tag_column.rs index 7f10c912138..d3ded3d0b05 100644 --- a/src/mito2/benches/bench_pk_tag_column.rs +++ b/src/mito2/benches/bench_pk_tag_column.rs @@ -95,6 +95,10 @@ fn metadata(tags: u32) -> RegionMetadataRef { /// Builds a sparse flat batch whose primary key dictionary holds /// `ROWS / rows_per_key` distinct keys with `tags` labels each. fn input(tags: u32, rows_per_key: usize) -> RecordBatch { + input_with_label_len(tags, rows_per_key, 24) +} + +fn input_with_label_len(tags: u32, rows_per_key: usize, label_len: usize) -> RecordBatch { let codec = SparsePrimaryKeyCodec::schemaless(); let mut keys = BinaryDictionaryBuilder::::new(); for series in 0..ROWS / rows_per_key { @@ -103,7 +107,14 @@ fn input(tags: u32, rows_per_key: usize) -> RecordBatch { .encode_internal((series / 128) as u32, series as u64, &mut key) .unwrap(); let labels: Vec<_> = (0..tags) - .map(|id| (id, format!("tag-{id:03}-value-{series:010}"))) + .map(|id| { + let mut value = format!("tag-{id:03}-value-{series:010}"); + value.extend(std::iter::repeat_n( + 'x', + label_len.saturating_sub(value.len()), + )); + (id, value) + }) .collect(); codec .encode_raw_tag_value( @@ -203,8 +214,8 @@ fn bench_pk_tag_filters(c: &mut Criterion) { const TAGS: u32 = 40; let metadata = metadata(TAGS); let mut group = c.benchmark_group("pk_tag_filters"); - for rows_per_key in [1, 32] { - let batch = input(TAGS, rows_per_key); + for (rows_per_key, label_len) in [(1, 24), (32, 24), (1, 1024), (32, 1024)] { + let batch = input_with_label_len(TAGS, rows_per_key, label_len); for predicate_count in [1, 2, 4, 8, 16, 32] { // Use distinct tags at the end of the key to exercise offset discovery. // Every row matches, so increasing the predicate count does not change selectivity. @@ -215,7 +226,10 @@ fn bench_pk_tag_filters(c: &mut Criterion) { // Validate the workload outside the timed section. assert_eq!(filter(batch.clone()).unwrap().unwrap().num_rows(), ROWS); group.bench_function( - BenchmarkId::new(format!("{rows_per_key}rpk"), predicate_count), + BenchmarkId::new( + format!("{rows_per_key}rpk_{label_len}bytes"), + predicate_count, + ), |b| b.iter(|| black_box(filter(black_box(batch.clone())).unwrap())), ); } diff --git a/src/mito2/src/series_index/writer.rs b/src/mito2/src/series_index/writer.rs index b6330aac044..0863b3b8519 100644 --- a/src/mito2/src/series_index/writer.rs +++ b/src/mito2/src/series_index/writer.rs @@ -26,7 +26,6 @@ use datatypes::arrow::datatypes::{DataType, Field, Schema, SchemaRef, UInt32Type use datatypes::arrow::record_batch::RecordBatch; use datatypes::prelude::ConcreteDataType; use datatypes::timestamp::timestamp_array_to_primitive; -use mito_codec::index::IndexValueCodec; use mito_codec::row_converter::SparseOffsetsCache; use mito_codec::row_converter::sparse::SparsePrimaryKeyView; use object_store::ObjectStore; @@ -603,25 +602,8 @@ fn decode_primary_key( let mut tags = Vec::with_capacity(tag_columns.len()); for (column_id, _) in tag_columns { - // `encode_sparse_value` returns None for missing and null labels and - // validates UTF-8 for string labels. - let value = IndexValueCodec::encode_sparse_value(&mut view, *column_id, buf) - .context(DecodeSnafu)?; - let tag = match value { - None => None, - Some(bytes) => { - let value = std::str::from_utf8(bytes).map_err(|_| { - InvalidRecordBatchSnafu { - reason: format!( - "sparse tag value of column {column_id} is not valid UTF-8" - ), - } - .build() - })?; - Some(value.to_string()) - } - }; - tags.push(tag); + let value = view.label(*column_id, buf).context(DecodeSnafu)?; + tags.push(value.map(str::to_owned)); } Ok(SeriesIndexRow { diff --git a/src/mito2/src/sst/parquet/flat_format.rs b/src/mito2/src/sst/parquet/flat_format.rs index c692dae0908..a99fc201baf 100644 --- a/src/mito2/src/sst/parquet/flat_format.rs +++ b/src/mito2/src/sst/parquet/flat_format.rs @@ -42,7 +42,6 @@ use datatypes::arrow::record_batch::RecordBatch; use datatypes::prelude::{ConcreteDataType, DataType}; use datatypes::value::ValueRef; use datatypes::vectors::MutableVector; -use mito_codec::index::IndexValueCodec; use mito_codec::row_converter::sparse::{ RESERVED_COLUMN_ID_TABLE_ID, RESERVED_COLUMN_ID_TSID, SparsePrimaryKeyView, }; @@ -798,23 +797,10 @@ fn push_sparse_tag_value_in_view( RESERVED_COLUMN_ID_TABLE_ID => builder.push_value_ref(&ValueRef::UInt32(view.table_id())), RESERVED_COLUMN_ID_TSID => builder.push_value_ref(&ValueRef::UInt64(view.tsid())), _ => { - // `encode_sparse_value` returns None for missing and null labels - // and validates UTF-8 for string labels. - let value = IndexValueCodec::encode_sparse_value(view, column_id, value_buf) - .context(DecodeSnafu)?; + let value = view.label(column_id, value_buf).context(DecodeSnafu)?; match value { None => builder.push_null(), - Some(bytes) => { - let value = std::str::from_utf8(bytes).map_err(|_| { - InvalidRecordBatchSnafu { - reason: format!( - "sparse tag value of column {column_id} is not valid UTF-8" - ), - } - .build() - })?; - builder.push_value_ref(&ValueRef::String(value)); - } + Some(value) => builder.push_value_ref(&ValueRef::String(value)), } } }