perf(mito-codec): streamline sparse label extraction and buffer sizing (#9217)

* perf(mito-codec): validate sparse labels once

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* perf(mito-codec): reserve complete sparse label chunks

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

---------

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>
This commit is contained in:
Lei, HUANG
2026-09-18 09:51:01 +00:00
committed by GitHub
parent 27090496d1
commit cb6c7a5c6d
6 changed files with 126 additions and 67 deletions
+4 -26
View File
@@ -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<u8>,
) -> Result<Option<&'a [u8]>> {
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
+21 -1
View File
@@ -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();
@@ -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<u8>,
) -> Result<Option<&'buf str>> {
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());
}
}
}
+18 -4
View File
@@ -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::<UInt32Type>::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())),
);
}
+2 -20
View File
@@ -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 {
+2 -16
View File
@@ -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)),
}
}
}