perf(mito2): lazily decode dense primary key columns (#9226)

* perf(mito2): lazily decode dense primary key columns

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

* perf(mito2): bypass lazy decoding for full primary keys

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

* fix(mito-codec): preserve prefix decoding errors

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

* 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 <ratuthomm@gmail.com>

* 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 <ratuthomm@gmail.com>

* 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 <ratuthomm@gmail.com>

---------

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>
This commit is contained in:
Lei, HUANG
2026-09-22 15:52:54 +00:00
committed by GitHub
parent 723da69b21
commit e91faa9df8
13 changed files with 1167 additions and 229 deletions
+30
View File
@@ -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<String>, 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<usize> {
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)),
}
}
}
+378 -103
View File
@@ -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<usize> {
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<usize> {
let invalid = || {
error::InvalidDensePrimaryKeySnafu {
reason: "truncated field or invalid encoding",
fn encoded_len(data_type: &ConcreteDataType, bytes: &[u8]) -> Result<usize> {
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<i64>
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<B: Buf>(
data_type: &ConcreteDataType,
deserializer: &mut Deserializer<B>,
) -> Option<Result<ValueRef<'static>>> {
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<usize> {
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<Item = Result<(ColumnId, Value)>> + '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<u8>,
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<usize>,
deserializer: &mut Deserializer<&[u8]>,
) -> Result<usize> {
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<usize>,
) -> Result<Value> {
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::<Result<Vec<_>>>()
.unwrap();
assert_eq!(
sequential
.iter()
.map(|(_, value)| value)
.collect::<Vec<_>>(),
row.iter().collect::<Vec<_>>()
);
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::<Result<Vec<_>>>()
.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::<Result<Vec<_>>>()
.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::<Result<Vec<_>>>()
.is_err()
);
assert!(
codec
.decode_value_at(&invalid_utf8, 0, &mut Vec::new())
.is_err()
);
}
#[test]
fn test_memcmp_dictionary() {
// Test Dictionary<i32, string>
@@ -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<String> 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)]
+112 -50
View File
@@ -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<Batch> = if encoded_only {
let pk = batch
.column_by_name("__primary_key")
.unwrap()
.as_any()
.downcast_ref::<DictionaryArray<UInt32Type>>()
.unwrap();
let values = pk.values().as_any().downcast_ref::<BinaryArray>().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();
+114 -3
View File
@@ -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::<UInt32Type>::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);
+125 -2
View File
@@ -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<u8>,
/// Possibly decoded `primary_key` values. Some places would decode it in advance.
pk_values: Option<CompositeValues>,
/// Lazily decoded Dense columns, shared by consumers such as index creators.
dense_pk_cache: Option<Box<DensePkCache>>,
/// Timestamps of rows, should be sorted and not null.
timestamps: VectorRef,
/// Sequences of rows
@@ -135,6 +137,13 @@ pub struct Batch {
fields_idx: Option<HashMap<ColumnId, usize>>,
}
/// Only used when no fully decoded primary key was supplied by the reader.
#[derive(Debug, PartialEq, Clone)]
struct DensePkCache {
offsets: Vec<usize>,
values: Vec<Option<Value>>,
}
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<u8>) {
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<Option<&Value>> {
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],
+53
View File
@@ -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<BloomFilterIndexer>,
last_mem_bloom_filter: usize,
/// Present only when the active creators together need every Dense PK field.
dense_pk_decoder: Option<DensePrimaryKeyCodec>,
#[cfg(feature = "vector_index")]
vector_indexer: Option<VectorIndexer>,
#[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<InvertedIndexer>,
bloom_filter_indexer: Option<BloomFilterIndexer>,
) -> 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());
+85 -5
View File
@@ -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" => {
+1
View File
@@ -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;
@@ -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 {
+22
View File
@@ -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 {
+220 -27
View File
@@ -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<Column
/// Decodes primary keys from a batch and returns decoded primary key information.
///
/// The batch must contain a primary key column at the expected index.
/// Sparse primary keys stay encoded: tag values are extracted lazily per column
/// Primary keys stay encoded: tag values are extracted lazily per column
/// in [`DecodedPrimaryKeys::get_tag_column`] without decoding every label.
/// The codec must describe the source key's field order and types.
pub fn decode_primary_keys(
codec: &dyn PrimaryKeyCodec,
batch: &RecordBatch,
@@ -729,15 +730,17 @@ pub fn decode_primary_keys(
values: pk_values_array.clone(),
distinct_keys,
},
PrimaryKeyEncoding::Dense => {
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<CompositeValues>),
/// Dense keys share positional offsets across projected tag columns.
Dense {
codec: DensePrimaryKeyCodec,
values: BinaryArray,
distinct_keys: Vec<u32>,
/// Initialized by positional access; full extraction needs no offsets.
offsets: Vec<Vec<usize>>,
},
/// 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<Vec<ArrayRef>> {
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::<UInt32Type>::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::<UInt32Type>::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,
+19 -1
View File
@@ -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<ColumnId>,
file_schema: SchemaRef,
) -> impl Fn(RecordBatch) -> crate::error::Result<RecordBatch> {
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,