diff --git a/python/Cargo.lock b/python/Cargo.lock index 39affa28041..aad67e708cb 100644 --- a/python/Cargo.lock +++ b/python/Cargo.lock @@ -5882,7 +5882,7 @@ dependencies = [ [[package]] name = "pylance" -version = "7.0.0" +version = "7.0.0+exa.1" dependencies = [ "arrow", "arrow-array", diff --git a/python/Cargo.toml b/python/Cargo.toml index 73fbeb623a9..57a758e7825 100644 --- a/python/Cargo.toml +++ b/python/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "pylance" -version = "7.0.0" +version = "7.0.0+exa.1" edition = "2024" authors = ["Lance Devs "] license = "Apache-2.0" diff --git a/rust/lance/src/dataset/transaction.rs b/rust/lance/src/dataset/transaction.rs index 3f96b9964d5..02052c64ba1 100644 --- a/rust/lance/src/dataset/transaction.rs +++ b/rust/lance/src/dataset/transaction.rs @@ -1890,13 +1890,15 @@ impl Transaction { .. } => { // Remove the deleted fragments + // Hash lookups keep this linear on tables with many fragments. + let deleted_ids: HashSet = deleted_fragment_ids.iter().copied().collect(); + let updated_by_id: HashMap = + updated_fragments.iter().map(|f| (f.id, f)).collect(); final_fragments.extend(maybe_existing_fragments?.clone()); - final_fragments.retain(|f| !deleted_fragment_ids.contains(&f.id)); + final_fragments.retain(|f| !deleted_ids.contains(&f.id)); final_fragments.iter_mut().for_each(|f| { - for updated in updated_fragments { - if updated.id == f.id { - *f = updated.clone(); - } + if let Some(updated) = updated_by_id.get(&f.id) { + *f = (*updated).clone(); } }); Self::retain_relevant_indices(&mut final_indices, &schema, &final_fragments) @@ -1916,13 +1918,20 @@ impl Transaction { let existing_fragments = maybe_existing_fragments?; // Apply updates to existing fragments + // Hash lookups keep this linear on tables with many fragments. + let removed_ids: HashSet = removed_fragment_ids.iter().copied().collect(); + let mut updated_by_id: HashMap = + HashMap::with_capacity(updated_fragments.len()); + for fragment in updated_fragments { + updated_by_id.entry(fragment.id).or_insert(fragment); + } let updated_frags: Vec = existing_fragments .iter() .filter_map(|f| { - if removed_fragment_ids.contains(&f.id) { + if removed_ids.contains(&f.id) { return None; } - if let Some(updated) = updated_fragments.iter().find(|uf| uf.id == f.id) { + if let Some(&updated) = updated_by_id.get(&f.id) { Some(updated.clone()) } else { Some(f.clone()) @@ -3851,10 +3860,14 @@ mod tests { use crate::session::Session; fn sample_manifest() -> Manifest { + sample_manifest_with_fragments(0..1) + } + + fn sample_manifest_with_fragments(ids: std::ops::Range) -> Manifest { let schema = ArrowSchema::new(vec![ArrowField::new("id", DataType::Int32, false)]); Manifest::new( LanceSchema::try_from(&schema).unwrap(), - Arc::new(vec![Fragment::new(0)]), + Arc::new(ids.map(Fragment::new).collect()), DataStorageFormat::new(LanceFileVersion::V2_0), HashMap::new(), ) @@ -4080,6 +4093,88 @@ mod tests { ); } + #[test] + fn test_update_build_manifest_replaces_and_removes_fragments() { + let manifest = sample_manifest_with_fragments(0..5); + + let mut updated2 = Fragment::new(2); + updated2.physical_rows = Some(42); + let mut updated4 = Fragment::new(4); + updated4.physical_rows = Some(43); + + let transaction = Transaction::new( + manifest.version, + Operation::Update { + removed_fragment_ids: vec![1], + // Fragment 99 does not exist in the dataset; it must be ignored, + // not appended. + updated_fragments: vec![updated2, updated4, Fragment::new(99)], + new_fragments: vec![], + fields_modified: vec![], + merged_generations: vec![], + fields_for_preserving_frag_bitmap: vec![], + update_mode: None, + inserted_rows_filter: None, + updated_fragment_offsets: None, + }, + None, + ); + + let (new_manifest, _) = transaction + .build_manifest( + Some(&manifest), + vec![], + "txn", + &ManifestWriteConfig::default(), + ) + .unwrap(); + + let ids: Vec = new_manifest.fragments.iter().map(|f| f.id).collect(); + assert_eq!(ids, vec![0, 2, 3, 4]); + let rows: Vec> = new_manifest + .fragments + .iter() + .map(|f| f.physical_rows) + .collect(); + assert_eq!(rows, vec![None, Some(42), None, Some(43)]); + } + + #[test] + fn test_delete_build_manifest_replaces_and_removes_fragments() { + let manifest = sample_manifest_with_fragments(0..5); + + let mut updated2 = Fragment::new(2); + updated2.physical_rows = Some(42); + + let transaction = Transaction::new( + manifest.version, + Operation::Delete { + updated_fragments: vec![updated2], + deleted_fragment_ids: vec![1, 3], + predicate: "id > 0".to_string(), + }, + None, + ); + + let (new_manifest, _) = transaction + .build_manifest( + Some(&manifest), + vec![], + "txn", + &ManifestWriteConfig::default(), + ) + .unwrap(); + + let ids: Vec = new_manifest.fragments.iter().map(|f| f.id).collect(); + assert_eq!(ids, vec![0, 2, 4]); + let rows: Vec> = new_manifest + .fragments + .iter() + .map(|f| f.physical_rows) + .collect(); + assert_eq!(rows, vec![None, Some(42), None]); + } + #[test] fn test_remove_tombstoned_data_files() { // Create a fragment with mixed data files: some normal, some fully tombstoned