fix(mito2): compare primary key ranges across schema versions (#9205)

* fix(mito2): compare primary key ranges across schema versions

Bind FileHandle ranges to the pinned region schema and append cached constant defaults to historical Dense keys. Preserve raw SST statistics, reject inexact bounds, and avoid invalidating views for unrelated metadata changes.

Cover schema evolution, default changes, and tombstone retention through real compaction and reopen regressions.

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

* refactor(mito2): report invalid primary key ranges

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

* refactor(mito2): scope primary key ranges to comparisons

Keep raw PK bounds only in FileHandleInner and move schema-aware mapping and caching into task-local comparison contexts.

Use explicit contexts for compaction overlap checks, window aggregation, and series scans. Preserve pinned-schema isolation, late statistics, and shared file lifecycle state without rebinding every handle.

Cover cache isolation across region owners and adapt range fixtures to real Dense encodings. All 1530 mito2 tests and Clippy for all targets with the testing feature pass.

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

* refactor(mito2): cache aligned primary key ranges per file

Replace task-local range maps with a single-slot cache in FileHandleInner, keyed by the target schema version. Preserve raw bounds for realignment across snapshots and default changes.

Share schema mappers across comparison paths and use copy-on-write SST lists for metadata updates. Simplify range mapping to accept encoded bounds and assert the same-table contract at the file accessor.

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

* docs(mito2): clarify primary key mapper schema snapshot

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

* test(mito2): align primary key range fixtures with table contract

Remove obsolete cross-table fallback expectations after region validation became a caller contract. Give compaction fixtures matching table identities, including the active-window L1 scenario.

Clarify the mapper precondition and format the simplified alignment call. All 1530 mito2 tests and Clippy for all targets with the testing feature pass.

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

* test(mito2): enable filesystem GC in release unit tests

Let unit tests use the filesystem-backed object-store GC path regardless of optimization profile. Keep the production release GC selection unchanged.

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

* perf(mito-codec): skip release value decoding in PK prefix counts

Validate field values only in debug builds and unit tests while keeping boundary, truncation, and trailing-byte checks in every build.

Cover the linked library in debug and release integration tests, and verify that release unit tests still perform value validation.

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

* test(mito-codec): remove redundant prefix integration tests

Retain the codec unit tests and cross-schema compaction regressions while dropping the standalone build-profile test file and its release-only invalid-value expectation.

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

* ci: wait for MySQL to accept authenticated TCP queries

Add a healthcheck using the configured test account and database. Docker Compose --wait previously only observed container startup because the fixture image had no healthcheck, allowing metasrv to connect before MySQL initialization completed.

Verify readiness with SELECT 1 over TCP rather than the initialization socket.

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

---------

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>
This commit is contained in:
Lei, HUANG
2026-09-20 14:30:12 +00:00
committed by GitHub
parent fcb165cf77
commit e9dd79d1d2
24 changed files with 1929 additions and 211 deletions
+10 -1
View File
@@ -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(),
}
}
+173 -51
View File
@@ -290,65 +290,84 @@ impl SortField {
bytes: &[u8],
deserializer: &mut Deserializer<&[u8]>,
) -> Result<usize> {
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<usize> {
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<usize> {
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<i64>
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<usize> {
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<Vec<Value>> {
let mut deserializer = Deserializer::new(bytes);
@@ -594,6 +642,25 @@ mod tests {
let value_ref = row.iter().map(|v| v.as_value_ref()).collect::<Vec<_>>();
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![
+7 -2
View File
@@ -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,
);
});
});
}
+3 -3
View File
@@ -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,
},
+9 -9
View File
@@ -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<CompactionOutput>,
files: Vec<FileHandle>,
metadata: &RegionMetadata,
metadata: &RegionMetadataRef,
) -> Vec<CompactionOutput> {
// 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!(),
};
+32 -19
View File
@@ -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<FileHandle>,
@@ -54,7 +56,7 @@ struct OverlapQuery<'a> {
}
impl<'a> FileOverlapIndex<'a> {
pub(super) fn new(files: Vec<FileHandle>, metadata: &'a RegionMetadata) -> Self {
pub(super) fn new(files: Vec<FileHandle>, 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();
+2 -1
View File
@@ -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<dyn FilePurger>,
+59 -39
View File
@@ -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<T: Item + Clone>(
l: &mut SortedRun<T>,
r: &mut SortedRun<T>,
result: &mut Vec<T>,
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<T: Item + Clone>(
}
// 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<T>(items: &mut [T]) -> Vec<SortedRun<T>>
/// Finds sorted runs using the caller's overlap test within each active time range.
pub fn find_sorted_runs<T>(items: &mut [T], overlaps: impl Fn(&T, &T) -> bool) -> Vec<SortedRun<T>>
where
T: Item,
{
@@ -296,7 +303,7 @@ where
let mut overlaps_any = false;
for idx in &active_run_item_indices {
let run_item = &current_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<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::<MockFile>::new()),
&mut SortedRun::from(Vec::<MockFile>::new()),
&mut result,
Ranged::overlap_inclusive,
);
assert_eq!(result, Vec::<MockFile>::new());
@@ -669,6 +675,7 @@ mod tests {
&mut SortedRun::from(files1.clone()),
&mut SortedRun::from(Vec::<MockFile>::new()),
&mut result,
Ranged::overlap_inclusive,
);
assert_eq!(result, Vec::<MockFile>::new());
@@ -676,6 +683,7 @@ mod tests {
&mut SortedRun::from(Vec::<MockFile>::new()),
&mut SortedRun::from(files1.clone()),
&mut result,
Ranged::overlap_inclusive,
);
assert_eq!(result, Vec::<MockFile>::new());
@@ -686,6 +694,7 @@ mod tests {
&mut SortedRun::from(files1),
&mut SortedRun::from(files2),
&mut result,
Ranged::overlap_inclusive,
);
assert_eq!(result, Vec::<MockFile>::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()));
}
}
+25 -2
View File
@@ -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(),
+232 -48
View File
@@ -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<i64, Window>,
phase: PickPhase,
primary_key_mapper: &'a PrimaryKeyRangeMapper,
}
struct WindowOutputContext {
active_window: Option<i64>,
time_window_size: Option<i64>,
max_outputs: Option<usize>,
primary_key_mapper: Arc<PrimaryKeyRangeMapper>,
}
/// 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<FileHandle>,
l1_files: Vec<FileHandle>,
pick: InactiveWindowPick,
pick: InactiveWindowPick<'_>,
max_output_file_size: Option<u64>,
) -> (Vec<FileHandle>, 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<FileHandle>,
max_output_file_size: Option<u64>,
phase: PickPhase,
mapper: &PrimaryKeyRangeMapper,
) -> (Vec<FileHandle>, 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<SortedRun<FileHandle>>, Option<u64>, &PrimaryKeyRangeMapper) -> Vec<FileHandle>;
fn pick_candidate_files(
mut files: Vec<FileHandle>,
max_output_file_size: Option<u64>,
picker: fn(Vec<SortedRun<FileHandle>>, Option<u64>) -> Vec<FileHandle>,
picker: CandidatePicker,
mapper: &PrimaryKeyRangeMapper,
) -> (Vec<FileHandle>, 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<bool>,
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<SortedRun<FileHandle>>,
max_output_file_size: Option<u64>,
mapper: &PrimaryKeyRangeMapper,
) -> Vec<FileHandle> {
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<SortedRun<FileHandle>>,
max_output_file_size: Option<u64>,
mapper: &PrimaryKeyRangeMapper,
) -> Vec<FileHandle> {
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<SortedRun<FileHandle>>,
max_output_file_size: Option<u64>,
mapper: &PrimaryKeyRangeMapper,
) -> Vec<FileHandle> {
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<SortedRun<FileHandle>>,
max_output_file_size: Option<u64>,
mapper: &PrimaryKeyRangeMapper,
) -> Vec<FileHandle> {
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<SortedRun<FileHandle>>,
max_output_file_size: Option<u64>,
mapper: &PrimaryKeyRangeMapper,
) -> Vec<FileHandle> {
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<SortedRun<FileHandle>>,
max_output_file_size: Option<u64>,
mapper: &PrimaryKeyRangeMapper,
) -> Vec<FileHandle> {
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<SortedRun<FileHandle>>,
max_output_file_size: Option<u64>,
mapper: &PrimaryKeyRangeMapper,
is_eligible: impl Fn(&Candidate) -> bool,
) -> Vec<FileHandle> {
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<bool> = 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<FileHandle>]) -> Vec<OrderedFile<'_>>
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<Item = &'a FileHandle>,
time_window_size: i64,
mapper: &PrimaryKeyRangeMapper,
) -> BTreeMap<i64, Window> {
let mut windows: HashMap<i64, Window> = 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<Item = &'a FileHandle>,
time_window_size: i64,
) -> BTreeMap<i64, Window> {
super::assign_to_windows(files, time_window_size, &primary_key_mapper_for_test())
}
fn pick_count_first(
sorted_runs: Vec<SortedRun<FileHandle>>,
max_output_file_size: Option<u64>,
) -> Vec<FileHandle> {
super::pick_count_first(
sorted_runs,
max_output_file_size,
&primary_key_mapper_for_test(),
)
}
fn pick_mixed_count_first(
sorted_runs: Vec<SortedRun<FileHandle>>,
max_output_file_size: Option<u64>,
) -> Vec<FileHandle> {
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,
);
+1 -1
View File
@@ -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,
+319 -2
View File
@@ -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::<Vec<_>>();
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::<Vec<_>>();
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 = <rand::rngs::StdRng as rand::SeedableRng>::seed_from_u64(0);
for _ in 0..3000 {
old_rows.extend(build_rows_for_key(
&format!("a{:032x}", rand::Rng::random::<u128>(&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::<Vec<_>>();
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<api::v1::Row>| {
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::<Vec<_>>();
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::<HashSet<_>>();
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::<Vec<_>>();
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::<UInt8Type>()
.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<i64> {
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<Arc<CompactionListener>>);
impl CompactionListenerGuard {
@@ -1790,6 +2098,15 @@ async fn env_for_manual_compaction(
env: &mut TestEnv,
region_id: RegionId,
flat_format: bool,
) -> (MitoEngine, Vec<ColumnSchema>) {
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<ColumnSchema>) {
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
+19 -1
View File
@@ -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,
+17 -1
View File
@@ -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<MemRangeBuilder>,
/// Handles to SST files to scan.
pub(crate) files: Vec<FileHandle>,
/// Shares the pinned schema's encoded defaults across parallel range readers.
primary_key_mapper: OnceLock<Arc<PrimaryKeyRangeMapper>>,
/// 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<PrimaryKeyRangeMapper>) -> 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 {
+1 -1
View File
@@ -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)
) {
+6 -4
View File
@@ -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(),
&region_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(),
&region_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]
+6 -3
View File
@@ -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
}
+1
View File
@@ -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;
+98 -6
View File
@@ -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<S>(bytes: &Option<Bytes>, serializer: S) -> Result<S::Ok, S::Error>
@@ -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<PrimaryKeyBounds>,
},
}
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<PrimaryKeyBounds>> {
match self {
Self::Aligned {
schema_version,
bounds,
..
} if *schema_version == target_version => Some(bounds),
_ => None,
}
}
fn align(
&mut self,
mapper: &PrimaryKeyRangeMapper,
) -> crate::error::Result<Option<PrimaryKeyBounds>> {
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<Option<(Bytes, Bytes)>>,
primary_key_range: RwLock<PrimaryKeyRange>,
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,
}
}
+4 -4
View File
@@ -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
}
+15
View File
@@ -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(
+847
View File
@@ -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<OnceLock<Option<Bytes>>>,
suffixes: Vec<OnceLock<Option<Bytes>>>,
}
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<Option<(Bytes, Bytes)>> {
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<usize> {
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<Bytes> {
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<Bytes> {
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<FileMeta>) -> 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<Value>,
) {
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<Value>) {
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<AtomicUsize>);
impl<S: Subscriber> Layer<S> 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));
}
}
+30 -13
View File
@@ -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<LevelMetaArray>,
primary_key_mapper: Arc<PrimaryKeyRangeMapper>,
}
pub(crate) type SstVersionRef = Arc<SstVersion>;
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<PrimaryKeyRangeMapper> {
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<Item = FileMeta>,
) {
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<Item = FileMeta>) {
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::<Vec<_>>();
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));
@@ -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: