diff --git a/src/mito-codec/src/error.rs b/src/mito-codec/src/error.rs index 3d3b0efcc07..b266ea9b04f 100644 --- a/src/mito-codec/src/error.rs +++ b/src/mito-codec/src/error.rs @@ -74,6 +74,13 @@ pub enum Error { location: Location, }, + #[snafu(display("Invalid dense primary key: {}", reason))] + InvalidDensePrimaryKey { + reason: String, + #[snafu(implicit)] + location: Location, + }, + #[snafu(display("Encode null value"))] IndexEncodeNull { #[snafu(implicit)] @@ -101,7 +108,9 @@ impl ErrorExt for Error { StatusCode::InvalidArguments } NotSupportedField { .. } | UnsupportedOperation { .. } => StatusCode::Unsupported, - InvalidSparsePrimaryKey { .. } => StatusCode::InvalidArguments, + InvalidSparsePrimaryKey { .. } | InvalidDensePrimaryKey { .. } => { + StatusCode::InvalidArguments + } EvaluateFilter { source, .. } => source.status_code(), } } diff --git a/src/mito-codec/src/row_converter/dense.rs b/src/mito-codec/src/row_converter/dense.rs index 6cc70feaeaf..56289a08b11 100644 --- a/src/mito-codec/src/row_converter/dense.rs +++ b/src/mito-codec/src/row_converter/dense.rs @@ -290,65 +290,84 @@ impl SortField { bytes: &[u8], deserializer: &mut Deserializer<&[u8]>, ) -> Result { - let pos = deserializer.position(); - if bytes[pos] == 0 { - deserializer.advance(1); - return Ok(1); - } - - Self::skip_deserialize_by_type(self.encode_data_type(), bytes, deserializer) + let len = self.encoded_field_len(&bytes[deserializer.position()..])?; + deserializer.advance(len); + Ok(len) } - fn skip_deserialize_by_type( - data_type: &ConcreteDataType, - bytes: &[u8], - deserializer: &mut Deserializer<&[u8]>, - ) -> Result { - let to_skip = match data_type { - ConcreteDataType::Boolean(_) => 2, - ConcreteDataType::Int8(_) | ConcreteDataType::UInt8(_) => 2, + /// 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", + } + .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, ConcreteDataType::Int16(_) | ConcreteDataType::UInt16(_) => 3, - ConcreteDataType::Int32(_) | ConcreteDataType::UInt32(_) => 5, - ConcreteDataType::Int64(_) | ConcreteDataType::UInt64(_) => 9, - ConcreteDataType::Float32(_) => 5, - ConcreteDataType::Float64(_) => 9, - ConcreteDataType::Binary(_) - | ConcreteDataType::Json(_) - | ConcreteDataType::Vector(_) => { - // Now the encoder encode binary as a list of bytes so we can't use - // skip bytes. - let pos_before = deserializer.position(); - let mut current = pos_before + 1; - while bytes[current] == 1 { - current += 2; - } - let to_skip = current - pos_before + 1; - deserializer.advance(to_skip); - return Ok(to_skip); - } - ConcreteDataType::String(_) => { - let pos_before = deserializer.position(); - deserializer.advance(1); - deserializer - .skip_bytes() - .context(error::DeserializeFieldSnafu)?; - return Ok(deserializer.position() - pos_before); - } - ConcreteDataType::Date(_) => 5, - ConcreteDataType::Timestamp(_) => 9, // We treat timestamp as Option - ConcreteDataType::Time(_) => 10, // i64 and 1 byte time unit - ConcreteDataType::Duration(_) => 10, + ConcreteDataType::Int32(_) + | ConcreteDataType::UInt32(_) + | ConcreteDataType::Float32(_) + | ConcreteDataType::Date(_) => 5, + ConcreteDataType::Int64(_) + | ConcreteDataType::UInt64(_) + | ConcreteDataType::Float64(_) + | ConcreteDataType::Timestamp(_) => 9, + ConcreteDataType::Time(_) | ConcreteDataType::Duration(_) => 10, ConcreteDataType::Interval(IntervalType::YearMonth(_)) => 5, ConcreteDataType::Interval(IntervalType::DayTime(_)) => 9, ConcreteDataType::Interval(IntervalType::MonthDayNano(_)) => 17, ConcreteDataType::Decimal128(_) => 19, - ConcreteDataType::Null(_) - | ConcreteDataType::List(_) - | ConcreteDataType::Struct(_) - | ConcreteDataType::Dictionary(_) => 0, + 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 { + data_type: data_type.clone(), + } + .fail(); + } }; - deserializer.advance(to_skip); - Ok(to_skip) + snafu::ensure!( + bytes.len() >= len, + error::InvalidDensePrimaryKeySnafu { + reason: "truncated field", + } + ); + Ok(len) } } @@ -411,6 +430,35 @@ impl DensePrimaryKeyCodec { Ok(values) } + /// Counts complete fields in a key encoded with a prefix of this codec's schema. + /// + /// Dense keys contain neither column ids nor types: callers must ensure that + /// existing fields have the same order and types. Only EOF between fields is + /// accepted; truncated fields and bytes beyond the schema return an error. + /// Field boundaries are checked in all builds; full value validation is only + /// enabled in debug builds and this crate's unit tests. + pub fn decode_prefix_len(&self, bytes: &[u8]) -> Result { + let mut deserializer = Deserializer::new(bytes); + for (index, (_, field)) in self.ordered_primary_key_columns.iter().enumerate() { + if !deserializer.has_remaining() { + return Ok(index); + } + let start = deserializer.position(); + let len = field.encoded_field_len(&bytes[start..])?; + // Production range mapping only needs boundaries, not allocated field values. + #[cfg(any(debug_assertions, test))] + field.deserialize(&mut Deserializer::new(&bytes[start..start + len]))?; + deserializer.advance(len); + } + snafu::ensure!( + !deserializer.has_remaining(), + error::InvalidDensePrimaryKeySnafu { + reason: "key contains bytes beyond the primary key schema", + } + ); + Ok(self.num_fields()) + } + /// Decode primary key values from bytes without column id. pub fn decode_dense_without_column_id(&self, bytes: &[u8]) -> Result> { let mut deserializer = Deserializer::new(bytes); @@ -594,6 +642,25 @@ mod tests { let value_ref = row.iter().map(|v| v.as_value_ref()).collect::>(); let result = encoder.encode(value_ref.iter().cloned()).unwrap(); + let boundaries: Vec<_> = (0..=row.len()) + .map(|count| { + encoder + .encode(value_ref[..count].iter().cloned()) + .unwrap() + .len() + }) + .collect(); + for end in 0..=result.len() { + match boundaries.iter().position(|boundary| *boundary == end) { + Some(count) => { + assert_eq!(count, encoder.decode_prefix_len(&result[..end]).unwrap()) + } + None => assert!( + encoder.decode_prefix_len(&result[..end]).is_err(), + "truncated {data_types:?} at {end}" + ), + } + } let decoded = encoder.decode(&result).unwrap().into_dense(); assert_eq!(decoded, row); let mut decoded = Vec::new(); @@ -610,6 +677,61 @@ mod tests { } } + #[test] + fn test_decode_prefix_len_accepts_only_complete_fields() { + let codec = DensePrimaryKeyCodec::with_fields(vec![ + (0, SortField::new(ConcreteDataType::string_datatype())), + (1, SortField::new(ConcreteDataType::int64_datatype())), + (2, SortField::new(ConcreteDataType::binary_datatype())), + (3, SortField::new(ConcreteDataType::string_datatype())), + ]); + let values = [ + Value::from("abcdefghijk"), + Value::Int64(42), + Value::Binary(vec![0, 255, 1].into()), + Value::Null, + ]; + let mut boundaries = vec![0]; + let mut encoded = Vec::new(); + for count in 1..=values.len() { + encoded.clear(); + codec + .encode_dense( + values[..count].iter().map(Value::as_value_ref), + &mut encoded, + ) + .unwrap(); + boundaries.push(encoded.len()); + assert_eq!(count, codec.decode_prefix_len(&encoded).unwrap()); + } + for end in 0..=encoded.len() { + match boundaries.iter().position(|boundary| *boundary == end) { + Some(count) => assert_eq!(count, codec.decode_prefix_len(&encoded[..end]).unwrap()), + None => assert!( + codec.decode_prefix_len(&encoded[..end]).is_err(), + "truncated at {end}" + ), + } + } + encoded.push(0); + assert!(codec.decode_prefix_len(&encoded).is_err()); + assert!(codec.decode_prefix_len(&[2]).is_err()); + let empty = DensePrimaryKeyCodec::with_fields(vec![]); + assert_eq!(0, empty.decode_prefix_len(&[]).unwrap()); + assert!(empty.decode_prefix_len(&[0]).is_err()); + } + + #[test] + fn test_prefix_len_validates_values_in_unit_tests() { + let codec = DensePrimaryKeyCodec::with_fields(vec![( + 0, + SortField::new(ConcreteDataType::boolean_datatype()), + )]); + // The field is complete, but 2 is not a boolean. Even release unit tests + // retain value validation through cfg(test). + assert!(codec.decode_prefix_len(&[1, 2]).is_err()); + } + #[test] fn test_memcmp() { let encoder = DensePrimaryKeyCodec::with_fields(vec![ diff --git a/src/mito2/benches/bench_compaction_picker.rs b/src/mito2/benches/bench_compaction_picker.rs index 439b9e6f226..99f02831536 100644 --- a/src/mito2/benches/bench_compaction_picker.rs +++ b/src/mito2/benches/bench_compaction_picker.rs @@ -78,7 +78,7 @@ fn bench_find_sorted_runs(c: &mut Criterion) { group.bench_function(format!("{}_new", case_name), |b| { let mut files = generate_same_timestamp_files(total_files, files_per_timestamp); b.iter(|| { - find_sorted_runs(black_box(&mut files)); + find_sorted_runs(black_box(&mut files), Ranged::overlap); }); }); @@ -121,7 +121,12 @@ fn bench_find_overlapping_items(c: &mut Criterion) { let mut r2 = SortedRun::from(files2); b.iter(|| { let mut result = vec![]; - find_overlapping_items(black_box(&mut r1), black_box(&mut r2), &mut result); + find_overlapping_items( + black_box(&mut r1), + black_box(&mut r2), + &mut result, + Ranged::overlap_inclusive, + ); }); }); } diff --git a/src/mito2/src/compaction/compactor.rs b/src/mito2/src/compaction/compactor.rs index 0a7bb61813b..c7d55fb8a28 100644 --- a/src/mito2/src/compaction/compactor.rs +++ b/src/mito2/src/compaction/compactor.rs @@ -208,7 +208,7 @@ pub async fn open_compaction_region( }; let current_version = { - let mut ssts = SstVersion::new(); + let mut ssts = SstVersion::new(region_metadata.clone()); ssts.add_files(file_purger.clone(), manifest.files.values().cloned()); CompactionVersion { metadata: region_metadata.clone(), @@ -1170,9 +1170,9 @@ mod tests { access_layer: env.access_layer.clone(), manifest_ctx, current_version: CompactionVersion { - metadata, + metadata: metadata.clone(), options: RegionOptions::default(), - ssts: Arc::new(SstVersion::new()), + ssts: Arc::new(SstVersion::new(metadata)), memtable_min_sequence: None, compaction_time_window: None, }, diff --git a/src/mito2/src/compaction/last_non_null.rs b/src/mito2/src/compaction/last_non_null.rs index 60f9bc254af..afcdc87e90f 100644 --- a/src/mito2/src/compaction/last_non_null.rs +++ b/src/mito2/src/compaction/last_non_null.rs @@ -15,7 +15,7 @@ use std::collections::{HashMap, HashSet}; use snafu::ResultExt; -use store_api::metadata::RegionMetadata; +use store_api::metadata::RegionMetadataRef; use store_api::storage::SequenceNumber; use crate::compaction::CompactionOutput; @@ -122,7 +122,7 @@ pub(super) fn inputs_precede_memtables( fn build_closed_outputs( seeds: Vec, files: Vec, - metadata: &RegionMetadata, + metadata: &RegionMetadataRef, ) -> Vec { // Maps every seed input to its seed, so a closure that reaches that input // can absorb the rest of the seed instead of letting it schedule separately. @@ -406,12 +406,12 @@ mod tests { #[case("dense_with_external", 1024)] #[case("time_dense_pk_disjoint", 1)] fn test_large_snapshot_closure(#[case] shape: &str, #[case] expected_count: usize) { - use bytes::Bytes; + use crate::compaction::test_util::{ + new_file_handle_with_size_sequence_and_primary_key_range, pk_range, + primary_key_metadata_for_test, + }; - use crate::compaction::test_util::new_file_handle_with_size_sequence_and_primary_key_range; - - let mut metadata = (*metadata_for_test()).clone(); - metadata.region_id = 0.into(); + let metadata = primary_key_metadata_for_test(); let files = (0..2048) .map(|i| { let (start, end, pk) = match shape { @@ -419,8 +419,8 @@ mod tests { "dense_with_external" if i < 1024 => (0, 1, None), "dense_with_external" => (i, i, None), "time_dense_pk_disjoint" => { - let pk = Bytes::copy_from_slice(&(i as u32).to_be_bytes()); - (0, 1, Some((pk.clone(), pk))) + let key = i.to_string(); + (0, 1, pk_range(key.as_bytes(), key.as_bytes())) } _ => unreachable!(), }; diff --git a/src/mito2/src/compaction/overlap.rs b/src/mito2/src/compaction/overlap.rs index 38cb7ad8a4a..f06c9058455 100644 --- a/src/mito2/src/compaction/overlap.rs +++ b/src/mito2/src/compaction/overlap.rs @@ -16,10 +16,11 @@ use std::collections::HashMap; use std::ops::Range; use common_time::Timestamp; -use store_api::metadata::RegionMetadata; +use store_api::metadata::{RegionMetadata, RegionMetadataRef}; use crate::compaction::run::primary_key_ranges_overlap; use crate::sst::file::{FileHandle, RegionFileId}; +use crate::sst::primary_key::PrimaryKeyRangeMapper; /// Snapshot-local overlap candidates, ordered by start time. Each subtree stores /// its maximum active end. Removing visited files prunes dense internal overlaps. @@ -28,6 +29,7 @@ use crate::sst::file::{FileHandle, RegionFileId}; pub(super) struct FileOverlapIndex<'a> { /// Current region metadata used to decide whether PK bounds can safely exclude overlaps. metadata: &'a RegionMetadata, + primary_key_mapper: PrimaryKeyRangeMapper, /// Candidates in original snapshot order, retained after removal so their /// indices remain stable and drained matches can preserve merge input order. files: Vec, @@ -54,7 +56,7 @@ struct OverlapQuery<'a> { } impl<'a> FileOverlapIndex<'a> { - pub(super) fn new(files: Vec, metadata: &'a RegionMetadata) -> Self { + pub(super) fn new(files: Vec, metadata: &'a RegionMetadataRef) -> Self { let mut by_start: Vec<_> = (0..files.len()).collect(); by_start.sort_unstable_by_key(|&i| (files[i].time_range().0, i)); // A single empty leaf also handles an empty snapshot. @@ -70,6 +72,7 @@ impl<'a> FileOverlapIndex<'a> { } Self { metadata, + primary_key_mapper: PrimaryKeyRangeMapper::new(metadata.clone()), files, by_start, positions, @@ -127,7 +130,13 @@ impl<'a> FileOverlapIndex<'a> { } if leaves.len() == 1 { let i = self.by_start[leaves.start]; - return files_may_overlap(query.input, &self.files[i], self.metadata).then_some(i); + return files_may_overlap( + query.input, + &self.files[i], + self.metadata, + &self.primary_key_mapper, + ) + .then_some(i); } let mid = leaves.start + leaves.len() / 2; self.find_in_subtree(2 * node, leaves.start..mid, query) @@ -137,7 +146,12 @@ impl<'a> FileOverlapIndex<'a> { /// SST bounds are inclusive. Missing statistics, foreign encodings and schema /// evolution must not exclude a possible logical-key dependency. -fn files_may_overlap(lhs: &FileHandle, rhs: &FileHandle, metadata: &RegionMetadata) -> bool { +fn files_may_overlap( + lhs: &FileHandle, + rhs: &FileHandle, + metadata: &RegionMetadata, + mapper: &PrimaryKeyRangeMapper, +) -> bool { let (lhs_start, lhs_end) = lhs.time_range(); let (rhs_start, rhs_end) = rhs.time_range(); if lhs_start.max(rhs_start) > lhs_end.min(rhs_end) { @@ -155,7 +169,7 @@ fn files_may_overlap(lhs: &FileHandle, rhs: &FileHandle, metadata: &RegionMetada { return true; } - match (lhs.primary_key_range(), rhs.primary_key_range()) { + match (lhs.primary_key_range(mapper), rhs.primary_key_range(mapper)) { (Some(lhs), Some(rhs)) if lhs.0 <= lhs.1 && rhs.0 <= rhs.1 => { primary_key_ranges_overlap(&lhs, &rhs) } @@ -166,12 +180,13 @@ fn files_may_overlap(lhs: &FileHandle, rhs: &FileHandle, metadata: &RegionMetada #[cfg(test)] mod tests { use std::collections::HashSet; + use std::sync::Arc; - use bytes::Bytes; use rand::{Rng, SeedableRng}; use store_api::storage::FileId; use super::*; + use crate::compaction::test_util::{pk_range, primary_key_metadata_for_test}; use crate::sst::file::FileMeta; use crate::test_util::memtable_util::metadata_for_test; use crate::test_util::new_noop_file_purger; @@ -187,12 +202,7 @@ mod tests { ..Default::default() }, new_noop_file_purger(), - pk.map(|(start, end)| { - ( - Bytes::copy_from_slice(start.as_bytes()), - Bytes::copy_from_slice(end.as_bytes()), - ) - }), + pk.and_then(|(start, end)| pk_range(start.as_bytes(), end.as_bytes())), ) } @@ -212,9 +222,11 @@ mod tests { #[case] foreign: bool, #[case] expected: bool, ) { - let mut metadata = (*metadata_for_test()).clone(); + let mut metadata = (*primary_key_metadata_for_test()).clone(); metadata.region_id = 0.into(); metadata.schema_version = schema_version; + let metadata = Arc::new(metadata); + let ranges = PrimaryKeyRangeMapper::new(metadata.clone()); let lhs = file(0, 10, Some(("a", "b"))); let mut rhs = file(start, start + 10, pk); if foreign { @@ -223,11 +235,11 @@ mod tests { rhs = FileHandle::new_with_primary_key_range( meta, new_noop_file_purger(), - rhs.primary_key_range(), + rhs.raw_primary_key_range(), ); } - assert_eq!(expected, files_may_overlap(&lhs, &rhs, &metadata)); - assert_eq!(expected, files_may_overlap(&rhs, &lhs, &metadata)); + assert_eq!(expected, files_may_overlap(&lhs, &rhs, &metadata, &ranges)); + assert_eq!(expected, files_may_overlap(&rhs, &lhs, &metadata, &ranges)); } #[test] @@ -247,8 +259,8 @@ mod tests { #[test] fn test_index_matches_linear_scan_after_removals() { - let mut metadata = (*metadata_for_test()).clone(); - metadata.region_id = 0.into(); + let metadata = primary_key_metadata_for_test(); + let ranges = PrimaryKeyRangeMapper::new(metadata.clone()); let mut rng = rand::rngs::StdRng::seed_from_u64(9146); for count in [0, 1, 7, 32, 127] { let files: Vec<_> = (0..count) @@ -276,7 +288,8 @@ mod tests { let expected: Vec<_> = files .iter() .filter(|f| { - !removed.contains(&f.file_id()) && files_may_overlap(&query, f, &metadata) + !removed.contains(&f.file_id()) + && files_may_overlap(&query, f, &metadata, &ranges) }) .map(FileHandle::file_id) .collect(); diff --git a/src/mito2/src/compaction/picker.rs b/src/mito2/src/compaction/picker.rs index de2e9d9d310..ec9fa91f4c8 100644 --- a/src/mito2/src/compaction/picker.rs +++ b/src/mito2/src/compaction/picker.rs @@ -91,7 +91,8 @@ impl From<&PickerOutput> for SerializedPickerOutput { } impl PickerOutput { - /// Converts a [SerializedPickerOutput] to a [PickerOutput]. + /// Converts a [SerializedPickerOutput] to a [PickerOutput]. File statistics + /// retain their original encoding; comparisons supply their own schema context. pub fn from_serialized( input: SerializedPickerOutput, file_purger: Arc, diff --git a/src/mito2/src/compaction/run.rs b/src/mito2/src/compaction/run.rs index 5698a6ce667..2ad9f2f1eba 100644 --- a/src/mito2/src/compaction/run.rs +++ b/src/mito2/src/compaction/run.rs @@ -23,6 +23,7 @@ use common_base::BitVec; use common_time::Timestamp; use crate::sst::file::FileHandle; +use crate::sst::primary_key::PrimaryKeyRangeMapper; /// Trait for any items with specific range (both boundaries are inclusive). pub trait Ranged { @@ -64,10 +65,12 @@ pub(crate) fn merge_primary_key_ranges( } } +/// Uses the caller's inclusive overlap test after pruning disjoint time ranges. pub fn find_overlapping_items( l: &mut SortedRun, r: &mut SortedRun, result: &mut Vec, + overlaps: impl Fn(&T, &T) -> bool, ) { if l.items.is_empty() || r.items.is_empty() { return; @@ -114,7 +117,7 @@ pub fn find_overlapping_items( } // We have an overlap (inclusive: touching boundaries count) - if lhs.overlap_inclusive(&r.items[j]) { + if overlaps(lhs, &r.items[j]) { if !selected[lhs_idx] { result.push(lhs.clone()); selected.set(lhs_idx, true); @@ -147,37 +150,41 @@ pub trait Item: Ranged + Clone { fn size(&self) -> usize; } +// Physical handles only supply time ranges. PK-aware comparisons need a target schema. impl Ranged for FileHandle { type BoundType = Timestamp; fn range(&self) -> (Self::BoundType, Self::BoundType) { self.time_range() } +} - fn overlap(&self, other: &Self) -> bool { - let (lhs_start, lhs_end) = self.range(); - let (rhs_start, rhs_end) = other.range(); - if lhs_start.max(rhs_start) >= lhs_end.min(rhs_end) { - return false; - } +/// Tests exclusive time overlap and inclusive PK overlap in one pinned schema. +pub(crate) fn files_overlap( + lhs: &FileHandle, + rhs: &FileHandle, + mapper: &PrimaryKeyRangeMapper, +) -> bool { + lhs.overlap(rhs) && file_primary_keys_overlap(lhs, rhs, mapper) +} - match (&self.primary_key_range(), &other.primary_key_range()) { - (Some(lhs), Some(rhs)) => primary_key_ranges_overlap(lhs, rhs), - _ => true, - } - } +/// Includes touching time and PK boundaries when checking file dependencies. +pub(crate) fn files_overlap_inclusive( + lhs: &FileHandle, + rhs: &FileHandle, + mapper: &PrimaryKeyRangeMapper, +) -> bool { + lhs.overlap_inclusive(rhs) && file_primary_keys_overlap(lhs, rhs, mapper) +} - fn overlap_inclusive(&self, other: &Self) -> bool { - let (lhs_start, lhs_end) = self.range(); - let (rhs_start, rhs_end) = other.range(); - if lhs_start.max(rhs_start) > lhs_end.min(rhs_end) { - return false; - } - - match (&self.primary_key_range(), &other.primary_key_range()) { - (Some(lhs), Some(rhs)) => primary_key_ranges_overlap(lhs, rhs), - _ => true, - } +fn file_primary_keys_overlap( + lhs: &FileHandle, + rhs: &FileHandle, + mapper: &PrimaryKeyRangeMapper, +) -> bool { + match (lhs.primary_key_range(mapper), rhs.primary_key_range(mapper)) { + (Some(lhs), Some(rhs)) => primary_key_ranges_overlap(&lhs, &rhs), + _ => true, } } @@ -251,8 +258,8 @@ where } } -/// Finds sorted runs in given items. -pub fn find_sorted_runs(items: &mut [T]) -> Vec> +/// Finds sorted runs using the caller's overlap test within each active time range. +pub fn find_sorted_runs(items: &mut [T], overlaps: impl Fn(&T, &T) -> bool) -> Vec> where T: Item, { @@ -296,7 +303,7 @@ where let mut overlaps_any = false; for idx in &active_run_item_indices { let run_item = ¤t_run.items[*idx]; - if run_item.overlap(item) { + if overlaps(run_item, item) { overlaps_any = true; break; } @@ -490,11 +497,13 @@ where #[cfg(test)] mod tests { - use bytes::Bytes; use store_api::storage::FileId; use super::*; - use crate::compaction::test_util::new_file_handle_with_size_sequence_and_primary_key_range; + use crate::compaction::test_util::{ + new_file_handle_with_size_sequence_and_primary_key_range, pk_range, + primary_key_mapper_for_test, + }; #[derive(Clone, Debug, PartialEq)] struct MockFile { @@ -528,10 +537,6 @@ mod tests { .collect() } - fn pk_range(min: &'static [u8], max: &'static [u8]) -> Option<(Bytes, Bytes)> { - Some((Bytes::from_static(min), Bytes::from_static(max))) - } - fn check_sorted_runs( ranges: &[(i64, i64)], expected_runs: &[Vec<(i64, i64)>], @@ -539,7 +544,7 @@ mod tests { let mut files = build_items(ranges); let mut files_clone = files.clone(); - let runs = find_sorted_runs(&mut files); + let runs = find_sorted_runs(&mut files, Ranged::overlap); let result_file_ranges: Vec> = runs .iter() @@ -574,7 +579,7 @@ mod tests { let mut files = build_items(ranges); let mut files_for_original = files.clone(); - let runs = find_sorted_runs(&mut files); + let runs = find_sorted_runs(&mut files, Ranged::overlap); let original_runs = find_sorted_runs_original(&mut files_for_original); assert_eq!(sorted_run_ranges(&original_runs), sorted_run_ranges(&runs)); @@ -661,6 +666,7 @@ mod tests { &mut SortedRun::from(Vec::::new()), &mut SortedRun::from(Vec::::new()), &mut result, + Ranged::overlap_inclusive, ); assert_eq!(result, Vec::::new()); @@ -669,6 +675,7 @@ mod tests { &mut SortedRun::from(files1.clone()), &mut SortedRun::from(Vec::::new()), &mut result, + Ranged::overlap_inclusive, ); assert_eq!(result, Vec::::new()); @@ -676,6 +683,7 @@ mod tests { &mut SortedRun::from(Vec::::new()), &mut SortedRun::from(files1.clone()), &mut result, + Ranged::overlap_inclusive, ); assert_eq!(result, Vec::::new()); @@ -686,6 +694,7 @@ mod tests { &mut SortedRun::from(files1), &mut SortedRun::from(files2), &mut result, + Ranged::overlap_inclusive, ); assert_eq!(result, Vec::::new()); @@ -696,6 +705,7 @@ mod tests { &mut SortedRun::from(files1), &mut SortedRun::from(files2), &mut result, + Ranged::overlap_inclusive, ); assert_eq!(result.len(), 2); assert_eq!(result[0].range(), (1, 5)); @@ -708,6 +718,7 @@ mod tests { &mut SortedRun::from(files1), &mut SortedRun::from(files2), &mut result, + Ranged::overlap_inclusive, ); assert_eq!(result.len(), 6); @@ -718,6 +729,7 @@ mod tests { &mut SortedRun::from(files1), &mut SortedRun::from(files2), &mut result, + Ranged::overlap_inclusive, ); assert_eq!(result.len(), 2); // Should overlap since ranges are inclusive @@ -728,6 +740,7 @@ mod tests { &mut SortedRun::from(files1), &mut SortedRun::from(files2), &mut result, + Ranged::overlap_inclusive, ); assert_eq!(result.len(), 2); @@ -738,6 +751,7 @@ mod tests { &mut SortedRun::from(files1), &mut SortedRun::from(files2), &mut result, + Ranged::overlap_inclusive, ); assert_eq!(result.len(), 2); @@ -748,6 +762,7 @@ mod tests { &mut SortedRun::from(files1), &mut SortedRun::from(files2), &mut result, + Ranged::overlap_inclusive, ); assert_eq!(result.len(), 4); // Should find both overlaps } @@ -773,7 +788,7 @@ mod tests { pk_range(b"x", b"z"), ); - assert!(!lhs.overlap(&rhs)); + assert!(!files_overlap(&lhs, &rhs, &primary_key_mapper_for_test())); } #[test] @@ -799,7 +814,8 @@ mod tests { ), ]; - let runs = find_sorted_runs(&mut files); + let ranges = primary_key_mapper_for_test(); + let runs = find_sorted_runs(&mut files, |lhs, rhs| files_overlap(lhs, rhs, &ranges)); assert_eq!(1, runs.len()); assert_eq!(2, runs[0].items().len()); @@ -837,7 +853,8 @@ mod tests { ), ]; - let runs = find_sorted_runs(&mut files); + let ranges = primary_key_mapper_for_test(); + let runs = find_sorted_runs(&mut files, |lhs, rhs| files_overlap(lhs, rhs, &ranges)); assert_eq!(2, runs.len()); assert_eq!(2, runs[0].items().len()); @@ -870,7 +887,10 @@ mod tests { ]); let mut result = Vec::new(); - find_overlapping_items(&mut left, &mut right, &mut result); + let ranges = primary_key_mapper_for_test(); + find_overlapping_items(&mut left, &mut right, &mut result, |lhs, rhs| { + files_overlap_inclusive(lhs, rhs, &ranges) + }); assert!(result.is_empty()); } @@ -896,6 +916,6 @@ mod tests { pk_range(b"a", b"f"), ); - assert!(!lhs.overlap(&rhs)); + assert!(!files_overlap(&lhs, &rhs, &primary_key_mapper_for_test())); } } diff --git a/src/mito2/src/compaction/test_util.rs b/src/mito2/src/compaction/test_util.rs index 453af8a8a68..1744b37c95b 100644 --- a/src/mito2/src/compaction/test_util.rs +++ b/src/mito2/src/compaction/test_util.rs @@ -19,6 +19,7 @@ use std::time::Duration; use bytes::Bytes; use common_base::Plugins; use common_time::Timestamp; +use store_api::metadata::RegionMetadataRef; use store_api::storage::FileId; use crate::cache::CacheManager; @@ -26,11 +27,30 @@ use crate::compaction::compactor::{CompactionRegion, CompactionVersion}; use crate::config::MitoConfig; use crate::region::options::RegionOptions; use crate::sst::file::{FileHandle, FileMeta, Level}; +use crate::sst::primary_key::PrimaryKeyRangeMapper; use crate::sst::version::SstVersion; use crate::test_util::memtable_util::metadata_for_test; use crate::test_util::new_noop_file_purger; use crate::test_util::scheduler_util::SchedulerEnv; +pub(crate) fn primary_key_metadata_for_test() -> RegionMetadataRef { + let mut metadata = crate::test_util::memtable_util::metadata_with_primary_key(vec![0], false); + metadata.region_id = 0.into(); + Arc::new(metadata) +} + +pub(crate) fn primary_key_mapper_for_test() -> PrimaryKeyRangeMapper { + PrimaryKeyRangeMapper::new(primary_key_metadata_for_test()) +} + +/// Encodes the single string tag used by compaction range fixtures. +pub(crate) fn pk_range(min: &[u8], max: &[u8]) -> Option<(Bytes, Bytes)> { + let encode = |key| { + crate::test_util::sst_util::new_primary_key(&[std::str::from_utf8(key).unwrap()]).into() + }; + Some((encode(min), encode(max))) +} + /// Test util to create file handles. pub fn new_file_handle( file_id: FileId, @@ -144,9 +164,12 @@ pub(crate) async fn compaction_region_with_ssts( ttl: Duration, ) -> CompactionRegion { let env = SchedulerEnv::new().await; - let metadata = metadata_for_test(); + let mut metadata = (*metadata_for_test()).clone(); + // Match the table used by new_file_handle* and default FileMeta fixtures. + metadata.region_id = 0.into(); + let metadata = Arc::new(metadata); let manifest_ctx = env.mock_manifest_context(metadata.clone()).await; - let mut ssts = SstVersion::new(); + let mut ssts = SstVersion::new(metadata.clone()); ssts.add_files( Arc::new(crate::sst::file_purger::NoopFilePurger), files.into_iter(), diff --git a/src/mito2/src/compaction/twcs.rs b/src/mito2/src/compaction/twcs.rs index b2e3274b062..dd407b30382 100644 --- a/src/mito2/src/compaction/twcs.rs +++ b/src/mito2/src/compaction/twcs.rs @@ -31,11 +31,12 @@ use crate::compaction::buckets::infer_time_bucket; use crate::compaction::compactor::CompactionRegion; use crate::compaction::picker::{Picker, PickerOutput, get_expired_ssts}; use crate::compaction::run::{ - Ranged, SortedRun, find_sorted_runs, find_sorted_runs_by_time_range, merge_primary_key_ranges, - primary_key_ranges_overlap, + Ranged, SortedRun, files_overlap, files_overlap_inclusive, find_sorted_runs, + find_sorted_runs_by_time_range, merge_primary_key_ranges, primary_key_ranges_overlap, }; use crate::error::{JoinSnafu, Result}; use crate::sst::file::{FileHandle, Level, overlaps}; +use crate::sst::primary_key::PrimaryKeyRangeMapper; use crate::sst::version::LevelMeta; const LEVEL_COMPACTED: Level = 1; @@ -69,12 +70,14 @@ struct WindowPickContext<'a> { files: &'a Window, windows: &'a BTreeMap, phase: PickPhase, + primary_key_mapper: &'a PrimaryKeyRangeMapper, } struct WindowOutputContext { active_window: Option, time_window_size: Option, max_outputs: Option, + primary_key_mapper: Arc, } /// A mixed L0/L1 compaction may rewrite at most this many L1 rows per L0 row. @@ -135,6 +138,7 @@ impl TwcsPicker { active_window, time_window_size, max_outputs, + primary_key_mapper, } = context; let mut output = vec![]; let windows = time_windows @@ -175,6 +179,7 @@ impl TwcsPicker { let picker = self.clone(); let time_windows = time_windows.clone(); let window = *window; + let primary_key_mapper = primary_key_mapper.clone(); handles.push(common_runtime::spawn_blocking_compact(move || { time_windows.get(&window).map(|files| { ( @@ -186,6 +191,7 @@ impl TwcsPicker { files, windows: &time_windows, phase, + primary_key_mapper: &primary_key_mapper, }, ), ) @@ -238,6 +244,7 @@ impl TwcsPicker { files, windows, phase, + primary_key_mapper, } = context; let is_active_window = active_window == Some(files.time_window); let window = &files.time_window; @@ -275,13 +282,21 @@ impl TwcsPicker { } let (inputs, found_runs) = if is_active_window { match phase { - PickPhase::HasL0 if num_l0_files >= self.trigger_file_num => { - pick_candidate_files(l0_files, self.max_output_file_size, pick_count_first) - } + PickPhase::HasL0 if num_l0_files >= self.trigger_file_num => pick_candidate_files( + l0_files, + self.max_output_file_size, + pick_count_first, + primary_key_mapper, + ), PickPhase::L1FileReduction | PickPhase::L1OverlapOnly if num_l1_files >= self.active_window_l1_merge_trigger => { - pick_l1_candidate_files(l1_files, self.max_output_file_size, phase) + pick_l1_candidate_files( + l1_files, + self.max_output_file_size, + phase, + primary_key_mapper, + ) } _ => (vec![], 0), } @@ -293,6 +308,7 @@ impl TwcsPicker { l0_file_num: self.inactive_window_trigger_file_num, l1_file_num: self.inactive_window_l1_merge_trigger, phase, + primary_key_mapper, }, self.max_output_file_size, ) @@ -303,7 +319,7 @@ impl TwcsPicker { let filter_deleted = !self.append_mode && !window_has_overlap(files, windows) - && !selected_overlaps_unselected(&inputs, files); + && !selected_overlaps_unselected(&inputs, files, primary_key_mapper); if inputs.len() > 1 { // If we have more than one file to compact. @@ -336,72 +352,106 @@ impl TwcsPicker { /// The rewrite is bounded by the L0 bytes and leaves large compacted files /// untouched. If nothing qualifies, the window is left as-is. #[derive(Debug, Clone, Copy)] -struct InactiveWindowPick { +struct InactiveWindowPick<'a> { l0_file_num: usize, l1_file_num: usize, phase: PickPhase, + primary_key_mapper: &'a PrimaryKeyRangeMapper, } fn pick_inactive_window_files( l0_files: Vec, l1_files: Vec, - pick: InactiveWindowPick, + pick: InactiveWindowPick<'_>, max_output_file_size: Option, ) -> (Vec, usize) { if !matches!(pick.phase, PickPhase::HasL0) { return if l1_files.len() >= pick.l1_file_num { - pick_l1_candidate_files(l1_files, max_output_file_size, pick.phase) + pick_l1_candidate_files( + l1_files, + max_output_file_size, + pick.phase, + pick.primary_key_mapper, + ) } else { (vec![], 0) }; } if l0_files.len() >= pick.l0_file_num { - let pick = pick_candidate_files(l0_files.clone(), max_output_file_size, pick_count_first); - if !pick.0.is_empty() { - return pick; + let candidate = pick_candidate_files( + l0_files.clone(), + max_output_file_size, + pick_count_first, + pick.primary_key_mapper, + ); + if !candidate.0.is_empty() { + return candidate; } } - let pick = pick_candidate_files(l0_files.clone(), max_output_file_size, pick_count_first); - if !pick.0.is_empty() { - return pick; + let candidate = pick_candidate_files( + l0_files.clone(), + max_output_file_size, + pick_count_first, + pick.primary_key_mapper, + ); + if !candidate.0.is_empty() { + return candidate; } let mut all_files = l0_files.clone(); all_files.extend(l1_files); - let pick = pick_candidate_files(all_files, max_output_file_size, pick_mixed_within_budget); - if !pick.0.is_empty() { - return pick; + let candidate = pick_candidate_files( + all_files, + max_output_file_size, + pick_mixed_within_budget, + pick.primary_key_mapper, + ); + if !candidate.0.is_empty() { + return candidate; } - pick_candidate_files(l0_files, max_output_file_size, pick_unbalanced_count_first) + pick_candidate_files( + l0_files, + max_output_file_size, + pick_unbalanced_count_first, + pick.primary_key_mapper, + ) } fn pick_l1_candidate_files( l1_files: Vec, max_output_file_size: Option, phase: PickPhase, + mapper: &PrimaryKeyRangeMapper, ) -> (Vec, usize) { let picker = match phase { PickPhase::L1FileReduction => pick_l1_file_reduction, PickPhase::L1OverlapOnly => pick_l1_overlap_only, PickPhase::HasL0 => return (vec![], 0), }; - pick_candidate_files(l1_files, max_output_file_size, picker) + pick_candidate_files(l1_files, max_output_file_size, picker, mapper) } +type CandidatePicker = + fn(Vec>, Option, &PrimaryKeyRangeMapper) -> Vec; + fn pick_candidate_files( mut files: Vec, max_output_file_size: Option, - picker: fn(Vec>, Option) -> Vec, + picker: CandidatePicker, + mapper: &PrimaryKeyRangeMapper, ) -> (Vec, usize) { let sorted_runs = if files.len() < 1024 { - find_sorted_runs(&mut files) + find_sorted_runs(&mut files, |lhs, rhs| files_overlap(lhs, rhs, mapper)) } else { find_sorted_runs_by_time_range(&mut files) }; let found_runs = sorted_runs.len(); - (picker(sorted_runs, max_output_file_size), found_runs) + ( + picker(sorted_runs, max_output_file_size, mapper), + found_runs, + ) } #[derive(Debug)] @@ -443,6 +493,7 @@ impl Candidate { file: &OrderedFile, preceding: &[OrderedFile], participations: &mut Vec, + mapper: &PrimaryKeyRangeMapper, ) { self.num_files += 1; let file_size = file.file.size() as usize; @@ -452,7 +503,8 @@ impl Candidate { let mut participates = false; for (offset, other) in preceding.iter().enumerate() { - if file.run_id != other.run_id && file.file.overlap_inclusive(other.file) { + if file.run_id != other.run_id && files_overlap_inclusive(file.file, other.file, mapper) + { if !participations[offset] { participations[offset] = true; self.overlap_participants += 1; @@ -582,15 +634,22 @@ impl CandidateScore { fn pick_count_first( sorted_runs: Vec>, max_output_file_size: Option, + mapper: &PrimaryKeyRangeMapper, ) -> Vec { - pick_count_first_where(sorted_runs, max_output_file_size, is_balanced_candidate) + pick_count_first_where( + sorted_runs, + max_output_file_size, + mapper, + is_balanced_candidate, + ) } fn pick_l1_file_reduction( sorted_runs: Vec>, max_output_file_size: Option, + mapper: &PrimaryKeyRangeMapper, ) -> Vec { - pick_count_first_where(sorted_runs, max_output_file_size, |candidate| { + pick_count_first_where(sorted_runs, max_output_file_size, mapper, |candidate| { is_balanced_candidate(candidate) && candidate.file_reduction(max_output_file_size) > 0 }) } @@ -598,8 +657,9 @@ fn pick_l1_file_reduction( fn pick_l1_overlap_only( sorted_runs: Vec>, max_output_file_size: Option, + mapper: &PrimaryKeyRangeMapper, ) -> Vec { - pick_count_first_where(sorted_runs, max_output_file_size, |candidate| { + pick_count_first_where(sorted_runs, max_output_file_size, mapper, |candidate| { is_balanced_candidate(candidate) && candidate.file_reduction(max_output_file_size) == 0 && candidate.overlap_participants > 0 @@ -610,8 +670,9 @@ fn pick_l1_overlap_only( fn pick_mixed_count_first( sorted_runs: Vec>, max_output_file_size: Option, + mapper: &PrimaryKeyRangeMapper, ) -> Vec { - pick_count_first_where(sorted_runs, max_output_file_size, |candidate| { + pick_count_first_where(sorted_runs, max_output_file_size, mapper, |candidate| { candidate.has_mixed_levels() && is_balanced_candidate(candidate) }) } @@ -621,8 +682,9 @@ fn pick_mixed_count_first( fn pick_mixed_within_budget( sorted_runs: Vec>, max_output_file_size: Option, + mapper: &PrimaryKeyRangeMapper, ) -> Vec { - pick_count_first_where(sorted_runs, max_output_file_size, |candidate| { + pick_count_first_where(sorted_runs, max_output_file_size, mapper, |candidate| { candidate.has_mixed_levels() && candidate.within_rewrite_budget(max_output_file_size) }) } @@ -633,8 +695,9 @@ fn pick_mixed_within_budget( fn pick_unbalanced_count_first( sorted_runs: Vec>, max_output_file_size: Option, + mapper: &PrimaryKeyRangeMapper, ) -> Vec { - pick_count_first_where(sorted_runs, max_output_file_size, |_| true) + pick_count_first_where(sorted_runs, max_output_file_size, mapper, |_| true) } fn is_balanced_candidate(candidate: &Candidate) -> bool { @@ -644,6 +707,7 @@ fn is_balanced_candidate(candidate: &Candidate) -> bool { fn pick_count_first_where( sorted_runs: Vec>, max_output_file_size: Option, + mapper: &PrimaryKeyRangeMapper, is_eligible: impl Fn(&Candidate) -> bool, ) -> Vec { let files = ordered_files(&sorted_runs); @@ -654,7 +718,12 @@ fn pick_count_first_where( let right_bound = left.saturating_add(*MAX_INPUT_FILES).min(files.len()); let mut participations: Vec = Vec::with_capacity(right_bound - left); for right in left..right_bound { - candidate.absorb(&files[right], &files[left..right], &mut participations); + candidate.absorb( + &files[right], + &files[left..right], + &mut participations, + mapper, + ); if candidate.num_files < 2 || !is_eligible(&candidate) || !candidate.makes_progress(max_output_file_size) @@ -707,7 +776,11 @@ fn ordered_files(sorted_runs: &[SortedRun]) -> Vec> files } -fn selected_overlaps_unselected(selected: &[FileHandle], window: &Window) -> bool { +fn selected_overlaps_unselected( + selected: &[FileHandle], + window: &Window, + mapper: &PrimaryKeyRangeMapper, +) -> bool { // The overall time span of the selection: a file outside it cannot overlap any // selected file (ranges are inclusive), so it needs no precise overlap check. let Some((span_start, span_end)) = selected @@ -731,7 +804,7 @@ fn selected_overlaps_unselected(selected: &[FileHandle], window: &Window) -> boo .any(|unselected| { selected .iter() - .any(|selected| selected.overlap_inclusive(unselected)) + .any(|selected| files_overlap_inclusive(selected, unselected, mapper)) }) } @@ -798,6 +871,8 @@ impl TwcsPicker { let region_id = compaction_region.region_id; let picker = self.clone(); let compaction_region = compaction_region.clone(); + let primary_key_mapper = compaction_region.current_version.ssts.primary_key_mapper(); + let window_mapper = primary_key_mapper.clone(); let (expired_ssts, time_window_size, active_window, windows) = common_runtime::spawn_blocking_compact(move || { let levels = compaction_region.current_version.ssts.levels(); @@ -835,6 +910,7 @@ impl TwcsPicker { .flat_map(LevelMeta::files) .filter(|file| !expired_file_ids.contains(&file.file_id())), time_window_size, + &window_mapper, ); // Compute activity from the candidate files so expired or // compacting files cannot identify a window absent from `windows`. @@ -865,6 +941,7 @@ impl TwcsPicker { active_window, time_window_size: Some(time_window_size), max_outputs, + primary_key_mapper, }, ) .await?; @@ -898,9 +975,9 @@ struct Window { impl Window { /// Creates a new [Window] with given file. - fn new_with_file(file: FileHandle) -> Self { + fn new_with_file(file: FileHandle, mapper: &PrimaryKeyRangeMapper) -> Self { let (start, end) = file.time_range(); - let primary_key_range = file.primary_key_range(); + let primary_key_range = file.primary_key_range(mapper); Self { start, end, @@ -916,12 +993,14 @@ impl Window { } /// Adds a new file to window and updates time range. - fn add_file(&mut self, file: FileHandle) { + fn add_file(&mut self, file: FileHandle, mapper: &PrimaryKeyRangeMapper) { let (start, end) = file.time_range(); self.start = self.start.min(start); self.end = self.end.max(end); - self.primary_key_range = - merge_primary_key_ranges(self.primary_key_range.take(), file.primary_key_range()); + self.primary_key_range = merge_primary_key_ranges( + self.primary_key_range.take(), + file.primary_key_range(mapper), + ); self.files.push(file); } @@ -934,6 +1013,7 @@ impl Window { fn assign_to_windows<'a>( files: impl Iterator, time_window_size: i64, + mapper: &PrimaryKeyRangeMapper, ) -> BTreeMap { let mut windows: HashMap = HashMap::new(); // Iterates all files and assign to time windows according to max timestamp @@ -951,10 +1031,10 @@ fn assign_to_windows<'a>( match windows.entry(time_window) { Entry::Occupied(mut e) => { - e.get_mut().add_file(f.clone()); + e.get_mut().add_file(f.clone(), mapper); } Entry::Vacant(e) => { - let mut window = Window::new_with_file(f.clone()); + let mut window = Window::new_with_file(f.clone(), mapper); window.time_window = time_window; e.insert(window); } @@ -1064,7 +1144,6 @@ mod tests { use std::sync::Arc; use std::time::Duration; - use bytes::Bytes; use common_base::Plugins; use common_time::range::TimestampRange; use store_api::storage::FileId; @@ -1075,7 +1154,8 @@ mod tests { use crate::compaction::test_util::{ compaction_region_with_ssts, new_file_handle, new_file_handle_with_sequence, new_file_handle_with_size_and_sequence, - new_file_handle_with_size_sequence_and_primary_key_range, + new_file_handle_with_size_sequence_and_primary_key_range, pk_range, + primary_key_mapper_for_test, }; use crate::config::MitoConfig; use crate::region::options::RegionOptions; @@ -1084,6 +1164,35 @@ mod tests { use crate::test_util::memtable_util::metadata_for_test; use crate::test_util::scheduler_util::SchedulerEnv; + fn assign_to_windows<'a>( + files: impl Iterator, + time_window_size: i64, + ) -> BTreeMap { + super::assign_to_windows(files, time_window_size, &primary_key_mapper_for_test()) + } + + fn pick_count_first( + sorted_runs: Vec>, + max_output_file_size: Option, + ) -> Vec { + super::pick_count_first( + sorted_runs, + max_output_file_size, + &primary_key_mapper_for_test(), + ) + } + + fn pick_mixed_count_first( + sorted_runs: Vec>, + max_output_file_size: Option, + ) -> Vec { + super::pick_mixed_count_first( + sorted_runs, + max_output_file_size, + &primary_key_mapper_for_test(), + ) + } + impl TwcsPicker { async fn build_output_with_time_range( &self, @@ -1099,12 +1208,88 @@ mod tests { active_window, time_window_size, max_outputs: self.max_background_tasks, + primary_key_mapper: Arc::new(primary_key_mapper_for_test()), }, ) .await } } + #[test] + fn test_cross_schema_pk_overlap_before_window_aggregation() { + use crate::test_util::sst_util::{new_primary_key, sst_region_metadata}; + + let metadata = Arc::new(sst_region_metadata()); + let mut old = FileMeta { + region_id: metadata.region_id, + file_id: FileId::random(), + time_range: ( + Timestamp::new_millisecond(0), + Timestamp::new_millisecond(1000), + ), + primary_key_min: Some(new_primary_key(&["a"]).into()), + primary_key_max: Some(new_primary_key(&["b"]).into()), + ..Default::default() + }; + let mut completed = new_primary_key(&["b"]); + completed.push(0); // The appended nullable tag's default. + let new = FileMeta { + file_id: FileId::random(), + time_range: ( + Timestamp::new_millisecond(500), + Timestamp::new_millisecond(2000), + ), + primary_key_min: Some(completed.clone().into()), + primary_key_max: Some(completed.clone().into()), + ..old.clone() + }; + assert!(old.primary_key_max < new.primary_key_min); + let mut ssts = SstVersion::new(metadata); + let old_id = old.file_id; + let new_id = new.file_id; + ssts.add_files( + crate::test_util::new_noop_file_purger(), + [old.clone(), new].into_iter(), + ); + let old_file = &ssts.levels()[0].files[&old_id]; + let new_file = &ssts.levels()[0].files[&new_id]; + let ranges = ssts.primary_key_mapper(); + let windows = super::assign_to_windows([old_file, new_file].into_iter(), 1, &ranges); + assert_eq!(2, windows.len()); + assert!( + windows + .values() + .all(|window| window_has_overlap(window, &windows)) + ); + + let mut window = Window::new_with_file(old_file.clone(), &ranges); + window.add_file(new_file.clone(), &ranges); + assert!(selected_overlaps_unselected( + std::slice::from_ref(new_file), + &window, + &ranges, + )); + assert_eq!( + Some(completed.into()), + window + .primary_key_range + .as_ref() + .map(|range| range.1.clone()) + ); + + // In the same target schema, a genuinely disjoint range still prunes. + old.file_id = FileId::random(); + old.primary_key_max = old.primary_key_min.clone(); + let disjoint_id = old.file_id; + ssts.add_files(crate::test_util::new_noop_file_purger(), [old].into_iter()); + let files = &ssts.levels()[0].files; + assert!(!files_overlap_inclusive( + &files[&disjoint_id], + &files[&new_id], + &ranges + )); + } + #[test] fn test_valid_max_input_files_env_overrides_default() { assert_eq!(64, parse_max_input_files(Some("64"))); @@ -1252,10 +1437,11 @@ mod tests { let env = SchedulerEnv::new().await; let metadata = metadata_for_test(); let manifest_ctx = env.mock_manifest_context(metadata.clone()).await; - let mut ssts = SstVersion::new(); + let mut ssts = SstVersion::new(metadata.clone()); ssts.add_files( Arc::new(crate::sst::file_purger::NoopFilePurger), [100, 101].into_iter().map(|sequence| FileMeta { + region_id: metadata.region_id, file_id: FileId::random(), time_range: ( Timestamp::new_millisecond(0), @@ -1530,10 +1716,6 @@ mod tests { /// (Window value, overlapping, files' time ranges in window) type ExpectedWindowSpec = (i64, bool, Vec<(i64, i64)>); - fn pk_range(min: &'static [u8], max: &'static [u8]) -> Option<(Bytes, Bytes)> { - Some((Bytes::from_static(min), Bytes::from_static(max))) - } - fn check_assign_to_windows_with_overlapping( file_time_ranges: &[(i64, i64)], time_window: i64, @@ -2554,6 +2736,7 @@ mod tests { l0_file_num: 8, l1_file_num: 2, phase: PickPhase::L1FileReduction, + primary_key_mapper: &primary_key_mapper_for_test(), }, None, ); @@ -2584,6 +2767,7 @@ mod tests { l0_file_num: 2, l1_file_num: 8, phase: PickPhase::L1FileReduction, + primary_key_mapper: &primary_key_mapper_for_test(), }, None, ); diff --git a/src/mito2/src/compaction/window.rs b/src/mito2/src/compaction/window.rs index b719543f9af..ab7ac0a5f9a 100644 --- a/src/mito2/src/compaction/window.rs +++ b/src/mito2/src/compaction/window.rs @@ -332,7 +332,7 @@ mod tests { let metadata = metadata_for_test(); let file_purger_ref = Arc::new(NoopFilePurger); - let mut ssts = SstVersion::new(); + let mut ssts = SstVersion::new(metadata.clone()); ssts.add_files( file_purger_ref, diff --git a/src/mito2/src/engine/compaction_test.rs b/src/mito2/src/engine/compaction_test.rs index 547d3f213f9..092abc296a6 100644 --- a/src/mito2/src/engine/compaction_test.rs +++ b/src/mito2/src/engine/compaction_test.rs @@ -208,6 +208,7 @@ async fn test_strict_window_output_file_size( let version = engine.get_region(region_id).unwrap().version(); assert!(version.ssts.levels()[0].files.is_empty()); + let primary_key_mapper = version.ssts.primary_key_mapper(); let files = version.ssts.levels()[1].files().collect::>(); assert_eq!( 6000, @@ -231,7 +232,7 @@ async fn test_strict_window_output_file_size( } let mut ranges = window_files .iter() - .map(|file| file.primary_key_range().unwrap()) + .map(|file| file.primary_key_range(&primary_key_mapper).unwrap()) .collect::>(); ranges.sort_unstable(); assert!(ranges.windows(2).all(|pair| pair[0].1 < pair[1].0)); @@ -541,6 +542,313 @@ async fn assert_partial_compaction_preserves_delete_order(flat_format: bool) { ); } +/// Real-engine regression for the tombstone resurrection chain described in +/// https://greptime.feishu.cn/wiki/L2WYwIM40iG9tckowrnctyw8nQc. +/// No synthetic file sizes, PK ranges, or manually selected compaction inputs. +#[rstest::rstest] +#[tokio::test] +async fn test_cross_schema_compaction_keeps_rows_deleted( + #[values(false, true)] flat_format: bool, + #[values(false, true)] cross_window: bool, + #[values(None, Some(""))] tag_default: Option<&str>, +) { + use api::v1::SemanticType; + use api::v1::value::ValueData; + use datatypes::prelude::{ConcreteDataType, Value}; + use datatypes::schema::{ColumnDefaultConstraint, ColumnSchema as DataColumnSchema}; + use store_api::metadata::ColumnMetadata; + use store_api::region_request::{AddColumn, AddColumnLocation, AlterKind}; + + let mut env = TestEnv::new().await; + let region_id = RegionId::new(1, 1); + let (engine, mut columns) = env_for_manual_compaction_with_window( + &mut env, + region_id, + flat_format, + if cross_window { "2h" } else { "1h" }, + ) + .await; + let (target_ts, other_ts) = if cross_window { (3500, 3700) } else { (10, 1) }; + + // Incompressible, distinct a* keys keep the old SST large enough for a + // partial pick to leave it behind. Its maximum PK is exactly the victim b. + let mut old_rows = build_rows_for_key("b", target_ts, target_ts + 1, 0); + let mut rng = ::seed_from_u64(0); + for _ in 0..3000 { + old_rows.extend(build_rows_for_key( + &format!("a{:032x}", rand::Rng::random::(&mut rng)), + other_ts, + other_ts + 1, + 0, + )); + } + put_rows( + &engine, + region_id, + Rows { + schema: columns.clone(), + rows: old_rows, + }, + ) + .await; + flush(&engine, region_id).await; + let old_version = engine.get_region(region_id).unwrap().version(); + let old_files = old_version + .ssts + .levels() + .iter() + .flat_map(|level| level.files()) + .collect::>(); + assert_eq!(1, old_files.len()); + let old_file = old_files[0]; + if cross_window { + // This real SST crosses the new one-hour boundary. TWCS assigns it by + // max_ts=3700 to window 7200, unlike Delete(b,3500) in window 3600. + set_compaction_window(&engine, region_id, "1h").await; + } + + let added = ColumnMetadata { + column_id: 3, + semantic_type: SemanticType::Tag, + column_schema: DataColumnSchema::new("tag_1", ConcreteDataType::string_datatype(), true) + .with_default_constraint( + tag_default.map(|value| ColumnDefaultConstraint::Value(Value::from(value))), + ) + .unwrap(), + }; + engine + .handle_request( + region_id, + RegionRequest::Alter(RegionAlterRequest { + kind: AlterKind::AddColumns { + columns: vec![AddColumn { + column_metadata: added.clone(), + location: Some(AddColumnLocation::First), + }], + }, + }), + ) + .await + .unwrap(); + columns.push(column_metadata_to_column_schema(&added)); + let append_tag = |mut rows: Vec| { + for row in &mut rows { + row.values.push(api::v1::Value { + value_data: tag_default.map(|value| ValueData::StringValue(value.to_string())), + }); + } + Rows { + schema: columns.clone(), + rows, + } + }; + let deletes = ["b", "c"] + .into_iter() + .flat_map(|key| build_rows_for_key(key, target_ts, target_ts + 1, 0)) + .collect(); + engine + .handle_request( + region_id, + RegionRequest::Delete(RegionDeleteRequest { + rows: append_tag(deletes), + hint: None, + partition_expr_version: None, + }), + ) + .await + .unwrap(); + flush(&engine, region_id).await; + for ts in target_ts + 10..target_ts + 25 { + put_rows( + &engine, + region_id, + append_tag(build_rows_for_key("c", ts, ts + 1, 0)), + ) + .await; + flush(&engine, region_id).await; + } + + let region = engine.get_region(region_id).unwrap(); + let table_dir = region.table_dir().to_string(); + let version = region.version(); + let files = version + .ssts + .levels() + .iter() + .flat_map(|level| level.files()) + .collect::>(); + assert_eq!(17, files.len()); + assert_eq!(vec![0, 3], version.metadata.primary_key); + let new_ids = files + .iter() + .filter(|file| file.file_id() != old_file.file_id()) + .map(|file| file.file_id()) + .collect::>(); + for file in files + .iter() + .filter(|file| new_ids.contains(&file.file_id())) + { + assert!(old_file.meta_ref().primary_key_max < file.meta_ref().primary_key_min); + assert!(file.time_range().1 < Timestamp::new_second(3600)); + } + assert_eq!( + cross_window, + old_file.time_range().1 > Timestamp::new_second(3600) + ); + + let picked = preview_regular_compaction(&engine, region_id).await; + assert_eq!(1, picked.outputs.len()); + let output = &picked.outputs[0]; + assert_eq!( + new_ids, + output.inputs.iter().map(|file| file.file_id()).collect() + ); + let filter_deleted = output.filter_deleted; + let before = sorted_scan_timestamps(&engine, region_id).await; + assert_eq!(3015, before.len()); + assert!(!before.contains(&(target_ts as i64 * 1000))); + + compact(&engine, region_id).await; + + let after_version = region.version(); + let after_files = after_version + .ssts + .levels() + .iter() + .flat_map(|level| level.files()) + .collect::>(); + assert_eq!(2, after_files.len()); + assert!( + after_files + .iter() + .any(|file| file.file_id() == old_file.file_id()) + ); + let output_file = after_files + .into_iter() + .find(|file| file.file_id() != old_file.file_id()) + .unwrap(); + assert!(!new_ids.contains(&output_file.file_id())); + let output_rows = output_file.meta_ref().num_rows; + let output_deletes = count_sst_deletes(&engine, region_id, output_file).await; + let after = sorted_scan_timestamps(&engine, region_id).await; + + let config = MitoConfig { + default_flat_format: flat_format, + min_compaction_interval: Duration::from_secs(3600), + ..Default::default() + }; + let engine = env.reopen_engine(engine, config).await; + crate::test_util::reopen_region(&engine, region_id, table_dir, false, HashMap::new()).await; + let reopened = sorted_scan_timestamps(&engine, region_id).await; + // Do not fail on filter_deleted before actually executing compaction/reopen: + // the negative-control run must demonstrate data resurrection, not just a bad plan. + assert!( + !filter_deleted + && output_rows == 17 + && output_deletes == 2 + && after == before + && reopened == before, + "cross-schema tombstone loss: flat_format={flat_format}, cross_window={cross_window}, default={tag_default:?}, selected={}, filter_deleted={filter_deleted}, output_rows={output_rows}, output_deletes={output_deletes}, before={}, after={}, reopened={}, resurrected_after={}, resurrected_reopen={}", + new_ids.len(), + before.len(), + after.len(), + reopened.len(), + after.contains(&(target_ts as i64 * 1000)), + reopened.contains(&(target_ts as i64 * 1000)), + ); +} + +async fn count_sst_deletes( + engine: &MitoEngine, + region_id: RegionId, + file: &crate::sst::file::FileHandle, +) -> usize { + use api::v1::OpType; + use datatypes::arrow::datatypes::UInt8Type; + + let region = engine.get_region(region_id).unwrap(); + let mut reader = region + .access_layer + .read_sst(file.clone()) + .build() + .await + .unwrap() + .unwrap(); + let mut deletes = 0; + while let Some(batch) = reader.next_record_batch().await.unwrap() { + deletes += batch + .column(batch.num_columns() - 1) + .as_primitive::() + .values() + .iter() + .filter(|op| **op == OpType::Delete as u8) + .count(); + } + deletes +} + +async fn set_compaction_window(engine: &MitoEngine, region_id: RegionId, window: &str) { + engine + .handle_request( + region_id, + RegionRequest::Alter(RegionAlterRequest { + kind: SetRegionOptions { + options: vec![SetRegionOption::Twsc( + "compaction.twcs.time_window".into(), + window.into(), + )], + }, + }), + ) + .await + .unwrap(); +} + +async fn sorted_scan_timestamps(engine: &MitoEngine, region_id: RegionId) -> Vec { + let stream = engine + .scanner(region_id, ScanRequest::default()) + .await + .unwrap() + .scan() + .await + .unwrap(); + let mut timestamps = collect_stream_ts(stream).await; + timestamps.sort_unstable(); + timestamps +} + +async fn preview_regular_compaction( + engine: &MitoEngine, + region_id: RegionId, +) -> crate::compaction::picker::PickerOutput { + use crate::compaction::compactor::CompactionRegion; + use crate::compaction::picker::new_picker; + + let region = engine.get_region(region_id).unwrap(); + let version = region.version(); + let picker = new_picker( + &RegionCompactRequest::default().options, + &version.options, + None, + None, + ); + let compaction_region = CompactionRegion { + region_id, + region_options: version.options.clone(), + engine_config: Arc::new(MitoConfig::default()), + region_metadata: version.metadata.clone(), + cache_manager: engine.cache_manager(), + access_layer: region.access_layer.clone(), + manifest_ctx: region.manifest_ctx.clone(), + current_version: version.into(), + file_purger: None, + ttl: None, + max_parallelism: 1, + plugins: common_base::Plugins::new(), + }; + picker.pick(&compaction_region).await.unwrap().unwrap() +} + struct CompactionListenerGuard(Option>); impl CompactionListenerGuard { @@ -1790,6 +2098,15 @@ async fn env_for_manual_compaction( env: &mut TestEnv, region_id: RegionId, flat_format: bool, +) -> (MitoEngine, Vec) { + env_for_manual_compaction_with_window(env, region_id, flat_format, "1h").await +} + +async fn env_for_manual_compaction_with_window( + env: &mut TestEnv, + region_id: RegionId, + flat_format: bool, + time_window: &str, ) -> (MitoEngine, Vec) { let engine = env .create_engine(MitoConfig { @@ -1812,7 +2129,7 @@ async fn env_for_manual_compaction( let request = CreateRequestBuilder::new() .insert_option("compaction.type", "twcs") - .insert_option("compaction.twcs.time_window", "1h") + .insert_option("compaction.twcs.time_window", time_window) .build(); let column_schemas = request .column_metadatas diff --git a/src/mito2/src/error.rs b/src/mito2/src/error.rs index 21416bf3652..1592bfebd60 100644 --- a/src/mito2/src/error.rs +++ b/src/mito2/src/error.rs @@ -1138,6 +1138,21 @@ pub enum Error { location: Location, }, + #[snafu(display("Invalid SST primary key range: {reason}"))] + InvalidPrimaryKeyRange { + reason: String, + #[snafu(implicit)] + location: Location, + }, + + #[snafu(display("Failed to decode SST primary key range {endpoint} endpoint"))] + DecodePrimaryKeyRange { + endpoint: &'static str, + source: mito_codec::error::Error, + #[snafu(implicit)] + location: Location, + }, + #[snafu(display("Region {} is busy", region_id))] RegionBusy { region_id: RegionId, @@ -1586,7 +1601,10 @@ impl ErrorExt for Error { FulltextPushText { source, .. } | FulltextFinish { source, .. } | ApplyFulltextIndex { source, .. } => source.status_code(), - DecodeStats { .. } | StatsNotPresent { .. } => StatusCode::Internal, + DecodeStats { .. } + | StatsNotPresent { .. } + | InvalidPrimaryKeyRange { .. } + | DecodePrimaryKeyRange { .. } => StatusCode::Internal, RegionBusy { .. } => StatusCode::RegionBusy, GetSchemaMetadata { source, .. } => source.status_code(), Timeout { .. } => StatusCode::Cancelled, diff --git a/src/mito2/src/read/scan_region.rs b/src/mito2/src/read/scan_region.rs index 7c93d137c04..1d0991a6bfb 100644 --- a/src/mito2/src/read/scan_region.rs +++ b/src/mito2/src/read/scan_region.rs @@ -17,7 +17,7 @@ use std::collections::{BTreeMap, HashSet}; use std::fmt; use std::num::NonZeroU64; -use std::sync::Arc; +use std::sync::{Arc, OnceLock}; use std::time::Instant; use api::v1::SemanticType; @@ -90,6 +90,7 @@ use crate::sst::index::vector_index::applier::{VectorIndexApplier, VectorIndexAp use crate::sst::parquet::Json2RewriteTargets; use crate::sst::parquet::file_range::PreFilterMode; use crate::sst::parquet::reader::ReaderMetrics; +use crate::sst::primary_key::PrimaryKeyRangeMapper; #[cfg(feature = "vector_index")] const VECTOR_INDEX_OVERFETCH_MULTIPLIER: usize = 2; @@ -589,6 +590,7 @@ impl ScanRegion { .with_predicate(predicate) .with_memtables(mem_range_builders) .with_files(files) + .with_primary_key_mapper(self.version.ssts.primary_key_mapper()) .with_cache(self.cache_strategy) .with_inverted_index_appliers(inverted_index_appliers) .with_bloom_filter_index_appliers(bloom_filter_appliers) @@ -977,6 +979,8 @@ pub struct ScanInput { pub(crate) memtables: Vec, /// Handles to SST files to scan. pub(crate) files: Vec, + /// Shares the pinned schema's encoded defaults across parallel range readers. + primary_key_mapper: OnceLock>, /// Scan-wide hint for rows in an execution batch. batch_size: usize, /// Cache. @@ -1066,6 +1070,7 @@ impl ScanInput { region_partition_expr: None, memtables: Vec::new(), files: Vec::new(), + primary_key_mapper: OnceLock::new(), batch_size: crate::sst::parquet::DEFAULT_READ_BATCH_SIZE, cache_strategy: CacheStrategy::Disabled, ignore_file_not_found: false, @@ -1101,6 +1106,12 @@ impl ScanInput { self.batch_size } + /// Interprets file statistics using this scan's pinned schema. + pub(crate) fn primary_key_mapper(&self) -> &PrimaryKeyRangeMapper { + self.primary_key_mapper + .get_or_init(|| Arc::new(PrimaryKeyRangeMapper::new(self.region_metadata().clone()))) + } + /// Returns the range implied by the range-cache time filters. pub(crate) fn implied_time_range(&self) -> Option<&TimestampRange> { self.scan_analysis @@ -1121,6 +1132,11 @@ impl ScanInput { } impl ScanInputBuilder { + fn with_primary_key_mapper(mut self, mapper: Arc) -> Self { + self.input.primary_key_mapper = OnceLock::from(mapper); + self + } + /// Sets whether to ignore range indexes during scans. #[must_use] pub(crate) fn with_ignore_range_index(mut self, ignore: bool) -> Self { diff --git a/src/mito2/src/read/series_reader.rs b/src/mito2/src/read/series_reader.rs index 9e86ff1e5e9..5ac69d191a6 100644 --- a/src/mito2/src/read/series_reader.rs +++ b/src/mito2/src/read/series_reader.rs @@ -408,7 +408,7 @@ async fn build_series_partition_range( if stream_ctx.is_file_range_index(index) { let file = stream_ctx.input.file_from_index(index); if matches!( - file.primary_key_range() + file.primary_key_range(stream_ctx.input.primary_key_mapper()) .and_then(|(min, max)| { filter.overlaps_encoded_bounds(&codec, &min, &max) }), Some(false) ) { diff --git a/src/mito2/src/region/opener.rs b/src/mito2/src/region/opener.rs index 033f9916905..fb5d260fe6f 100644 --- a/src/mito2/src/region/opener.rs +++ b/src/mito2/src/region/opener.rs @@ -1232,7 +1232,7 @@ async fn preload_parquet_meta_cache_for_files( .get_compact_sst_meta_data(file_id, PageIndexPolicy::Optional) .await { - if file_handle.primary_key_range().is_none() + if file_handle.raw_primary_key_range().is_none() && let Some(primary_key_range) = extract_primary_key_range( metadata.parquet_metadata().as_ref(), ®ion_metadata, @@ -1251,7 +1251,7 @@ async fn preload_parquet_meta_cache_for_files( .await { let decoded = metadata.decoded(); - if file_handle.primary_key_range().is_none() + if file_handle.raw_primary_key_range().is_none() && let Some(primary_key_range) = extract_primary_key_range( decoded.parquet_metadata().as_ref(), ®ion_metadata, @@ -1664,7 +1664,9 @@ mod tests { let file_id = FileId::random(); let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef; - let primary_key = Arc::new(BinaryArray::from_iter_values([b"a", b"b", b"c"])) as ArrayRef; + let keys = + ["a", "b", "c"].map(|tag| crate::test_util::sst_util::new_primary_key(&[tag, ""])); + let primary_key = Arc::new(BinaryArray::from_iter_values(&keys)) as ArrayRef; let batch = RecordBatch::try_from_iter([ ("col", col), ( @@ -1744,7 +1746,7 @@ mod tests { .await .is_some() ); - assert!(file_handle.primary_key_range().is_some()); + assert!(file_handle.raw_primary_key_range().is_some()); } #[tokio::test] diff --git a/src/mito2/src/region/version.rs b/src/mito2/src/region/version.rs index e9ba4a1cdbf..a56cef1f544 100644 --- a/src/mito2/src/region/version.rs +++ b/src/mito2/src/region/version.rs @@ -404,9 +404,9 @@ impl VersionBuilder { /// Returns a new builder. pub(crate) fn new(metadata: RegionMetadataRef, mutable: TimePartitionsRef) -> Self { VersionBuilder { - metadata, + metadata: metadata.clone(), memtables: Arc::new(MemtableVersion::new(mutable)), - ssts: Arc::new(SstVersion::new()), + ssts: Arc::new(SstVersion::new(metadata)), flushed_entry_id: 0, flushed_sequence: 0, truncated_entry_id: None, @@ -437,6 +437,9 @@ impl VersionBuilder { /// Sets metadata. pub(crate) fn metadata(mut self, metadata: RegionMetadataRef) -> Self { + if !Arc::ptr_eq(&self.metadata, &metadata) { + Arc::make_mut(&mut self.ssts).set_metadata(metadata.clone()); + } self.metadata = metadata; self } @@ -525,7 +528,7 @@ impl VersionBuilder { /// Clear all files in the builder. pub(crate) fn clear_files(mut self) -> Self { - self.ssts = Arc::new(SstVersion::new()); + self.ssts = Arc::new(SstVersion::new(self.metadata.clone())); self } diff --git a/src/mito2/src/sst.rs b/src/mito2/src/sst.rs index e53c28e54c7..7a0e3603836 100644 --- a/src/mito2/src/sst.rs +++ b/src/mito2/src/sst.rs @@ -47,6 +47,7 @@ pub mod file_ref; pub mod index; pub mod location; pub mod parquet; +pub(crate) mod primary_key; pub mod range_index; pub(crate) mod version; diff --git a/src/mito2/src/sst/file.rs b/src/mito2/src/sst/file.rs index 9702477b07a..f160bbb3785 100644 --- a/src/mito2/src/sst/file.rs +++ b/src/mito2/src/sst/file.rs @@ -24,7 +24,7 @@ use std::sync::{Arc, Mutex, RwLock}; use base64::prelude::{BASE64_STANDARD, Engine}; use bytes::Bytes; use common_base::readable_size::ReadableSize; -use common_telemetry::{debug, error}; +use common_telemetry::{debug, error, warn}; use common_time::Timestamp; use partition::expr::PartitionExpr; use serde::{Deserialize, Serialize}; @@ -39,6 +39,7 @@ use crate::cache::file_cache::{FileType, IndexKey}; use crate::sst::file_purger::FilePurgerRef; use crate::sst::location; use crate::sst::parquet::SstInfo; +use crate::sst::primary_key::PrimaryKeyRangeMapper; /// Custom serde functions for Bytes fields serialized as base64 strings. fn serialize_bytes_option(bytes: &Option, serializer: S) -> Result @@ -497,6 +498,7 @@ impl fmt::Debug for FileHandle { } impl FileHandle { + /// Creates a handle sharing the file's physical state and original statistics. pub fn new(meta: FileMeta, file_purger: FilePurgerRef) -> FileHandle { let pk_range = meta.primary_key_range(); FileHandle { @@ -612,12 +614,100 @@ impl FileHandle { self.inner.deleted.load(Ordering::Relaxed) } - pub fn primary_key_range(&self) -> Option<(Bytes, Bytes)> { - self.inner.primary_key_range.read().unwrap().clone() + /// Returns bounds aligned to the caller's pinned schema, before any comparison + /// or aggregation. Unknown and invalid statistics cannot exclude possible data. + pub(crate) fn primary_key_range( + &self, + mapper: &PrimaryKeyRangeMapper, + ) -> Option<(Bytes, Bytes)> { + debug_assert_eq!(self.region_id().table_id(), mapper.region_id().table_id()); + if let Some(range) = self + .inner + .primary_key_range + .read() + .unwrap() + .aligned(mapper.schema_version()) + { + return range.clone(); + } + // Recheck under the write lock: another snapshot may have replaced the cached schema. + let aligned = self.inner.primary_key_range.write().unwrap().align(mapper); + match aligned { + Ok(range) => range, + Err(err) => { + warn!(err; "Invalid SST primary key range; using unknown bounds, region: {}, file: {}, schema version: {}", + self.region_id(), self.file_id(), mapper.schema_version()); + None + } + } + } + + /// Returns original statistics for metadata hydration, never schema-aligned bounds. + pub fn raw_primary_key_range(&self) -> Option<(Bytes, Bytes)> { + self.inner.primary_key_range.read().unwrap().raw().cloned() } pub(crate) fn set_primary_key_range(&self, primary_key_range: (Bytes, Bytes)) { - *self.inner.primary_key_range.write().unwrap() = Some(primary_key_range); + // SST contents are immutable. Hydrate missing raw statistics without + // replacing the source of already cached schema views. + let mut range = self.inner.primary_key_range.write().unwrap(); + if matches!(*range, PrimaryKeyRange::Missing) { + *range = PrimaryKeyRange::Raw(primary_key_range); + } + } +} + +type PrimaryKeyBounds = (Bytes, Bytes); + +/// A single-slot cache shared by file handles. Always retain the source bounds so +/// older snapshots and changed defaults can realign without interpreting padded values as stored. +enum PrimaryKeyRange { + Missing, + Raw(PrimaryKeyBounds), + Aligned { + raw: PrimaryKeyBounds, + schema_version: u64, + bounds: Option, + }, +} + +impl PrimaryKeyRange { + fn raw(&self) -> Option<&PrimaryKeyBounds> { + match self { + Self::Missing => None, + Self::Raw(raw) | Self::Aligned { raw, .. } => Some(raw), + } + } + + fn aligned(&self, target_version: u64) -> Option<&Option> { + match self { + Self::Aligned { + schema_version, + bounds, + .. + } if *schema_version == target_version => Some(bounds), + _ => None, + } + } + + fn align( + &mut self, + mapper: &PrimaryKeyRangeMapper, + ) -> crate::error::Result> { + if let Some(bounds) = self.aligned(mapper.schema_version()) { + return Ok(bounds.clone()); + } + let Some(raw) = self.raw().cloned() else { + return Ok(None); + }; + let aligned = mapper.map(raw.clone()); + // Failed mappings also occupy the cache, so repeated hits don't repeat the warning. + *self = Self::Aligned { + raw, + schema_version: mapper.schema_version(), + bounds: aligned.as_ref().ok().cloned().flatten(), + }; + aligned } } @@ -629,7 +719,7 @@ struct FileHandleInner { compacting: AtomicBool, deleted: AtomicBool, index_outdated: AtomicBool, - primary_key_range: RwLock>, + primary_key_range: RwLock, file_purger: FilePurgerRef, } @@ -656,7 +746,9 @@ impl FileHandleInner { compacting: AtomicBool::new(false), deleted: AtomicBool::new(false), index_outdated: AtomicBool::new(false), - primary_key_range: RwLock::new(primary_key_range), + primary_key_range: RwLock::new( + primary_key_range.map_or(PrimaryKeyRange::Missing, PrimaryKeyRange::Raw), + ), file_purger, } } diff --git a/src/mito2/src/sst/file_purger.rs b/src/mito2/src/sst/file_purger.rs index 5ff88c0874b..1b3b6239897 100644 --- a/src/mito2/src/sst/file_purger.rs +++ b/src/mito2/src/sst/file_purger.rs @@ -71,15 +71,15 @@ impl fmt::Debug for LocalFilePurger { } } -#[cfg(not(debug_assertions))] +#[cfg(not(any(debug_assertions, test)))] /// Whether to enable GC for the file purger. pub fn should_enable_gc(global_gc_enabled: bool, object_store_scheme: &'static str) -> bool { global_gc_enabled && object_store_scheme != object_store::services::FS_SCHEME } -#[cfg(debug_assertions)] -/// For debug build, we may use Fs as the object store scheme, -/// so we need to enable GC for local file system. +#[cfg(any(debug_assertions, test))] +/// Debug builds and unit tests use Fs to exercise object-store GC, including +/// unit tests compiled with the release profile. pub fn should_enable_gc(global_gc_enabled: bool, _object_store_scheme: &'static str) -> bool { global_gc_enabled } diff --git a/src/mito2/src/sst/parquet/metadata.rs b/src/mito2/src/sst/parquet/metadata.rs index 3c0d197c5c6..8e6a0a5998c 100644 --- a/src/mito2/src/sst/parquet/metadata.rs +++ b/src/mito2/src/sst/parquet/metadata.rs @@ -144,6 +144,11 @@ pub(crate) fn extract_primary_key_range( return None; }; + // Truncated bounds can end at a field boundary and masquerade as an + // older Dense schema. Only complete endpoints can be schema-normalized. + if !stats.min_is_exact() || !stats.max_is_exact() { + return None; + } let row_group_min = Bytes::copy_from_slice(stats.min_bytes_opt()?); let row_group_max = Bytes::copy_from_slice(stats.max_bytes_opt()?); min = Some(match min { @@ -321,6 +326,16 @@ mod tests { ); } + #[test] + fn test_extract_primary_key_range_rejects_truncated_statistics() { + let key = vec![b'a'; 1024]; + let metadata = build_test_metadata(true, &[&key], &[1], EnabledStatistics::Page); + assert_eq!( + None, + extract_primary_key_range(&metadata, &sst_region_metadata()) + ); + } + #[test] fn test_extract_primary_key_range_returns_none_when_any_rg_stats_missing() { let metadata = build_test_metadata( diff --git a/src/mito2/src/sst/primary_key.rs b/src/mito2/src/sst/primary_key.rs new file mode 100644 index 00000000000..5fb258052c7 --- /dev/null +++ b/src/mito2/src/sst/primary_key.rs @@ -0,0 +1,847 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Schema-bound views of persisted primary key ranges. + +use std::sync::OnceLock; + +use bytes::Bytes; +use datatypes::prelude::ConcreteDataType; +use mito_codec::row_converter::{DensePrimaryKeyCodec, PrimaryKeyCodec, SortField}; +use snafu::{ResultExt, ensure}; +use store_api::codec::PrimaryKeyEncoding; +use store_api::metadata::{ColumnMetadata, RegionMetadataRef}; +use store_api::storage::RegionId; + +use crate::error::{DecodePrimaryKeyRangeSnafu, InvalidPrimaryKeyRangeSnafu, Result}; + +/// Schema conversion and default encodings shared by a version's comparison contexts. +#[derive(Debug)] +pub(crate) struct PrimaryKeyRangeMapper { + /// Target metadata pinned by the owning version or scan, not the SST's write schema. + /// ALTER creates a new mapper; existing snapshots keep their original target schema. + metadata: RegionMetadataRef, + codec: DensePrimaryKeyCodec, + /// Encoding a constant is independent of the file's missing-prefix length. + encoded_defaults: Vec>>, + suffixes: Vec>>, +} + +impl PrimaryKeyRangeMapper { + pub(crate) fn new(metadata: RegionMetadataRef) -> Self { + let codec = DensePrimaryKeyCodec::new(&metadata); + let encoded_defaults = (0..codec.num_fields()).map(|_| OnceLock::new()).collect(); + let suffixes = (0..codec.num_fields()).map(|_| OnceLock::new()).collect(); + Self { + metadata, + codec, + encoded_defaults, + suffixes, + } + } + + /// Reuses encoded constants without retaining a chain of old schema snapshots. + pub(crate) fn with_metadata(&self, metadata: RegionMetadataRef) -> Self { + let mut mapper = Self::new(metadata); + for (index, (old, new)) in self + .metadata + .primary_key_columns() + .zip(mapper.metadata.primary_key_columns()) + .enumerate() + { + if same_pk_column(old, new) { + mapper.encoded_defaults[index] = self.encoded_defaults[index].clone(); + } + } + mapper + } + + /// Returns the schema version used to interpret the bounds. + pub(crate) fn schema_version(&self) -> u64 { + self.metadata.schema_version + } + + /// Returns the region id of mapper. + pub(crate) fn region_id(&self) -> RegionId { + self.metadata.region_id + } + + /// Maps bounds from the same table's encoding and a compatible schema prefix. + /// Returns unknown for defaults that cannot be safely encoded. + /// Invalid bounds return an error so callers can diagnose them before falling back + /// to unknown. Neither the error nor its source includes the encoded keys. + pub(crate) fn map(&self, (min, max): (Bytes, Bytes)) -> Result> { + ensure!( + min <= max, + InvalidPrimaryKeyRangeSnafu { + reason: "min is greater than max", + } + ); + // Sparse-to-sparse read compatibility preserves encoded keys verbatim. + if self.metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse { + return Ok(Some((min, max))); + } + // Ordinary same-table Dense ALTER only appends PK fields and preserves + // prefix order/types. SyncColumns violating that invariant is out of scope. + let prefix_len = self.range_prefix_len(&min, &max)?; + if prefix_len == self.codec.num_fields() { + return Ok(Some((min, max))); + } + let Some(suffix) = self.suffixes[prefix_len] + .get_or_init(|| self.encode_suffix(prefix_len)) + .as_ref() + else { + return Ok(None); + }; + // Fixed-layout Dense keys are prefix-free, so appending one constant + // suffix is monotone. Transform each file before aggregating any ranges. + Ok(Some(( + append_suffix(min, suffix), + append_suffix(max, suffix), + ))) + } + + fn range_prefix_len(&self, min: &[u8], max: &[u8]) -> Result { + let min_len = self + .codec + .decode_prefix_len(min) + .context(DecodePrimaryKeyRangeSnafu { endpoint: "min" })?; + let max_len = self + .codec + .decode_prefix_len(max) + .context(DecodePrimaryKeyRangeSnafu { endpoint: "max" })?; + ensure!( + min_len == max_len, + InvalidPrimaryKeyRangeSnafu { + reason: format!( + "endpoints have different field counts: min {min_len}, max {max_len}" + ), + } + ); + Ok(min_len) + } + + fn encode_suffix(&self, prefix_len: usize) -> Option { + let mut suffix = Vec::with_capacity(self.codec.num_fields() - prefix_len); + for (index, column) in self + .metadata + .primary_key_columns() + .enumerate() + .skip(prefix_len) + { + let encoded = self.encoded_defaults[index] + .get_or_init(|| encode_default(column)) + .as_ref()?; + suffix.extend_from_slice(encoded); + } + Some(suffix.into()) + } +} + +fn same_pk_column(left: &ColumnMetadata, right: &ColumnMetadata) -> bool { + left.column_id == right.column_id + && left.column_schema.data_type == right.column_schema.data_type + && left.column_schema.is_nullable() == right.column_schema.is_nullable() + && left.column_schema.default_constraint() == right.column_schema.default_constraint() +} + +/// Missing, impure or unencodable defaults cannot produce exact bounds. +fn encode_default(column: &ColumnMetadata) -> Option { + let schema = &column.column_schema; + if schema.is_default_impure() { + return None; + } + let default = schema.create_default().ok()??; + let field = SortField::new(schema.data_type.clone()); + if default.is_null() { + // Dense uses one NULL marker per field, but unsupported types must stay unknown. + return match field.encode_data_type() { + ConcreteDataType::Null(_) + | ConcreteDataType::List(_) + | ConcreteDataType::Struct(_) + | ConcreteDataType::Dictionary(_) => None, + _ => Some(Bytes::from_static(&[0])), + }; + } + let mut encoded = Vec::new(); + DensePrimaryKeyCodec::with_fields(vec![(column.column_id, field)]) + .encode_values(&[(column.column_id, default)], &mut encoded) + .ok()?; + Some(encoded.into()) +} + +fn append_suffix(key: Bytes, suffix: &Bytes) -> Bytes { + let mut completed = Vec::with_capacity(key.len() + suffix.len()); + completed.extend_from_slice(&key); + completed.extend_from_slice(suffix); + completed.into() +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use api::v1::SemanticType; + use common_time::Timestamp; + use datatypes::prelude::{ConcreteDataType, Value}; + use datatypes::schema::{ColumnDefaultConstraint, ColumnSchema}; + use rstest::rstest; + use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder}; + use store_api::storage::FileId; + + use super::*; + use crate::compaction::run::files_overlap_inclusive; + use crate::manifest::action::RegionEdit; + use crate::memtable::time_partition::TimePartitions; + use crate::memtable::time_series::TimeSeriesMemtableBuilder; + use crate::region::version::{VersionBuilder, VersionRef}; + use crate::sst::file::{FileHandle, FileMeta}; + use crate::test_util::new_noop_file_purger; + + fn metadata(defaults: &[Value]) -> RegionMetadataRef { + let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1)); + for (id, value) in defaults.iter().enumerate() { + let data_type = if value.is_null() { + ConcreteDataType::string_datatype() + } else { + value.data_type() + }; + builder.push_column_metadata(ColumnMetadata { + column_id: id as u32, + semantic_type: SemanticType::Tag, + column_schema: ColumnSchema::new(format!("tag_{id}"), data_type, true) + .with_default_constraint(Some(ColumnDefaultConstraint::Value(value.clone()))) + .unwrap(), + }); + } + builder.push_column_metadata(ColumnMetadata { + column_id: 100, + semantic_type: SemanticType::Timestamp, + column_schema: ColumnSchema::new( + "ts", + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ), + }); + builder.primary_key((0..defaults.len() as u32).collect()); + Arc::new(builder.build().unwrap()) + } + + fn metadata_at_version(defaults: &[Value], version: u64) -> RegionMetadataRef { + let mut metadata = (*metadata(defaults)).clone(); + metadata.schema_version = version; + Arc::new(metadata) + } + + fn encode(metadata: &RegionMetadataRef, values: &[Value]) -> Bytes { + let mut bytes = Vec::new(); + let values: Vec<_> = values + .iter() + .enumerate() + .map(|(id, value)| (id as u32, value.clone())) + .collect(); + DensePrimaryKeyCodec::new(metadata) + .encode_values(&values, &mut bytes) + .unwrap(); + bytes.into() + } + + fn file_meta(metadata: &RegionMetadataRef, values: &[Value]) -> FileMeta { + let key = encode(metadata, values); + FileMeta { + region_id: metadata.region_id, + file_id: FileId::random(), + time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)), + primary_key_min: Some(key.clone()), + primary_key_max: Some(key), + ..Default::default() + } + } + + fn version_with_files(metadata: RegionMetadataRef, files: Vec) -> VersionRef { + let mutable = Arc::new(TimePartitions::new( + metadata.clone(), + Arc::new(TimeSeriesMemtableBuilder::default()), + 0, + None, + )); + Arc::new( + VersionBuilder::new(metadata, mutable) + .add_files(new_noop_file_purger(), files.into_iter()) + .build(), + ) + } + + #[test] + fn test_field_addition_updates_cache_version_and_shares_sst_list() { + let metadata = metadata(&[Value::from(""), Value::from("default")]); + let raw = file_meta(&metadata, &[Value::from("a")]); + let file_id = raw.file_id; + let original = version_with_files(metadata.clone(), vec![raw]); + let mapper = original.ssts.primary_key_mapper(); + let range = original.ssts.levels()[0].files[&file_id] + .primary_key_range(&mapper) + .unwrap(); + let mut builder = RegionMetadataBuilder::from_existing((*metadata).clone()); + builder + .push_column_metadata(ColumnMetadata { + column_id: 101, + semantic_type: SemanticType::Field, + column_schema: ColumnSchema::new("field", ConcreteDataType::int64_datatype(), true), + }) + .bump_version(); + let changed = VersionBuilder::from_version(original.clone()) + .metadata(Arc::new(builder.build().unwrap())) + .build(); + + assert_eq!(1, changed.metadata.schema_version); + assert!(changed.metadata.column_by_id(101).is_some()); + assert_eq!( + original.ssts.levels().as_ptr(), + changed.ssts.levels().as_ptr() + ); + let changed_mapper = changed.ssts.primary_key_mapper(); + assert_eq!(1, changed_mapper.schema_version()); + let changed_range = changed.ssts.levels()[0].files[&file_id] + .primary_key_range(&changed_mapper) + .unwrap(); + assert_eq!(range, changed_range); + let cached = changed.ssts.levels()[0].files[&file_id] + .primary_key_range(&changed_mapper) + .unwrap(); + assert_eq!(cached.0.as_ptr(), changed_range.0.as_ptr()); + assert_eq!(cached.1.as_ptr(), changed_range.1.as_ptr()); + } + + // SET/DROP DEFAULT may change padding even when the number of PK columns stays fixed. + #[rstest] + #[case::set(Value::from("x"), Some(Value::from("y")))] + #[case::set_null(Value::from("x"), Some(Value::Null))] + #[case::drop(Value::from("x"), None)] + #[case::replace_null(Value::Null, Some(Value::from("y")))] + fn test_tag_default_changes_refresh_only_missing_values( + #[case] previous_default: Value, + #[case] new_default: Option, + ) { + let metadata = metadata(&[Value::from(""), previous_default.clone()]); + let missing = file_meta(&metadata, &[Value::from("a")]); + let complete = file_meta(&metadata, &[Value::from("b"), Value::from("stored")]); + let missing_id = missing.file_id; + let complete_id = complete.file_id; + let complete_range = complete.primary_key_range(); + let original = version_with_files(metadata.clone(), vec![missing, complete]); + let old_mapper = original.ssts.primary_key_mapper(); + let old_range = original.ssts.levels()[0].files[&missing_id].primary_key_range(&old_mapper); + + let mut changed_metadata = (*metadata).clone(); + changed_metadata.column_metadatas[1].column_schema = changed_metadata.column_metadatas[1] + .column_schema + .clone() + .with_default_constraint(new_default.clone().map(ColumnDefaultConstraint::Value)) + .unwrap(); + let mut builder = RegionMetadataBuilder::from_existing(changed_metadata); + builder.bump_version(); + let changed_metadata = Arc::new(builder.build().unwrap()); + let changed = VersionBuilder::from_version(original.clone()) + .metadata(changed_metadata.clone()) + .build(); + + let expected = encode( + &changed_metadata, + &[Value::from("a"), new_default.unwrap_or(Value::Null)], + ); + let changed_mapper = changed.ssts.primary_key_mapper(); + assert_eq!( + Some((expected.clone(), expected)), + changed.ssts.levels()[0].files[&missing_id].primary_key_range(&changed_mapper) + ); + assert_eq!( + complete_range, + changed.ssts.levels()[0].files[&complete_id].primary_key_range(&changed_mapper) + ); + let expected_old = encode(&metadata, &[Value::from("a"), previous_default]); + assert_eq!(Some((expected_old.clone(), expected_old)), old_range); + assert_eq!( + old_range, + original.ssts.levels()[0].files[&missing_id].primary_key_range(&old_mapper) + ); + } + + #[rstest] + #[case::mixed(vec![ + Value::from("default"), + Value::Null, + Value::Int64(42), + Value::Binary(vec![0, 255].into()), + ])] + #[case::nulls(vec![Value::Null; 4])] + #[case::empty(vec![])] + fn test_dense_ranges_complete_every_historical_prefix(#[case] defaults: Vec) { + let metadata = metadata(&defaults); + let mapper = PrimaryKeyRangeMapper::new(metadata.clone()); + let expected = encode(&metadata, &defaults); + for count in 0..=defaults.len() { + let prefix = encode(&metadata, &defaults[..count]); + assert_eq!( + Some((expected.clone(), expected.clone())), + mapper.map((prefix.clone(), prefix)).unwrap() + ); + } + } + + #[rstest] + #[case::cold(false)] + #[case::warm(true)] + fn test_successive_pk_appends_reuse_constants_and_complete_all_prefixes(#[case] warm: bool) { + let defaults = [ + Value::from(""), + Value::from("constant"), + Value::Null, + Value::Int64(42), + Value::Null, + ]; + let mut mapper = PrimaryKeyRangeMapper::new(metadata(&defaults[..2])); + let mut values = defaults.clone(); + values[0] = Value::from("stored"); + let raw = encode(&mapper.metadata, &values[..1]); + if warm { + assert!(mapper.map((raw.clone(), raw)).unwrap().is_some()); + } + + for count in 3..=defaults.len() { + let next_metadata = metadata(&defaults[..count]); + let next_mapper = mapper.with_metadata(next_metadata.clone()); + if let Some(Some(encoded)) = mapper.encoded_defaults[1].get() { + // An already encoded non-NULL constant must share its allocation across ALTER. + let reused = next_mapper.encoded_defaults[1] + .get() + .unwrap() + .as_ref() + .unwrap(); + assert_eq!(encoded.as_ptr(), reused.as_ptr()); + } + let expected = encode(&next_metadata, &values[..count]); + for prefix_len in 1..=count { + let prefix = encode(&next_metadata, &values[..prefix_len]); + assert_eq!( + Some((expected.clone(), expected.clone())), + next_mapper.map((prefix.clone(), prefix)).unwrap() + ); + } + let old_expected = encode(&mapper.metadata, &values[..count - 1]); + let raw = encode(&mapper.metadata, &values[..1]); + assert_eq!( + Some((old_expected.clone(), old_expected)), + mapper.map((raw.clone(), raw)).unwrap() + ); + mapper = next_mapper; + } + } + + #[rstest] + #[case::string(ConcreteDataType::string_datatype(), true, true)] + #[case::int(ConcreteDataType::int64_datatype(), true, true)] + #[case::binary(ConcreteDataType::binary_datatype(), true, true)] + #[case::timestamp(ConcreteDataType::timestamp_millisecond_datatype(), true, true)] + #[case::required(ConcreteDataType::string_datatype(), false, false)] + #[case::unsupported_list( + ConcreteDataType::list_datatype(Arc::new(ConcreteDataType::string_datatype())), + true, + false + )] + #[case::unsupported_null(ConcreteDataType::null_datatype(), true, false)] + fn test_null_padding_requires_a_nullable_encodable_column( + #[case] data_type: ConcreteDataType, + #[case] nullable: bool, + #[case] known: bool, + ) { + let mut metadata = (*metadata(&[Value::Null])).clone(); + metadata.column_metadatas[0].column_schema = + ColumnSchema::new("tag_0", data_type, nullable); + let metadata = Arc::new( + RegionMetadataBuilder::from_existing(metadata) + .build() + .unwrap(), + ); + let mapper = PrimaryKeyRangeMapper::new(metadata.clone()); + let expected = known.then(|| (Bytes::from_static(&[0]), Bytes::from_static(&[0]))); + assert_eq!(expected, mapper.map((Bytes::new(), Bytes::new())).unwrap()); + if known { + let encoded_null = encode(&metadata, &[Value::Null]); + assert_eq!(Some((encoded_null.clone(), encoded_null)), expected); + } + } + + // Invalid order/layout is an error; foreign bounds and unavailable defaults + // remain unknown. Exercise both endpoints without an inverted range masking decoding. + #[rstest] + #[case::dense(PrimaryKeyEncoding::Dense)] + #[case::sparse(PrimaryKeyEncoding::Sparse)] + fn test_reversed_pk_bounds_return_error(#[case] encoding: PrimaryKeyEncoding) { + use common_error::ext::ErrorExt; + use common_error::status_code::StatusCode; + use mito_codec::row_converter::build_primary_key_codec; + + let mut metadata = (*metadata(&[Value::from("")])).clone(); + metadata.primary_key_encoding = encoding; + let metadata = Arc::new(metadata); + let codec = build_primary_key_codec(&metadata); + let mut a = Vec::new(); + let mut b = Vec::new(); + codec + .encode_values(&[(0, Value::from("a"))], &mut a) + .unwrap(); + codec + .encode_values(&[(0, Value::from("b"))], &mut b) + .unwrap(); + let mapper = PrimaryKeyRangeMapper::new(metadata); + let err = mapper.map((b.into(), a.into())).unwrap_err(); + assert!(matches!( + err, + crate::error::Error::InvalidPrimaryKeyRange { .. } + )); + assert_eq!(StatusCode::Internal, err.status_code()); + } + + #[test] + fn test_different_pk_prefix_lengths_return_error() { + let metadata = metadata(&[Value::from(""), Value::Int64(42)]); + let mapper = PrimaryKeyRangeMapper::new(metadata.clone()); + let min = encode(&metadata, &[Value::from("a")]); + let max = encode(&metadata, &[Value::from("b"), Value::Int64(42)]); + assert!(matches!( + mapper.map((min, max)), + Err(crate::error::Error::InvalidPrimaryKeyRange { .. }) + )); + } + + #[rstest] + #[case::min("min")] + #[case::max("max")] + fn test_truncated_pk_endpoint_preserves_decode_error(#[case] endpoint: &'static str) { + use common_error::ext::ErrorExt; + use common_error::status_code::StatusCode; + + let metadata = metadata(&[Value::from("")]); + let mapper = PrimaryKeyRangeMapper::new(metadata.clone()); + let mut min = encode(&metadata, &[Value::from("a")]); + let mut max = encode(&metadata, &[Value::from("b")]); + if endpoint == "min" { + min.truncate(min.len() - 1); + } else { + max.truncate(max.len() - 1); + } + let err = mapper.map((min, max)).unwrap_err(); + assert_eq!(StatusCode::Internal, err.status_code()); + assert!(matches!(err, crate::error::Error::DecodePrimaryKeyRange { + endpoint: actual, + source: mito_codec::error::Error::InvalidDensePrimaryKey { .. }, + .. + } if actual == endpoint)); + } + + #[rstest] + #[case::source_first(false)] + #[case::destination_first(true)] + fn test_aligned_cache_uses_table_schema_across_region_migration( + #[case] destination_first: bool, + ) { + let source = metadata(&[Value::from("")]); + let file = FileHandle::new( + file_meta(&source, &[Value::from("a")]), + new_noop_file_purger(), + ); + let target = metadata_at_version(&[Value::from(""), Value::Int64(42)], 1); + let expected = encode(&target, &[Value::from("a"), Value::Int64(42)]); + let mut regions = [RegionId::new(1, 1), RegionId::new(1, 2)]; + if destination_first { + regions.reverse(); + } + let mut previous: Option<(Bytes, Bytes)> = None; + for region_id in regions { + let mut metadata = (*target).clone(); + metadata.region_id = region_id; + let mapper = PrimaryKeyRangeMapper::new(Arc::new(metadata)); + let aligned = file.primary_key_range(&mapper).unwrap(); + assert_eq!((expected.clone(), expected.clone()), aligned); + if let Some(previous) = previous { + // Different metadata allocations/region ids still share one schema-version cache. + assert_eq!(previous.0.as_ptr(), aligned.0.as_ptr()); + assert_eq!(previous.1.as_ptr(), aligned.1.as_ptr()); + } + previous = Some(aligned); + } + } + + #[test] + fn test_invalid_file_ranges_log_once_and_keep_possible_overlap() { + use std::sync::atomic::{AtomicUsize, Ordering}; + + use common_telemetry::tracing_subscriber::Layer; + use common_telemetry::tracing_subscriber::layer::Context; + use common_telemetry::tracing_subscriber::prelude::*; + use tracing::{Event, Level, Subscriber}; + + struct WarningCounter(Arc); + impl Layer for WarningCounter { + fn on_event(&self, event: &Event<'_>, _ctx: Context<'_, S>) { + if *event.metadata().level() == Level::WARN { + self.0.fetch_add(1, Ordering::Relaxed); + } + } + } + + let warnings = Arc::new(AtomicUsize::new(0)); + let subscriber = + common_telemetry::tracing_subscriber::registry().with(WarningCounter(warnings.clone())); + let metadata = metadata(&[Value::from("")]); + let mapper = Arc::new(PrimaryKeyRangeMapper::new(metadata.clone())); + let healthy_meta = file_meta(&metadata, &[Value::from("c")]); + let healthy = FileHandle::new(healthy_meta, new_noop_file_purger()); + let a = encode(&metadata, &[Value::from("a")]); + let b = encode(&metadata, &[Value::from("b")]); + + tracing::subscriber::with_default(subscriber, || { + for (min, max) in [(b.clone(), a.clone()), (a.slice(..a.len() - 1), b)] { + let mut meta = file_meta(&metadata, &[Value::from("a")]); + meta.primary_key_min = Some(min); + meta.primary_key_max = Some(max); + let file = FileHandle::new(meta, new_noop_file_purger()); + assert_eq!(None, file.primary_key_range(&mapper)); + assert_eq!(None, file.clone().primary_key_range(&mapper)); + // These raw ranges look disjoint from c; unknown must prevent pruning. + assert!(files_overlap_inclusive(&file, &healthy, &mapper)); + assert!(files_overlap_inclusive(&healthy, &file, &mapper)); + } + }); + assert_eq!(2, warnings.load(Ordering::Relaxed)); + } + + #[test] + fn test_missing_impure_default_is_unknown_but_complete_key_is_usable() { + let mut metadata = (*metadata(&[Value::Timestamp(Timestamp::new_millisecond(0))])).clone(); + metadata.column_metadatas[0].column_schema = metadata.column_metadatas[0] + .column_schema + .clone() + .with_default_constraint(Some(ColumnDefaultConstraint::Function( + "current_timestamp()".into(), + ))) + .unwrap(); + let metadata = Arc::new(metadata); + let mapper = PrimaryKeyRangeMapper::new(metadata.clone()); + assert_eq!(None, mapper.map((Bytes::new(), Bytes::new())).unwrap()); + let key = encode( + &metadata, + &[Value::Timestamp(Timestamp::new_millisecond(10))], + ); + assert_eq!( + Some((key.clone(), key.clone())), + mapper.map((key.clone(), key)).unwrap() + ); + } + + #[test] + fn test_sparse_ranges_keep_their_original_encoding() { + use mito_codec::row_converter::SparsePrimaryKeyCodec; + let mut metadata = (*metadata(&[Value::from(""), Value::Null])).clone(); + metadata.primary_key_encoding = PrimaryKeyEncoding::Sparse; + let metadata = Arc::new(metadata); + let mut bytes = Vec::new(); + SparsePrimaryKeyCodec::new(&metadata) + .encode_values(&[(0, Value::from("a"))], &mut bytes) + .unwrap(); + let key = Bytes::from(bytes); + let mapper = PrimaryKeyRangeMapper::new(metadata); + assert_eq!( + Some((key.clone(), key.clone())), + mapper.map((key.clone(), key)).unwrap() + ); + } + + #[test] + fn test_aligned_cache_isolates_schemas_and_shares_file_state() { + let old_metadata = metadata(&[Value::from("")]); + let new_metadata = + metadata_at_version(&[Value::from(""), Value::Null, Value::Int64(42)], 1); + let raw_meta = file_meta(&old_metadata, &[Value::from("b")]); + let original = version_with_files(old_metadata.clone(), vec![raw_meta.clone()]); + let old_file = original.ssts.levels()[0].files().next().unwrap(); + let old_mapper = original.ssts.primary_key_mapper(); + let changed = VersionBuilder::from_version(original.clone()) + .metadata(new_metadata.clone()) + .build(); + let new_file = changed.ssts.levels()[0].files().next().unwrap(); + let new_mapper = changed.ssts.primary_key_mapper(); + let expected = encode( + &new_metadata, + &[Value::from("b"), Value::Null, Value::Int64(42)], + ); + assert_eq!( + raw_meta.primary_key_range(), + old_file.primary_key_range(&old_mapper) + ); + assert_eq!( + Some((expected.clone(), expected)), + new_file.primary_key_range(&new_mapper) + ); + assert_eq!(raw_meta, *new_file.meta_ref()); + assert_eq!( + raw_meta.primary_key_range(), + new_file.raw_primary_key_range() + ); + assert_eq!( + old_file.primary_key_range(&old_mapper), + new_file.primary_key_range(&old_mapper) + ); + assert_eq!( + new_file.primary_key_range(&new_mapper), + old_file.primary_key_range(&new_mapper) + ); + old_file.set_compacting(true); + assert!(new_file.compacting()); + new_file.mark_deleted(); + assert!(old_file.is_deleted()); + + // Old-schema flushes can finish after ALTER; compare them with the target mapper. + let added_meta = file_meta(&old_metadata, &[Value::from("a")]); + let added_id = added_meta.file_id; + let changed = VersionBuilder::from_version(Arc::new(changed)) + .apply_edit( + RegionEdit { + files_to_add: vec![added_meta], + files_to_remove: vec![], + timestamp_ms: None, + compaction_time_window: None, + flushed_entry_id: None, + flushed_sequence: None, + committed_sequence: None, + }, + new_noop_file_purger(), + ) + .build(); + let added = &changed.ssts.levels()[0].files[&added_id]; + let expected = encode( + &new_metadata, + &[Value::from("a"), Value::Null, Value::Int64(42)], + ); + assert_eq!( + Some((expected.clone(), expected)), + added.primary_key_range(&new_mapper) + ); + } + + #[test] + fn test_late_statistics_and_changed_defaults_use_each_pinned_schema() { + let metadata_v1 = metadata(&[Value::from(""), Value::Int64(42)]); + let metadata_v2 = metadata_at_version(&[Value::from(""), Value::Int64(100)], 1); + let mut meta = file_meta(&metadata_v1, &[Value::from("a")]); + let raw = meta.primary_key_range().unwrap(); + meta.primary_key_min = None; + meta.primary_key_max = None; + let file = FileHandle::new(meta, new_noop_file_purger()); + let v1 = PrimaryKeyRangeMapper::new(metadata_v1.clone()); + let v2 = PrimaryKeyRangeMapper::new(metadata_v2.clone()); + assert_eq!(None, file.primary_key_range(&v1)); + assert_eq!(None, file.primary_key_range(&v2)); + file.set_primary_key_range(raw.clone()); + for (mapper, metadata, default) in [ + (&v2, &metadata_v2, 100), + (&v1, &metadata_v1, 42), + (&v2, &metadata_v2, 100), + ] { + let expected = encode(metadata, &[Value::from("a"), Value::Int64(default)]); + assert_eq!( + Some((expected.clone(), expected)), + file.primary_key_range(mapper) + ); + } + assert_eq!(Some(raw), file.raw_primary_key_range()); + } + + #[test] + fn test_concurrent_snapshots_do_not_return_each_others_aligned_bounds() { + let old_metadata = metadata(&[Value::from(""), Value::Int64(42)]); + let new_metadata = metadata_at_version(&[Value::from(""), Value::Int64(100)], 1); + let raw = file_meta(&old_metadata, &[Value::from("a")]); + let original_bounds = raw.primary_key_range(); + let file = FileHandle::new(raw, new_noop_file_purger()); + let barrier = std::sync::Barrier::new(2); + std::thread::scope(|scope| { + for (metadata, default) in [(old_metadata, 42), (new_metadata, 100)] { + let file = &file; + let barrier = &barrier; + scope.spawn(move || { + let expected = encode(&metadata, &[Value::from("a"), Value::Int64(default)]); + let mapper = PrimaryKeyRangeMapper::new(metadata); + let mut all_matched = true; + for _ in 0..64 { + barrier.wait(); + all_matched &= file.primary_key_range(&mapper) + == Some((expected.clone(), expected.clone())); + } + // Assert after all barriers so a failure cannot strand the other snapshot. + assert!( + all_matched, + "wrong bounds for schema version {}", + mapper.schema_version() + ); + }); + } + }); + assert_eq!(original_bounds, file.raw_primary_key_range()); + } + + #[test] + fn test_deserialized_compaction_files_use_the_comparison_schema() { + use crate::compaction::CompactionOutput; + use crate::compaction::picker::{PickerOutput, SerializedPickerOutput}; + + let metadata = metadata(&[Value::from(""), Value::Int64(42)]); + let raw = file_meta(&metadata, &[Value::from("a")]); + let picked = PickerOutput { + outputs: vec![CompactionOutput { + output_level: 1, + inputs: vec![FileHandle::new(raw.clone(), new_noop_file_purger())], + filter_deleted: false, + output_time_range: None, + }], + ..Default::default() + }; + let serialized = SerializedPickerOutput::from(&picked); + let output = PickerOutput::from_serialized(serialized, new_noop_file_purger()); + let mapper = PrimaryKeyRangeMapper::new(metadata.clone()); + let file = &output.outputs[0].inputs[0]; + let completed = encode(&metadata, &[Value::from("a"), Value::Int64(42)]); + assert_eq!( + Some((completed.clone(), completed)), + file.primary_key_range(&mapper) + ); + assert_eq!(raw, *file.meta_ref()); + } + + #[test] + fn test_compaction_overlap_uses_completed_endpoints() { + let metadata = metadata(&[Value::from(""), Value::Null]); + let mapper = Arc::new(PrimaryKeyRangeMapper::new(metadata.clone())); + let mut old = file_meta(&metadata, &[Value::from("a")]); + old.primary_key_max = Some(encode(&metadata, &[Value::from("b")])); + let mut new = file_meta(&metadata, &[Value::from("b"), Value::Null]); + new.primary_key_max = Some(encode(&metadata, &[Value::from("c"), Value::Null])); + assert!(old.primary_key_max < new.primary_key_min); + let old = FileHandle::new(old, new_noop_file_purger()); + let new = FileHandle::new(new, new_noop_file_purger()); + assert!(files_overlap_inclusive(&old, &new, &mapper)); + assert!(files_overlap_inclusive(&new, &old, &mapper)); + } +} diff --git a/src/mito2/src/sst/version.rs b/src/mito2/src/sst/version.rs index 172d8701a73..1ca0410db7b 100644 --- a/src/mito2/src/sst/version.rs +++ b/src/mito2/src/sst/version.rs @@ -18,31 +18,45 @@ use std::fmt; use std::sync::Arc; use common_time::{TimeToLive, Timestamp}; +use store_api::metadata::RegionMetadataRef; use store_api::storage::{FileId, RegionId}; use crate::sst::file::{FileHandle, FileMeta, FileTimeRange, Level, MAX_LEVEL}; use crate::sst::file_purger::FilePurgerRef; +use crate::sst::primary_key::PrimaryKeyRangeMapper; /// A version of all SSTs in a region. #[derive(Debug, Clone)] pub(crate) struct SstVersion { /// SST metadata organized by levels. - levels: LevelMetaArray, + levels: Arc, + primary_key_mapper: Arc, } pub(crate) type SstVersionRef = Arc; impl SstVersion { /// Returns a new [SstVersion]. - pub(crate) fn new() -> SstVersion { + pub(crate) fn new(metadata: RegionMetadataRef) -> SstVersion { SstVersion { - levels: new_level_meta_vec(), + levels: Arc::new(new_level_meta_vec()), + primary_key_mapper: Arc::new(PrimaryKeyRangeMapper::new(metadata)), } } + /// Changes the target schema without copying the SST list or rebinding handles. + pub(crate) fn set_metadata(&mut self, metadata: RegionMetadataRef) { + self.primary_key_mapper = Arc::new(self.primary_key_mapper.with_metadata(metadata)); + } + + /// Shares the target schema and encoded defaults with comparisons. + pub(crate) fn primary_key_mapper(&self) -> Arc { + self.primary_key_mapper.clone() + } + /// Returns a slice to metadatas of all levels. pub(crate) fn levels(&self) -> &[LevelMeta] { - &self.levels + self.levels.as_ref() } /// Returns the current handle matching the selected file's identity in its immutable level. @@ -65,11 +79,12 @@ impl SstVersion { file_purger: FilePurgerRef, files_to_add: impl Iterator, ) { + let levels = Arc::make_mut(&mut self.levels); for file in files_to_add { let level = file.level; let new_index_version = file.index_version; // If the file already exists, then we should only replace the handle when the index is outdated. - self.levels[level as usize] + levels[level as usize] .files .entry(file.file_id) .and_modify(|f| { @@ -99,9 +114,10 @@ impl SstVersion { /// # Panics /// Panics if level of [FileMeta] is greater than [MAX_LEVEL]. pub(crate) fn remove_files(&mut self, files_to_remove: impl Iterator) { + let levels = Arc::make_mut(&mut self.levels); for file in files_to_remove { let level = file.level; - if let Some(handle) = self.levels[level as usize].files.remove(&file.file_id) { + if let Some(handle) = levels[level as usize].files.remove(&file.file_id) { handle.mark_deleted(); } } @@ -109,7 +125,7 @@ impl SstVersion { /// Marks all SSTs in this version as deleted. pub(crate) fn mark_all_deleted(&self) { - for level_meta in &self.levels { + for level_meta in self.levels.iter() { for file_handle in level_meta.files.values() { file_handle.mark_deleted(); } @@ -289,7 +305,7 @@ mod tests { ..Default::default() }; - let mut version = SstVersion::new(); + let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test()); version.add_files(purger, [owned, referenced].into_iter()); assert_eq!( @@ -318,7 +334,7 @@ mod tests { ..Default::default() }; - let mut version = SstVersion::new(); + let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test()); version.add_files(purger, [seconds, millis].into_iter()); assert_eq!( @@ -332,7 +348,8 @@ mod tests { #[test] fn time_range_is_none_without_files() { - assert_eq!(SstVersion::new().time_range(), None); + let version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test()); + assert_eq!(version.time_range(), None); } #[test] @@ -346,7 +363,7 @@ mod tests { }) .collect::>(); - let mut version = SstVersion::new(); + let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test()); // files[1] is added multiple times, and that's ok. version.add_files(purger.clone(), files[..=1].iter().cloned()); version.add_files(purger, files[1..].iter().cloned()); @@ -370,7 +387,7 @@ mod tests { }, purger.clone(), ); - let mut version = SstVersion::new(); + let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test()); version.add_files( purger, [ @@ -422,7 +439,7 @@ mod tests { }, ]; - let mut version = SstVersion::new(); + let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test()); version.add_files(purger, files.iter().cloned()); assert_eq!(3, version.owned_num_rows(region_id)); diff --git a/tests-integration/fixtures/docker-compose.yml b/tests-integration/fixtures/docker-compose.yml index 283998cc28f..1e527c27019 100644 --- a/tests-integration/fixtures/docker-compose.yml +++ b/tests-integration/fixtures/docker-compose.yml @@ -134,6 +134,19 @@ services: - MYSQL_USER=greptimedb - MYSQL_PASSWORD=admin - MYSQL_ROOT_PASSWORD=admin + # Wait for authenticated TCP queries, not just the container or initialization socket. + healthcheck: + test: + - CMD-SHELL + - >- + MYSQL_PWD="$${MYSQL_PASSWORD}" mysql + --protocol=TCP --host=127.0.0.1 + --user="$${MYSQL_USER}" --database="$${MYSQL_DATABASE}" + --execute="SELECT 1" + interval: 2s + timeout: 5s + retries: 30 + start_period: 10s volumes: minio_data: