diff --git a/src/mito2/src/access_layer.rs b/src/mito2/src/access_layer.rs index 14c629a01b..ac18f838db 100644 --- a/src/mito2/src/access_layer.rs +++ b/src/mito2/src/access_layer.rs @@ -568,6 +568,8 @@ pub struct SstWriteRequest { pub cache_manager: CacheManagerRef, #[allow(dead_code)] pub storage: Option, + /// 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, pub sst_write_format: FormatType, diff --git a/src/mito2/src/compaction/compactor.rs b/src/mito2/src/compaction/compactor.rs index fb2172b487..2e59aae348 100644 --- a/src/mito2/src/compaction/compactor.rs +++ b/src/mito2/src/compaction/compactor.rs @@ -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, + 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 { let mut max: Option = 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::>(); @@ -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, 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(®ion, 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::() + .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 { diff --git a/src/mito2/src/engine/compaction_test.rs b/src/mito2/src/engine/compaction_test.rs index e427f3f34f..80530876ce 100644 --- a/src/mito2/src/engine/compaction_test.rs +++ b/src/mito2/src/engine/compaction_test.rs @@ -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 { 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 { + 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::() + .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 = ::seed_from_u64(0); + for _ in 0..3000 { + deletes.extend(build_rows_for_key( + &format!("{:032x}", rand::Rng::random::(&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::>(); + 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::>(); + 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::>(); + 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::(); + 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::>(); + 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>); impl CompactionListenerGuard { diff --git a/src/mito2/src/engine/merge_mode_test.rs b/src/mito2/src/engine/merge_mode_test.rs index afb984e77b..5430e6b246 100644 --- a/src/mito2/src/engine/merge_mode_test.rs +++ b/src/mito2/src/engine/merge_mode_test.rs @@ -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 = ::seed_from_u64(0); + for _ in 0..3000 { + b.extend(build_rows_with_fields( + &format!("{:032x}", rand::Rng::random::(&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::() + .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; diff --git a/src/mito2/src/engine/scan_test.rs b/src/mito2/src/engine/scan_test.rs index 8258792c35..eedcee031d 100644 --- a/src/mito2/src/engine/scan_test.rs +++ b/src/mito2/src/engine/scan_test.rs @@ -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::() - .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::() + .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( diff --git a/src/mito2/src/read/scan_region.rs b/src/mito2/src/read/scan_region.rs index 073d1635cf..8154d73b3b 100644 --- a/src/mito2/src/read/scan_region.rs +++ b/src/mito2/src/read/scan_region.rs @@ -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 diff --git a/src/mito2/src/region/options.rs b/src/mito2/src/region/options.rs index a2a165c12a..fe603f9d1b 100644 --- a/src/mito2/src/region/options.rs +++ b/src/mito2/src/region/options.rs @@ -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, - /// 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 diff --git a/src/mito2/src/sst/file.rs b/src/mito2/src/sst/file.rs index b604275ce9..9702477b07 100644 --- a/src/mito2/src/sst/file.rs +++ b/src/mito2/src/sst/file.rs @@ -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, /// Partition expression from the region metadata when the file is created. /// @@ -284,6 +287,8 @@ pub struct FileMeta { pub primary_key_max: Option, /// 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, }