From 8bb367f5bfe92eb6fd4a367b4a3baa7652fb0b16 Mon Sep 17 00:00:00 2001 From: tanish Date: Tue, 25 Aug 2026 19:02:28 +0000 Subject: [PATCH 1/2] perf: make Rewrite transaction fragment handling O(n) instead of O(groups * fragments) Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- rust/lance/src/dataset/transaction.rs | 115 ++++++++++++++++++-------- 1 file changed, 80 insertions(+), 35 deletions(-) diff --git a/rust/lance/src/dataset/transaction.rs b/rust/lance/src/dataset/transaction.rs index fc0d4f4dc19..96df008f094 100644 --- a/rust/lance/src/dataset/transaction.rs +++ b/rust/lance/src/dataset/transaction.rs @@ -2556,37 +2556,41 @@ impl Transaction { version: u64, _next_row_id: Option<&u64>, ) -> Result<()> { + // Index fragment positions up front so each rewrite group is located in + // O(1). Scanning and splicing/retaining the fragment list per group is + // O(groups * fragments) and dominates commit time on manifests with + // millions of fragments. + let positions: HashMap = final_fragments + .iter() + .enumerate() + .map(|(position, fragment)| (fragment.id, position)) + .collect(); + + // Contiguous groups have their new fragments spliced in at the position + // of their first old fragment; non-contiguous groups have their new + // fragments appended at the end. + let mut removed = vec![false; final_fragments.len()]; + let mut replacements: HashMap> = HashMap::new(); + let mut appended: Vec = Vec::new(); + for group in groups { - // If the old fragments are contiguous, find the range - let replace_range = { - let start = final_fragments + let start = *positions.get(&group.old_fragments[0].id).ok_or_else(|| { + Error::commit_conflict_source( + version, + format!( + "dataset does not contain a fragment a rewrite operation wants to replace: id={}", + group.old_fragments[0].id + ) + .into(), + ) + })?; + + let contiguous = start + group.old_fragments.len() <= final_fragments.len() + && group + .old_fragments .iter() .enumerate() - .find(|(_, f)| f.id == group.old_fragments[0].id) - .ok_or_else(|| { - Error::commit_conflict_source( - version, - format!( - "dataset does not contain a fragment a rewrite operation wants to replace: id={}", - group.old_fragments[0].id - ) - .into(), - ) - })? - .0; - - // Verify old_fragments matches contiguous range - let mut i = 1; - loop { - if i == group.old_fragments.len() { - break Some(start..start + i); - } - if final_fragments[start + i].id != group.old_fragments[i].id { - break None; - } - i += 1; - } - }; + .all(|(i, old)| final_fragments[start + i].id == old.id); let new_fragments = Self::fragments_with_ids(group.new_fragments.clone(), fragment_id) .collect::>(); @@ -2595,17 +2599,36 @@ impl Transaction { // (recalc_versions_for_rewritten_fragments) which preserves version information // from the original fragments. We don't modify it here. - if let Some(replace_range) = replace_range { - // Efficiently path using slice - final_fragments.splice(replace_range, new_fragments); + if contiguous { + for flag in removed + .iter_mut() + .skip(start) + .take(group.old_fragments.len()) + { + *flag = true; + } + replacements.insert(start, new_fragments); } else { - // Slower path for non-contiguous ranges - for fragment in group.old_fragments.iter() { - final_fragments.retain(|f| f.id != fragment.id); + for old_fragment in group.old_fragments.iter() { + if let Some(&position) = positions.get(&old_fragment.id) { + removed[position] = true; + } } - final_fragments.extend(new_fragments); + appended.extend(new_fragments); } } + + let mut rewritten = Vec::with_capacity(final_fragments.len() + appended.len()); + for (position, fragment) in std::mem::take(final_fragments).into_iter().enumerate() { + if let Some(new_fragments) = replacements.remove(&position) { + rewritten.extend(new_fragments); + } + if !removed[position] { + rewritten.push(fragment); + } + } + rewritten.extend(appended); + *final_fragments = rewritten; Ok(()) } @@ -3527,6 +3550,28 @@ mod tests { assert_eq!(final_fragments, expected_fragments); } + #[test] + fn test_rewrite_fragments_missing_fragment() { + let mut final_fragments: Vec = (0..4).map(Fragment::new).collect(); + let rewrite_groups = vec![RewriteGroup { + old_fragments: vec![Fragment::new(7)], + new_fragments: vec![Fragment::new(0)], + }]; + + let mut fragment_id = 10; + let result = Transaction::handle_rewrite_fragments( + &mut final_fragments, + &rewrite_groups, + &mut fragment_id, + 42, + None, + ); + + let error = result.unwrap_err(); + assert!(matches!(error, Error::CommitConflict { .. })); + assert!(error.to_string().contains("id=7")); + } + #[test] fn test_merge_fragments_valid() { use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema}; From 535ec125cfecfca7c6daf065eb98303c7574cf05 Mon Sep 17 00:00:00 2001 From: tanish Date: Tue, 25 Aug 2026 21:17:56 +0000 Subject: [PATCH 2/2] fix: reject rewrite groups that replace the same fragment more than once Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- rust/lance/src/dataset/transaction.rs | 54 ++++++++++++++++++++++++--- 1 file changed, 48 insertions(+), 6 deletions(-) diff --git a/rust/lance/src/dataset/transaction.rs b/rust/lance/src/dataset/transaction.rs index 96df008f094..2971c522a8e 100644 --- a/rust/lance/src/dataset/transaction.rs +++ b/rust/lance/src/dataset/transaction.rs @@ -2549,6 +2549,17 @@ impl Transaction { Ok(()) } + fn duplicate_rewrite_error(version: u64, fragment_id: u64) -> Error { + Error::commit_conflict_source( + version, + format!( + "a rewrite operation attempts to replace fragment id={} more than once", + fragment_id + ) + .into(), + ) + } + fn handle_rewrite_fragments( final_fragments: &mut Vec, groups: &[RewriteGroup], @@ -2600,17 +2611,20 @@ impl Transaction { // from the original fragments. We don't modify it here. if contiguous { - for flag in removed - .iter_mut() - .skip(start) - .take(group.old_fragments.len()) - { - *flag = true; + for (i, old_fragment) in group.old_fragments.iter().enumerate() { + let position = start + i; + if removed[position] { + return Err(Self::duplicate_rewrite_error(version, old_fragment.id)); + } + removed[position] = true; } replacements.insert(start, new_fragments); } else { for old_fragment in group.old_fragments.iter() { if let Some(&position) = positions.get(&old_fragment.id) { + if removed[position] { + return Err(Self::duplicate_rewrite_error(version, old_fragment.id)); + } removed[position] = true; } } @@ -3572,6 +3586,34 @@ mod tests { assert!(error.to_string().contains("id=7")); } + #[test] + fn test_rewrite_fragments_duplicate_fragment() { + let mut final_fragments: Vec = (0..4).map(Fragment::new).collect(); + let rewrite_groups = vec![ + RewriteGroup { + old_fragments: vec![Fragment::new(1), Fragment::new(2)], + new_fragments: vec![Fragment::new(0)], + }, + RewriteGroup { + old_fragments: vec![Fragment::new(2), Fragment::new(3)], + new_fragments: vec![Fragment::new(0)], + }, + ]; + + let mut fragment_id = 10; + let result = Transaction::handle_rewrite_fragments( + &mut final_fragments, + &rewrite_groups, + &mut fragment_id, + 42, + None, + ); + + let error = result.unwrap_err(); + assert!(matches!(error, Error::CommitConflict { .. })); + assert!(error.to_string().contains("id=2")); + } + #[test] fn test_merge_fragments_valid() { use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema};