Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion python/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion python/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "pylance"
version = "7.0.0"
version = "7.0.0+exa.1"
edition = "2024"
authors = ["Lance Devs <dev@lance.org>"]
license = "Apache-2.0"
Expand Down
111 changes: 103 additions & 8 deletions rust/lance/src/dataset/transaction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1890,13 +1890,15 @@ impl Transaction {
..
} => {
// Remove the deleted fragments
// Hash lookups keep this linear on tables with many fragments.
let deleted_ids: HashSet<u64> = deleted_fragment_ids.iter().copied().collect();
let updated_by_id: HashMap<u64, &Fragment> =
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)
Expand All @@ -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<u64> = removed_fragment_ids.iter().copied().collect();
let mut updated_by_id: HashMap<u64, &Fragment> =
HashMap::with_capacity(updated_fragments.len());
for fragment in updated_fragments {
updated_by_id.entry(fragment.id).or_insert(fragment);
}
let updated_frags: Vec<Fragment> = 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())
Expand Down Expand Up @@ -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<u64>) -> 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(),
)
Expand Down Expand Up @@ -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<u64> = new_manifest.fragments.iter().map(|f| f.id).collect();
assert_eq!(ids, vec![0, 2, 3, 4]);
let rows: Vec<Option<usize>> = 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<u64> = new_manifest.fragments.iter().map(|f| f.id).collect();
assert_eq!(ids, vec![0, 2, 4]);
let rows: Vec<Option<usize>> = 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
Expand Down
Loading