fix(mito2): preserve effective sequences during compaction (#9147)

* fix(mito2): preserve effective row sequences during compaction

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

* fix(mito2): inherit compaction input sequence bounds

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-15 07:43:16 +00:00
committed by GitHub
parent 32bf865ffb
commit e0552eee2d
8 changed files with 894 additions and 99 deletions
+2
View File
@@ -568,6 +568,8 @@ pub struct SstWriteRequest {
pub cache_manager: CacheManagerRef,
#[allow(dead_code)]
pub storage: Option<String>,
/// Optional uniform row sequence for writes that do not preserve sequences.
/// Compaction passes `None` to retain the reader's effective input sequences.
pub max_sequence: Option<SequenceNumber>,
pub sst_write_format: FormatType,
+296 -41
View File
@@ -404,10 +404,33 @@ pub trait SstMerger: Send + Sync + 'static {
#[derive(Clone)]
pub struct DefaultSstMerger;
/// File-level sequence metadata is independent of physical row encoding.
struct OutputSequenceMetadata {
file_sequence: Option<NonZeroU64>,
exact_sequence_trusted: bool,
}
/// Compaction does not admit new data, so it inherits the input sequence bounds
/// without allocating a new admission marker. Retaining physical row sequences
/// does not by itself restore exact-read capability for untrusted inputs.
fn output_sequence_metadata(
region: &CompactionRegion,
inputs: &[FileHandle],
) -> OutputSequenceMetadata {
let exact_sequence_trusted = region.region_options.preserve_row_sequence
&& inputs
.iter()
.all(|f| f.is_effective_target_sequence_trusted(region.region_id));
OutputSequenceMetadata {
file_sequence: known_max_input_sequence(inputs),
exact_sequence_trusted,
}
}
/// Computes the maximum target-domain sequence bound of the output of merging
/// `inputs`. Foreign files are described by their target-local barrier; local
/// trusted files use their physical sequence bound. Unknown bounds never become
/// trusted through a partial maximum.
/// `inputs`. A bound can be a row maximum or an admission marker; compaction
/// must preserve both for sequence-based pruning. Unknown bounds must remain
/// unknown rather than being replaced with a partial maximum.
fn known_max_input_sequence(inputs: &[FileHandle]) -> Option<NonZeroU64> {
let mut max: Option<NonZeroU64> = None;
for input in inputs {
@@ -452,40 +475,7 @@ impl SstMerger for DefaultSstMerger {
.iter()
.map(|f| f.file_id().to_string())
.join(",");
let input_max_sequence = known_max_input_sequence(&output.inputs);
// The output is trusted only when every input is trusted in the target
// domain. Foreign marker values describe the source and are ignored;
// their present FileMeta.sequence is the target-local barrier.
let output_preserves_sequence = compaction_region.region_options.preserve_row_sequence
&& output
.inputs
.iter()
.all(|f| f.is_effective_target_sequence_trusted(region_id));
let output_sequence = if output_preserves_sequence {
input_max_sequence
} else {
// The manifest records the region's latest accepted sequence for
// edit paths and its flushed frontier otherwise. Compaction only
// rewrites the immutable SST snapshot, so this is the same
// region-local admission barrier used when accepting files.
let manifest = compaction_region
.manifest_ctx
.manifest_manager
.read()
.await
.manifest();
NonZeroU64::new(
manifest
.committed_sequence
.unwrap_or(manifest.flushed_sequence)
+ 1,
)
};
// For untrusted output, keep the physical sequence override compatible
// with the main write path: use the known maximum sequence of the
// inputs. The FileMeta sequence remains the admission barrier, but it
// must not be written into the rows (or replaced with zero).
let write_max_sequence = input_max_sequence.map(NonZeroU64::get);
let sequence_metadata = output_sequence_metadata(&compaction_region, &output.inputs);
let builder = CompactionSstReaderBuilder {
metadata: compaction_region.region_metadata.clone(),
sst_layer: compaction_region.access_layer.clone(),
@@ -509,13 +499,16 @@ impl SstMerger for DefaultSstMerger {
source,
cache_manager: compaction_region.cache_manager.clone(),
storage,
max_sequence: write_max_sequence,
// Readers resolve file overrides before merge/dedup. Replacing
// their effective sequences here could promote old rows above
// versions in SSTs that were not part of this merge.
max_sequence: None,
sst_write_format: if flat_format {
FormatType::Flat
} else {
FormatType::PrimaryKey
},
preserve_row_sequence: output_preserves_sequence,
preserve_row_sequence: sequence_metadata.exact_sequence_trusted,
index_options,
index_config,
inverted_index_config,
@@ -564,12 +557,12 @@ impl SstMerger for DefaultSstMerger {
index_version: 0,
num_rows: sst_info.num_rows as u64,
num_row_groups: sst_info.num_row_groups,
sequence: output_sequence,
sequence: sequence_metadata.file_sequence,
partition_expr: partition_expr.clone(),
num_series: sst_info.num_series,
primary_key_min,
primary_key_max,
preserve_row_sequence: output_preserves_sequence,
preserve_row_sequence: sequence_metadata.exact_sequence_trusted,
}
})
.collect::<Vec<_>>();
@@ -888,6 +881,268 @@ mod tests {
assert_eq!(None, known_max_input_sequence(&[]));
}
#[tokio::test]
async fn test_output_sequence_metadata_preserves_trust_and_unknown_bounds() {
let mut region = new_test_compaction_region().await;
let local = region.region_id;
let foreign = RegionId::new(1, 2);
let file = |region_id, sequence: Option<u64>, preserve_row_sequence| {
new_file_handle(FileMeta {
region_id,
sequence: sequence.and_then(NonZeroU64::new),
preserve_row_sequence,
..dummy_file_meta()
})
};
// Bounds are independent of row trust: known, mixed trust, unknown in
// either sequence domain, empty inputs, and the largest possible bound.
let cases = [
(
vec![file(local, Some(3), true), file(foreign, Some(9), false)],
Some(9),
true,
),
(
vec![file(local, Some(3), true), file(local, Some(6), false)],
Some(6),
false,
),
(
vec![file(foreign, Some(9), false), file(foreign, Some(11), true)],
Some(11),
true,
),
(
vec![file(local, Some(3), true), file(local, None, true)],
None,
true,
),
(
vec![file(local, None, false), file(local, Some(6), true)],
None,
false,
),
(
vec![file(local, Some(3), true), file(foreign, None, true)],
None,
false,
),
(vec![], None, true),
(
vec![file(local, Some(u64::MAX), false)],
Some(u64::MAX),
false,
),
];
for committed_sequence in [100, 101] {
region
.manifest_ctx
.manifest_manager
.write()
.await
.update(
RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
committed_sequence: Some(committed_sequence),
files_to_add: vec![],
files_to_remove: vec![],
timestamp_ms: None,
compaction_time_window: None,
flushed_entry_id: None,
flushed_sequence: None,
})),
false,
)
.await
.unwrap();
for preserve in [false, true] {
region.region_options.preserve_row_sequence = preserve;
for (inputs, known_bound, inputs_trusted) in &cases {
let metadata = output_sequence_metadata(&region, inputs);
let trusted = preserve && *inputs_trusted;
assert_eq!(trusted, metadata.exact_sequence_trusted);
assert_eq!(
known_bound.and_then(NonZeroU64::new),
metadata.file_sequence
);
}
}
}
}
/// Independent outputs must be safe in either visibility order, including when
/// only one output succeeds. Merely committing both outputs together is not a
/// substitute for preserving version order across their separate merge streams.
#[rstest::rstest]
#[case(true)]
#[case(false)]
#[tokio::test]
async fn test_independent_outputs_preserve_delete_in_every_visibility_state(
#[case] flat_format: bool,
#[values(false, true)] unknown_put_bound: bool,
) {
use datatypes::arrow::array::AsArray;
use datatypes::arrow::datatypes::TimestampMillisecondType;
use store_api::region_engine::RegionEngine;
use store_api::region_request::RegionRequest;
use crate::compaction::CompactionOutput;
use crate::compaction::compactor::{CompactionRegion, Compactor, DefaultCompactor};
use crate::compaction::picker::PickerOutput;
use crate::compaction::reader::CompactionSstReaderBuilder;
use crate::engine::compaction_test::{delete_and_flush, put_and_flush};
use crate::region::options::MergeMode;
use crate::sst::file::FileHandle;
use crate::sst::file_purger::NoopFilePurger;
use crate::test_util::{CreateRequestBuilder, TestEnv, rows_schema};
let mut env = TestEnv::new().await;
let region_id = RegionId::new(1, 1);
let config = MitoConfig {
default_flat_format: flat_format,
min_compaction_interval: Duration::from_secs(3600),
..Default::default()
};
let engine = env.create_engine(config).await;
let request = CreateRequestBuilder::new().build();
let columns = rows_schema(&request);
engine
.handle_request(region_id, RegionRequest::Create(request))
.await
.unwrap();
put_and_flush(&engine, region_id, &columns, 0..1).await;
delete_and_flush(&engine, region_id, &columns, 0..1).await;
put_and_flush(&engine, region_id, &columns, 10..11).await;
let region = engine.get_region(region_id).unwrap();
let version = region.version();
let mut files: Vec<_> = version
.ssts
.levels()
.iter()
.flat_map(|level| level.files())
.cloned()
.collect();
assert_eq!(3, files.len());
files.sort_unstable_by_key(|file| file.meta_ref().sequence);
if unknown_put_bound {
// Model a legacy SST whose physical rows are readable but whose
// file-level boundary is unknown. Merging cannot invent that bound.
let mut meta = files[0].meta_ref().clone();
meta.sequence = None;
files[0] = FileHandle::new(meta, Arc::new(NoopFilePurger));
}
let compaction_region = CompactionRegion {
region_id,
region_options: version.options.clone(),
engine_config: Arc::new(MitoConfig {
default_flat_format: flat_format,
..Default::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: 2,
plugins: common_base::Plugins::new(),
};
let puts = vec![files[0].clone(), files[2].clone()];
let deletes = vec![files[1].clone()];
let picker_output = PickerOutput {
outputs: [puts.clone(), deletes.clone()]
.into_iter()
.map(|inputs| CompactionOutput {
output_level: 1,
inputs,
filter_deleted: false,
output_time_range: None,
})
.collect(),
expired_ssts: vec![],
time_window_size: 3600,
max_file_size: None,
};
let merged = DefaultCompactor::with_merger(DefaultSstMerger)
.merge_ssts(&compaction_region, picker_output)
.await
.unwrap();
assert_eq!(2, merged.files_to_add.len());
let outputs: Vec<_> = merged
.files_to_add
.into_iter()
.map(|meta| FileHandle::new(meta, Arc::new(NoopFilePurger)))
.collect();
let put_output = outputs
.iter()
.find(|file| file.meta_ref().num_rows == 2)
.unwrap();
let delete_output = outputs
.iter()
.find(|file| file.meta_ref().num_rows == 1)
.unwrap();
assert_eq!(
if unknown_put_bound {
None
} else {
NonZeroU64::new(3)
},
put_output.meta_ref().sequence,
);
assert_eq!(NonZeroU64::new(2), delete_output.meta_ref().sequence);
assert!(
outputs
.iter()
.all(|file| !file.meta_ref().preserve_row_sequence)
);
// Cover the original snapshot, each partial success/cancellation state,
// and both completed outputs using real persisted SSTs and merge readers.
for (put_done, delete_done) in [(false, false), (true, false), (false, true), (true, true)]
{
let mut visible = if put_done {
vec![put_output.clone()]
} else {
puts.clone()
};
visible.extend(if delete_done {
vec![delete_output.clone()]
} else {
deletes.clone()
});
let mut reader = CompactionSstReaderBuilder {
metadata: compaction_region.region_metadata.clone(),
sst_layer: region.access_layer.clone(),
cache: engine.cache_manager(),
inputs: &visible,
append_mode: false,
filter_deleted: true,
time_range: None,
merge_mode: MergeMode::LastRow,
}
.build_flat_sst_reader()
.await
.unwrap();
let mut timestamps = Vec::new();
while let Some(batch) = reader.next_batch().await.unwrap() {
timestamps.extend(
batch
.column_by_name("ts")
.unwrap()
.as_primitive::<TimestampMillisecondType>()
.values()
.iter()
.copied(),
);
}
assert_eq!(
vec![10_000],
timestamps,
"put_done={put_done}, delete_done={delete_done}"
);
}
}
/// Build a minimal [`CompactionRegion`] suitable for tests where the
/// [`SstMerger`] is mocked and never touches the access layer.
async fn new_test_compaction_region() -> CompactionRegion {
+290 -1
View File
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use std::collections::HashMap;
use std::collections::{HashMap, HashSet};
use std::ops::Range;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
@@ -137,6 +137,295 @@ async fn collect_stream_ts(stream: SendableRecordBatchStream) -> Vec<i64> {
res
}
/// Flush may collapse versions within one file, but compaction must not promote
/// them, even with no external overlaps or after rewriting a previous output.
#[rstest::rstest]
#[case(false, false)]
#[case(true, false)]
#[case(true, true)]
#[tokio::test]
async fn test_compaction_preserves_flush_sequences_across_generations(
#[values(false, true)] flat_format: bool,
#[case] append_mode: bool,
#[case] preserve_row_sequence: bool,
) {
let mut env = TestEnv::new().await;
let config = MitoConfig {
default_flat_format: flat_format,
min_compaction_interval: Duration::from_secs(3600),
..Default::default()
};
let engine = env.create_engine(config.clone()).await;
let region_id = RegionId::new(1, 1);
let request = CreateRequestBuilder::new()
.insert_option("append_mode", &append_mode.to_string())
.insert_option("preserve_row_sequence", &preserve_row_sequence.to_string())
.insert_option("compaction.type", "twcs")
.insert_option("compaction.twcs.time_window", "1h")
.build();
let table_dir = request.table_dir.clone();
let options = request.options.clone();
let columns = crate::test_util::rows_schema(&request);
engine
.handle_request(region_id, RegionRequest::Create(request))
.await
.unwrap();
async fn read_sequences(engine: &MitoEngine, region_id: RegionId) -> Vec<u64> {
let region = engine.get_region(region_id).unwrap();
let mut sequences = Vec::new();
for file in region
.version()
.ssts
.levels()
.iter()
.flat_map(|l| l.files())
{
let mut reader = region
.access_layer
.read_sst(file.clone())
.build()
.await
.unwrap()
.unwrap();
while let Some(batch) = reader.next_record_batch().await.unwrap() {
sequences.extend_from_slice(
batch
.column(batch.num_columns() - 2)
.as_primitive::<datatypes::arrow::datatypes::UInt64Type>()
.values(),
);
}
}
sequences.sort_unstable();
sequences
}
let mut expected = Vec::new();
for generation in 0..=2 {
// Interleaving time ranges exercise a merge without duplicate rows.
let rows = [generation, generation + 10]
.into_iter()
.flat_map(|ts| build_rows_for_key("a", ts, ts + 1, 0))
.collect();
put_rows(
&engine,
region_id,
Rows {
schema: columns.clone(),
rows,
},
)
.await;
flush(&engine, region_id).await;
let max = (generation * 2 + 2) as u64;
expected.extend(if preserve_row_sequence {
[max - 1, max]
} else {
[max, max]
});
assert_eq!(expected, read_sequences(&engine, region_id).await);
if generation == 0 {
continue;
}
engine
.handle_request(
region_id,
RegionRequest::Compact(RegionCompactRequest {
options: api::v1::region::compact_request::Options::StrictWindow(
api::v1::region::StrictWindow {
window_seconds: 3600,
},
),
..Default::default()
}),
)
.await
.unwrap();
let version = engine.get_region(region_id).unwrap().version();
let files: Vec<_> = version
.ssts
.levels()
.iter()
.flat_map(|l| l.files())
.collect();
assert_eq!(
1,
files.len(),
"generation {generation} must merge all inputs"
);
assert_eq!(
preserve_row_sequence,
files[0].meta_ref().preserve_row_sequence
);
assert_eq!(Some(max), files[0].meta_ref().sequence.map(|s| s.get()));
assert_eq!(
expected,
read_sequences(&engine, region_id).await,
"generation {generation}"
);
}
let engine = env.reopen_engine(engine, config).await;
crate::test_util::reopen_region(&engine, region_id, table_dir, false, options).await;
assert_eq!(expected, read_sequences(&engine, region_id).await);
}
#[tokio::test]
async fn test_partial_compaction_preserves_delete_order_flat() {
assert_partial_compaction_preserves_delete_order(true).await;
}
#[tokio::test]
async fn test_partial_compaction_preserves_delete_order_primary_key() {
assert_partial_compaction_preserves_delete_order(false).await;
}
/// A newer unrelated input must not promote an old Put above an unselected Delete.
/// All files fit in the same one-hour window; no cross-window SST is required.
async fn assert_partial_compaction_preserves_delete_order(flat_format: bool) {
let mut env = TestEnv::new().await;
let region_id = RegionId::new(1, 1);
let (engine, columns) = env_for_manual_compaction(&mut env, region_id, flat_format).await;
put_and_flush(&engine, region_id, &columns, 10..16).await;
let mut deletes = build_rows_for_key("a", 10, 16, 0);
// Distinct keys keep the Delete SST large even after compression, so the
// byte-balance rule excludes it without fabricating FileMeta sizes.
let mut rng = <rand::rngs::StdRng as rand::SeedableRng>::seed_from_u64(0);
for _ in 0..3000 {
deletes.extend(build_rows_for_key(
&format!("{:032x}", rand::Rng::random::<u128>(&mut rng)),
1,
2,
0,
));
}
engine
.handle_request(
region_id,
RegionRequest::Delete(RegionDeleteRequest {
rows: Rows {
schema: columns.clone(),
rows: deletes,
},
hint: None,
partition_expr_version: None,
}),
)
.await
.unwrap();
flush(&engine, region_id).await;
// These newer writes are legitimate, but must not change the version of 10..16.
for i in 10..25 {
put_and_flush(&engine, region_id, &columns, i * 10..i * 10 + 6).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!(files.iter().all(|file| {
file.time_range().0 >= Timestamp::new_millisecond(0)
&& file.time_range().1 < Timestamp::new_millisecond(3_600_000)
}));
let delete_file = files
.iter()
.find(|file| file.meta_ref().num_rows == 3006)
.unwrap();
let delete_id = delete_file.file_id();
let delete_sequence = delete_file.meta_ref().sequence;
let old_ids = files
.iter()
.map(|file| file.file_id())
.collect::<HashSet<_>>();
let before = collect_stream_ts(
engine
.scanner(region_id, ScanRequest::default())
.await
.unwrap()
.scan()
.await
.unwrap(),
)
.await;
let expected = (10..25)
.flat_map(|i| (i * 10..i * 10 + 6).map(|ts| ts * 1000))
.collect::<Vec<i64>>();
assert_eq!(
expected, before,
"the Delete must hide old rows before compaction"
);
compact(&engine, region_id).await;
let scanner = engine
.scanner(region_id, ScanRequest::default())
.await
.unwrap();
let current_ids = scanner.file_ids();
assert_eq!(
2,
current_ids.len(),
"the 16 small inputs must merge, leaving the large Delete SST"
);
assert!(current_ids.contains(&delete_id));
assert_eq!(
1,
current_ids.iter().filter(|id| old_ids.contains(id)).count()
);
let after = collect_stream_ts(scanner.scan().await.unwrap()).await;
let current = region.version();
let output = current
.ssts
.levels()
.iter()
.flat_map(|level| level.files())
.find(|file| !old_ids.contains(&file.file_id()))
.unwrap();
let mut reader = region
.access_layer
.read_sst(output.clone())
.build()
.await
.unwrap()
.unwrap();
let batch = reader.next_record_batch().await.unwrap().unwrap();
let output_sequences = batch
.column(batch.num_columns() - 2)
.as_primitive::<datatypes::arrow::datatypes::UInt64Type>();
let first_output_sequence = output_sequences.value(0);
crate::test_util::reopen_region(&engine, region_id, table_dir, true, HashMap::new()).await;
let reopened = collect_stream_ts(
engine
.scanner(region_id, ScanRequest::default())
.await
.unwrap()
.scan()
.await
.unwrap(),
)
.await;
let resurrected = after
.iter()
.filter(|ts| (10_000..16_000).contains(*ts))
.collect::<Vec<_>>();
assert!(
after == expected && reopened == expected,
"partial compaction changed delete ordering (flat_format={flat_format}, unselected Delete sequence={delete_sequence:?}, first output row sequence={first_output_sequence}); expected rows={}, after rows={}, reopened rows={}, resurrected timestamps={resurrected:?}",
expected.len(),
after.len(),
reopened.len()
);
}
struct CompactionListenerGuard(Option<Arc<CompactionListener>>);
impl CompactionListenerGuard {
+189
View File
@@ -16,6 +16,9 @@
use api::v1::Rows;
use common_recordbatch::RecordBatches;
use datafusion_expr::{col, lit};
use datatypes::arrow::array::AsArray;
use datatypes::arrow::datatypes::UInt64Type;
use store_api::region_engine::RegionEngine;
use store_api::region_request::{RegionCompactRequest, RegionRequest};
use store_api::storage::{RegionId, ScanRequest};
@@ -27,6 +30,192 @@ use crate::test_util::{
delete_rows_schema, flush_region, put_rows, reopen_region, rows_schema,
};
#[rstest::rstest]
#[tokio::test]
async fn test_partial_compaction_preserves_last_row_put_versions(
#[values(false, true)] flat_format: bool,
) {
let mut env = TestEnv::new().await;
let config = MitoConfig {
default_flat_format: flat_format,
min_compaction_interval: std::time::Duration::from_secs(3600),
..Default::default()
};
let engine = env.create_engine(config.clone()).await;
let region_id = RegionId::new(1, 1);
env.get_schema_metadata_manager()
.register_region_table_info(
region_id.table_id(),
"test_table",
"test_catalog",
"test_schema",
None,
env.get_kv_backend(),
)
.await;
let request = CreateRequestBuilder::new()
.field_num(2)
.insert_option("compaction.type", "twcs")
.insert_option("compaction.twcs.time_window", "1h")
.insert_option("merge_mode", "last_row")
.build();
let table_dir = request.table_dir.clone();
let region_opts = request.options.clone();
let schema = rows_schema(&request);
engine
.handle_request(region_id, RegionRequest::Create(request))
.await
.unwrap();
// A(x=1), B(x=2) share a key and timestamp; C is unrelated.
// Promoting A when merging A+C must not hide the unselected B.
let a = build_rows_with_fields("a", &[10], &[(Some(1), None)]);
put_rows(
&engine,
region_id,
Rows {
schema: schema.clone(),
rows: a,
},
)
.await;
flush_region(&engine, region_id, None).await;
let mut b = build_rows_with_fields("a", &[10], &[(Some(2), None)]);
let mut rng = <rand::rngs::StdRng as rand::SeedableRng>::seed_from_u64(0);
for _ in 0..3000 {
b.extend(build_rows_with_fields(
&format!("{:032x}", rand::Rng::random::<u128>(&mut rng)),
&[1],
&[(Some(2), None)],
));
}
put_rows(
&engine,
region_id,
Rows {
schema: schema.clone(),
rows: b,
},
)
.await;
flush_region(&engine, region_id, None).await;
let b_id = engine
.get_region(region_id)
.unwrap()
.version()
.ssts
.levels()
.iter()
.flat_map(|level| level.files())
.find(|file| file.meta_ref().num_rows == 3001)
.unwrap()
.file_id();
let c = build_rows_with_fields("z", &[10], &[(None, Some(3))]);
put_rows(
&engine,
region_id,
Rows {
schema: schema.clone(),
rows: c,
},
)
.await;
flush_region(&engine, region_id, None).await;
// Sixteen small inputs plus the large middle-version file B. The normal
// byte-balanced candidate excludes B.
for i in 0..14 {
let rows = build_rows_with_fields("z", &[100 + i], &[(Some(4), None)]);
put_rows(
&engine,
region_id,
Rows {
schema: schema.clone(),
rows,
},
)
.await;
flush_region(&engine, region_id, None).await;
}
let scan = || ScanRequest {
filters: vec![col("tag_0").eq(lit("a"))],
..Default::default()
};
let before =
RecordBatches::try_collect(engine.scan_to_stream(region_id, scan()).await.unwrap())
.await
.unwrap()
.pretty_print()
.unwrap();
let expected = "\
+-------+---------+---------+---------------------+
| tag_0 | field_0 | field_1 | ts |
+-------+---------+---------+---------------------+
| a | 2.0 | | 1970-01-01T00:00:10 |
+-------+---------+---------+---------------------+";
assert_eq!(expected, before);
engine
.handle_request(
region_id,
RegionRequest::Compact(RegionCompactRequest::default()),
)
.await
.unwrap();
let after = RecordBatches::try_collect(engine.scan_to_stream(region_id, scan()).await.unwrap())
.await
.unwrap()
.pretty_print()
.unwrap();
assert_eq!(expected, after, "flat_format={flat_format}");
let version = engine.get_region(region_id).unwrap().version();
let files: Vec<_> = version
.ssts
.levels()
.iter()
.flat_map(|level| level.files())
.collect();
assert_eq!(2, files.len());
assert!(files.iter().any(|file| file.file_id() == b_id));
assert!(
files
.iter()
.all(|file| !file.meta_ref().preserve_row_sequence)
);
let region = engine.get_region(region_id).unwrap();
let mut output_sequences = std::collections::HashSet::new();
for file in files.iter().filter(|file| file.file_id() != b_id) {
let mut reader = region
.access_layer
.read_sst((*file).clone())
.build()
.await
.unwrap()
.unwrap();
while let Some(batch) = reader.next_record_batch().await.unwrap() {
output_sequences.extend(
batch
.column(batch.num_columns() - 2)
.as_primitive::<UInt64Type>()
.values()
.iter()
.copied(),
);
}
}
assert!(
output_sequences.len() > 1,
"compaction must retain effective input sequences"
);
let engine = env.reopen_engine(engine, config).await;
reopen_region(&engine, region_id, table_dir, false, region_opts).await;
let reopened =
RecordBatches::try_collect(engine.scan_to_stream(region_id, scan()).await.unwrap())
.await
.unwrap()
.pretty_print()
.unwrap();
assert_eq!(expected, reopened);
}
#[tokio::test]
async fn test_merge_mode_write_query() {
test_merge_mode_write_query_with_format(false).await;
+103 -50
View File
@@ -2907,9 +2907,8 @@ async fn test_bulk_write_sequence_not_committed_before_install() {
}
}
/// Non-preserving compaction must keep the physical input sequence below the
/// admission barrier. Otherwise a later row at the barrier can collide with
/// the compacted row when ordinary reads deduplicate overlapping SSTs.
/// Compaction must not allocate a new sequence that can collide with the next
/// write. Preserve the input's effective row sequence and file-level bound.
#[tokio::test]
async fn test_non_preserve_compaction_sequence_collision() {
for flat_format in [false, true] {
@@ -2979,8 +2978,7 @@ async fn test_non_preserve_compaction_sequence_collision_with_format(flat_format
);
// Both current SSTs are compacted through the real picker, merger, writer,
// and manifest update. With committed/flushed sequence 2, the untrusted
// output's admission barrier is exactly 3.
// and manifest update. Both the row version and inherited file bound stay 2.
engine
.handle_request(
region_id,
@@ -3003,14 +3001,13 @@ async fn test_non_preserve_compaction_sequence_collision_with_format(flat_format
assert!(!compacted.meta_ref().preserve_row_sequence);
assert_eq!(1, compacted.num_rows());
assert_eq!(
Some(std::num::NonZeroU64::new(3).unwrap()),
Some(std::num::NonZeroU64::new(2).unwrap()),
compacted.meta_ref().sequence,
"admission barrier is input max 2 plus one"
"compaction must inherit the input bound without allocating a sequence"
);
// Hide FileMeta.sequence so this is a direct physical read, not a reader
// admission-barrier override. The output must contain physical sequence 2,
// not zero and not the output barrier 3.
// override. The output must contain physical sequence 2, not zero or 3.
let mut physical_meta = compacted.meta_ref().clone();
physical_meta.sequence = None;
let mut reader = region
@@ -3037,10 +3034,8 @@ async fn test_non_preserve_compaction_sequence_collision_with_format(flat_format
"after compaction: {after_compaction}"
);
// The next write receives physical sequence 3, equal to the compacted
// output's admission barrier. It must still win because the compacted row
// remains physically at sequence 2; an old zero/barrier encoding would
// collide here and incorrectly retain value 2.
// The next write receives sequence 3 and must win. The old zero/barrier
// encoding promoted the compacted row to 3 and could incorrectly retain 2.
test_util::put_rows(
&engine,
region_id,
@@ -3075,7 +3070,7 @@ async fn test_non_preserve_compaction_sequence_collision_with_format(flat_format
.find(|file| file.meta_ref().level == 1)
.expect("compacted L1 SST");
assert_eq!(
Some(std::num::NonZeroU64::new(3).unwrap()),
Some(std::num::NonZeroU64::new(2).unwrap()),
compacted.meta_ref().sequence
);
@@ -3086,24 +3081,27 @@ async fn test_non_preserve_compaction_sequence_collision_with_format(flat_format
);
}
/// Compaction rewrites a legacy (unmarked) input as an untrusted output with
/// the region-local admission barrier. Its physical rows retain the known
/// input maximum, while exact scans skip it once C reaches that barrier and
/// fail closed while C is below it.
/// A legacy input must retain its admission boundary without gaining row trust.
/// A later flush can advance beyond manifest.committed_sequence, so compaction
/// must inherit the maximum input bound rather than sample that stale frontier.
#[rstest::rstest]
#[tokio::test]
async fn test_compaction_output_non_preserve_not_laundered_from_legacy_input() {
async fn test_compaction_output_non_preserve_not_laundered_from_legacy_input(
#[values(false, true)] flat_format: bool,
#[values(false, true)] flush_after_edit: bool,
) {
let mut env =
TestEnv::with_prefix("test_compaction_output_not_laundered_from_legacy_input").await;
// Suppress automatic edit-triggered compactions. The high TWCS trigger below
// also prevents flush-triggered compaction from consuming the inputs, so the
// explicit Compact below is the only compaction in flight (deterministic).
let engine = env
.create_engine(MitoConfig {
min_compaction_interval: std::time::Duration::from_secs(60 * 60),
schedule_compaction_after_edit: false,
..Default::default()
})
.await;
let config = MitoConfig {
default_flat_format: flat_format,
min_compaction_interval: std::time::Duration::from_secs(60 * 60),
schedule_compaction_after_edit: false,
..Default::default()
};
let engine = env.create_engine(config.clone()).await;
let region_id = RegionId::new(1, 1);
let request = CreateRequestBuilder::new()
@@ -3113,6 +3111,8 @@ async fn test_compaction_output_non_preserve_not_laundered_from_legacy_input() {
.insert_option("compaction.twcs.trigger_file_num", "100")
.build();
let column_schemas = test_util::rows_schema(&request);
let table_dir = request.table_dir.clone();
let region_options = request.options.clone();
engine
.handle_request(region_id, RegionRequest::Create(request))
@@ -3172,6 +3172,30 @@ async fn test_compaction_output_non_preserve_not_laundered_from_legacy_input() {
"seeded file should be unmarked"
);
let mut expected_sequences = vec![1, 2, 3, 4, 5, 6];
let expected_bound = if flush_after_edit {
test_util::put_rows(
&engine,
region_id,
Rows {
schema: column_schemas,
rows: test_util::build_rows(6, 9),
},
)
.await;
test_util::flush_region(&engine, region_id, None).await;
expected_sequences.extend([8, 9, 10]);
10
} else {
7
};
let manifest = region.manifest_ctx.manifest().await;
assert_eq!(Some(7), manifest.committed_sequence);
assert_eq!(
if flush_after_edit { 10 } else { 6 },
manifest.flushed_sequence
);
// Compact: the rewritten output must NOT be laundered back to marked.
engine
.handle_request(
@@ -3198,25 +3222,24 @@ async fn test_compaction_output_non_preserve_not_laundered_from_legacy_input() {
assert_eq!(
1,
outputs.len(),
"two inputs should rewrite into one output file"
"all inputs should rewrite into one output file"
);
assert!(
!outputs[0].meta_ref().preserve_row_sequence,
"legacy input must not be laundered into a marked output"
);
// The compacted output is sequence-less and carries the current region's
// admission barrier rather than any source-domain sequence.
// Inherit the admitted legacy bound and, if present, the newer flush bound.
let barrier = outputs[0]
.meta_ref()
.sequence
.expect("compaction output barrier")
.get();
assert_eq!(8, barrier);
assert_eq!(expected_bound, barrier);
// Reinstalling the legacy input assigns it sequence 7, so the physical
// parquet retains that known input maximum rather than encoding either
// zero or the output admission barrier 8. The manifest marker stays false.
// Reinstalling the legacy input assigns its metadata sequence 7, but the
// existing rows keep sequences 1..=6. Compaction must not rewrite them to
// the metadata bound, including when newer flushed rows join the merge.
let mut sequence_meta = outputs[0].meta_ref().clone();
sequence_meta.sequence = None;
let sequence_handle = FileHandle::new(
@@ -3230,22 +3253,34 @@ async fn test_compaction_output_non_preserve_not_laundered_from_legacy_input() {
.await
.unwrap()
.expect("compaction output reader");
let batch = reader
.next_record_batch()
.await
.unwrap()
.expect("compaction output batch");
let sequence = batch
.column(batch.num_columns() - 2)
.as_any()
.downcast_ref::<datatypes::arrow::array::UInt64Array>()
.expect("sequence column");
assert!(
sequence
.values()
.iter()
.all(|sequence| *sequence == barrier - 1)
let mut sequences = Vec::new();
while let Some(batch) = reader.next_record_batch().await.unwrap() {
let sequence = batch
.column(batch.num_columns() - 2)
.as_any()
.downcast_ref::<datatypes::arrow::array::UInt64Array>()
.expect("sequence column");
sequences.extend_from_slice(sequence.values());
}
sequences.sort_unstable();
assert_eq!(expected_sequences, sequences);
// The inherited bound and untrusted marker must survive manifest recovery.
let engine = env.reopen_engine(engine, config).await;
test_util::reopen_region(&engine, region_id, table_dir, false, region_options).await;
let version = engine.get_region(region_id).unwrap().version();
let files: Vec<_> = version
.ssts
.levels()
.iter()
.flat_map(|level| level.files())
.collect();
assert_eq!(1, files.len());
assert_eq!(
Some(expected_bound),
files[0].meta_ref().sequence.map(|s| s.get())
);
assert!(!files[0].meta_ref().preserve_row_sequence);
// A cursor before the barrier must fail closed.
let err = engine
@@ -3261,9 +3296,27 @@ async fn test_compaction_output_non_preserve_not_laundered_from_legacy_input() {
.await
.err()
.expect("newer barrier must disable exact scanning");
assert!(matches!(err, Error::SequenceRangeUnsupported { .. }));
if flush_after_edit {
// Without exact capability, the lower-bound fence is checked first.
assert!(
matches!(
err,
Error::IncrementalQueryStale {
given_seq: 9,
min_readable_seq: 10,
..
}
),
"unexpected error: {err:?}"
);
} else {
assert!(
matches!(err, Error::SequenceRangeUnsupported { .. }),
"unexpected error: {err:?}"
);
}
// Once C reaches the barrier the sequence-less file is skipped, so exact
// Once C reaches the barrier the untrusted file is skipped, so exact
// capability is restored without attempting row-level filtering.
let scanner = engine
.scanner(
+4 -3
View File
@@ -1836,9 +1836,10 @@ fn pre_filter_mode(append_mode: bool, merge_mode: MergeMode) -> PreFilterMode {
/// contribute a row to `(C, H]`.
///
/// Unmarked local files use `FileMeta.sequence` as an admission barrier rather
/// than a row maximum. Compaction and edit assign it as `committed_sequence + 1`;
/// `C >= barrier` proves Flow has already consumed the entire file, so such a
/// file is excluded before the capability check.
/// than a row maximum. Region edits allocate this barrier; compaction inherits
/// the maximum input bound without assigning a new barrier. `C >= barrier`
/// proves Flow has already consumed the entire file, so such a file is excluded
/// before the capability check. An unknown bound cannot prove this exclusion.
///
/// A foreign file is different: the parquet reader virtualizes every row to its
/// target-local `FileMeta.sequence`. Consequently, a present sequence is the
+2 -1
View File
@@ -121,7 +121,8 @@ pub struct RegionOptions {
/// the configured size; zero disables both limits.
#[serde(skip_serializing_if = "Option::is_none")]
pub write_buffer_size: Option<ReadableSize>,
/// Whether to preserve per-row sequence numbers through flush and compaction.
/// Whether to preserve original per-row sequence numbers when flushing.
/// Compaction always retains effective input sequences, regardless of this option.
///
/// Only meaningful for append-only tables (`append_mode = true`): when enabled,
/// every row keeps its exact sequence number in memtables, flushed SSTs and
+8 -3
View File
@@ -243,10 +243,13 @@ pub struct FileMeta {
/// the default value `0` doesn't means the file doesn't contains any rows,
/// but instead means the number of rows is unknown.
pub num_row_groups: u64,
/// Sequence in this file.
/// File-level sequence bound or admission marker in the target region.
///
/// This sequence is the only sequence in this file. And it's retrieved from the max
/// sequence of the rows on generating this file.
/// Flush records the maximum input row sequence. Compaction outputs inherit
/// the maximum input bound, or remain unknown if any input bound is unknown,
/// independently of per-row sequence trust. This does not imply that every
/// physical row has this sequence.
/// Readers also use it to normalize foreign and legacy all-zero files.
pub sequence: Option<NonZeroU64>,
/// Partition expression from the region metadata when the file is created.
///
@@ -284,6 +287,8 @@ pub struct FileMeta {
pub primary_key_max: Option<Bytes>,
/// Whether the file preserves per-row sequence numbers usable for exact
/// row-level sequence filtering.
/// Merely retaining physical input sequences during compaction does not
/// restore this capability for untrusted inputs.
#[serde(default, skip_serializing_if = "is_false")]
pub preserve_row_sequence: bool,
}