From b9fdc8e6e9e445c1aba416879330970d48890827 Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Fri, 7 Aug 2026 13:35:39 +0000 Subject: [PATCH 01/10] fix(compaction): preserve task order during commit --- rust/lance/src/dataset/optimize.rs | 106 ++++++++++++++++++++--------- 1 file changed, 72 insertions(+), 34 deletions(-) diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index b487c803c78..fbe1abc99e2 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -2087,6 +2087,19 @@ pub async fn commit_compaction( .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(), + )); + } + + // Rewrite tasks finish in an arbitrary order, but new fragment ids determine + // their order in the manifest. Reserve ids in the original fragment order so + // parallel and distributed compaction preserve the dataset's row order. + completed_tasks.sort_unstable_by_key(|task| task.original_fragments[0].id); // Collect the rewritten fragments' file paths up front so every failure // path below can clean them up (or deliberately keep them). Fragment ids @@ -2505,14 +2518,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] @@ -2950,9 +2955,7 @@ mod tests { .unwrap(); let first_new_frag_idx = 7; - // Predicting the remap is difficult. One task will remap to fragments 7/8 and the other - // will remap to fragments 9/10 but we don't know which is which and so we just allow ourselves - // to expect both possibilities. + // Tasks may finish in any order, but commit assigns fragment ids in plan order. let remap_a = expect_remap( &[ vec![ @@ -2976,26 +2979,6 @@ mod tests { ], first_new_frag_idx, ); - let remap_b = expect_remap( - &[ - // Frags 4, 5, and 6 are rewritten to frags 7 & 8 - vec![ - (row_addrs(4, 0..200), true), - (row_addrs(4, 200..400), false), - (row_addrs(4, 400..1000), true), - (row_addrs(5, 0..200), true), - ], - vec![(row_addrs(5, 200..300), true), (row_addrs(6, 0..300), true)], - // 3 small fragments rewritten to frags 9 & 10 - vec![ - (row_addrs(0, 0..400), true), - (row_addrs(1, 0..400), true), - (row_addrs(2, 0..200), true), - ], - vec![(row_addrs(2, 200..400), true)], - ], - first_new_frag_idx, - ); // Create compaction plan let options = CompactionOptions { @@ -3024,10 +3007,8 @@ mod tests { vec![4, 5, 6] ); - let mock_remapper = MockIndexRemapper::in_any_order(&[remap_a, remap_b]); - // 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(); @@ -3045,6 +3026,63 @@ mod tests { assert_eq!(fragment_ids, vec![3, 7, 8, 9, 10]); } + #[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] #[tokio::test] async fn test_compact_data_files( From 46c174434cd4eef750c13655b6b08ccc9307b052 Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Fri, 7 Aug 2026 14:34:13 +0000 Subject: [PATCH 02/10] fix(compaction): preserve bounded compaction order --- python/python/lance/dataset.py | 12 +- rust/lance/src/dataset/optimize.rs | 117 +++++++++++++++++- rust/lance/src/io/commit/conflict_resolver.rs | 14 ++- 3 files changed, 127 insertions(+), 16 deletions(-) diff --git a/python/python/lance/dataset.py b/python/python/lance/dataset.py index d387880be48..1edffd669df 100644 --- a/python/python/lance/dataset.py +++ b/python/python/lance/dataset.py @@ -7108,12 +7108,12 @@ def compact_files( * Removes dropped columns from fragments * Merges small fragments into larger ones - This method preserves the insertion order of the dataset. This may mean - it leaves small fragments in the dataset if they are not adjacent to - other fragments that need compaction. For example, if you have fragments - with row counts 5 million, 100, and 5 million, the middle fragment will - not be compacted because the fragments it is adjacent to do not need - compaction. + Row order is preserved when compaction rewrites a contiguous suffix, + including the entire dataset. In other partial or gapped plans, each + rewritten range keeps its internal order, but the ranges move after + untouched fragments because replacement fragment IDs are allocated + monotonically. ``max_source_fragments`` limits work to a suffix so that + bounded compaction preserves row order. Default values for these options can be stored in the dataset manifest config using keys prefixed with ``lance.compaction.``. For example, diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index fbe1abc99e2..ac9493f5ca4 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -808,9 +808,30 @@ impl CompactionPlanner for DefaultCompactionPlanner { candidate_bins.push(bin); } - let all_tasks: Vec = candidate_bins + let compactable_bins = if self.options.max_source_fragments.is_some() { + // New fragment ids are allocated above the manifest high-water mark, + // so a bounded replacement before an untouched fragment would move + // behind it. Select a contiguous suffix when bounding the work. + let mut suffix_start = i; + let mut suffix_bins = Vec::new(); + for bin in candidate_bins.into_iter().rev() { + if bin.pos_range.end != suffix_start || bin.is_noop() { + break; + } + suffix_start = bin.pos_range.start; + suffix_bins.push(bin); + } + suffix_bins.reverse(); + suffix_bins + } else { + candidate_bins + .into_iter() + .filter(|bin| !bin.is_noop()) + .collect() + }; + + let all_tasks: Vec = compactable_bins .into_iter() - .filter(|bin| !bin.is_noop()) .flat_map(|bin| bin.split_for_size(self.options.target_rows_per_fragment)) .map(|bin| TaskData { fragments: bin.fragments, @@ -819,13 +840,16 @@ impl CompactionPlanner for DefaultCompactionPlanner { let tasks = if let Some(max_frags) = self.options.max_source_fragments { let mut total_frags = 0; - all_tasks + let mut tasks = all_tasks .into_iter() + .rev() .take_while(|task| { total_frags += task.fragments.len(); total_frags <= max_frags }) - .collect() + .collect::>(); + tasks.reverse(); + tasks } else { all_tasks }; @@ -838,14 +862,17 @@ impl CompactionPlanner for DefaultCompactionPlanner { } } -/// Compacts the files in the dataset without reordering them. +/// Compacts files in the dataset. /// /// By default, this does a few things: /// * Removes deleted rows from fragments. /// * Removes dropped columns from fragments. /// * Merges fragments that are too small. /// -/// This method tries to preserve the insertion order of rows in the dataset. +/// Row order is preserved when the selected fragments form a contiguous suffix, +/// including a full-dataset compaction. Within other partial or gapped plans, +/// rewritten ranges retain their relative order but move after untouched +/// fragments because fragment ids are monotonically allocated. /// /// If no compaction is needed, this method will not make a new version of the table. pub async fn compact_files( @@ -2045,6 +2072,11 @@ async fn recalc_versions_for_rewritten_fragments( /// 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 fragments. This preserves row +/// order when they collectively replace a contiguous suffix. A partial or +/// gapped set is appended after untouched fragments because replacement ids are +/// monotonically allocated. pub async fn commit_compaction( dataset: &mut Dataset, completed_tasks: Vec, @@ -7086,6 +7118,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 ({})", @@ -7098,6 +7141,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, @@ -7117,6 +7162,8 @@ 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, @@ -7124,6 +7171,64 @@ mod tests { ); } + #[rstest] + #[tokio::test] + async fn test_bounded_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(); + let options = CompactionOptions { + target_rows_per_fragment: 250, + max_source_fragments: Some(2), + ..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![3, 4]); + + compact_files(&mut dataset, options, None).await.unwrap(); + + let compacted = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(compacted, data.slice(0, offset)); + } + #[tokio::test] async fn test_compaction_uses_manifest_config() { let test_dir = TempStrDir::default(); diff --git a/rust/lance/src/io/commit/conflict_resolver.rs b/rust/lance/src/io/commit/conflict_resolver.rs index d28b1cc9882..00fcb5d243d 100644 --- a/rust/lance/src/io/commit/conflict_resolver.rs +++ b/rust/lance/src/io/commit/conflict_resolver.rs @@ -773,8 +773,14 @@ impl<'a> TransactionRebase<'a> { 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 { .. } + // A rewrite allocates replacement fragment ids above the current + // high-water mark. If an append landed after the rewrite was + // planned, rebasing would place the replacement after the newly + // appended rows and violate insertion order. + Operation::Append { .. } => { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } + Operation::ReserveFragments { .. } | Operation::Project { .. } | Operation::Clone { .. } | Operation::UpdateConfig { .. } @@ -2837,7 +2843,7 @@ mod tests { frag_reuse_index: None, }, [ - Compatible, // append + Retryable, // append Retryable, // create index Compatible, // delete Retryable, // merge @@ -2859,7 +2865,7 @@ mod tests { frag_reuse_index: None, }, [ - Compatible, // append + Retryable, // append Retryable, // create index Retryable, // delete Retryable, // merge From e47e8df74f858e136ddef4e462c692df261b97b2 Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Fri, 7 Aug 2026 17:09:26 +0000 Subject: [PATCH 03/10] fix(compaction): preserve logical fragment order --- python/python/lance/dataset.py | 12 +- rust/lance/src/dataset.rs | 68 +++++----- rust/lance/src/dataset/optimize.rs | 119 ++++++++++-------- rust/lance/src/dataset/transaction.rs | 114 ++++++++++++++++- rust/lance/src/dataset/write/commit.rs | 6 +- rust/lance/src/io/commit.rs | 12 +- rust/lance/src/io/commit/conflict_resolver.rs | 51 ++++++-- 7 files changed, 268 insertions(+), 114 deletions(-) diff --git a/python/python/lance/dataset.py b/python/python/lance/dataset.py index 1edffd669df..d387880be48 100644 --- a/python/python/lance/dataset.py +++ b/python/python/lance/dataset.py @@ -7108,12 +7108,12 @@ def compact_files( * Removes dropped columns from fragments * Merges small fragments into larger ones - Row order is preserved when compaction rewrites a contiguous suffix, - including the entire dataset. In other partial or gapped plans, each - rewritten range keeps its internal order, but the ranges move after - untouched fragments because replacement fragment IDs are allocated - monotonically. ``max_source_fragments`` limits work to a suffix so that - bounded compaction preserves row order. + This method preserves the insertion order of the dataset. This may mean + it leaves small fragments in the dataset if they are not adjacent to + other fragments that need compaction. For example, if you have fragments + with row counts 5 million, 100, and 5 million, the middle fragment will + not be compacted because the fragments it is adjacent to do not need + compaction. Default values for these options can be stored in the dataset manifest config using keys prefixed with ``lance.compaction.``. For example, diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index b98a2d1ee11..e81019f3264 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -183,6 +183,9 @@ pub struct Dataset { // Bitmap of fragment ids in this dataset. pub(crate) fragment_bitmap: Arc, + // Manifest positions indexed by ascending fragment-id rank. Fragment ids are + // stable identities, while manifest position defines logical row order. + fragment_indices_by_id: Arc>, // These are references to session caches, but with the dataset URI as a prefix. pub(crate) index_cache: Arc, @@ -210,6 +213,22 @@ impl std::fmt::Debug for Dataset { } } +fn build_fragment_lookup(fragments: &[Fragment]) -> (Arc, Arc>) { + let fragment_bitmap: Arc = Arc::new( + fragments + .iter() + .map(|fragment| fragment.id as u32) + .collect(), + ); + let mut fragment_indices_by_id = vec![0; fragments.len()]; + for (manifest_index, fragment) in fragments.iter().enumerate() { + let id_rank = fragment_bitmap.rank(fragment.id as u32) as usize; + debug_assert!(id_rank > 0); + fragment_indices_by_id[id_rank - 1] = manifest_index; + } + (fragment_bitmap, Arc::new(fragment_indices_by_id)) +} + /// Dataset Version #[derive(Deserialize, Serialize, Debug)] pub struct Version { @@ -492,13 +511,8 @@ impl Dataset { let (manifest, manifest_location) = self.latest_manifest().await?; self.manifest = manifest; self.manifest_location = manifest_location; - self.fragment_bitmap = Arc::new( - self.manifest - .fragments - .iter() - .map(|f| f.id as u32) - .collect(), - ); + (self.fragment_bitmap, self.fragment_indices_by_id) = + build_fragment_lookup(&self.manifest.fragments); Ok(()) } @@ -831,7 +845,7 @@ impl Dataset { ); let metadata_cache = Arc::new(session.metadata_cache.for_dataset(&uri)); let index_cache = Arc::new(session.index_cache.for_dataset(&uri)); - let fragment_bitmap = Arc::new(manifest.fragments.iter().map(|f| f.id as u32).collect()); + let (fragment_bitmap, fragment_indices_by_id) = build_fragment_lookup(&manifest.fragments); write::log_unregistered_base_scoped_options( store_params.as_ref(), &manifest.base_paths, @@ -847,6 +861,7 @@ impl Dataset { session, refs, fragment_bitmap, + fragment_indices_by_id, metadata_cache, index_cache, file_reader_options, @@ -1626,13 +1641,8 @@ impl Dataset { self.manifest = Arc::new(manifest); self.manifest_location = manifest_location; - self.fragment_bitmap = Arc::new( - self.manifest - .fragments - .iter() - .map(|f| f.id as u32) - .collect(), - ); + (self.fragment_bitmap, self.fragment_indices_by_id) = + build_fragment_lookup(&self.manifest.fragments); Ok(()) } @@ -2714,7 +2724,9 @@ impl Dataset { Projection::full(self.clone()) } - /// Get fragments. + /// Get fragments in logical row order. + /// + /// Fragment ids are stable identities and are not guaranteed to be sorted. pub fn get_fragments(&self) -> Vec { let dataset = Arc::new(self.clone()); self.manifest @@ -2724,7 +2736,8 @@ impl Dataset { .collect() } - /// Iterate over manifest fragments without allocating [`FileFragment`] wrappers. + /// Iterate over manifest fragments in logical row order without allocating + /// [`FileFragment`] wrappers. pub fn iter_fragments(&self) -> impl Iterator { self.manifest.fragments.iter() } @@ -2835,11 +2848,12 @@ impl Dataset { if !self.fragment_bitmap.contains(*id) { return None; } - let fragment_index = self.fragment_bitmap.rank(*id) as usize - 1; + let id_rank = self.fragment_bitmap.rank(*id) as usize - 1; + let fragment_index = self.fragment_indices_by_id[id_rank]; let fragment = self.manifest.fragments.get(fragment_index)?; debug_assert_eq!( fragment.id, *id as u64, - "fragment_bitmap rank({id}) resolved to fragment {}, but fragment_bitmap and manifest.fragments are expected to stay in sync", + "fragment lookup for id {id} resolved to fragment {}, but the lookup and manifest.fragments are expected to stay in sync", fragment.id ); Some(FileFragment::new(dataset.clone(), fragment.clone())) @@ -3025,22 +3039,6 @@ impl Dataset { } } - // Fragments are sorted in increasing fragment id order - self.manifest - .fragments - .iter() - .map(|f| f.id) - .try_fold(0, |prev, id| { - if id < prev { - Err(Error::corrupt_file(self.base.clone(), format!( - "Fragment ids are not sorted in increasing fragment-id order. Found {} after {} in dataset {:?}", - id, prev, self.base - ))) - } else { - Ok(id) - } - })?; - // All fragments have equal lengths futures::stream::iter(self.get_fragments()) .map(|f| async move { f.validate().await }) diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index ac9493f5ca4..647451205ea 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -703,14 +703,9 @@ impl CompactionPlanner for DefaultCompactionPlanner { )); } - // get_fragments should be returning fragments in sorted order (by id) - // and fragment ids should be unique + // Manifest order is logical row order. Fragment ids are stable identities + // and may be out of order after a rewrite. let fragments = dataset.get_fragments(); - - debug_assert!( - fragments.windows(2).all(|w| w[0].id() < w[1].id()), - "fragments in manifest are not sorted" - ); let mut fragment_metrics = futures::stream::iter(fragments) .map(|fragment| async move { match collect_metrics(&fragment).await { @@ -808,30 +803,9 @@ impl CompactionPlanner for DefaultCompactionPlanner { candidate_bins.push(bin); } - let compactable_bins = if self.options.max_source_fragments.is_some() { - // New fragment ids are allocated above the manifest high-water mark, - // so a bounded replacement before an untouched fragment would move - // behind it. Select a contiguous suffix when bounding the work. - let mut suffix_start = i; - let mut suffix_bins = Vec::new(); - for bin in candidate_bins.into_iter().rev() { - if bin.pos_range.end != suffix_start || bin.is_noop() { - break; - } - suffix_start = bin.pos_range.start; - suffix_bins.push(bin); - } - suffix_bins.reverse(); - suffix_bins - } else { - candidate_bins - .into_iter() - .filter(|bin| !bin.is_noop()) - .collect() - }; - - let all_tasks: Vec = compactable_bins + let all_tasks: Vec = candidate_bins .into_iter() + .filter(|bin| !bin.is_noop()) .flat_map(|bin| bin.split_for_size(self.options.target_rows_per_fragment)) .map(|bin| TaskData { fragments: bin.fragments, @@ -840,16 +814,13 @@ impl CompactionPlanner for DefaultCompactionPlanner { let tasks = if let Some(max_frags) = self.options.max_source_fragments { let mut total_frags = 0; - let mut tasks = all_tasks + all_tasks .into_iter() - .rev() .take_while(|task| { total_frags += task.fragments.len(); total_frags <= max_frags }) - .collect::>(); - tasks.reverse(); - tasks + .collect() } else { all_tasks }; @@ -862,17 +833,14 @@ impl CompactionPlanner for DefaultCompactionPlanner { } } -/// Compacts files in the dataset. +/// Compacts the files in the dataset without reordering them. /// /// By default, this does a few things: /// * Removes deleted rows from fragments. /// * Removes dropped columns from fragments. /// * Merges fragments that are too small. /// -/// Row order is preserved when the selected fragments form a contiguous suffix, -/// including a full-dataset compaction. Within other partial or gapped plans, -/// rewritten ranges retain their relative order but move after untouched -/// fragments because fragment ids are monotonically allocated. +/// This method tries to preserve the insertion order of rows in the dataset. /// /// If no compaction is needed, this method will not make a new version of the table. pub async fn compact_files( @@ -2073,10 +2041,8 @@ async fn recalc_versions_for_rewritten_fragments( /// 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 fragments. This preserves row -/// order when they collectively replace a contiguous suffix. A partial or -/// gapped set is appended after untouched fragments because replacement ids are -/// monotonically allocated. +/// Completed tasks are ordered by their source fragments' current logical +/// positions so partial and gapped result sets replace their ranges in place. pub async fn commit_compaction( dataset: &mut Dataset, completed_tasks: Vec, @@ -2128,10 +2094,21 @@ pub async fn commit_compaction( )); } - // Rewrite tasks finish in an arbitrary order, but new fragment ids determine - // their order in the manifest. Reserve ids in the original fragment order so - // parallel and distributed compaction preserve the dataset's row order. - completed_tasks.sort_unstable_by_key(|task| task.original_fragments[0].id); + // Rewrite tasks finish in an arbitrary order. Apply their groups in current + // manifest order so every replacement is spliced into its logical position. + let fragment_positions = dataset + .manifest + .fragments + .iter() + .enumerate() + .map(|(position, fragment)| (fragment.id, position)) + .collect::>(); + completed_tasks.sort_by_key(|task| { + fragment_positions + .get(&task.original_fragments[0].id) + .copied() + .map_or((1, usize::MAX), |position| (0, position)) + }); // Collect the rewritten fragments' file paths up front so every failure // path below can clean them up (or deliberately keep them). Fragment ids @@ -3038,6 +3015,7 @@ mod tests { .collect::>(), vec![4, 5, 6] ); + let expected_data = dataset.scan().try_into_batch().await.unwrap(); // Run compaction let metrics = compact_files(&mut dataset, options, Some(Arc::new(remap_a))) @@ -3055,7 +3033,9 @@ mod tests { .iter() .map(|f| f.id()) .collect::>(); - assert_eq!(fragment_ids, vec![3, 7, 8, 9, 10]); + assert_eq!(fragment_ids, vec![7, 8, 3, 9, 10]); + let compacted_data = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(compacted_data, expected_data); } #[rstest] @@ -7123,8 +7103,7 @@ mod tests { .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..] + let expected_fragment_ids = dataset.manifest.fragments[..bounded_source_frags] .iter() .map(|fragment| fragment.id) .collect::>(); @@ -7166,8 +7145,8 @@ mod tests { 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}" ); } @@ -7221,12 +7200,44 @@ mod tests { .iter() .flat_map(|task| task.fragments.iter().map(|fragment| fragment.id)) .collect::>(); - assert_eq!(planned_fragment_ids, vec![3, 4]); + assert_eq!(planned_fragment_ids, vec![0, 1]); + + compact_files(&mut dataset, options.clone(), None) + .await + .unwrap(); + + assert_eq!(dataset.get_fragments().len(), 4); + let logical_fragment_ids = dataset + .manifest + .fragments + .iter() + .map(|fragment| fragment.id) + .collect::>(); + assert_eq!(logical_fragment_ids, vec![5, 2, 3, 4]); + let fragments_by_id = dataset + .get_fragments_from_ids(&[5, 2, 4, 3]) + .unwrap() + .into_iter() + .map(|fragment| fragment.id()) + .collect::>(); + assert_eq!(fragments_by_id, vec![2, 3, 4, 5]); + dataset.validate().await.unwrap(); + let after_first = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(after_first, data.slice(0, offset)); + + let second_plan = plan_compaction(&dataset, &options).await.unwrap(); + let second_fragment_ids = second_plan + .tasks() + .iter() + .flat_map(|task| task.fragments.iter().map(|fragment| fragment.id)) + .collect::>(); + assert_eq!(second_fragment_ids, vec![3, 4]); compact_files(&mut dataset, options, None).await.unwrap(); let compacted = dataset.scan().try_into_batch().await.unwrap(); assert_eq!(compacted, data.slice(0, offset)); + assert_eq!(dataset.get_fragments().len(), 3); } #[tokio::test] diff --git a/rust/lance/src/dataset/transaction.rs b/rust/lance/src/dataset/transaction.rs index be757867d0e..76257b06ff0 100644 --- a/rust/lance/src/dataset/transaction.rs +++ b/rust/lance/src/dataset/transaction.rs @@ -2397,6 +2397,28 @@ impl Transaction { final_fragments.extend(unmodified_fragments); + // Data replacement changes files, not row placement. Reassemble + // the fragments in their existing logical order without using + // fragment ids as ordering keys. + let mut fragments_by_id = HashMap::with_capacity(final_fragments.len()); + for fragment in final_fragments.drain(..) { + if fragments_by_id.insert(fragment.id, fragment).is_some() { + return Err(Error::invalid_input( + "DataReplacement contains multiple replacements for the same fragment" + .to_string(), + )); + } + } + for existing_fragment in existing_fragments { + let Some(fragment) = fragments_by_id.remove(&existing_fragment.id) else { + return Err(Error::internal(format!( + "DataReplacement lost fragment {} while preserving logical order", + existing_fragment.id + ))); + }; + final_fragments.push(fragment); + } + // 5. Invalidate index bitmaps for replaced fields let modified_fragments: Vec = final_fragments .iter() @@ -2487,9 +2509,6 @@ impl Transaction { } }; - // If a fragment was reserved then it may not belong at the end of the fragments list. - final_fragments.sort_by_key(|frag| frag.id); - // Clean up data files that only contain tombstoned fields Self::remove_tombstoned_data_files(&mut final_fragments); @@ -4307,6 +4326,95 @@ mod tests { assert_eq!(final_fragments, expected_fragments); } + #[test] + fn test_rewrite_build_manifest_preserves_logical_fragment_order() { + 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", + &ManifestWriteConfig::default(), + ) + .unwrap(); + + let fragment_ids = rewritten + .fragments + .iter() + .map(|fragment| fragment.id) + .collect::>(); + assert_eq!(fragment_ids, vec![0, 4, 3]); + } + + #[test] + fn test_rewrite_after_row_adding_update_preserves_logical_order() { + 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", + &ManifestWriteConfig::default(), + ) + .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", + &ManifestWriteConfig::default(), + ) + .unwrap(); + + let fragment_ids = rewritten + .fragments + .iter() + .map(|fragment| fragment.id) + .collect::>(); + assert_eq!(fragment_ids, vec![4, 1, 2, 3]); + } + #[test] fn test_merge_fragments_valid() { // Create a simple schema for testing diff --git a/rust/lance/src/dataset/write/commit.rs b/rust/lance/src/dataset/write/commit.rs index fd57f68f9e7..6a3649c5533 100644 --- a/rust/lance/src/dataset/write/commit.rs +++ b/rust/lance/src/dataset/write/commit.rs @@ -17,7 +17,7 @@ use crate::io::commit::DEFAULT_COMMIT_RETRY_TIMEOUT; use crate::{ Dataset, Error, Result, dataset::{ - ManifestWriteConfig, ReadParams, + ManifestWriteConfig, ReadParams, build_fragment_lookup, builder::DatasetBuilder, commit_detached_transaction, commit_new_dataset, commit_transaction, refs::Refs, @@ -473,7 +473,7 @@ impl<'a> CommitBuilder<'a> { operation=&transaction.operation.name() ); - let fragment_bitmap = Arc::new(manifest.fragments.iter().map(|f| f.id as u32).collect()); + let (fragment_bitmap, fragment_indices_by_id) = build_fragment_lookup(&manifest.fragments); match &self.dest { WriteDestination::Dataset(dataset) => Ok(Dataset { @@ -481,6 +481,7 @@ impl<'a> CommitBuilder<'a> { manifest_location, session, fragment_bitmap, + fragment_indices_by_id, ..dataset.as_ref().clone() }), WriteDestination::Uri(uri) => { @@ -505,6 +506,7 @@ impl<'a> CommitBuilder<'a> { refs, index_cache, fragment_bitmap, + fragment_indices_by_id, metadata_cache, file_reader_options: None, store_params: self.store_params.clone().map(Box::new), diff --git a/rust/lance/src/io/commit.rs b/rust/lance/src/io/commit.rs index 186afff3468..c0e57a066e5 100644 --- a/rust/lance/src/io/commit.rs +++ b/rust/lance/src/io/commit.rs @@ -706,16 +706,20 @@ fn check_storage_version(manifest: &mut Manifest) -> Result<()> { /// 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 ))); } Ok(()) diff --git a/rust/lance/src/io/commit/conflict_resolver.rs b/rust/lance/src/io/commit/conflict_resolver.rs index 00fcb5d243d..e3408da045c 100644 --- a/rust/lance/src/io/commit/conflict_resolver.rs +++ b/rust/lance/src/io/commit/conflict_resolver.rs @@ -773,14 +773,8 @@ impl<'a> TransactionRebase<'a> { match &other_transaction.operation { // Rewrite is only compatible with operations that don't touch // existing fragments or update fragments we don't touch. - // A rewrite allocates replacement fragment ids above the current - // high-water mark. If an append landed after the rewrite was - // planned, rebasing would place the replacement after the newly - // appended rows and violate insertion order. - Operation::Append { .. } => { - Err(self.retryable_conflict_err(other_transaction, other_version)) - } - Operation::ReserveFragments { .. } + Operation::Append { .. } + | Operation::ReserveFragments { .. } | Operation::Project { .. } | Operation::Clone { .. } | Operation::UpdateConfig { .. } @@ -2649,6 +2643,43 @@ mod tests { Retryable, } + #[test] + fn test_rewrite_is_compatible_with_row_adding_update() { + 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(), + }; + let other = Transaction::new( + 0, + Operation::Update { + removed_fragment_ids: vec![], + updated_fragments: vec![], + new_fragments: vec![Fragment::new(1)], + 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, + }, + None, + ); + + assert!(rebase.check_txn(&other, 1).is_ok()); + } + #[test] fn test_conflicts() { use io::commit::conflict_resolver::tests::{ConflictResult::*, modified_fragment_ids}; @@ -2843,7 +2874,7 @@ mod tests { frag_reuse_index: None, }, [ - Retryable, // append + Compatible, // append Retryable, // create index Compatible, // delete Retryable, // merge @@ -2865,7 +2896,7 @@ mod tests { frag_reuse_index: None, }, [ - Retryable, // append + Compatible, // append Retryable, // create index Retryable, // delete Retryable, // merge From 3e8e101fb641261f396e789335bbd14306ebb699 Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Fri, 7 Aug 2026 18:12:46 +0000 Subject: [PATCH 04/10] fix(compaction): gate logical fragment ordering --- docs/src/format/table/versioning.md | 5 +- protos/table.proto | 2 + rust/lance-table/src/feature_flags.rs | 58 ++++++++++++- rust/lance-table/src/format/manifest.rs | 15 +++- rust/lance/src/dataset.rs | 20 +++++ rust/lance/src/dataset/optimize.rs | 10 +++ rust/lance/src/dataset/transaction.rs | 104 +++++++++++++++++++----- rust/lance/src/io/commit.rs | 37 ++++++++- 8 files changed, 224 insertions(+), 27 deletions(-) diff --git a/docs/src/format/table/versioning.md b/docs/src/format/table/versioning.md index 745dd1ccd87..92e110ea725 100644 --- a/docs/src/format/table/versioning.md +++ b/docs/src/format/table/versioning.md @@ -28,7 +28,10 @@ they should return an "unsupported" error on any read or write operation. | 4 | `FLAG_USE_V2_FORMAT_DEPRECATED` | No | No | Files are written with the new v2 format. This flag is deprecated and no longer used. | | 8 | `FLAG_TABLE_CONFIG` | No | Yes | Table config is present in the manifest. | | 16 | `FLAG_BASE_PATHS` | Yes | Yes | Dataset uses multiple base paths (for shallow clones or multi-base datasets). | +| 32 | `FLAG_DISABLE_TRANSACTION_FILE` | No | Yes | Transaction files are stored inline in manifests instead of under `_transaction`. | +| 64 | `FLAG_UNSTABLE_DATA_OVERLAY_FILES` | Yes | Yes | Fragments contain data overlay files; this feature is not yet released. | +| 128 | `FLAG_LOGICAL_FRAGMENT_ORDER` | Yes | Yes | Manifest position defines logical row order, so fragments must not be reordered by ID. | -Flags with bit values 32 and above are unknown and will cause implementations to reject the dataset with an "unsupported" error. +Flags with bit values 256 and above are unknown and will cause implementations to reject the dataset with an "unsupported" error. diff --git a/protos/table.proto b/protos/table.proto index 9a64230f40f..157ffba40ec 100644 --- a/protos/table.proto +++ b/protos/table.proto @@ -118,6 +118,8 @@ message Manifest { // * 1 << 6: data overlay files are present (see DataOverlayFile). Readers that do // not understand overlays must refuse the dataset, since ignoring an overlay // would silently return stale base values. + // * 1 << 7: fragments are stored in logical row order rather than fragment-ID + // order. Readers and writers must preserve their manifest positions. uint64 reader_feature_flags = 9; // Feature flags for writers. diff --git a/rust/lance-table/src/feature_flags.rs b/rust/lance-table/src/feature_flags.rs index 41b8e415f8e..fc31f2a9db5 100644 --- a/rust/lance-table/src/feature_flags.rs +++ b/rust/lance-table/src/feature_flags.rs @@ -30,8 +30,12 @@ pub const FLAG_DISABLE_TRANSACTION_FILE: u64 = 32; /// unless [`ENABLE_UNSTABLE_DATA_OVERLAY_FILES_ENV`] is set, which lets benchmarks opt in. /// Debug builds always understand it so tests exercise the path. pub const FLAG_UNSTABLE_DATA_OVERLAY_FILES: u64 = 64; +/// Manifest fragments are stored in logical row order instead of fragment-ID +/// order. Readers must preserve the list order, and writers must not reorder the +/// fragments by ID. +pub const FLAG_LOGICAL_FRAGMENT_ORDER: u64 = 128; /// The first bit that is unknown as a feature flag -pub const FLAG_UNKNOWN: u64 = 128; +pub const FLAG_UNKNOWN: u64 = 256; /// Environment variable that opts a release build into reading and writing data /// overlay files before the feature is generally released. @@ -97,6 +101,15 @@ pub fn apply_feature_flags( manifest.writer_feature_flags |= FLAG_UNSTABLE_DATA_OVERLAY_FILES; } + let has_logical_fragment_order = manifest + .fragments + .windows(2) + .any(|fragments| fragments[0].id > fragments[1].id); + if has_logical_fragment_order { + manifest.reader_feature_flags |= FLAG_LOGICAL_FRAGMENT_ORDER; + manifest.writer_feature_flags |= FLAG_LOGICAL_FRAGMENT_ORDER; + } + if disable_transaction_file { manifest.writer_feature_flags |= FLAG_DISABLE_TRANSACTION_FILE; } @@ -161,6 +174,7 @@ mod tests { assert!(can_read_dataset(super::FLAG_TABLE_CONFIG)); assert!(can_read_dataset(super::FLAG_BASE_PATHS)); assert!(can_read_dataset(super::FLAG_DISABLE_TRANSACTION_FILE)); + assert!(can_read_dataset(super::FLAG_LOGICAL_FRAGMENT_ORDER)); // Overlay support is gated on the build profile / env opt-in, so the // flag is readable exactly when overlays are enabled (see // test_data_overlay_flag_release_gating for the full policy). @@ -237,6 +251,7 @@ mod tests { assert!(can_write_dataset(super::FLAG_TABLE_CONFIG)); assert!(can_write_dataset(super::FLAG_BASE_PATHS)); assert!(can_write_dataset(super::FLAG_DISABLE_TRANSACTION_FILE)); + assert!(can_write_dataset(super::FLAG_LOGICAL_FRAGMENT_ORDER)); // Overlay support is gated on the build profile / env opt-in, so the // flag is writable exactly when overlays are enabled (see // test_data_overlay_flag_release_gating for the full policy). @@ -305,4 +320,45 @@ mod tests { 0 ); } + + #[test] + fn test_logical_fragment_order_is_mixed_version_gated() { + use crate::format::{DataStorageFormat, Fragment}; + use arrow_schema::{Field as ArrowField, Schema as ArrowSchema}; + use lance_core::datatypes::Schema; + use std::collections::HashMap; + use std::sync::Arc; + + let arrow_schema = ArrowSchema::new(vec![ArrowField::new( + "id", + arrow_schema::DataType::Int64, + false, + )]); + let schema = Schema::try_from(&arrow_schema).unwrap(); + let mut manifest = Manifest::new( + schema, + Arc::new(vec![Fragment::new(5), Fragment::new(2)]), + DataStorageFormat::default(), + HashMap::new(), + ); + + apply_feature_flags(&mut manifest, false, false).unwrap(); + + assert_ne!( + manifest.reader_feature_flags & FLAG_LOGICAL_FRAGMENT_ORDER, + 0 + ); + assert_ne!( + manifest.writer_feature_flags & FLAG_LOGICAL_FRAGMENT_ORDER, + 0 + ); + assert!(can_read_dataset(manifest.reader_feature_flags)); + assert!(can_write_dataset(manifest.writer_feature_flags)); + + // Lance releases predating logical fragment order support only recognize + // lower feature bits, so they reject instead of misreading this manifest. + let released_supported_flags = FLAG_LOGICAL_FRAGMENT_ORDER - 1; + assert_ne!(manifest.reader_feature_flags & !released_supported_flags, 0); + assert_ne!(manifest.writer_feature_flags & !released_supported_flags, 0); + } } diff --git a/rust/lance-table/src/format/manifest.rs b/rust/lance-table/src/format/manifest.rs index 5543511bb95..1c9e9e48cea 100644 --- a/rust/lance-table/src/format/manifest.rs +++ b/rust/lance-table/src/format/manifest.rs @@ -18,7 +18,9 @@ use std::ops::Range; use std::sync::Arc; use super::Fragment; -use crate::feature_flags::{FLAG_STABLE_ROW_IDS, has_deprecated_v2_feature_flag}; +use crate::feature_flags::{ + FLAG_LOGICAL_FRAGMENT_ORDER, FLAG_STABLE_ROW_IDS, has_deprecated_v2_feature_flag, +}; use crate::format::fragment::DataFileFieldInterner; use crate::format::pb; use lance_core::cache::LanceCache; @@ -49,8 +51,9 @@ pub struct Manifest { /// Fragments, the pieces to build the dataset. /// - /// This list is stored in order, sorted by fragment id. However, the fragment id - /// sequence may have gaps. + /// This list is stored in logical row order when + /// [`FLAG_LOGICAL_FRAGMENT_ORDER`] is set. Otherwise it is sorted by fragment + /// id, though the fragment id sequence may have gaps. pub fragments: Arc>, /// The file position of the version aux data. @@ -496,6 +499,12 @@ impl Manifest { self.reader_feature_flags & FLAG_STABLE_ROW_IDS != 0 } + /// Whether manifest position, rather than fragment ID, defines logical row + /// order. + pub fn uses_logical_fragment_order(&self) -> bool { + self.reader_feature_flags & FLAG_LOGICAL_FRAGMENT_ORDER != 0 + } + /// Creates a serialized copy of the manifest, suitable for IPC or temp storage /// and can be used to create a dataset pub fn serialized(&self) -> Vec { diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index e81019f3264..ce564cae6be 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -3039,6 +3039,26 @@ impl Dataset { } } + if !self.manifest.uses_logical_fragment_order() { + self.manifest + .fragments + .iter() + .map(|fragment| fragment.id) + .try_fold(0, |previous_id, fragment_id| { + if fragment_id < previous_id { + Err(Error::corrupt_file( + self.base.clone(), + format!( + "Fragment ids are not sorted in increasing fragment-id order, but the logical fragment order feature is not set. Found {fragment_id} after {previous_id} in dataset {:?}", + self.base + ), + )) + } else { + Ok(fragment_id) + } + })?; + } + // All fragments have equal lengths futures::stream::iter(self.get_fragments()) .map(|f| async move { f.validate().await }) diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index 647451205ea..b6ceb63843e 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -7214,6 +7214,16 @@ mod tests { .map(|fragment| fragment.id) .collect::>(); assert_eq!(logical_fragment_ids, vec![5, 2, 3, 4]); + assert_ne!( + dataset.manifest.reader_feature_flags + & lance_table::feature_flags::FLAG_LOGICAL_FRAGMENT_ORDER, + 0 + ); + assert_ne!( + dataset.manifest.writer_feature_flags + & lance_table::feature_flags::FLAG_LOGICAL_FRAGMENT_ORDER, + 0 + ); let fragments_by_id = dataset .get_fragments_from_ids(&[5, 2, 4, 3]) .unwrap() diff --git a/rust/lance/src/dataset/transaction.rs b/rust/lance/src/dataset/transaction.rs index 76257b06ff0..dbc34835ed3 100644 --- a/rust/lance/src/dataset/transaction.rs +++ b/rust/lance/src/dataset/transaction.rs @@ -1750,6 +1750,48 @@ impl Transaction { }) } + /// Keep replacements for existing fragment IDs in current manifest order, + /// then append genuinely new fragments in their supplied order. + fn preserve_existing_fragment_order( + existing_fragments: &[Fragment], + fragments: Vec, + operation_name: &str, + ) -> Result> { + let existing_fragment_ids = existing_fragments + .iter() + .map(|fragment| fragment.id) + .collect::>(); + let mut seen_fragment_ids = HashSet::with_capacity(fragments.len()); + let mut existing_fragments_by_id = HashMap::with_capacity(existing_fragments.len()); + let mut new_fragments = Vec::with_capacity(fragments.len()); + for fragment in fragments { + if !seen_fragment_ids.insert(fragment.id) { + return Err(Error::invalid_input(format!( + "{operation_name} contains multiple fragments with id {}", + fragment.id + ))); + } + if existing_fragment_ids.contains(&fragment.id) { + existing_fragments_by_id.insert(fragment.id, fragment); + } else { + new_fragments.push(fragment); + } + } + + let mut ordered_fragments = Vec::with_capacity(seen_fragment_ids.len()); + for existing_fragment in existing_fragments { + let Some(fragment) = existing_fragments_by_id.remove(&existing_fragment.id) else { + return Err(Error::internal(format!( + "{operation_name} lost fragment {} while preserving logical order", + existing_fragment.id + ))); + }; + ordered_fragments.push(fragment); + } + ordered_fragments.extend(new_fragments); + Ok(ordered_fragments) + } + fn data_storage_format_from_files( fragments: &[Fragment], user_requested: Option, @@ -2239,7 +2281,16 @@ impl Transaction { } } } - final_fragments.extend(merged_fragments); + + // Merge changes columns in existing fragments without changing + // their row placement. Reassemble those fragments in current + // manifest order, then append genuinely new fragments in the + // caller-provided order. + final_fragments = Self::preserve_existing_fragment_order( + existing_fragments, + merged_fragments, + "Merge", + )?; // A Merge can rewrite a column's data file in place; the field stays // in the schema, so the index is retained -- prune its now-stale @@ -2400,24 +2451,11 @@ impl Transaction { // Data replacement changes files, not row placement. Reassemble // the fragments in their existing logical order without using // fragment ids as ordering keys. - let mut fragments_by_id = HashMap::with_capacity(final_fragments.len()); - for fragment in final_fragments.drain(..) { - if fragments_by_id.insert(fragment.id, fragment).is_some() { - return Err(Error::invalid_input( - "DataReplacement contains multiple replacements for the same fragment" - .to_string(), - )); - } - } - for existing_fragment in existing_fragments { - let Some(fragment) = fragments_by_id.remove(&existing_fragment.id) else { - return Err(Error::internal(format!( - "DataReplacement lost fragment {} while preserving logical order", - existing_fragment.id - ))); - }; - final_fragments.push(fragment); - } + final_fragments = Self::preserve_existing_fragment_order( + existing_fragments, + final_fragments, + "DataReplacement", + )?; // 5. Invalidate index bitmaps for replaced fields let modified_fragments: Vec = final_fragments @@ -4415,6 +4453,34 @@ mod tests { assert_eq!(fragment_ids, vec![4, 1, 2, 3]); } + #[test] + fn test_merge_build_manifest_preserves_logical_fragment_order() { + 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(), + }; + validate_operation(Some(&manifest), &operation).unwrap(); + let transaction = Transaction::new(manifest.version, operation, None); + + let (merged, _) = transaction + .build_manifest( + Some(&manifest), + vec![], + "txn", + &ManifestWriteConfig::default(), + ) + .unwrap(); + + let fragment_ids = merged + .fragments + .iter() + .map(|fragment| fragment.id) + .collect::>(); + assert_eq!(fragment_ids, vec![0, 1, 2]); + } + #[test] fn test_merge_fragments_valid() { // Create a simple schema for testing diff --git a/rust/lance/src/io/commit.rs b/rust/lance/src/io/commit.rs index c0e57a066e5..b3a6ac12736 100644 --- a/rust/lance/src/io/commit.rs +++ b/rust/lance/src/io/commit.rs @@ -701,9 +701,10 @@ fn check_storage_version(manifest: &mut Manifest) -> Result<()> { Ok(()) } -/// 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 their order +/// requires a feature flag that is not set. 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. @@ -722,6 +723,18 @@ fn check_fragment_ids(manifest: &Manifest) -> Result<()> { fragment.id ))); } + if !manifest.uses_logical_fragment_order() + && 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 {}, but the logical \ + fragment order feature is not set", + fragments[0].id, fragments[1].id + ))); + } Ok(()) } @@ -1686,6 +1699,24 @@ mod tests { use crate::index::vector::VectorIndexParams; use crate::utils::test::{DatagenExt, FragmentCount, FragmentRowCount}; + #[test] + fn test_check_fragment_ids_requires_logical_order_feature() { + let mut 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("logical fragment order feature")); + + manifest.reader_feature_flags |= lance_table::feature_flags::FLAG_LOGICAL_FRAGMENT_ORDER; + manifest.writer_feature_flags |= lance_table::feature_flags::FLAG_LOGICAL_FRAGMENT_ORDER; + check_fragment_ids(&manifest).unwrap(); + } + async fn test_commit_handler(handler: Arc, should_succeed: bool) { // Create a dataset, passing handler as commit handler let schema = Arc::new(ArrowSchema::new(vec![ArrowField::new( From 93b75efe78c0132c59e58d924d0b3442fec89873 Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Fri, 7 Aug 2026 19:19:09 +0000 Subject: [PATCH 05/10] fix(compaction): preserve order without format changes --- docs/src/format/table/versioning.md | 5 +- protos/table.proto | 2 - rust/lance-table/src/feature_flags.rs | 58 +-- rust/lance-table/src/format/manifest.rs | 15 +- rust/lance/src/dataset.rs | 86 ++-- rust/lance/src/dataset/optimize.rs | 461 +++++++++++++----- rust/lance/src/dataset/transaction.rs | 90 +--- rust/lance/src/dataset/write/commit.rs | 6 +- rust/lance/src/io/commit.rs | 35 +- rust/lance/src/io/commit/conflict_resolver.rs | 44 +- 10 files changed, 465 insertions(+), 337 deletions(-) diff --git a/docs/src/format/table/versioning.md b/docs/src/format/table/versioning.md index 92e110ea725..745dd1ccd87 100644 --- a/docs/src/format/table/versioning.md +++ b/docs/src/format/table/versioning.md @@ -28,10 +28,7 @@ they should return an "unsupported" error on any read or write operation. | 4 | `FLAG_USE_V2_FORMAT_DEPRECATED` | No | No | Files are written with the new v2 format. This flag is deprecated and no longer used. | | 8 | `FLAG_TABLE_CONFIG` | No | Yes | Table config is present in the manifest. | | 16 | `FLAG_BASE_PATHS` | Yes | Yes | Dataset uses multiple base paths (for shallow clones or multi-base datasets). | -| 32 | `FLAG_DISABLE_TRANSACTION_FILE` | No | Yes | Transaction files are stored inline in manifests instead of under `_transaction`. | -| 64 | `FLAG_UNSTABLE_DATA_OVERLAY_FILES` | Yes | Yes | Fragments contain data overlay files; this feature is not yet released. | -| 128 | `FLAG_LOGICAL_FRAGMENT_ORDER` | Yes | Yes | Manifest position defines logical row order, so fragments must not be reordered by ID. | -Flags with bit values 256 and above are unknown and will cause implementations to reject the dataset with an "unsupported" error. +Flags with bit values 32 and above are unknown and will cause implementations to reject the dataset with an "unsupported" error. diff --git a/protos/table.proto b/protos/table.proto index 157ffba40ec..9a64230f40f 100644 --- a/protos/table.proto +++ b/protos/table.proto @@ -118,8 +118,6 @@ message Manifest { // * 1 << 6: data overlay files are present (see DataOverlayFile). Readers that do // not understand overlays must refuse the dataset, since ignoring an overlay // would silently return stale base values. - // * 1 << 7: fragments are stored in logical row order rather than fragment-ID - // order. Readers and writers must preserve their manifest positions. uint64 reader_feature_flags = 9; // Feature flags for writers. diff --git a/rust/lance-table/src/feature_flags.rs b/rust/lance-table/src/feature_flags.rs index fc31f2a9db5..41b8e415f8e 100644 --- a/rust/lance-table/src/feature_flags.rs +++ b/rust/lance-table/src/feature_flags.rs @@ -30,12 +30,8 @@ pub const FLAG_DISABLE_TRANSACTION_FILE: u64 = 32; /// unless [`ENABLE_UNSTABLE_DATA_OVERLAY_FILES_ENV`] is set, which lets benchmarks opt in. /// Debug builds always understand it so tests exercise the path. pub const FLAG_UNSTABLE_DATA_OVERLAY_FILES: u64 = 64; -/// Manifest fragments are stored in logical row order instead of fragment-ID -/// order. Readers must preserve the list order, and writers must not reorder the -/// fragments by ID. -pub const FLAG_LOGICAL_FRAGMENT_ORDER: u64 = 128; /// The first bit that is unknown as a feature flag -pub const FLAG_UNKNOWN: u64 = 256; +pub const FLAG_UNKNOWN: u64 = 128; /// Environment variable that opts a release build into reading and writing data /// overlay files before the feature is generally released. @@ -101,15 +97,6 @@ pub fn apply_feature_flags( manifest.writer_feature_flags |= FLAG_UNSTABLE_DATA_OVERLAY_FILES; } - let has_logical_fragment_order = manifest - .fragments - .windows(2) - .any(|fragments| fragments[0].id > fragments[1].id); - if has_logical_fragment_order { - manifest.reader_feature_flags |= FLAG_LOGICAL_FRAGMENT_ORDER; - manifest.writer_feature_flags |= FLAG_LOGICAL_FRAGMENT_ORDER; - } - if disable_transaction_file { manifest.writer_feature_flags |= FLAG_DISABLE_TRANSACTION_FILE; } @@ -174,7 +161,6 @@ mod tests { assert!(can_read_dataset(super::FLAG_TABLE_CONFIG)); assert!(can_read_dataset(super::FLAG_BASE_PATHS)); assert!(can_read_dataset(super::FLAG_DISABLE_TRANSACTION_FILE)); - assert!(can_read_dataset(super::FLAG_LOGICAL_FRAGMENT_ORDER)); // Overlay support is gated on the build profile / env opt-in, so the // flag is readable exactly when overlays are enabled (see // test_data_overlay_flag_release_gating for the full policy). @@ -251,7 +237,6 @@ mod tests { assert!(can_write_dataset(super::FLAG_TABLE_CONFIG)); assert!(can_write_dataset(super::FLAG_BASE_PATHS)); assert!(can_write_dataset(super::FLAG_DISABLE_TRANSACTION_FILE)); - assert!(can_write_dataset(super::FLAG_LOGICAL_FRAGMENT_ORDER)); // Overlay support is gated on the build profile / env opt-in, so the // flag is writable exactly when overlays are enabled (see // test_data_overlay_flag_release_gating for the full policy). @@ -320,45 +305,4 @@ mod tests { 0 ); } - - #[test] - fn test_logical_fragment_order_is_mixed_version_gated() { - use crate::format::{DataStorageFormat, Fragment}; - use arrow_schema::{Field as ArrowField, Schema as ArrowSchema}; - use lance_core::datatypes::Schema; - use std::collections::HashMap; - use std::sync::Arc; - - let arrow_schema = ArrowSchema::new(vec![ArrowField::new( - "id", - arrow_schema::DataType::Int64, - false, - )]); - let schema = Schema::try_from(&arrow_schema).unwrap(); - let mut manifest = Manifest::new( - schema, - Arc::new(vec![Fragment::new(5), Fragment::new(2)]), - DataStorageFormat::default(), - HashMap::new(), - ); - - apply_feature_flags(&mut manifest, false, false).unwrap(); - - assert_ne!( - manifest.reader_feature_flags & FLAG_LOGICAL_FRAGMENT_ORDER, - 0 - ); - assert_ne!( - manifest.writer_feature_flags & FLAG_LOGICAL_FRAGMENT_ORDER, - 0 - ); - assert!(can_read_dataset(manifest.reader_feature_flags)); - assert!(can_write_dataset(manifest.writer_feature_flags)); - - // Lance releases predating logical fragment order support only recognize - // lower feature bits, so they reject instead of misreading this manifest. - let released_supported_flags = FLAG_LOGICAL_FRAGMENT_ORDER - 1; - assert_ne!(manifest.reader_feature_flags & !released_supported_flags, 0); - assert_ne!(manifest.writer_feature_flags & !released_supported_flags, 0); - } } diff --git a/rust/lance-table/src/format/manifest.rs b/rust/lance-table/src/format/manifest.rs index 1c9e9e48cea..5543511bb95 100644 --- a/rust/lance-table/src/format/manifest.rs +++ b/rust/lance-table/src/format/manifest.rs @@ -18,9 +18,7 @@ use std::ops::Range; use std::sync::Arc; use super::Fragment; -use crate::feature_flags::{ - FLAG_LOGICAL_FRAGMENT_ORDER, FLAG_STABLE_ROW_IDS, has_deprecated_v2_feature_flag, -}; +use crate::feature_flags::{FLAG_STABLE_ROW_IDS, has_deprecated_v2_feature_flag}; use crate::format::fragment::DataFileFieldInterner; use crate::format::pb; use lance_core::cache::LanceCache; @@ -51,9 +49,8 @@ pub struct Manifest { /// Fragments, the pieces to build the dataset. /// - /// This list is stored in logical row order when - /// [`FLAG_LOGICAL_FRAGMENT_ORDER`] is set. Otherwise it is sorted by fragment - /// id, though the fragment id sequence may have gaps. + /// This list is stored in order, sorted by fragment id. However, the fragment id + /// sequence may have gaps. pub fragments: Arc>, /// The file position of the version aux data. @@ -499,12 +496,6 @@ impl Manifest { self.reader_feature_flags & FLAG_STABLE_ROW_IDS != 0 } - /// Whether manifest position, rather than fragment ID, defines logical row - /// order. - pub fn uses_logical_fragment_order(&self) -> bool { - self.reader_feature_flags & FLAG_LOGICAL_FRAGMENT_ORDER != 0 - } - /// Creates a serialized copy of the manifest, suitable for IPC or temp storage /// and can be used to create a dataset pub fn serialized(&self) -> Vec { diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index ce564cae6be..b98a2d1ee11 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -183,9 +183,6 @@ pub struct Dataset { // Bitmap of fragment ids in this dataset. pub(crate) fragment_bitmap: Arc, - // Manifest positions indexed by ascending fragment-id rank. Fragment ids are - // stable identities, while manifest position defines logical row order. - fragment_indices_by_id: Arc>, // These are references to session caches, but with the dataset URI as a prefix. pub(crate) index_cache: Arc, @@ -213,22 +210,6 @@ impl std::fmt::Debug for Dataset { } } -fn build_fragment_lookup(fragments: &[Fragment]) -> (Arc, Arc>) { - let fragment_bitmap: Arc = Arc::new( - fragments - .iter() - .map(|fragment| fragment.id as u32) - .collect(), - ); - let mut fragment_indices_by_id = vec![0; fragments.len()]; - for (manifest_index, fragment) in fragments.iter().enumerate() { - let id_rank = fragment_bitmap.rank(fragment.id as u32) as usize; - debug_assert!(id_rank > 0); - fragment_indices_by_id[id_rank - 1] = manifest_index; - } - (fragment_bitmap, Arc::new(fragment_indices_by_id)) -} - /// Dataset Version #[derive(Deserialize, Serialize, Debug)] pub struct Version { @@ -511,8 +492,13 @@ impl Dataset { let (manifest, manifest_location) = self.latest_manifest().await?; self.manifest = manifest; self.manifest_location = manifest_location; - (self.fragment_bitmap, self.fragment_indices_by_id) = - build_fragment_lookup(&self.manifest.fragments); + self.fragment_bitmap = Arc::new( + self.manifest + .fragments + .iter() + .map(|f| f.id as u32) + .collect(), + ); Ok(()) } @@ -845,7 +831,7 @@ impl Dataset { ); let metadata_cache = Arc::new(session.metadata_cache.for_dataset(&uri)); let index_cache = Arc::new(session.index_cache.for_dataset(&uri)); - let (fragment_bitmap, fragment_indices_by_id) = build_fragment_lookup(&manifest.fragments); + let fragment_bitmap = Arc::new(manifest.fragments.iter().map(|f| f.id as u32).collect()); write::log_unregistered_base_scoped_options( store_params.as_ref(), &manifest.base_paths, @@ -861,7 +847,6 @@ impl Dataset { session, refs, fragment_bitmap, - fragment_indices_by_id, metadata_cache, index_cache, file_reader_options, @@ -1641,8 +1626,13 @@ impl Dataset { self.manifest = Arc::new(manifest); self.manifest_location = manifest_location; - (self.fragment_bitmap, self.fragment_indices_by_id) = - build_fragment_lookup(&self.manifest.fragments); + self.fragment_bitmap = Arc::new( + self.manifest + .fragments + .iter() + .map(|f| f.id as u32) + .collect(), + ); Ok(()) } @@ -2724,9 +2714,7 @@ impl Dataset { Projection::full(self.clone()) } - /// Get fragments in logical row order. - /// - /// Fragment ids are stable identities and are not guaranteed to be sorted. + /// Get fragments. pub fn get_fragments(&self) -> Vec { let dataset = Arc::new(self.clone()); self.manifest @@ -2736,8 +2724,7 @@ impl Dataset { .collect() } - /// Iterate over manifest fragments in logical row order without allocating - /// [`FileFragment`] wrappers. + /// Iterate over manifest fragments without allocating [`FileFragment`] wrappers. pub fn iter_fragments(&self) -> impl Iterator { self.manifest.fragments.iter() } @@ -2848,12 +2835,11 @@ impl Dataset { if !self.fragment_bitmap.contains(*id) { return None; } - let id_rank = self.fragment_bitmap.rank(*id) as usize - 1; - let fragment_index = self.fragment_indices_by_id[id_rank]; + let fragment_index = self.fragment_bitmap.rank(*id) as usize - 1; let fragment = self.manifest.fragments.get(fragment_index)?; debug_assert_eq!( fragment.id, *id as u64, - "fragment lookup for id {id} resolved to fragment {}, but the lookup and manifest.fragments are expected to stay in sync", + "fragment_bitmap rank({id}) resolved to fragment {}, but fragment_bitmap and manifest.fragments are expected to stay in sync", fragment.id ); Some(FileFragment::new(dataset.clone(), fragment.clone())) @@ -3039,25 +3025,21 @@ impl Dataset { } } - if !self.manifest.uses_logical_fragment_order() { - self.manifest - .fragments - .iter() - .map(|fragment| fragment.id) - .try_fold(0, |previous_id, fragment_id| { - if fragment_id < previous_id { - Err(Error::corrupt_file( - self.base.clone(), - format!( - "Fragment ids are not sorted in increasing fragment-id order, but the logical fragment order feature is not set. Found {fragment_id} after {previous_id} in dataset {:?}", - self.base - ), - )) - } else { - Ok(fragment_id) - } - })?; - } + // Fragments are sorted in increasing fragment id order + self.manifest + .fragments + .iter() + .map(|f| f.id) + .try_fold(0, |prev, id| { + if id < prev { + Err(Error::corrupt_file(self.base.clone(), format!( + "Fragment ids are not sorted in increasing fragment-id order. Found {} after {} in dataset {:?}", + id, prev, self.base + ))) + } else { + Ok(id) + } + })?; // All fragments have equal lengths futures::stream::iter(self.get_fragments()) diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index b6ceb63843e..e5595922f8a 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -114,6 +114,7 @@ use lance_core::Error; use lance_core::datatypes::{ BLOB_V2_LOGICAL_FIELDS, BLOB_V2_LOGICAL_TYPE, BlobHandling, BlobKind, BlobV2Layout, }; +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_index::frag_reuse::{FRAG_REUSE_INDEX_NAME, FragReuseGroup}; @@ -129,6 +130,7 @@ pub mod remapping; use crate::index::frag_reuse::build_new_frag_reuse_index; use crate::io::deletion::read_dataset_deletion_file; use binary_copy::rewrite_files_binary_copy; +use lance_table::io::deletion::deletion_file_path; pub use remapping::{IgnoreRemap, IndexRemapper, IndexRemapperOptions, RemappedIndex}; /// Controls how data is rewritten during compaction. @@ -703,9 +705,12 @@ impl CompactionPlanner for DefaultCompactionPlanner { )); } - // Manifest order is logical row order. Fragment ids are stable identities - // and may be out of order after a rewrite. + // Manifest order follows fragment ID order, which is also row order. let fragments = dataset.get_fragments(); + debug_assert!( + fragments.windows(2).all(|pair| pair[0].id() < pair[1].id()), + "fragments in manifest are not sorted" + ); let mut fragment_metrics = futures::stream::iter(fragments) .map(|fragment| async move { match collect_metrics(&fragment).await { @@ -1608,6 +1613,275 @@ 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, + capture_row_addrs: bool, + 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 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(); + let row_addrs = if capture_row_addrs { + Some(serialize_fragment_row_addrs(&fragment)?) + } else { + None + }; + completed_suffix.push(RewriteResult { + metrics: CompactionMetrics::default(), + new_fragments: vec![fragment.clone()], + read_version, + original_fragments: vec![fragment], + row_addrs, + }); + 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 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. @@ -2041,8 +2315,9 @@ async fn recalc_versions_for_rewritten_fragments( /// 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 fragments' current logical -/// positions so partial and gapped result sets replace their ranges in place. +/// 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, @@ -2053,18 +2328,6 @@ pub async fn commit_compaction( return Ok(CompactionMetrics::default()); } - 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 @@ -2084,7 +2347,6 @@ 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()) @@ -2094,43 +2356,48 @@ pub async fn commit_compaction( )); } - // Rewrite tasks finish in an arbitrary order. Apply their groups in current - // manifest order so every replacement is spliced into its logical position. - let fragment_positions = dataset - .manifest - .fragments - .iter() - .enumerate() - .map(|(position, fragment)| (fragment.id, position)) - .collect::>(); - completed_tasks.sort_by_key(|task| { - fragment_positions - .get(&task.original_fragments[0].id) - .copied() - .map_or((1, usize::MAX), |position| (0, position)) - }); - // 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, + capture_row_addrs, + tasks_read_version, + dataset.manifest.version, + )?; + 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(); @@ -2260,19 +2527,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() }; @@ -2974,8 +3228,10 @@ mod tests { (row_addrs(2, 0..200), true), ], vec![(row_addrs(2, 200..400), true)], - // frag 3 is skipped since it does not have enough missing data - // Frags 4, 5, and 6 are rewritten to frags 9 & 10 + // Frag 3 keeps its files but is relabeled as fragment 9 so the + // manifest can remain ID-sorted. + vec![(row_addrs(3, 0..1000), true)], + // Frags 4, 5, and 6 are rewritten to frags 10 & 11 vec![ // Only 800 of the 1000 rows taken from frag 4 (row_addrs(4, 0..200), true), @@ -3033,7 +3289,7 @@ mod tests { .iter() .map(|f| f.id()) .collect::>(); - assert_eq!(fragment_ids, vec![7, 8, 3, 9, 10]); + assert_eq!(fragment_ids, vec![7, 8, 9, 10, 11]); let compacted_data = dataset.scan().try_into_batch().await.unwrap(); assert_eq!(compacted_data, expected_data); } @@ -3419,10 +3675,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, ) @@ -3447,6 +3704,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] @@ -4617,11 +4883,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]), @@ -4640,16 +4901,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]), @@ -4668,7 +4926,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 @@ -4692,15 +4952,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); } @@ -4785,11 +5039,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]), @@ -4812,10 +5061,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)); @@ -7189,6 +7435,8 @@ mod tests { } 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, max_source_fragments: Some(2), @@ -7207,33 +7455,24 @@ mod tests { .unwrap(); assert_eq!(dataset.get_fragments().len(), 4); - let logical_fragment_ids = dataset + let fragment_ids = dataset .manifest .fragments .iter() .map(|fragment| fragment.id) .collect::>(); - assert_eq!(logical_fragment_ids, vec![5, 2, 3, 4]); - assert_ne!( - dataset.manifest.reader_feature_flags - & lance_table::feature_flags::FLAG_LOGICAL_FRAGMENT_ORDER, - 0 - ); - assert_ne!( - dataset.manifest.writer_feature_flags - & lance_table::feature_flags::FLAG_LOGICAL_FRAGMENT_ORDER, - 0 - ); + assert_eq!(fragment_ids, vec![5, 6, 7, 8]); let fragments_by_id = dataset - .get_fragments_from_ids(&[5, 2, 4, 3]) + .get_fragments_from_ids(&[8, 5, 7, 6]) .unwrap() .into_iter() .map(|fragment| fragment.id()) .collect::>(); - assert_eq!(fragments_by_id, vec![2, 3, 4, 5]); + 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, data.slice(0, offset)); + assert_eq!(after_first, expected); let second_plan = plan_compaction(&dataset, &options).await.unwrap(); let second_fragment_ids = second_plan @@ -7241,12 +7480,12 @@ mod tests { .iter() .flat_map(|task| task.fragments.iter().map(|fragment| fragment.id)) .collect::>(); - assert_eq!(second_fragment_ids, vec![3, 4]); + assert_eq!(second_fragment_ids, vec![7, 8]); compact_files(&mut dataset, options, None).await.unwrap(); let compacted = dataset.scan().try_into_batch().await.unwrap(); - assert_eq!(compacted, data.slice(0, offset)); + assert_eq!(compacted, expected); assert_eq!(dataset.get_fragments().len(), 3); } diff --git a/rust/lance/src/dataset/transaction.rs b/rust/lance/src/dataset/transaction.rs index dbc34835ed3..f9df669b6c5 100644 --- a/rust/lance/src/dataset/transaction.rs +++ b/rust/lance/src/dataset/transaction.rs @@ -1750,48 +1750,6 @@ impl Transaction { }) } - /// Keep replacements for existing fragment IDs in current manifest order, - /// then append genuinely new fragments in their supplied order. - fn preserve_existing_fragment_order( - existing_fragments: &[Fragment], - fragments: Vec, - operation_name: &str, - ) -> Result> { - let existing_fragment_ids = existing_fragments - .iter() - .map(|fragment| fragment.id) - .collect::>(); - let mut seen_fragment_ids = HashSet::with_capacity(fragments.len()); - let mut existing_fragments_by_id = HashMap::with_capacity(existing_fragments.len()); - let mut new_fragments = Vec::with_capacity(fragments.len()); - for fragment in fragments { - if !seen_fragment_ids.insert(fragment.id) { - return Err(Error::invalid_input(format!( - "{operation_name} contains multiple fragments with id {}", - fragment.id - ))); - } - if existing_fragment_ids.contains(&fragment.id) { - existing_fragments_by_id.insert(fragment.id, fragment); - } else { - new_fragments.push(fragment); - } - } - - let mut ordered_fragments = Vec::with_capacity(seen_fragment_ids.len()); - for existing_fragment in existing_fragments { - let Some(fragment) = existing_fragments_by_id.remove(&existing_fragment.id) else { - return Err(Error::internal(format!( - "{operation_name} lost fragment {} while preserving logical order", - existing_fragment.id - ))); - }; - ordered_fragments.push(fragment); - } - ordered_fragments.extend(new_fragments); - Ok(ordered_fragments) - } - fn data_storage_format_from_files( fragments: &[Fragment], user_requested: Option, @@ -2281,16 +2239,7 @@ impl Transaction { } } } - - // Merge changes columns in existing fragments without changing - // their row placement. Reassemble those fragments in current - // manifest order, then append genuinely new fragments in the - // caller-provided order. - final_fragments = Self::preserve_existing_fragment_order( - existing_fragments, - merged_fragments, - "Merge", - )?; + final_fragments.extend(merged_fragments); // A Merge can rewrite a column's data file in place; the field stays // in the schema, so the index is retained -- prune its now-stale @@ -2448,15 +2397,6 @@ impl Transaction { final_fragments.extend(unmodified_fragments); - // Data replacement changes files, not row placement. Reassemble - // the fragments in their existing logical order without using - // fragment ids as ordering keys. - final_fragments = Self::preserve_existing_fragment_order( - existing_fragments, - final_fragments, - "DataReplacement", - )?; - // 5. Invalidate index bitmaps for replaced fields let modified_fragments: Vec = final_fragments .iter() @@ -2547,6 +2487,9 @@ impl Transaction { } }; + // If a fragment was reserved then it may not belong at the end of the fragments list. + final_fragments.sort_by_key(|frag| frag.id); + // Clean up data files that only contain tombstoned fields Self::remove_tombstoned_data_files(&mut final_fragments); @@ -2956,6 +2899,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 { @@ -4365,7 +4323,7 @@ mod tests { } #[test] - fn test_rewrite_build_manifest_preserves_logical_fragment_order() { + fn test_rewrite_build_manifest_keeps_fragment_ids_sorted() { let manifest = sample_manifest_with_fragments(0..4); let transaction = Transaction::new( manifest.version, @@ -4394,11 +4352,11 @@ mod tests { .iter() .map(|fragment| fragment.id) .collect::>(); - assert_eq!(fragment_ids, vec![0, 4, 3]); + assert_eq!(fragment_ids, vec![0, 3, 4]); } #[test] - fn test_rewrite_after_row_adding_update_preserves_logical_order() { + 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, @@ -4450,11 +4408,11 @@ mod tests { .iter() .map(|fragment| fragment.id) .collect::>(); - assert_eq!(fragment_ids, vec![4, 1, 2, 3]); + assert_eq!(fragment_ids, vec![1, 2, 3, 4]); } #[test] - fn test_merge_build_manifest_preserves_logical_fragment_order() { + 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 { diff --git a/rust/lance/src/dataset/write/commit.rs b/rust/lance/src/dataset/write/commit.rs index 6a3649c5533..fd57f68f9e7 100644 --- a/rust/lance/src/dataset/write/commit.rs +++ b/rust/lance/src/dataset/write/commit.rs @@ -17,7 +17,7 @@ use crate::io::commit::DEFAULT_COMMIT_RETRY_TIMEOUT; use crate::{ Dataset, Error, Result, dataset::{ - ManifestWriteConfig, ReadParams, build_fragment_lookup, + ManifestWriteConfig, ReadParams, builder::DatasetBuilder, commit_detached_transaction, commit_new_dataset, commit_transaction, refs::Refs, @@ -473,7 +473,7 @@ impl<'a> CommitBuilder<'a> { operation=&transaction.operation.name() ); - let (fragment_bitmap, fragment_indices_by_id) = build_fragment_lookup(&manifest.fragments); + let fragment_bitmap = Arc::new(manifest.fragments.iter().map(|f| f.id as u32).collect()); match &self.dest { WriteDestination::Dataset(dataset) => Ok(Dataset { @@ -481,7 +481,6 @@ impl<'a> CommitBuilder<'a> { manifest_location, session, fragment_bitmap, - fragment_indices_by_id, ..dataset.as_ref().clone() }), WriteDestination::Uri(uri) => { @@ -506,7 +505,6 @@ impl<'a> CommitBuilder<'a> { refs, index_cache, fragment_bitmap, - fragment_indices_by_id, metadata_cache, file_reader_options: None, store_params: self.store_params.clone().map(Box::new), diff --git a/rust/lance/src/io/commit.rs b/rust/lance/src/io/commit.rs index b3a6ac12736..d5ee0508c2a 100644 --- a/rust/lance/src/io/commit.rs +++ b/rust/lance/src/io/commit.rs @@ -701,10 +701,10 @@ fn check_storage_version(manifest: &mut Manifest) -> Result<()> { Ok(()) } -/// Reject a manifest in which fragment identities are ambiguous or their order -/// requires a feature flag that is not set. 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. @@ -723,15 +723,14 @@ fn check_fragment_ids(manifest: &Manifest) -> Result<()> { fragment.id ))); } - if !manifest.uses_logical_fragment_order() - && let Some(fragments) = manifest - .fragments - .windows(2) - .find(|fragments| fragments[0].id > fragments[1].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 {}, but the logical \ - fragment order feature is not set", + "The commit would place fragment {} before lower fragment id {}. Fragment ids must \ + remain sorted in increasing order", fragments[0].id, fragments[1].id ))); } @@ -1700,8 +1699,8 @@ mod tests { use crate::utils::test::{DatagenExt, FragmentCount, FragmentRowCount}; #[test] - fn test_check_fragment_ids_requires_logical_order_feature() { - let mut manifest = Manifest::new( + 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(), @@ -1710,11 +1709,11 @@ mod tests { let error = check_fragment_ids(&manifest).unwrap_err(); assert!(matches!(error, Error::InvalidInput { .. })); - assert!(error.to_string().contains("logical fragment order feature")); - - manifest.reader_feature_flags |= lance_table::feature_flags::FLAG_LOGICAL_FRAGMENT_ORDER; - manifest.writer_feature_flags |= lance_table::feature_flags::FLAG_LOGICAL_FRAGMENT_ORDER; - check_fragment_ids(&manifest).unwrap(); + assert!( + error + .to_string() + .contains("Fragment ids must remain sorted in increasing order") + ); } async fn test_commit_handler(handler: Arc, should_succeed: bool) { diff --git a/rust/lance/src/io/commit/conflict_resolver.rs b/rust/lance/src/io/commit/conflict_resolver.rs index e3408da045c..5f84dc0da0f 100644 --- a/rust/lance/src/io/commit/conflict_resolver.rs +++ b/rust/lance/src/io/commit/conflict_resolver.rs @@ -773,8 +773,14 @@ impl<'a> TransactionRebase<'a> { 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 { .. } + // Rewrites allocate a consecutive suffix of fragment IDs for the + // replacement range and every following fragment. If a row-adding + // operation landed before that reservation, ID sorting would put + // its rows before the rewritten range, so replan from the new manifest. + Operation::Append { .. } => { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } + Operation::ReserveFragments { .. } | Operation::Project { .. } | Operation::Clone { .. } | Operation::UpdateConfig { .. } @@ -784,11 +790,6 @@ impl<'a> TransactionRebase<'a> { updated_fragments, deleted_fragment_ids, .. - } - | Operation::Update { - updated_fragments, - removed_fragment_ids: deleted_fragment_ids, - .. } => { if updated_fragments .iter() @@ -801,6 +802,24 @@ impl<'a> TransactionRebase<'a> { Ok(()) } } + Operation::Update { + updated_fragments, + removed_fragment_ids: deleted_fragment_ids, + new_fragments, + .. + } => { + if !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 @@ -2644,7 +2663,7 @@ mod tests { } #[test] - fn test_rewrite_is_compatible_with_row_adding_update() { + fn test_rewrite_conflicts_with_row_adding_update() { let operation = Operation::Rewrite { groups: vec![RewriteGroup { old_fragments: vec![Fragment::new(0)], @@ -2677,7 +2696,10 @@ mod tests { None, ); - assert!(rebase.check_txn(&other, 1).is_ok()); + assert!(matches!( + rebase.check_txn(&other, 1), + Err(Error::RetryableCommitConflict { .. }) + )); } #[test] @@ -2874,7 +2896,7 @@ mod tests { frag_reuse_index: None, }, [ - Compatible, // append + Retryable, // append Retryable, // create index Compatible, // delete Retryable, // merge @@ -2896,7 +2918,7 @@ mod tests { frag_reuse_index: None, }, [ - Compatible, // append + Retryable, // append Retryable, // create index Retryable, // delete Retryable, // merge From 840b7d85d75532a21ad572564481ca2e6fcf14f6 Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Fri, 7 Aug 2026 19:38:27 +0000 Subject: [PATCH 06/10] test: align index coverage with compaction relabeling --- rust/lance/src/index/create.rs | 17 +++++----- rust/lance/src/index/scalar_logical.rs | 46 ++++++++++++++------------ 2 files changed, 32 insertions(+), 31 deletions(-) diff --git a/rust/lance/src/index/create.rs b/rust/lance/src/index/create.rs index 16bb655a871..8e7dafbe6c2 100644 --- a/rust/lance/src/index/create.rs +++ b/rust/lance/src/index/create.rs @@ -3036,10 +3036,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, @@ -3058,16 +3058,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 3b7dd4fcbda..f28524a1755 100644 --- a/rust/lance/src/index/scalar_logical.rs +++ b/rust/lance/src/index/scalar_logical.rs @@ -907,8 +907,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( @@ -1001,6 +1003,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 @@ -1019,7 +1025,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); } @@ -1428,7 +1434,8 @@ mod tests { "compaction should retire fragment 0" ); - // Merge: the retired fragment should be dropped from coverage + // Merge: stable-row-ID index coverage should follow both the compacted + // fragment and the metadata-only relabeled trailing fragment. let segments = dataset .load_indices_by_name("text_fmindex_compact") .await @@ -1439,14 +1446,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, &live_frags); // Commit the merged segment and verify search works dataset @@ -1851,6 +1851,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(); @@ -1875,9 +1876,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 @@ -1886,14 +1893,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" ); } } From 56ee82491990d183c12b4fa517e5d512699dc3d0 Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Fri, 7 Aug 2026 20:43:58 +0000 Subject: [PATCH 07/10] fix(compaction): bound suffix relabel work --- .../lance/compaction/CompactionOptions.java | 6 +- python/python/lance/dataset.py | 10 +- python/python/lance/optimize.py | 8 +- rust/lance/src/dataset/optimize.rs | 214 ++++++++++++++---- 4 files changed, 181 insertions(+), 57 deletions(-) diff --git a/java/src/main/java/org/lance/compaction/CompactionOptions.java b/java/src/main/java/org/lance/compaction/CompactionOptions.java index 7c3d65ffc3f..841990928d4 100644 --- a/java/src/main/java/org/lance/compaction/CompactionOptions.java +++ b/java/src/main/java/org/lance/compaction/CompactionOptions.java @@ -236,9 +236,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. */ public Builder withMaxSourceFragments(long maxSourceFragments) { this.maxSourceFragments = Optional.of(maxSourceFragments); diff --git a/python/python/lance/dataset.py b/python/python/lance/dataset.py index d387880be48..09e4d80545f 100644 --- a/python/python/lance/dataset.py +++ b/python/python/lance/dataset.py @@ -7185,11 +7185,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. Returns ------- diff --git a/python/python/lance/optimize.py b/python/python/lance/optimize.py index 3ac7547960b..4d0121d291e 100644 --- a/python/python/lance/optimize.py +++ b/python/python/lance/optimize.py @@ -91,9 +91,9 @@ class CompactionOptions(TypedDict): """ 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) """ diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index e5595922f8a..6c3043d7008 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -260,10 +260,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 data overlay files a fragment may carry before it is @@ -808,7 +809,27 @@ impl CompactionPlanner for DefaultCompactionPlanner { candidate_bins.push(bin); } - let all_tasks: Vec = candidate_bins + if let Some(max_frags) = self.options.max_source_fragments { + let suffix_start = dataset.manifest.fragments.len().saturating_sub(max_frags); + 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 tasks: Vec = candidate_bins .into_iter() .filter(|bin| !bin.is_noop()) .flat_map(|bin| bin.split_for_size(self.options.target_rows_per_fragment)) @@ -817,19 +838,6 @@ impl CompactionPlanner for DefaultCompactionPlanner { }) .collect(); - let tasks = if let Some(max_frags) = self.options.max_source_fragments { - let mut total_frags = 0; - all_tasks - .into_iter() - .take_while(|task| { - total_frags += task.fragments.len(); - total_frags <= max_frags - }) - .collect() - } else { - all_tasks - }; - let mut compaction_plan = CompactionPlan::new(dataset.manifest.version, self.options.clone()); compaction_plan.extend_tasks(tasks); @@ -1734,7 +1742,7 @@ fn rebase_task_source_fragments( fn complete_rewrite_suffix( manifest_fragments: &[Fragment], completed_tasks: Vec, - capture_row_addrs: bool, + max_affected_fragments: Option, read_version: u64, current_version: u64, ) -> Result> { @@ -1795,6 +1803,21 @@ fn complete_rewrite_suffix( 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), ); @@ -1816,17 +1839,12 @@ fn complete_rewrite_suffix( } let fragment = manifest_fragments[position].clone(); - let row_addrs = if capture_row_addrs { - Some(serialize_fragment_row_addrs(&fragment)?) - } else { - None - }; completed_suffix.push(RewriteResult { metrics: CompactionMetrics::default(), new_fragments: vec![fragment.clone()], read_version, original_fragments: vec![fragment], - row_addrs, + row_addrs: None, }); position += 1; } @@ -1851,6 +1869,47 @@ fn is_metadata_only_relabel(task: &RewriteResult) -> bool { 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], @@ -2374,10 +2433,13 @@ pub async fn commit_compaction( let mut completed_tasks = complete_rewrite_suffix( &dataset.manifest.fragments, completed_tasks, - capture_row_addrs, + options.max_source_fragments, tasks_read_version, dataset.manifest.version, )?; + 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. @@ -2606,7 +2668,7 @@ mod tests { use crate::dataset::optimize::remapping::{transpose_row_addrs, transpose_row_ids_from_digest}; 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, @@ -7326,8 +7388,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), @@ -7349,7 +7411,8 @@ mod tests { .iter() .flat_map(|task| task.fragments.iter().map(|fragment| fragment.id)) .collect::>(); - let expected_fragment_ids = dataset.manifest.fragments[..bounded_source_frags] + let expected_fragment_ids = dataset.manifest.fragments + [dataset.manifest.fragments.len() - bounded_source_frags..] .iter() .map(|fragment| fragment.id) .collect::>(); @@ -7398,7 +7461,7 @@ mod tests { #[rstest] #[tokio::test] - async fn test_bounded_compaction_preserves_order_across_candidate_gap( + async fn test_partial_compaction_preserves_order_across_candidate_gap( #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)] data_storage_version: LanceFileVersion, ) { @@ -7439,7 +7502,6 @@ mod tests { let expected = dataset.scan().try_into_batch().await.unwrap(); let options = CompactionOptions { target_rows_per_fragment: 250, - max_source_fragments: Some(2), ..Default::default() }; let plan = plan_compaction(&dataset, &options).await.unwrap(); @@ -7448,11 +7510,20 @@ mod tests { .iter() .flat_map(|task| task.fragments.iter().map(|fragment| fragment.id)) .collect::>(); - assert_eq!(planned_fragment_ids, vec![0, 1]); + 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(); - compact_files(&mut dataset, options.clone(), None) - .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 @@ -7474,21 +7545,74 @@ mod tests { let after_first = dataset.scan().try_into_batch().await.unwrap(); assert_eq!(after_first, expected); - let second_plan = plan_compaction(&dataset, &options).await.unwrap(); - let second_fragment_ids = second_plan - .tasks() - .iter() - .flat_map(|task| task.fragments.iter().map(|fragment| fragment.id)) - .collect::>(); - assert_eq!(second_fragment_ids, vec![7, 8]); - - compact_files(&mut dataset, options, None).await.unwrap(); + 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); + } + #[tokio::test] async fn test_compaction_uses_manifest_config() { let test_dir = TempStrDir::default(); From abf1119236bc4238464659bf59f48491f5261a3f Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Wed, 19 Aug 2026 18:37:02 +0000 Subject: [PATCH 08/10] fix(compaction): preserve excluded fragment identities --- .../lance/compaction/CompactionOptions.java | 5 +-- python/python/lance/dataset.py | 4 ++- python/python/lance/optimize.py | 5 +-- rust/lance-table/src/transaction.rs | 18 +++-------- rust/lance/src/dataset/optimize.rs | 32 ++++++++++++++----- 5 files changed, 37 insertions(+), 27 deletions(-) diff --git a/java/src/main/java/org/lance/compaction/CompactionOptions.java b/java/src/main/java/org/lance/compaction/CompactionOptions.java index 222437bf57a..878adb9a7e5 100644 --- a/java/src/main/java/org/lance/compaction/CompactionOptions.java +++ b/java/src/main/java/org/lance/compaction/CompactionOptions.java @@ -347,8 +347,9 @@ public Builder withMaxSourceBytes(long maxSourceBytes) { /** * 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 8df214d3051..0a04e5ed722 100644 --- a/python/python/lance/dataset.py +++ b/python/python/lance/dataset.py @@ -7260,7 +7260,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. Returns ------- diff --git a/python/python/lance/optimize.py b/python/python/lance/optimize.py index e54e9746f7e..cc597eb861e 100644 --- a/python/python/lance/optimize.py +++ b/python/python/lance/optimize.py @@ -118,6 +118,7 @@ class CompactionOptions(TypedDict): """ 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) """ diff --git a/rust/lance-table/src/transaction.rs b/rust/lance-table/src/transaction.rs index 7a87e2b6ede..41bf4a43190 100644 --- a/rust/lance-table/src/transaction.rs +++ b/rust/lance-table/src/transaction.rs @@ -5222,12 +5222,7 @@ mod tests { ); let (rewritten, _) = transaction - .build_manifest( - Some(&manifest), - vec![], - "txn", - &ManifestWriteConfig::default(), - ) + .build_manifest(Some(&manifest), vec![], "txn", &default_build_config()) .unwrap(); let fragment_ids = rewritten @@ -5261,7 +5256,7 @@ mod tests { Some(&manifest), vec![], "update-txn", - &ManifestWriteConfig::default(), + &default_build_config(), ) .unwrap(); @@ -5282,7 +5277,7 @@ mod tests { Some(&updated), vec![], "rewrite-txn", - &ManifestWriteConfig::default(), + &default_build_config(), ) .unwrap(); @@ -5307,12 +5302,7 @@ mod tests { let transaction = Transaction::new(manifest.version, operation, None); let (merged, _) = transaction - .build_manifest( - Some(&manifest), - vec![], - "txn", - &ManifestWriteConfig::default(), - ) + .build_manifest(Some(&manifest), vec![], "txn", &default_build_config()) .unwrap(); let fragment_ids = merged diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index efa1ec272d6..456074bbe6c 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -292,7 +292,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, @@ -753,6 +755,17 @@ impl CompactionPlanner for DefaultCompactionPlanner { fragments.windows(2).all(|pair| pair[0].id() < pair[1].id()), "fragments in manifest are not sorted" ); + // 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 exclusion so excluded identities stay unchanged. + let exclusion_suffix_start = fragments + .iter() + .rposition(|fragment| { + u32::try_from(fragment.id()) + .is_ok_and(|fragment_id| self.excluded_fragment_ids.contains(fragment_id)) + }) + .map_or(0, |position| position + 1); let mut fragment_metrics = futures::stream::iter(fragments) .map(|fragment| async { if u32::try_from(fragment.id()) @@ -864,8 +877,11 @@ impl CompactionPlanner for DefaultCompactionPlanner { candidate_bins.push(bin); } - if let Some(max_frags) = self.options.max_source_fragments { - let suffix_start = dataset.manifest.fragments.len().saturating_sub(max_frags); + 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| { @@ -8269,10 +8285,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() @@ -8280,8 +8293,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() From 5ce3b3f4311075ab8652be3aa975dc783434756e Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Mon, 14 Sep 2026 15:36:17 +0000 Subject: [PATCH 09/10] test(compaction): align storage version coverage --- rust/lance/src/dataset/optimize.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index cec19b4d0bc..e274ef21363 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -3777,7 +3777,7 @@ mod tests { .unwrap(); let options = CompactionOptions { data_storage_version: selector, - excluded_fragment_ids: vec![2], + excluded_fragment_ids: vec![0], ..Default::default() }; @@ -3794,7 +3794,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(); From 4f9a19fbcaa6fcfb3d2a5a261e9b004d85efbc6f Mon Sep 17 00:00:00 2001 From: Gatefixer <312823363+lance-gatefixer[bot]@users.noreply.github.com> Date: Sun, 27 Sep 2026 00:42:32 +0000 Subject: [PATCH 10/10] test(compaction): align index fixtures with suffix relabeling --- rust/lance/src/dataset/index/frag_reuse.rs | 12 +++++----- rust/lance/src/index/create.rs | 21 ++++++++--------- rust/lance/src/index/scalar_logical.rs | 26 +++++++++++++--------- 3 files changed, 33 insertions(+), 26 deletions(-) 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/index/create.rs b/rust/lance/src/index/create.rs index f6659a85755..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 { diff --git a/rust/lance/src/index/scalar_logical.rs b/rust/lance/src/index/scalar_logical.rs index dc5af206e07..37913a75ea9 100644 --- a/rust/lance/src/index/scalar_logical.rs +++ b/rust/lance/src/index/scalar_logical.rs @@ -1482,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(), @@ -1502,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); @@ -1527,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( @@ -1547,12 +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: stable-row-ID index coverage should follow both the compacted - // fragment and the metadata-only relabeled trailing fragment. + // 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 @@ -1563,7 +1567,7 @@ mod tests { .unwrap(); let coverage = merged.fragment_bitmap.as_ref().unwrap(); - assert_eq!(coverage, &live_frags); + 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 {