diff --git a/java/src/main/java/org/lance/compaction/CompactionOptions.java b/java/src/main/java/org/lance/compaction/CompactionOptions.java index d83e618da5b..3b0bbfbc629 100644 --- a/java/src/main/java/org/lance/compaction/CompactionOptions.java +++ b/java/src/main/java/org/lance/compaction/CompactionOptions.java @@ -334,9 +334,9 @@ public Builder withBinaryCopyReadBatchBytes(long binaryCopyReadBatchBytes) { } /** - * Maximum number of source fragments to compact in a single run. Tasks are included until - * adding the next task would exceed this limit, allowing for incremental compaction. Fragments - * are processed oldest first. + * Maximum number of fragments whose identities may change in a single run. This includes + * compacted source fragments and following fragments relabeled to preserve row order. The + * planner selects candidates from a suffix within this limit. * * @throws IllegalArgumentException if {@code maxSourceFragments} is not positive */ @@ -384,8 +384,9 @@ public Builder withDataStorageVersion(DataStorageVersion version) { /** * Fragment IDs to exclude from compaction planning. Excluded fragments remain unchanged and act - * as boundaries, so fragments on opposite sides are not combined into the same task. Duplicate - * and unknown IDs are ignored. + * as boundaries, so fragments on opposite sides are not combined into the same task. To + * preserve row order, only candidates after the last present excluded fragment are eligible. + * Duplicate and unknown IDs are ignored. * * @throws IllegalArgumentException if an ID is negative or exceeds the unsigned 32-bit range */ diff --git a/python/python/lance/dataset.py b/python/python/lance/dataset.py index 9013e35781a..42ab3b67b1e 100644 --- a/python/python/lance/dataset.py +++ b/python/python/lance/dataset.py @@ -7696,11 +7696,11 @@ def compact_files( Controls how much data is read at once when performing binary copy. Defaults to 16MB. max_source_fragments: int, optional - Maximum number of source fragments to compact in a single run. - Compaction tasks are included until adding the next task would - exceed this limit, allowing compaction to proceed incrementally. - Fragments are processed oldest first. If not specified, uses the - manifest config value, or applies no limit. + Maximum number of fragments whose identities may change in a + single run. This includes compacted source fragments and following + fragments relabeled to preserve row order. The planner selects + candidates from a suffix within this limit. If not specified, uses + the manifest config value, or applies no limit. max_source_rows: int, optional Maximum number of source rows to compact in a single run. Rows are counted as live rows (physical rows minus soft-deleted rows). @@ -7717,7 +7717,9 @@ def compact_files( Fragment IDs to exclude from compaction planning. Excluded fragments remain unchanged and act as boundaries, so fragments on opposite sides are not combined into the same compaction task. - Duplicate and unknown IDs are ignored. + To preserve row order, only candidates after the last present + excluded fragment are eligible. Duplicate and unknown IDs are + ignored. data_storage_version: str, optional Output data file version, such as "2.2", "stable", or "next". Uses the compaction config target when set, otherwise the dataset's diff --git a/python/python/lance/optimize.py b/python/python/lance/optimize.py index ab18255c853..080ca1b0f1d 100644 --- a/python/python/lance/optimize.py +++ b/python/python/lance/optimize.py @@ -91,10 +91,10 @@ class CompactionOptions(TypedDict, total=False): """ max_source_fragments: Optional[int] """ - Maximum number of source fragments to compact in a single run. Tasks - are included until adding the next task would exceed this limit, - allowing for incremental compaction (e.g., compact 20 fragments at a - time). Fragments are processed oldest first. + Maximum number of fragments whose identities may change in a single run. + This includes compacted source fragments and following fragments relabeled + to preserve row order. The planner selects candidates from a suffix within + this limit. (default: None, no limit) """ max_source_rows: Optional[int] @@ -118,8 +118,9 @@ class CompactionOptions(TypedDict, total=False): """ Fragment IDs to exclude from compaction planning. Excluded fragments remain unchanged and act as boundaries, so fragments on opposite sides - are not combined into the same task. Duplicate and unknown IDs are - ignored. (default: None) + are not combined into the same task. To preserve row order, only candidates + after the last present excluded fragment are eligible. Duplicate and + unknown IDs are ignored. (default: None) """ data_storage_version: Optional[str] """ diff --git a/rust/lance-table/src/transaction/index_maintenance.rs b/rust/lance-table/src/transaction/index_maintenance.rs index 6c43eb58e51..96a7caf34ac 100644 --- a/rust/lance-table/src/transaction/index_maintenance.rs +++ b/rust/lance-table/src/transaction/index_maintenance.rs @@ -476,6 +476,21 @@ impl Transaction { groups: &[RewriteGroup], ) { for group in groups { + // Order-preserving partial compaction represents untouched trailing + // fragments as exact metadata copies with a new ID. Their overlays + // are still live and must not invalidate index coverage. + if let ([old_fragment], [new_fragment]) = ( + group.old_fragments.as_slice(), + group.new_fragments.as_slice(), + ) { + let mut relabeled_fragment = new_fragment.clone(); + relabeled_fragment.id = old_fragment.id; + relabeled_fragment.deletion_file = old_fragment.deletion_file.clone(); + if relabeled_fragment == *old_fragment { + continue; + } + } + // field id -> newest overlay committed_version supplying that field let mut overlaid_field_versions: HashMap = HashMap::new(); for old_frag in &group.old_fragments { diff --git a/rust/lance-table/src/transaction/manifest_build.rs b/rust/lance-table/src/transaction/manifest_build.rs index 5320e8f1543..e792f78fe9f 100644 --- a/rust/lance-table/src/transaction/manifest_build.rs +++ b/rust/lance-table/src/transaction/manifest_build.rs @@ -3202,6 +3202,114 @@ mod tests { assert_eq!(last_updated_at_versions(&new_manifest, 2), vec![2; 42]); } + #[test] + fn test_rewrite_build_manifest_keeps_fragment_ids_sorted() { + let manifest = sample_manifest_with_fragments(0..4); + let transaction = Transaction::new( + manifest.version, + Operation::Rewrite { + groups: vec![RewriteGroup { + old_fragments: vec![Fragment::new(1), Fragment::new(2)], + new_fragments: vec![Fragment::new(0)], + }], + rewritten_indices: vec![], + frag_reuse_index: None, + }, + None, + ); + + let (rewritten, _) = transaction + .build_manifest(Some(&manifest), vec![], "txn", &default_build_config()) + .unwrap(); + + let fragment_ids = rewritten + .fragments + .iter() + .map(|fragment| fragment.id) + .collect::>(); + assert_eq!(fragment_ids, vec![0, 3, 4]); + } + + #[test] + fn test_rewrite_after_row_adding_update_keeps_fragment_ids_sorted() { + let manifest = sample_manifest_with_fragments(0..3); + let update = Transaction::new( + manifest.version, + Operation::Update { + removed_fragment_ids: vec![], + updated_fragments: vec![], + new_fragments: vec![Fragment::new(0)], + fields_modified: vec![], + compacted_sstables: vec![], + fields_for_preserving_frag_bitmap: vec![], + update_mode: Some(UpdateMode::RewriteRows), + inserted_rows_filter: None, + updated_fragment_offsets: None, + }, + None, + ); + let (updated, _) = update + .build_manifest( + Some(&manifest), + vec![], + "update-txn", + &default_build_config(), + ) + .unwrap(); + + let rewrite = Transaction::new( + manifest.version, + Operation::Rewrite { + groups: vec![RewriteGroup { + old_fragments: vec![Fragment::new(0)], + new_fragments: vec![Fragment::new(0)], + }], + rewritten_indices: vec![], + frag_reuse_index: None, + }, + None, + ); + let (rewritten, _) = rewrite + .build_manifest( + Some(&updated), + vec![], + "rewrite-txn", + &default_build_config(), + ) + .unwrap(); + + let fragment_ids = rewritten + .fragments + .iter() + .map(|fragment| fragment.id) + .collect::>(); + assert_eq!(fragment_ids, vec![1, 2, 3, 4]); + } + + #[test] + fn test_merge_build_manifest_keeps_fragment_ids_sorted() { + let manifest = sample_manifest_with_fragments(0..3); + let fragments = manifest.fragments.iter().rev().cloned().collect::>(); + let operation = Operation::Merge { + fragments, + schema: manifest.schema.clone(), + preserves_nullability: true, + }; + validate_operation(Some(&manifest), &operation).unwrap(); + let transaction = Transaction::new(manifest.version, operation, None); + + let (merged, _) = transaction + .build_manifest(Some(&manifest), vec![], "txn", &default_build_config()) + .unwrap(); + + let fragment_ids = merged + .fragments + .iter() + .map(|fragment| fragment.id) + .collect::>(); + assert_eq!(fragment_ids, vec![0, 1, 2]); + } + #[test] fn test_remove_tombstoned_data_files() { // Create a fragment with mixed data files: some normal, some fully tombstoned diff --git a/rust/lance/src/dataset/index/frag_reuse.rs b/rust/lance/src/dataset/index/frag_reuse.rs index 251c9d76819..308258e6baa 100644 --- a/rust/lance/src/dataset/index/frag_reuse.rs +++ b/rust/lance/src/dataset/index/frag_reuse.rs @@ -725,7 +725,7 @@ mod tests { RowAddress::from(before[&1250]).fragment_id(), RowAddress::from(before[&2250]).fragment_id(), ]; - let untouched_addr = before[&5000]; + let trailing_addr = before[&5000]; // Offset 0 of the first rewritten fragment is i=1000, deleted above. let deleted_addr = u64::from(RowAddress::new_from_parts(rewritten_frags[0], 0)); @@ -777,11 +777,13 @@ mod tests { } assert_eq!(remap.get(deleted_addr), Some(None)); assert_eq!(frag_reuse_index.remap_row_id(deleted_addr), None); - assert_eq!(remap.get(untouched_addr), None); - assert_eq!(after[&5000], untouched_addr); + // The trailing fragment keeps its data but receives a new ID so the + // ID-sorted manifest preserves row order after the earlier rewrites. + assert_ne!(after[&5000], trailing_addr); + assert_eq!(remap.get(trailing_addr), Some(Some(after[&5000]))); assert_eq!( - frag_reuse_index.remap_row_id(untouched_addr), - Some(untouched_addr) + frag_reuse_index.remap_row_id(trailing_addr), + Some(after[&5000]) ); let pre_compaction = dataset diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index 55e930bc160..623e3c34ab9 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -128,6 +128,7 @@ use lance_core::datatypes::{ BLOB_V2_LOGICAL_FIELDS, BLOB_V2_LOGICAL_TYPE, BlobHandling, BlobKind, BlobV2Layout, Field as LanceField, }; +use lance_core::utils::address::RowAddress; use lance_core::utils::tokio::get_num_compute_intensive_cpus; use lance_core::utils::tracing::{DATASET_COMPACTING_EVENT, TRACE_DATASET_EVENTS}; use lance_file::version::{ConcreteFileVersion, LanceFileVersion}; @@ -149,6 +150,7 @@ pub mod remapping; use crate::index::frag_reuse::build_new_frag_reuse_index; use crate::io::deletion::read_dataset_deletion_file; +use lance_table::io::deletion::deletion_file_path; pub use remapping::{IgnoreRemap, IndexRemapper, IndexRemapperOptions, RemappedIndex}; /// Controls how data is rewritten during compaction. @@ -278,10 +280,11 @@ pub struct CompactionOptions { /// Controls how much data is read at once when performing binary copy. /// Defaults to 16MB (16 * 1024 * 1024). pub binary_copy_read_batch_bytes: Option, - /// Maximum number of source fragments to compact in a single run. When set, - /// tasks are included in the plan until adding the next task would exceed - /// this limit. This allows for incremental compaction (e.g., compact 20 - /// fragments at a time). + /// Maximum number of fragments whose identities may change in a single run. + /// This includes compacted source fragments and any following fragments that + /// must be relabeled to preserve row order. The planner selects candidates + /// from a suffix within this limit, allowing bounded incremental compaction + /// without unbounded index remapping or metadata churn. /// Defaults to `None` (no limit, all eligible fragments are compacted). pub max_source_fragments: Option, /// Maximum number of source rows to compact in a single run. Rows are @@ -303,7 +306,9 @@ pub struct CompactionOptions { /// /// Excluded fragments act as boundaries between adjacent compaction candidates, /// so fragments on opposite sides of an exclusion are never combined into the - /// same task. IDs that are duplicated or absent from the dataset are ignored. + /// same task. To preserve row order without changing an excluded fragment's + /// identity, only candidates after the last present excluded fragment are + /// eligible. IDs that are duplicated or absent from the dataset are ignored. /// Defaults to an empty list. #[serde(default)] pub excluded_fragment_ids: Vec, @@ -829,12 +834,10 @@ impl CompactionPlanner for DefaultCompactionPlanner { )); } - // get_fragments should be returning fragments in sorted order (by id) - // and fragment ids should be unique + // Manifest order follows fragment ID order, which is also row order. let fragments = dataset.get_fragments(); - debug_assert!( - fragments.windows(2).all(|w| w[0].id() < w[1].id()), + fragments.windows(2).all(|pair| pair[0].id() < pair[1].id()), "fragments in manifest are not sorted" ); // Without stable row ids a rewrite moves every row address, so the @@ -863,8 +866,19 @@ impl CompactionPlanner for DefaultCompactionPlanner { } excluded_fragment_ids |= unremappable; } + // Compaction replacements use fresh IDs, and stable manifests are ID-sorted. + // Rewriting at or before an excluded fragment would therefore require + // relabeling that fragment to keep its rows in place. Restrict planning to + // the suffix after the last caller or index exclusion so excluded identities + // stay unchanged. + let exclusion_suffix_start = fragments + .iter() + .rposition(|fragment| { + u32::try_from(fragment.id()) + .is_ok_and(|fragment_id| excluded_fragment_ids.contains(fragment_id)) + }) + .map_or(0, |position| position + 1); let excluded_fragment_ids = &excluded_fragment_ids; - let mut fragment_metrics = futures::stream::iter(fragments) .map(|fragment| async move { if u32::try_from(fragment.id()) @@ -995,6 +1009,29 @@ impl CompactionPlanner for DefaultCompactionPlanner { candidate_bins.push(bin); } + let budget_suffix_start = self.options.max_source_fragments.map_or(0, |max_frags| { + dataset.manifest.fragments.len().saturating_sub(max_frags) + }); + let suffix_start = exclusion_suffix_start.max(budget_suffix_start); + if suffix_start > 0 { + candidate_bins = candidate_bins + .into_iter() + .filter_map(|mut bin| { + if bin.pos_range.end <= suffix_start { + return None; + } + let skip = suffix_start.saturating_sub(bin.pos_range.start); + if skip > 0 { + bin.fragments.drain(..skip); + bin.candidacy.drain(..skip); + bin.row_counts.drain(..skip); + bin.pos_range.start += skip; + } + (!bin.is_noop()).then_some(bin) + }) + .collect(); + } + let all_tasks: Vec<(TaskData, usize)> = candidate_bins .into_iter() .filter(|bin| !bin.is_noop()) @@ -1171,14 +1208,14 @@ fn pins_dead_user_columns_for_lineage( }) } -/// Truncates a planned task list to the configured per-run source budgets -/// (`max_source_fragments`, `max_source_rows`, `max_source_bytes`). +/// Truncates a planned task list to the configured per-run row and byte budgets. /// -/// All configured budgets apply together: tasks are kept, in order, until -/// adding the next task would exceed any one of them. The budgets are hard -/// upper bounds, so if the first task already exceeds one of them the -/// returned plan is empty and a warning is logged, since compaction would -/// otherwise stall silently. +/// Both configured budgets apply together: tasks are kept, in order, until adding +/// the next task would exceed either one. The budgets are hard upper bounds, so +/// if the first task already exceeds one of them the returned plan is empty and a +/// warning is logged, since compaction would otherwise stall silently. +/// `max_source_fragments` is enforced separately against the entire suffix whose +/// fragment identities may change, including metadata-only relabels. /// /// Each task is paired with the number of live rows in its source fragments. fn limit_tasks_to_source_budget( @@ -1186,10 +1223,7 @@ fn limit_tasks_to_source_budget( schema: &lance_core::datatypes::Schema, all_tasks: Vec<(TaskData, usize)>, ) -> Result> { - if options.max_source_fragments.is_none() - && options.max_source_rows.is_none() - && options.max_source_bytes.is_none() - { + if options.max_source_rows.is_none() && options.max_source_bytes.is_none() { return Ok(all_tasks.into_iter().map(|(task, _)| task).collect()); } @@ -1203,21 +1237,16 @@ fn limit_tasks_to_source_budget( }; let num_candidate_tasks = all_tasks.len(); - let mut total_fragments = 0_usize; let mut total_rows = 0_usize; let mut total_bytes = 0_u64; let mut tasks = Vec::with_capacity(all_tasks.len()); for (task, live_rows) in all_tasks { - total_fragments += task.fragments.len(); total_rows = total_rows.saturating_add(live_rows); if options.max_source_bytes.is_some() { total_bytes = total_bytes.saturating_add(task_source_bytes(&task, &schema_field_ids)?); } - let over_budget = options - .max_source_fragments - .is_some_and(|max| total_fragments > max) - || options.max_source_rows.is_some_and(|max| total_rows > max) + let over_budget = options.max_source_rows.is_some_and(|max| total_rows > max) || options .max_source_bytes .is_some_and(|max| total_bytes > max); @@ -1231,12 +1260,9 @@ fn limit_tasks_to_source_budget( if tasks.is_empty() && num_candidate_tasks > 0 { warn!( "Compaction plan is empty: the first of {} candidate tasks already exceeds a source \ - budget (max_source_fragments={:?}, max_source_rows={:?}, max_source_bytes={:?}); \ + budget (max_source_rows={:?}, max_source_bytes={:?}); \ compaction cannot make progress until the budget is raised", - num_candidate_tasks, - options.max_source_fragments, - options.max_source_rows, - options.max_source_bytes + num_candidate_tasks, options.max_source_rows, options.max_source_bytes ); } @@ -2437,6 +2463,326 @@ async fn reserve_fragment_ids( Ok(()) } +fn serialize_fragment_row_addrs(fragment: &Fragment) -> Result> { + let physical_rows = fragment.physical_rows.ok_or_else(|| { + Error::invalid_input(format!( + "cannot preserve compaction order because trailing fragment {} is missing physical_rows", + fragment.id + )) + })?; + let fragment_id = u32::try_from(fragment.id).map_err(|_| { + Error::invalid_input(format!( + "cannot preserve compaction order because trailing fragment id {} exceeds the row-address limit", + fragment.id + )) + })?; + let start = u64::from(RowAddress::first_row(fragment_id)); + let physical_rows = u64::try_from(physical_rows).map_err(|_| { + Error::invalid_input(format!( + "cannot preserve compaction order because fragment {} physical_rows does not fit in a row address", + fragment.id + )) + })?; + let end = start.checked_add(physical_rows).ok_or_else(|| { + Error::invalid_input(format!( + "cannot preserve compaction order because fragment {} row-address range overflows", + fragment.id + )) + })?; + let mut row_addrs = RoaringTreemap::new(); + row_addrs.insert_range(start..end); + let mut serialized = Vec::with_capacity(row_addrs.serialized_size()); + row_addrs.serialize_into(&mut serialized)?; + Ok(serialized) +} + +fn same_fragment_with_different_id(left: &Fragment, right: &Fragment) -> bool { + let mut right = right.clone(); + right.id = left.id; + right == *left +} + +fn rebase_task_source_fragments( + manifest_fragments: &[Fragment], + task: &mut RewriteResult, + current_version: u64, +) -> Result<()> { + let mut relabeled_ids = HashMap::new(); + for source_fragment in &mut task.original_fragments { + if manifest_fragments + .iter() + .any(|fragment| fragment.id == source_fragment.id) + { + continue; + } + + let mut matching_fragments = manifest_fragments + .iter() + .filter(|fragment| same_fragment_with_different_id(source_fragment, fragment)); + let Some(matching_fragment) = matching_fragments.next() else { + return Err(Error::retryable_commit_conflict_source( + current_version, + format!( + "compaction source fragment {} is no longer present in the current manifest", + source_fragment.id + ) + .into(), + )); + }; + if matching_fragments.next().is_some() { + return Err(Error::invalid_input(format!( + "compaction source fragment {} matches multiple relabeled fragments", + source_fragment.id + ))); + } + let old_id = u32::try_from(source_fragment.id).map_err(|_| { + Error::invalid_input(format!( + "compaction source fragment id {} exceeds the row-address limit", + source_fragment.id + )) + })?; + let new_id = u32::try_from(matching_fragment.id).map_err(|_| { + Error::invalid_input(format!( + "relabeled compaction source fragment id {} exceeds the row-address limit", + matching_fragment.id + )) + })?; + relabeled_ids.insert(old_id, new_id); + *source_fragment = matching_fragment.clone(); + } + + if relabeled_ids.is_empty() { + return Ok(()); + } + let Some(serialized_row_addrs) = task.row_addrs.as_mut() else { + return Ok(()); + }; + let row_addrs = RoaringTreemap::deserialize_from(&mut Cursor::new(&*serialized_row_addrs))?; + let mut rebased_row_addrs = RoaringTreemap::new(); + for row_addr in row_addrs { + let row_addr = RowAddress::from(row_addr); + let fragment_id = relabeled_ids + .get(&row_addr.fragment_id()) + .copied() + .unwrap_or_else(|| row_addr.fragment_id()); + rebased_row_addrs.insert(u64::from(RowAddress::new_from_parts( + fragment_id, + row_addr.row_offset(), + ))); + } + serialized_row_addrs.clear(); + serialized_row_addrs.reserve(rebased_row_addrs.serialized_size()); + rebased_row_addrs.serialize_into(serialized_row_addrs)?; + Ok(()) +} + +/// Complete a partial rewrite with metadata-only replacements for every +/// following fragment. Replacement IDs are reserved in this returned order, so +/// the manifest remains ID-sorted without moving compacted rows behind an +/// untouched suffix. The trailing replacements keep their existing data files +/// and only change physical fragment identity. +fn complete_rewrite_suffix( + manifest_fragments: &[Fragment], + completed_tasks: Vec, + max_affected_fragments: Option, + read_version: u64, + current_version: u64, +) -> Result> { + let fragment_positions = manifest_fragments + .iter() + .enumerate() + .map(|(position, fragment)| (fragment.id, position)) + .collect::>(); + let mut positioned_tasks = Vec::with_capacity(completed_tasks.len()); + + for mut task in completed_tasks { + rebase_task_source_fragments(manifest_fragments, &mut task, current_version)?; + let first_fragment = task.original_fragments.first().ok_or_else(|| { + Error::invalid_input( + "compaction results must replace at least one original fragment".to_string(), + ) + })?; + let start = fragment_positions.get(&first_fragment.id).copied().ok_or_else(|| { + Error::invalid_input(format!( + "compaction result references fragment {} which is not present in the current manifest", + first_fragment.id + )) + })?; + let end = start + .checked_add(task.original_fragments.len()) + .ok_or_else(|| Error::invalid_input("compaction source range overflow".to_string()))?; + let current_range = manifest_fragments.get(start..end).ok_or_else(|| { + Error::invalid_input(format!( + "compaction result starting at fragment {} extends beyond the current manifest", + first_fragment.id + )) + })?; + if let Some((expected, actual)) = task + .original_fragments + .iter() + .zip(current_range) + .find(|(expected, actual)| expected.id != actual.id) + { + return Err(Error::invalid_input(format!( + "compaction source fragments must be contiguous and in manifest order: expected fragment {} but found {}", + expected.id, actual.id + ))); + } + positioned_tasks.push((start, end, task)); + } + + positioned_tasks.sort_unstable_by_key(|(start, _, _)| *start); + if let Some(tasks) = positioned_tasks + .windows(2) + .find(|tasks| tasks[0].1 > tasks[1].0) + { + return Err(Error::invalid_input(format!( + "compaction results overlap at manifest position {}", + tasks[1].0 + ))); + } + + let Some(first_position) = positioned_tasks.first().map(|(start, _, _)| *start) else { + return Ok(Vec::new()); + }; + let affected_fragments = manifest_fragments.len() - first_position; + if let Some(max_frags) = max_affected_fragments + && affected_fragments > max_frags + { + let message = format!( + "compaction would change {affected_fragments} fragment identities, exceeding max_source_fragments={max_frags}" + ); + if read_version < current_version { + return Err(Error::retryable_commit_conflict_source( + current_version, + message.into(), + )); + } + return Err(Error::invalid_input(message)); + } + let mut completed_suffix = Vec::with_capacity( + positioned_tasks.len() + manifest_fragments.len().saturating_sub(first_position), + ); + let mut positioned_tasks = positioned_tasks.into_iter().peekable(); + let mut position = first_position; + while position < manifest_fragments.len() { + if positioned_tasks + .peek() + .is_some_and(|(start, _, _)| *start == position) + { + let Some((_, end, task)) = positioned_tasks.next() else { + return Err(Error::internal( + "compaction task ordering changed while completing rewrite suffix", + )); + }; + completed_suffix.push(task); + position = end; + continue; + } + + let fragment = manifest_fragments[position].clone(); + completed_suffix.push(RewriteResult { + metrics: CompactionMetrics::default(), + new_fragments: vec![fragment.clone()], + read_version, + original_fragments: vec![fragment], + row_addrs: None, + }); + position += 1; + } + + if let Some((start, _, _)) = positioned_tasks.next() { + return Err(Error::invalid_input(format!( + "compaction result starts at manifest position {start}, before the completed suffix position {position}" + ))); + } + Ok(completed_suffix) +} + +fn is_metadata_only_relabel(task: &RewriteResult) -> bool { + let ([old_fragment], [new_fragment]) = ( + task.original_fragments.as_slice(), + task.new_fragments.as_slice(), + ) else { + return false; + }; + let mut relabeled_fragment = new_fragment.clone(); + relabeled_fragment.id = old_fragment.id; + relabeled_fragment == *old_fragment +} + +async fn prepare_metadata_only_relabels( + dataset: &Dataset, + completed_tasks: &mut [RewriteResult], +) -> Result<()> { + let relabel_task_indices = completed_tasks + .iter() + .enumerate() + .filter_map(|(index, task)| is_metadata_only_relabel(task).then_some(index)) + .collect::>(); + if relabel_task_indices.is_empty() { + return Ok(()); + } + + let relabeled_fragments = relabel_task_indices + .iter() + .map(|index| completed_tasks[*index].original_fragments[0].clone()) + .collect::>(); + // Manifests written before writer-version metadata may contain missing or + // inaccurate row counts. Recompute those bounded suffix fragments before + // their physical addresses are serialized. + let recompute_stats = dataset.manifest.writer_version.is_none(); + let migrated_fragments = + migrate_fragments(dataset, &relabeled_fragments, recompute_stats).await?; + if migrated_fragments.len() != relabel_task_indices.len() { + return Err(Error::invalid_input( + "cannot preserve compaction order because a trailing fragment became empty during metadata migration" + .to_string(), + )); + } + + for (task_index, migrated_fragment) in relabel_task_indices.into_iter().zip(migrated_fragments) + { + let row_addrs = serialize_fragment_row_addrs(&migrated_fragment)?; + let task = &mut completed_tasks[task_index]; + task.original_fragments[0] = migrated_fragment.clone(); + task.new_fragments[0] = migrated_fragment; + task.row_addrs = Some(row_addrs); + } + Ok(()) +} + +async fn copy_relabeled_deletion_files( + dataset: &Dataset, + completed_tasks: &mut [RewriteResult], +) -> Result<()> { + for task in completed_tasks { + if !is_metadata_only_relabel(task) { + continue; + } + let Some(old_fragment) = task.original_fragments.first() else { + return Err(Error::internal( + "metadata-only compaction relabel is missing its source fragment", + )); + }; + let Some(old_deletion_file) = old_fragment.deletion_file.as_ref() else { + continue; + }; + let Some(new_fragment) = task.new_fragments.first_mut() else { + return Err(Error::internal( + "metadata-only compaction relabel is missing its replacement fragment", + )); + }; + let deletion_base = dataset.dataset_dir_for_deletion(old_deletion_file)?; + let deletion_store = dataset.object_store_for_deletion(old_deletion_file).await?; + let old_path = deletion_file_path(&deletion_base, old_fragment.id, old_deletion_file); + let new_path = deletion_file_path(&deletion_base, new_fragment.id, old_deletion_file); + let deletion_data = deletion_store.read_one_all(&old_path).await?; + deletion_store.put(&new_path, &deletion_data).await?; + } + Ok(()) +} + /// Rewrite the files in a single task. /// /// This assumes that the dataset is the correct read version to be compacted. @@ -3153,6 +3499,10 @@ fn append_row_lineage_columns( /// they can be omitted and the successful tasks can be committed. However, once /// some of the tasks have been committed, the remainder of the tasks will not /// be able to be committed and should be considered cancelled. +/// +/// Completed tasks are ordered by their source positions. Partial and gapped +/// result sets also relabel the untouched trailing fragments so replacement IDs +/// remain monotonic without changing row order. pub async fn commit_compaction( dataset: &mut Dataset, completed_tasks: Vec, @@ -3180,26 +3530,6 @@ pub async fn commit_compaction( )); } - // Before anything is written or committed. The condition is the planner's, - // not `has_address_style`: a dataset whose only index is one this build - // cannot read captures no row addresses at all, which is exactly the plan - // that has to be refused here. - if !dataset.manifest.uses_stable_row_ids() && !options.defer_index_remap { - reject_unremappable_rewrite(dataset, &completed_tasks).await?; - } - - let has_address_style = completed_tasks.iter().any(|t| t.row_addrs.is_some()); - // Address-style results require immediate index remapping unless it is deferred. - let needs_remapping = - !dataset.manifest.uses_stable_row_ids() && !options.defer_index_remap && has_address_style; - - // Confirm there is a remapper before materializing the potentially very large row address map. - let index_remapper = if needs_remapping { - remap_options.create_remapper(dataset).await? - } else { - None - }; - // Determine the earliest version at which compaction tasks were planned/executed. // // In distributed mode (e.g. Spark) the caller opens *two separate* Dataset @@ -3219,29 +3549,67 @@ pub async fn commit_compaction( .min() .unwrap_or(dataset.manifest.version); - let mut completed_tasks = completed_tasks; + if completed_tasks + .iter() + .any(|task| task.original_fragments.is_empty()) + { + return Err(Error::invalid_input( + "compaction results must replace at least one original fragment".to_string(), + )); + } // Collect the rewritten fragments' file paths up front so every failure - // path below can clean them up (or deliberately keep them). Fragment ids - // may still be reassigned by reserve_fragment_ids; cleanup only needs the - // file paths, which never change. + // path below can clean them up (or deliberately keep them). This deliberately + // excludes the metadata-only trailing replacements added below because their + // files remain live. let all_new_fragments: Vec = completed_tasks .iter() .flat_map(|t| t.new_fragments.iter().cloned()) .collect(); - // Single reserve_fragment_ids for all address-style tasks - if has_address_style { - let frags: Vec<&mut Fragment> = completed_tasks - .iter_mut() - .filter(|t| t.row_addrs.is_some()) - .flat_map(|t| t.new_fragments.iter_mut()) - .collect(); - if let Err(e) = reserve_fragment_ids(dataset, frags.into_iter()).await { - cleanup_compaction_files_after_reservation_failure(dataset, &all_new_fragments).await; - return Err(e); - } + let uses_stable_row_ids = dataset.manifest.uses_stable_row_ids(); + let capture_row_addrs = completed_tasks.iter().any(|task| task.row_addrs.is_some()); + + // Fragment IDs are the manifest ordering key. For a partial or gapped + // compaction, relabel every untouched fragment after the first rewritten + // range so all replacement IDs can be reserved as one ordered suffix. + let mut completed_tasks = complete_rewrite_suffix( + &dataset.manifest.fragments, + completed_tasks, + options.max_source_fragments, + tasks_read_version, + dataset.manifest.version, + )?; + // A metadata-only suffix relabel moves physical row addresses too. Check + // the completed suffix, not only the originally submitted results, so a + // custom or distributed plan cannot move rows covered by an index this + // build cannot open and remap. + if !uses_stable_row_ids && !options.defer_index_remap { + reject_unremappable_rewrite(dataset, &completed_tasks).await?; + } + if capture_row_addrs { + prepare_metadata_only_relabels(dataset, &mut completed_tasks).await?; + } + let has_address_style = completed_tasks.iter().any(|task| task.row_addrs.is_some()); + let needs_remapping = !uses_stable_row_ids && !options.defer_index_remap && has_address_style; + // Confirm there is a remapper before materializing the potentially very large row address map. + let index_remapper = if needs_remapping { + remap_options.create_remapper(dataset).await? + } else { + None + }; + + // Reserve one consecutive ID range for compacted outputs and metadata-only + // trailing replacements. The returned task order is the desired row order. + let new_fragments = completed_tasks + .iter_mut() + .flat_map(|task| task.new_fragments.iter_mut()) + .collect::>(); + if let Err(error) = reserve_fragment_ids(dataset, new_fragments.into_iter()).await { + cleanup_compaction_files_after_reservation_failure(dataset, &all_new_fragments).await; + return Err(error); } + copy_relabeled_deletion_files(dataset, &mut completed_tasks).await?; let mut rewrite_groups = Vec::with_capacity(completed_tasks.len()); let mut metrics = CompactionMetrics::default(); @@ -3446,19 +3814,6 @@ pub async fn commit_compaction( new_index_files: rewritten.files, }) .collect() - } else if !options.defer_index_remap && !has_address_style { - // We need to reserve fragment ids here so that the fragment bitmap - // can be updated for each index. Only needed for stable row IDs - // since address-style IDs were already reserved above. - let new_fragments = rewrite_groups - .iter_mut() - .flat_map(|group| group.new_fragments.iter_mut()) - .collect::>(); - if let Err(e) = reserve_fragment_ids(dataset, new_fragments.into_iter()).await { - cleanup_compaction_files_after_reservation_failure(dataset, &all_new_fragments).await; - return Err(e); - } - Vec::new() } else { Vec::new() }; @@ -3674,7 +4029,7 @@ mod tests { use crate::index::DatasetIndexExt; use crate::index::frag_reuse::{load_frag_reuse_index_details, open_frag_reuse_index}; use crate::index::vector::{StageParams, VectorIndexParams}; - use crate::utils::test::{DatagenExt, FragmentCount, FragmentRowCount}; + use crate::utils::test::{DatagenExt, FragmentCount, FragmentRowCount, copy_test_data_to_tmp}; use arrow_array::types::{Float32Type, Float64Type, Int32Type, Int64Type}; use arrow_array::{ ArrayRef, Float32Array, Int32Array, Int64Array, LargeBinaryArray, LargeStringArray, @@ -3852,14 +4207,6 @@ mod tests { result_str.push_str(&first_keys); result_str } - - fn in_any_order(expectations: &[Self]) -> Self { - let expectations = expectations - .iter() - .flat_map(|item| item.expectations.clone()) - .collect::>(); - Self { expectations } - } } #[async_trait] @@ -3967,7 +4314,7 @@ mod tests { .unwrap(); let options = CompactionOptions { data_storage_version: selector, - excluded_fragment_ids: vec![2], + excluded_fragment_ids: vec![0], ..Default::default() }; @@ -3984,7 +4331,7 @@ mod tests { plan ); assert_eq!(plan.num_tasks(), 1); - let retained = dataset.manifest.fragments[2].clone(); + let retained = dataset.manifest.fragments[0].clone(); let task = plan.compaction_tasks().next().unwrap(); let task: CompactionTask = serde_json::from_slice(&serde_json::to_vec(&task).unwrap()).unwrap(); @@ -4240,12 +4587,26 @@ mod tests { let ours = CompactionOptions { target_rows_per_fragment: 16, defer_index_remap: true, - excluded_fragment_ids: vec![2, 3], ..Default::default() }; let planned = fresh(uri).await.unwrap(); let read_version = planned.manifest.version; - let completed = execute_compaction_plan(&planned, &ours).await; + let plan = plan_compaction(&planned, &ours).await.unwrap(); + assert_eq!(plan.num_tasks(), 2); + let first_task = plan.tasks().first().unwrap().clone(); + assert_eq!( + first_task + .fragments + .iter() + .map(|f| f.id) + .collect::>(), + vec![0, 1] + ); + let completed = vec![ + rewrite_files(Cow::Borrowed(&planned), first_task, &ours) + .await + .unwrap(), + ]; assert_eq!(completed.len(), 1); assert_eq!(completed[0].read_version, read_version); @@ -4282,14 +4643,15 @@ mod tests { let ledger = crate::index::frag_reuse::decode_frag_reuse_ledger(&committed, entry) .await .unwrap(); - assert_eq!(ledger.transitions().len(), 2); + // The earlier rewrite also relabels the later writer's replacement to + // keep fragment IDs in row order. + assert_eq!(ledger.transitions().len(), 3); let batch = committed.scan().try_into_batch().await.unwrap(); - let mut values: Vec = batch["i"] + let values: Vec = batch["i"] .as_primitive::() .iter() .map(|v| v.unwrap()) .collect(); - values.sort_unstable(); assert_eq!(values, (0..32).collect::>()); for predicate in ["i = 3", "i = 20", "i > 10"] { let indexed = committed @@ -4666,8 +5028,7 @@ mod tests { .unwrap(); let first_new_frag_idx = 7; - // The tasks execute concurrently, so either one may reserve the first - // output fragment id. + // Tasks may finish in any order, but commit assigns fragment ids in plan order. let remap_a = expect_remap( &[ vec![ @@ -4676,8 +5037,10 @@ mod tests { (row_addrs(1, 0..400), true), (row_addrs(2, 0..400), true), ], - // frag 3 is skipped since it does not have enough missing data - // Frags 4, 5, and 6 are rewritten to frag 8 + // Frag 3 keeps its files but is relabeled as fragment 8 so the + // manifest can remain ID-sorted. + vec![(row_addrs(3, 0..1000), true)], + // Frags 4, 5, and 6 are rewritten to frag 9. vec![ (row_addrs(4, 0..200), true), (row_addrs(4, 200..400), false), @@ -4688,25 +5051,6 @@ mod tests { ], first_new_frag_idx, ); - let remap_b = expect_remap( - &[ - // Frags 4, 5, and 6 are rewritten to frag 7 - vec![ - (row_addrs(4, 0..200), true), - (row_addrs(4, 200..400), false), - (row_addrs(4, 400..1000), true), - (row_addrs(5, 0..300), true), - (row_addrs(6, 0..300), true), - ], - // 3 small fragments rewritten to frag 8 - vec![ - (row_addrs(0, 0..400), true), - (row_addrs(1, 0..400), true), - (row_addrs(2, 0..400), true), - ], - ], - first_new_frag_idx, - ); // Create compaction plan let options = CompactionOptions { @@ -4734,11 +5078,10 @@ mod tests { .collect::>(), vec![4, 5, 6] ); - - let mock_remapper = MockIndexRemapper::in_any_order(&[remap_a, remap_b]); + let expected_data = dataset.scan().try_into_batch().await.unwrap(); // Run compaction - let metrics = compact_files(&mut dataset, options, Some(Arc::new(mock_remapper))) + let metrics = compact_files(&mut dataset, options, Some(Arc::new(remap_a))) .await .unwrap(); @@ -4753,7 +5096,66 @@ mod tests { .iter() .map(|f| f.id()) .collect::>(); - assert_eq!(fragment_ids, vec![3, 7, 8]); + assert_eq!(fragment_ids, vec![7, 8, 9]); + let compacted_data = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(compacted_data, expected_data); + } + + #[rstest] + #[tokio::test] + async fn test_commit_compaction_preserves_plan_order( + #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)] + data_storage_version: LanceFileVersion, + ) { + let data = sample_data(); + let reader = RecordBatchIterator::new(vec![Ok(data.clone())], data.schema()); + let mut dataset = Dataset::write( + reader, + "memory://", + Some(WriteParams { + max_rows_per_file: 2_500, + max_rows_per_group: 2_500, + data_storage_version: Some(data_storage_version), + ..Default::default() + }), + ) + .await + .unwrap(); + + let options = CompactionOptions { + target_rows_per_fragment: 5_000, + ..Default::default() + }; + let plan = plan_compaction(&dataset, &options).await.unwrap(); + let tasks = plan.compaction_tasks().collect::>(); + assert_eq!(tasks.len(), 2); + + // Compaction tasks may complete in any order when executed concurrently + // or by distributed workers. + let completed_tasks = vec![ + tasks[1].execute(&dataset).await.unwrap(), + tasks[0].execute(&dataset).await.unwrap(), + ]; + commit_compaction( + &mut dataset, + completed_tasks, + Arc::new(DatasetIndexRemapperOptions::default()), + &options, + ) + .await + .unwrap(); + + let compacted = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(compacted, data); + + let taken = dataset + .take(&[0, 4_999, 5_000, 9_999], dataset.schema().clone()) + .await + .unwrap(); + assert_eq!( + taken["a"].as_primitive::().values(), + &[0, 4_999, 5_000, 9_999] + ); } #[rstest] @@ -5545,10 +5947,11 @@ mod tests { assert_eq!(results[0].metrics.files_removed, 3); assert_eq!(results[0].metrics.files_added, 1); - // Just commit the last task + // Commit the first task. This relabels the untouched suffix, so the + // remaining worker results refer to stale source fragment IDs. commit_compaction( &mut dataset, - vec![results.pop().unwrap()], + vec![results.remove(0)], Arc::new(IgnoreRemap::default()), &options, ) @@ -5573,6 +5976,15 @@ mod tests { assert_eq!(dataset.manifest.version, 5); assert_eq!(dataset.manifest.uses_stable_row_ids(), use_stable_row_id,); + assert!( + dataset + .manifest + .fragments + .windows(2) + .all(|fragments| fragments[0].id < fragments[1].id) + ); + let compacted = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(compacted, data.slice(0, 9_000)); } #[tokio::test] @@ -6889,11 +7301,6 @@ mod tests { let rewrite_result2 = rewrite_files(Cow::Borrowed(&dataset), tasks[1].clone(), &options) .await .unwrap(); - let rewritten_frags2 = rewrite_result2 - .original_fragments - .iter() - .map(|f| f.id) - .collect::>(); commit_compaction( &mut dataset, Vec::from([rewrite_result2]), @@ -6912,16 +7319,13 @@ mod tests { let frag_reuse_details2 = load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta2) .await .unwrap(); - let new_frags2 = frag_reuse_details2.versions.last().unwrap().new_frag_ids(); + let reuse_version2 = frag_reuse_details2.versions.last().unwrap(); + let old_frags2 = reuse_version2.old_frag_ids(); + let new_frags2 = reuse_version2.new_frag_ids(); let rewrite_result3 = rewrite_files(Cow::Borrowed(&dataset), tasks[2].clone(), &options) .await .unwrap(); - let rewritten_frags3 = rewrite_result3 - .original_fragments - .iter() - .map(|f| f.id) - .collect::>(); commit_compaction( &mut dataset, Vec::from([rewrite_result3]), @@ -6940,7 +7344,9 @@ mod tests { let frag_reuse_details3 = load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta3) .await .unwrap(); - let new_frags3 = frag_reuse_details3.versions.last().unwrap().new_frag_ids(); + let reuse_version3 = frag_reuse_details3.versions.last().unwrap(); + let old_frags3 = reuse_version3.old_frag_ids(); + let new_frags3 = reuse_version3.new_frag_ids(); // Concurrently commit a frag_reuse_index cleanup operation. dataset_clone // only knows the first reuse version; catch its index up so the cleanup @@ -6964,15 +7370,9 @@ mod tests { .await .unwrap(); assert_eq!(frag_reuse_details.versions.len(), 2); - assert_eq!( - frag_reuse_details.versions[0].old_frag_ids(), - rewritten_frags2 - ); + assert_eq!(frag_reuse_details.versions[0].old_frag_ids(), old_frags2); assert_eq!(frag_reuse_details.versions[0].new_frag_ids(), new_frags2); - assert_eq!( - frag_reuse_details.versions[1].old_frag_ids(), - rewritten_frags3 - ); + assert_eq!(frag_reuse_details.versions[1].old_frag_ids(), old_frags3); assert_eq!(frag_reuse_details.versions[1].new_frag_ids(), new_frags3); } @@ -7057,11 +7457,6 @@ mod tests { rewrite_files(Cow::Borrowed(&dataset_clone), tasks[1].clone(), &options) .await .unwrap(); - let rewritten_frags2 = rewrite_result2 - .original_fragments - .iter() - .map(|f| f.id) - .collect::>(); commit_compaction( &mut dataset_clone, Vec::from([rewrite_result2]), @@ -7084,10 +7479,7 @@ mod tests { .await .unwrap(); assert_eq!(frag_reuse_details.versions.len(), 1); - assert_eq!( - frag_reuse_details.versions[0].old_frag_ids(), - rewritten_frags2 - ); + assert_eq!(frag_reuse_details.versions[0].groups.len(), 3); // Verify new fragment IDs are non-zero (allocated by commit_compaction) let new_frags2 = frag_reuse_details.versions[0].new_frag_ids(); assert!(new_frags2.iter().all(|id| *id != 0)); @@ -9518,7 +9910,37 @@ mod tests { #[tokio::test] async fn test_max_source_fragments() { let test_dir = TempStrDir::default(); - let dataset = dataset_with_ten_small_fragments(&test_dir).await; + let test_uri = &test_dir; + + let data = sample_data(); + let schema = data.schema(); + + // Create 10 small fragments (100 rows each) via 10 appends. + let write_params = WriteParams { + max_rows_per_file: 100, + ..Default::default() + }; + Dataset::write( + RecordBatchIterator::new(vec![Ok(data.slice(0, 100))], schema.clone()), + test_uri, + Some(write_params.clone()), + ) + .await + .unwrap(); + for i in 1..10 { + let mut append_params = write_params.clone(); + append_params.mode = WriteMode::Append; + Dataset::write( + RecordBatchIterator::new(vec![Ok(data.slice(i * 100, 100))], schema.clone()), + test_uri, + Some(append_params), + ) + .await + .unwrap(); + } + + let dataset = Dataset::open(test_uri).await.unwrap(); + assert_eq!(dataset.get_fragments().len(), 10); // Plan without limit - all 10 fragments should be candidates. // Use a target that splits the 10 fragments into multiple tasks. @@ -9535,8 +9957,8 @@ mod tests { plan_all.num_tasks() ); - // Plan with max_source_fragments=4 should include tasks covering <= 4 - // source fragments + // Plan with max_source_fragments=4 should only affect a suffix of at + // most four fragments, including any required metadata relabels. let opts_bounded = CompactionOptions { target_rows_per_fragment: 250, max_source_fragments: Some(4), @@ -9553,6 +9975,17 @@ mod tests { bounded_source_frags > 0, "expected at least 1 source fragment in bounded plan" ); + let bounded_fragment_ids = plan_bounded + .tasks() + .iter() + .flat_map(|task| task.fragments.iter().map(|fragment| fragment.id)) + .collect::>(); + let expected_fragment_ids = dataset.manifest.fragments + [dataset.manifest.fragments.len() - bounded_source_frags..] + .iter() + .map(|fragment| fragment.id) + .collect::>(); + assert_eq!(bounded_fragment_ids, expected_fragment_ids); assert!( plan_bounded.num_tasks() < plan_all.num_tasks(), "bounded plan ({}) should have fewer tasks than unbounded ({})", @@ -9565,6 +9998,8 @@ mod tests { compact_files(&mut dataset, opts_bounded, None) .await .unwrap(); + let after_first_data = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(after_first_data, data.slice(0, 1_000)); let after_first = dataset.get_fragments().len(); assert!( after_first < 10, @@ -9584,17 +10019,173 @@ mod tests { compact_files(&mut dataset, opts_bounded, None) .await .unwrap(); + let after_second_data = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(after_second_data, data.slice(0, 1_000)); let after_second = dataset.get_fragments().len(); assert!( - after_second <= after_first, - "expected progress: {after_second} should be <= {after_first}" + after_second < after_first, + "expected progress: {after_second} should be < {after_first}" ); } + #[rstest] + #[tokio::test] + async fn test_partial_compaction_preserves_order_across_candidate_gap( + #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)] + data_storage_version: LanceFileVersion, + ) { + let test_dir = TempStrDir::default(); + let test_uri = &test_dir; + let data = sample_data(); + let schema = data.schema(); + let fragment_rows = [100, 100, 1_000, 100, 100]; + let write_params = WriteParams { + max_rows_per_file: 1_000, + data_storage_version: Some(data_storage_version), + ..Default::default() + }; + + Dataset::write( + RecordBatchIterator::new(vec![Ok(data.slice(0, fragment_rows[0]))], schema.clone()), + test_uri, + Some(write_params.clone()), + ) + .await + .unwrap(); + let mut offset = fragment_rows[0]; + for row_count in fragment_rows.iter().copied().skip(1) { + let mut append_params = write_params.clone(); + append_params.mode = WriteMode::Append; + Dataset::write( + RecordBatchIterator::new(vec![Ok(data.slice(offset, row_count))], schema.clone()), + test_uri, + Some(append_params), + ) + .await + .unwrap(); + offset += row_count; + } + + let mut dataset = Dataset::open(test_uri).await.unwrap(); + dataset.delete("a = 1250").await.unwrap(); + let expected = dataset.scan().try_into_batch().await.unwrap(); + let options = CompactionOptions { + target_rows_per_fragment: 250, + ..Default::default() + }; + let plan = plan_compaction(&dataset, &options).await.unwrap(); + let planned_fragment_ids = plan + .tasks() + .iter() + .flat_map(|task| task.fragments.iter().map(|fragment| fragment.id)) + .collect::>(); + assert_eq!(planned_fragment_ids, vec![0, 1, 3, 4]); + let tasks = plan.compaction_tasks().collect::>(); + assert_eq!(tasks.len(), 2); + let first_result = tasks[0].execute(&dataset).await.unwrap(); + let second_result = tasks[1].execute(&dataset).await.unwrap(); + + commit_compaction( + &mut dataset, + vec![first_result], + Arc::new(DatasetIndexRemapperOptions::default()), + &options, + ) + .await + .unwrap(); + + assert_eq!(dataset.get_fragments().len(), 4); + let fragment_ids = dataset + .manifest + .fragments + .iter() + .map(|fragment| fragment.id) + .collect::>(); + assert_eq!(fragment_ids, vec![5, 6, 7, 8]); + let fragments_by_id = dataset + .get_fragments_from_ids(&[8, 5, 7, 6]) + .unwrap() + .into_iter() + .map(|fragment| fragment.id()) + .collect::>(); + assert_eq!(fragments_by_id, vec![5, 6, 7, 8]); + assert!(dataset.manifest.fragments[2].deletion_file.is_some()); + dataset.validate().await.unwrap(); + let after_first = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(after_first, expected); + + commit_compaction( + &mut dataset, + vec![second_result], + Arc::new(DatasetIndexRemapperOptions::default()), + &options, + ) + .await + .unwrap(); + + let compacted = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(compacted, expected); + assert_eq!(dataset.get_fragments().len(), 3); + } + + #[test] + fn test_complete_rewrite_suffix_respects_max_source_fragments() { + let mut source = Fragment::new(0); + source.physical_rows = Some(1); + let mut trailing = Fragment::new(1); + trailing.physical_rows = Some(100_000); + let task = RewriteResult { + metrics: CompactionMetrics::default(), + new_fragments: vec![Fragment::new(2)], + read_version: 1, + original_fragments: vec![source.clone()], + row_addrs: None, + }; + + let error = + complete_rewrite_suffix(&[source, trailing], vec![task], Some(1), 1, 1).unwrap_err(); + let message = error.to_string(); + assert!(matches!(error, Error::InvalidInput { .. })); + assert!( + message.contains( + "compaction would change 2 fragment identities, exceeding max_source_fragments=1" + ), + "unexpected error: {message}" + ); + } + + #[tokio::test] + async fn test_metadata_only_relabel_migrates_released_fragment_stats() { + let test_dir = copy_test_data_to_tmp("v0.7.5/with_deletions").unwrap(); + let dataset = Dataset::open(&test_dir.path_str()).await.unwrap(); + assert!(dataset.manifest.writer_version.is_none()); + let trailing = dataset.manifest.fragments[0].clone(); + assert_eq!(trailing.physical_rows, None); + let mut tasks = vec![RewriteResult { + metrics: CompactionMetrics::default(), + new_fragments: vec![trailing.clone()], + read_version: dataset.version().version, + original_fragments: vec![trailing], + row_addrs: None, + }]; + + prepare_metadata_only_relabels(&dataset, &mut tasks) + .await + .unwrap(); + + assert_eq!(tasks[0].original_fragments[0].physical_rows, Some(100)); + assert_eq!(tasks[0].new_fragments[0].physical_rows, Some(100)); + let row_addrs = RoaringTreemap::deserialize_from(&mut Cursor::new( + tasks[0].row_addrs.as_ref().unwrap(), + )) + .unwrap(); + assert_eq!(row_addrs.len(), 100); + } + /// Writes `sample_data` as 10 fragments of 100 rows each, deletes half of - /// fragment 0's rows, and puts a one-cell data overlay on it, so every - /// `max_source_*` planning budget is exercised against a dataset that - /// carries deleted rows and overlay files. + /// fragment 0's rows, and puts a one-cell data overlay on it, so the row and + /// byte planning budgets are exercised against a dataset that carries + /// deleted rows and overlay files. async fn dataset_with_ten_small_fragments(test_uri: &str) -> Dataset { let data = sample_data(); let schema = data.schema(); @@ -9664,10 +10255,7 @@ mod tests { }) .collect::>(); - assert_eq!( - planned_fragment_ids, - vec![vec![0, 1, 2, 3], vec![5, 6, 7, 8, 9]] - ); + assert_eq!(planned_fragment_ids, vec![vec![5, 6, 7, 8, 9]]); assert!( planned_fragment_ids .iter() @@ -9675,8 +10263,11 @@ mod tests { .all(|fragment_id| *fragment_id != excluded_fragment_id) ); + let before = dataset.scan().try_into_batch().await.unwrap(); let metrics = compact_files(&mut dataset, options, None).await.unwrap(); - assert_eq!(metrics.fragments_removed, 9); + assert_eq!(metrics.fragments_removed, 5); + let after = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(after, before); let remaining_fragment_ids = dataset .get_fragments() .iter() diff --git a/rust/lance/src/index/create.rs b/rust/lance/src/index/create.rs index fc63fe6114b..5af79a7e692 100644 --- a/rust/lance/src/index/create.rs +++ b/rust/lance/src/index/create.rs @@ -3108,13 +3108,14 @@ mod tests { /// surviving half made the coverage look healthy. #[tokio::test] async fn test_merge_uncommitted_segments_partly_retired_by_compaction() { - // Two undersized fragments and one already at the compaction target, so - // compaction rewrites the pair and leaves the third alone. + // One fragment at the compaction target followed by two undersized + // fragments. Rewriting the trailing pair leaves the first fragment's + // staged coverage live without relabeling it. let reader = gen_batch() .col("id", lance_datagen::array::step::()) .into_reader_rows( - lance_datagen::RowCount::from(2), - lance_datagen::BatchCount::from(2), + lance_datagen::RowCount::from(4), + lance_datagen::BatchCount::from(1), ); let test_dir = tempfile::tempdir().unwrap(); let dataset_uri = test_dir.path().to_str().unwrap(); @@ -3122,7 +3123,7 @@ mod tests { reader, dataset_uri, Some(WriteParams { - max_rows_per_file: 2, + max_rows_per_file: 4, enable_stable_row_ids: false, ..Default::default() }), @@ -3132,15 +3133,15 @@ mod tests { let reader = gen_batch() .col("id", lance_datagen::array::step::()) .into_reader_rows( - lance_datagen::RowCount::from(4), - lance_datagen::BatchCount::from(1), + lance_datagen::RowCount::from(2), + lance_datagen::BatchCount::from(2), ); let mut dataset = Dataset::write( reader, dataset_uri, Some(WriteParams { mode: WriteMode::Append, - max_rows_per_file: 4, + max_rows_per_file: 2, enable_stable_row_ids: false, ..Default::default() }), @@ -3172,8 +3173,8 @@ mod tests { ) .await .unwrap(); - // Two rows per fragment against a four-row target pairs some fragments - // and leaves at least one alone. + // Two trailing two-row fragments compact while the leading four-row + // fragment keeps its identity. crate::dataset::optimize::compact_files( &mut dataset, crate::dataset::optimize::CompactionOptions { @@ -3856,10 +3857,10 @@ mod tests { ); assert_eq!(count_in_range(&dataset, &merged, 50, 100).await, 50); - // Phase 2 — retire fragment 0: delete >10% of its rows so compaction - // rewrites only frag 0 (frag 1 has no deletions and is at target size). - // The committed per-fragment segment now claims a fragment the dataset - // no longer has. + // Phase 2 — delete >10% of fragment 0's rows so compaction rewrites its + // data. Preserving row order in the ID-sorted manifest also relabels the + // untouched trailing fragment, so both committed segments must follow + // their stable row-id coverage to fresh fragment IDs. dataset.delete("id < 16").await.unwrap(); crate::dataset::optimize::compact_files( &mut dataset, @@ -3878,16 +3879,15 @@ mod tests { .collect(); assert!(!live_frags.contains(0), "compaction should retire frag 0"); - // Filtered merge: coverage drops the retired fragment but keeps the - // live one, and the merged page data does not leak the retired row ids - // (ids < 16 lived only in frag 0, so the range now returns nothing). + // Filtered merge: coverage follows the current fragments, and the + // merged page data does not leak the retired row ids (ids < 16 lived + // only in frag 0, so the range now returns nothing). let merged = dataset .merge_existing_index_segments(dataset.load_indices_by_name("id_btree").await.unwrap()) .await .unwrap(); let coverage = merged.fragment_bitmap.as_ref().unwrap(); - assert!(!coverage.contains(0), "must drop retired frag 0"); - assert!(coverage.contains(1), "must keep live frag 1"); + assert_eq!(coverage, &live_frags); assert_eq!( count_in_range(&dataset, &merged, 0, 16).await, 0, diff --git a/rust/lance/src/index/scalar_logical.rs b/rust/lance/src/index/scalar_logical.rs index 18cc840b823..37913a75ea9 100644 --- a/rust/lance/src/index/scalar_logical.rs +++ b/rust/lance/src/index/scalar_logical.rs @@ -1024,8 +1024,10 @@ mod tests { .await .unwrap(); let coverage = merged.fragment_bitmap.as_ref().unwrap(); - assert!(!coverage.contains(0), "must drop retired frag 0"); - assert!(coverage.contains(1), "must keep live indexed frag 1"); + assert!( + coverage.is_empty(), + "physical-address coverage must be invalidated by fragment relabeling" + ); let field_path = dataset.schema().field_path(merged.fields[0]).unwrap(); let index = crate::index::scalar::open_scalar_index( @@ -1118,6 +1120,10 @@ mod tests { .await .unwrap(); let merged_coverage = merged.fragment_bitmap.as_ref().unwrap().clone(); + assert!( + merged_coverage.is_empty(), + "physical-address coverage must be invalidated by fragment relabeling" + ); let merged_uuid = merged.uuid; dataset @@ -1136,7 +1142,7 @@ mod tests { scalar_index_fragment_bitmap(&dataset, "value", "value_zonemap_replace_retired") .await .unwrap() - .unwrap(); + .unwrap_or_default(); assert_eq!(combined_bitmap, merged_coverage); } @@ -1476,14 +1482,14 @@ mod tests { arrow_array::RecordBatch::try_new( schema.clone(), vec![Arc::new(arrow_array::StringArray::from(vec![ - "alpha beta gamma", - "beta gamma delta", - "gamma delta epsilon", - "delta epsilon zeta", "epsilon zeta eta", "zeta eta theta", "eta theta iota", "theta iota kappa", + "alpha beta gamma", + "beta gamma delta", + "gamma delta epsilon", + "delta epsilon zeta", ]))], ) .unwrap(), @@ -1496,6 +1502,8 @@ mod tests { let fragments = dataset.get_fragments(); assert_eq!(fragments.len(), 2); + let surviving_fragment_id = fragments[0].id() as u32; + let rewritten_fragment_id = fragments[1].id() as u32; // Build per-fragment FM-Index segments and commit let params = ScalarIndexParams::for_builtin(BuiltinIndexType::Fm); @@ -1521,7 +1529,8 @@ mod tests { .unwrap(); assert_eq!(committed.len(), 2); - // Delete rows from fragment 0 to trigger compaction retirement + // Delete rows from the trailing fragment so compaction leaves the + // leading indexed fragment's physical addresses intact. dataset.delete("text = 'alpha beta gamma'").await.unwrap(); dataset.delete("text = 'beta gamma delta'").await.unwrap(); crate::dataset::optimize::compact_files( @@ -1541,11 +1550,13 @@ mod tests { .map(|f| f.id() as u32) .collect(); assert!( - !live_frags.contains(0), - "compaction should retire fragment 0" + !live_frags.contains(rewritten_fragment_id), + "compaction should retire the trailing fragment" ); + assert!(live_frags.contains(surviving_fragment_id)); - // Merge: the retired fragment should be dropped from coverage + // The surviving leading segment remains usable; the rewritten + // trailing segment cannot claim its new physical addresses. let segments = dataset .load_indices_by_name("text_fmindex_compact") .await @@ -1556,14 +1567,7 @@ mod tests { .unwrap(); let coverage = merged.fragment_bitmap.as_ref().unwrap(); - assert!( - !coverage.contains(0), - "merged coverage must drop retired fragment 0" - ); - assert!( - coverage.contains(1), - "merged coverage must keep live fragment 1" - ); + assert_eq!(coverage, &RoaringBitmap::from_iter([surviving_fragment_id])); // Commit the merged segment and verify search works dataset @@ -1599,7 +1603,7 @@ mod tests { "deleted rows from retired fragment should not appear in merged index" ); - // "theta" exists in fragment 1 rows only + // "theta" exists in the surviving leading fragment. let query = lance_index::scalar::TextQuery::StringContains("theta".to_string()); let result = logical.search(&query, &NoOpMetricsCollector).await.unwrap(); let row_addrs = match result { @@ -1968,6 +1972,7 @@ mod tests { .await .unwrap(); let source_uuid = segment.uuid; + let source_coverage = segment.fragment_bitmap.as_ref().unwrap().clone(); // Retire fragment 0: delete its rows and compact it away. dataset.delete("text = 'alpha beta gamma'").await.unwrap(); @@ -1992,9 +1997,15 @@ mod tests { !live_frags.contains(0), "compaction should retire fragment 0" ); - assert!(live_frags.contains(1), "fragment 1 should stay live"); + assert_eq!( + source_coverage.intersection_len(&live_frags), + 0, + "compaction should retire every fragment covered by the stale segment" + ); - // Coverage shrank, so even a single segment must be rebuilt. + // The uncommitted segment was not present for compaction to relabel its + // coverage, so even a single segment must be rebuilt without claiming + // any current fragment. let merged = dataset .merge_existing_index_segments(vec![segment]) .await @@ -2003,14 +2014,9 @@ mod tests { merged.uuid, source_uuid, "shrunk coverage must trigger a rebuild" ); - assert_eq!( - merged - .fragment_bitmap - .as_ref() - .unwrap() - .iter() - .collect::>(), - vec![1] + assert!( + merged.fragment_bitmap.as_ref().unwrap().is_empty(), + "rebuilt coverage must exclude every retired fragment" ); } } diff --git a/rust/lance/src/io/commit.rs b/rust/lance/src/io/commit.rs index 107180ba1d9..817580a6911 100644 --- a/rust/lance/src/io/commit.rs +++ b/rust/lance/src/io/commit.rs @@ -774,21 +774,37 @@ fn check_storage_version(manifest: &mut Manifest) -> Result<()> { crate::dataset::versions::check_manifest_storage_version(manifest) } -/// Reject a manifest in which two fragments share an id. Per-fragment state is -/// keyed by fragment id — deletion file paths, cached row id sequences, row -/// addresses — so a duplicate makes it ambiguous which rows that state describes. +/// Reject a manifest in which fragment identities are ambiguous or out of order. +/// Per-fragment state is keyed by fragment id — deletion file paths, cached row +/// id sequences, row addresses — so a duplicate makes it ambiguous which rows +/// that state describes. /// /// Runs after the legacy fixups above, so a dataset that needs a rollback for some -/// other reason is diagnosed with that first. Relies on `build_manifest` leaving -/// the fragments sorted by id. +/// other reason is diagnosed with that first. fn check_fragment_ids(manifest: &Manifest) -> Result<()> { - if let Some(pair) = manifest.fragments.windows(2).find(|p| p[0].id == p[1].id) { + let mut fragment_ids = HashSet::with_capacity(manifest.fragments.len()); + if let Some(fragment) = manifest + .fragments + .iter() + .find(|fragment| !fragment_ids.insert(fragment.id)) + { return Err(Error::invalid_input(format!( "The commit would produce two fragments with id {}. Fragment ids must be \ unique. Datasets written by Lance 0.16 and earlier may already contain \ duplicate ids; those have to be rewritten, or rolled back to a version \ without the duplicate.", - pair[0].id + fragment.id + ))); + } + if let Some(fragments) = manifest + .fragments + .windows(2) + .find(|fragments| fragments[0].id > fragments[1].id) + { + return Err(Error::invalid_input(format!( + "The commit would place fragment {} before lower fragment id {}. Fragment ids must \ + remain sorted in increasing order", + fragments[0].id, fragments[1].id ))); } Ok(()) @@ -1952,6 +1968,24 @@ mod tests { use crate::index::vector::VectorIndexParams; use crate::utils::test::{DatagenExt, FragmentCount, FragmentRowCount}; + #[test] + fn test_check_fragment_ids_requires_sorted_order() { + let manifest = Manifest::new( + Schema::try_from(&ArrowSchema::empty()).unwrap(), + Arc::new(vec![Fragment::new(5), Fragment::new(2)]), + DataStorageFormat::default(), + HashMap::new(), + ); + + let error = check_fragment_ids(&manifest).unwrap_err(); + assert!(matches!(error, Error::InvalidInput { .. })); + assert!( + error + .to_string() + .contains("Fragment ids must remain sorted in increasing order") + ); + } + #[rstest::rstest] #[case::unset(None, DEFAULT_COMMIT_RETRY_TIMEOUT)] #[case::whole_seconds(Some("600"), Duration::from_secs(600))] diff --git a/rust/lance/src/io/commit/conflict_resolver.rs b/rust/lance/src/io/commit/conflict_resolver.rs index 08f978f6a9d..e393a9705ac 100644 --- a/rust/lance/src/io/commit/conflict_resolver.rs +++ b/rust/lance/src/io/commit/conflict_resolver.rs @@ -1249,11 +1249,32 @@ impl<'a> TransactionRebase<'a> { .as_ref() .is_some_and(|entry| !is_tagged(entry)) && self.reuse.added_transitions.is_none(); + // Compaction reserves a new ID suffix after planning. A concurrent + // row addition can then sort before its replacements, so retry the + // compaction. Stable-partition rewrites carry their own tagged row + // mapping and keep the existing append rebase behavior. + let is_stable_partition_rewrite = self.reuse.added_transitions.as_ref().is_some_and( + |transitions| { + !transitions.is_empty() + && transitions.iter().all(|transition| { + matches!( + transition.mapping, + Some(lance_table::format::pb::fragment_reuse_index_details::transition::Mapping::StablePartition(_)) + ) + }) + }, + ); match &other_transaction.operation { // Rewrite is only compatible with operations that don't touch // existing fragments or update fragments we don't touch. - Operation::Append { .. } - | Operation::ReserveFragments { .. } + Operation::Append { .. } => { + if !is_stable_partition_rewrite { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } + Operation::ReserveFragments { .. } | Operation::Project { .. } | Operation::Clone { .. } | Operation::UpdateConfig { .. } @@ -1263,11 +1284,6 @@ impl<'a> TransactionRebase<'a> { updated_fragments, deleted_fragment_ids, .. - } - | Operation::Update { - updated_fragments, - removed_fragment_ids: deleted_fragment_ids, - .. } => { if updated_fragments .iter() @@ -1280,6 +1296,24 @@ impl<'a> TransactionRebase<'a> { Ok(()) } } + Operation::Update { + updated_fragments, + removed_fragment_ids: deleted_fragment_ids, + new_fragments, + .. + } => { + if (!is_stable_partition_rewrite && !new_fragments.is_empty()) + || updated_fragments + .iter() + .map(|f| f.id) + .chain(deleted_fragment_ids.iter().copied()) + .any(|id| self.modified_fragment_ids.contains(&id)) + { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } Operation::DataOverlay { groups } => { // Rewriting a fragment changes its physical row addresses, so // an overlay addressed by physical offset on that fragment is @@ -3679,6 +3713,57 @@ mod tests { Retryable, } + #[rstest::rstest] + #[case::append(false)] + #[case::update(true)] + fn test_compaction_rewrite_retries_row_additions(#[case] is_update: bool) { + let operation = Operation::Rewrite { + groups: vec![RewriteGroup { + old_fragments: vec![Fragment::new(0)], + new_fragments: vec![Fragment::new(2)], + }], + rewritten_indices: vec![], + frag_reuse_index: None, + }; + let mut rebase = TransactionRebase { + transaction: Transaction::new(0, operation.clone(), None), + initial_fragments: HashMap::new(), + modified_fragment_ids: modified_fragment_ids(&operation).collect::>(), + affected_rows: None, + conflicting_frag_reuse_indices: Vec::new(), + conflicting_mem_wal_compacted_sstables: Vec::new(), + current_lineage: None, + current_live: None, + current_schema: None, + read_fragments: None, + read_schema: None, + reuse: Default::default(), + }; + let other_operation = if is_update { + Operation::Update { + removed_fragment_ids: vec![], + updated_fragments: vec![], + new_fragments: vec![Fragment::new(3)], + fields_modified: vec![], + compacted_sstables: Vec::new(), + fields_for_preserving_frag_bitmap: vec![], + update_mode: Some(RewriteRows), + inserted_rows_filter: None, + updated_fragment_offsets: None, + } + } else { + Operation::Append { + fragments: vec![Fragment::new(3)], + } + }; + let other = Transaction::new(0, other_operation, None); + + assert!(matches!( + rebase.check_txn(&other, 1), + Err(Error::RetryableCommitConflict { .. }) + )); + } + #[test] fn test_conflicts() { use io::commit::conflict_resolver::tests::{ConflictResult::*, modified_fragment_ids}; @@ -3875,7 +3960,7 @@ mod tests { frag_reuse_index: None, }, [ - Compatible, // append + Retryable, // append Retryable, // create index Compatible, // delete Retryable, // merge @@ -3897,7 +3982,7 @@ mod tests { frag_reuse_index: None, }, [ - Compatible, // append + Retryable, // append Retryable, // create index Retryable, // delete Retryable, // merge @@ -7586,16 +7671,24 @@ mod tests { /// An on-disk tagged fixture (fragments 10 and 11, four rows each, /// values 0..8) plus two two-row appended fragments (100..102 and - /// 102..104): with `target_rows_per_fragment: 4` only the appended - /// fragments are compaction candidates, so a real `compact_files` - /// stays disjoint from a stable-partition rewrite of 10 and 11. - async fn disk_tagged_fixture_with_small_fragments(uri: &str) -> Dataset { + /// 102..104). The appended fragments have higher IDs, so compacting + /// their suffix leaves the tagged fragments' identities intact. + async fn disk_tagged_fixture_with_small_fragments( + uri: &str, + extra_rows: Option>, + ) -> Dataset { // Order matters: the two-row fragments are appended and indexed // BEFORE the tagging rewrite, so deferred compaction of them // must record a transition (uncovered fragments would commit a // plain rewrite instead; see // `uncovered_deferred_compaction_commits_plain_rewrite`). - let dataset = disk_fixture(uri, 2, 4).await; + let mut dataset = disk_fixture(uri, 2, 4).await; + reserve(&mut dataset, 40).await; + let dataset = if let Some(rows) = extra_rows { + append_rows(&dataset, rows).await + } else { + dataset + }; let dataset = append_rows(&dataset, 100..102).await; let mut dataset = append_rows(&dataset, 102..104).await; dataset @@ -7608,7 +7701,6 @@ mod tests { ) .await .unwrap(); - reserve(&mut dataset, 40).await; let old_fragments: Vec = dataset .fragments() .iter() @@ -7918,7 +8010,7 @@ mod tests { #[tokio::test] async fn real_deferred_compaction_rebases_over_concurrent_sp_fresh_session() { let dir = TempStrDir::default(); - disk_tagged_fixture_with_small_fragments(dir.as_str()).await; + disk_tagged_fixture_with_small_fragments(dir.as_str(), None).await; // The compactor opens BEFORE the rewrite commits: its plan and // its commit read-version anchor at the pre-rewrite snapshot. @@ -7931,7 +8023,7 @@ mod tests { .filter(|f| f.id == 10 || f.id == 11) .cloned() .collect(); - let (b_transition, b_destinations) = prepare_partition(&sp_writer, &[10, 11], 50).await; + let (b_transition, b_destinations) = prepare_partition(&sp_writer, &[10, 11], 20).await; let read_version = sp_writer.manifest.version; commit_sp( &sp_writer, @@ -7967,7 +8059,7 @@ mod tests { #[tokio::test] async fn sp_rebases_over_real_deferred_compaction_fresh_session() { let dir = TempStrDir::default(); - disk_tagged_fixture_with_small_fragments(dir.as_str()).await; + disk_tagged_fixture_with_small_fragments(dir.as_str(), None).await; let sp_writer = fresh_session(dir.as_str()).await; let b_old: Vec = sp_writer @@ -8012,7 +8104,7 @@ mod tests { #[tokio::test] async fn tagged_compactions_disjoint_both_land_fresh_session() { let dir = TempStrDir::default(); - disk_tagged_fixture_with_small_fragments(dir.as_str()).await; + disk_tagged_fixture_with_small_fragments(dir.as_str(), None).await; let oc_writer = fresh_session(dir.as_str()).await; let b_old: Vec = oc_writer @@ -8058,7 +8150,7 @@ mod tests { #[tokio::test] async fn tagged_compactions_overlapping_rejected_fresh_session() { let dir = TempStrDir::default(); - disk_tagged_fixture_with_small_fragments(dir.as_str()).await; + disk_tagged_fixture_with_small_fragments(dir.as_str(), None).await; let mut loser = fresh_session(dir.as_str()).await; let mut winner = fresh_session(dir.as_str()).await; @@ -8153,7 +8245,7 @@ mod tests { #[tokio::test] async fn overlapping_compaction_and_sp_rejected_fresh_session() { let dir = TempStrDir::default(); - disk_tagged_fixture_with_small_fragments(dir.as_str()).await; + disk_tagged_fixture_with_small_fragments(dir.as_str(), None).await; let sp_writer = fresh_session(dir.as_str()).await; let small_ids: Vec = sp_writer @@ -8204,11 +8296,25 @@ mod tests { #[tokio::test] async fn sp_compaction_sp_three_way_disjoint_merge() { let dir = TempStrDir::default(); - let dataset = disk_tagged_fixture_with_small_fragments(dir.as_str()).await; - // A fourth region for the second stable-partition writer. - let mut dataset = append_rows(&dataset, 200..204).await; + // A fourth region before the compactable suffix for the second + // stable-partition writer. + let mut dataset = + disk_tagged_fixture_with_small_fragments(dir.as_str(), Some(200..204)).await; reserve(&mut dataset, 100).await; - let extra_id = dataset.fragments().last().unwrap().id; + let small_id = dataset + .fragments() + .iter() + .find(|fragment| fragment.physical_rows == Some(2)) + .unwrap() + .id; + let extra_id = dataset + .fragments() + .iter() + .find(|fragment| { + fragment.id != 10 && fragment.id != 11 && fragment.physical_rows == Some(4) + }) + .unwrap() + .id; // Both racing writers open before the first rewrite commits. let mut compactor = fresh_session(dir.as_str()).await; @@ -8232,7 +8338,7 @@ mod tests { .filter(|f| f.id == 10 || f.id == 11) .cloned() .collect(); - let (w1_transition, w1_destinations) = prepare_partition(&dataset, &[10, 11], 80).await; + let (w1_transition, w1_destinations) = prepare_partition(&dataset, &[10, 11], 20).await; let w1_version = dataset.manifest.version; commit_sp( &dataset, @@ -8268,7 +8374,7 @@ mod tests { // Fixture rewrite + writers 1 and 3, plus the compaction. assert_eq!(count_mappings(&ledger), (3, 1)); assert!(ledger.consumer(10).is_some()); - assert!(ledger.consumer(2).is_some()); + assert!(ledger.consumer(u32::try_from(small_id).unwrap()).is_some()); assert!(ledger.consumer(extra_id as u32).is_some()); let mut expected: Vec = (0..8).collect(); expected.extend(100..104);