Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
b9fdc8e
fix(compaction): preserve task order during commit
lance-gatefixer[bot] Aug 7, 2026
46c1744
fix(compaction): preserve bounded compaction order
lance-gatefixer[bot] Aug 7, 2026
e47e8df
fix(compaction): preserve logical fragment order
lance-gatefixer[bot] Aug 7, 2026
3e8e101
fix(compaction): gate logical fragment ordering
lance-gatefixer[bot] Aug 7, 2026
93b75ef
fix(compaction): preserve order without format changes
lance-gatefixer[bot] Aug 7, 2026
840b7d8
test: align index coverage with compaction relabeling
lance-gatefixer[bot] Aug 7, 2026
56ee824
fix(compaction): bound suffix relabel work
lance-gatefixer[bot] Aug 7, 2026
d23b189
Merge remote-tracking branch 'origin/main' into gatekeeper/fix-3465-1
lance-gatefixer[bot] Aug 11, 2026
b7d6961
Merge remote-tracking branch 'origin/main' into gatekeeper/fix-3465-1
lance-gatefixer[bot] Aug 13, 2026
c5d7501
Merge remote-tracking branch 'origin/main' into gatekeeper/fix-3465-1
lance-gatefixer[bot] Aug 13, 2026
8f8da60
Merge remote-tracking branch 'origin/main' into gatekeeper/fix-3465-1
lance-gatefixer[bot] Aug 19, 2026
abf1119
fix(compaction): preserve excluded fragment identities
lance-gatefixer[bot] Aug 19, 2026
0392526
Merge remote-tracking branch 'origin/main' into gatekeeper/fix-3465-1
lance-gatefixer[bot] Aug 19, 2026
c68a81f
Merge remote-tracking branch 'refs/remotes/origin/main' into gatekeep…
lance-gatefixer[bot] Aug 31, 2026
69ec884
Merge remote-tracking branch 'refs/remotes/origin/main' into gatekeep…
lance-gatefixer[bot] Aug 31, 2026
58b67b4
Merge remote-tracking branch 'refs/remotes/origin/main' into gatekeep…
lance-gatefixer[bot] Aug 31, 2026
5f1d37c
Merge remote-tracking branch 'origin/main' into gatekeeper/fix-3465-1
lance-gatefixer[bot] Sep 14, 2026
a6b4a7a
Merge remote-tracking branch 'origin/main' into gatekeeper/fix-3465-1
lance-gatefixer[bot] Sep 14, 2026
5ce3b3f
test(compaction): align storage version coverage
lance-gatefixer[bot] Sep 14, 2026
848e002
Merge remote-tracking branch 'origin/main' into gatekeeper/fix-3465-1
lance-gatefixer[bot] Sep 14, 2026
3a2eeab
Merge main into compaction row-order fix
lance-gatefixer[bot] Sep 27, 2026
4f9a19f
test(compaction): align index fixtures with suffix relabeling
lance-gatefixer[bot] Sep 27, 2026
929761d
Merge main into gatekeeper/fix-3465-1
lance-gatefixer[bot] Sep 29, 2026
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
11 changes: 6 additions & 5 deletions java/src/main/java/org/lance/compaction/CompactionOptions.java
Original file line number Diff line number Diff line change
Expand Up @@ -334,9 +334,9 @@ public Builder withBinaryCopyReadBatchBytes(long binaryCopyReadBatchBytes) {
}

/**
* Maximum number of source fragments to compact in a single run. Tasks are included until
* adding the next task would exceed this limit, allowing for incremental compaction. Fragments
* are processed oldest first.
* Maximum number of fragments whose identities may change in a single run. This includes
* compacted source fragments and following fragments relabeled to preserve row order. The
* planner selects candidates from a suffix within this limit.
*
* @throws IllegalArgumentException if {@code maxSourceFragments} is not positive
*/
Expand Down Expand Up @@ -384,8 +384,9 @@ public Builder withDataStorageVersion(DataStorageVersion version) {

/**
* Fragment IDs to exclude from compaction planning. Excluded fragments remain unchanged and act
* as boundaries, so fragments on opposite sides are not combined into the same task. Duplicate
* and unknown IDs are ignored.
* as boundaries, so fragments on opposite sides are not combined into the same task. To
* preserve row order, only candidates after the last present excluded fragment are eligible.
* Duplicate and unknown IDs are ignored.
*
* @throws IllegalArgumentException if an ID is negative or exceeds the unsigned 32-bit range
*/
Expand Down
14 changes: 8 additions & 6 deletions python/python/lance/dataset.py
Original file line number Diff line number Diff line change
Expand Up @@ -7696,11 +7696,11 @@ def compact_files(
Controls how much data is read at once when performing binary copy.
Defaults to 16MB.
max_source_fragments: int, optional
Maximum number of source fragments to compact in a single run.
Compaction tasks are included until adding the next task would
exceed this limit, allowing compaction to proceed incrementally.
Fragments are processed oldest first. If not specified, uses the
manifest config value, or applies no limit.
Maximum number of fragments whose identities may change in a
single run. This includes compacted source fragments and following
fragments relabeled to preserve row order. The planner selects
candidates from a suffix within this limit. If not specified, uses
the manifest config value, or applies no limit.
max_source_rows: int, optional
Maximum number of source rows to compact in a single run. Rows are
counted as live rows (physical rows minus soft-deleted rows).
Expand All @@ -7717,7 +7717,9 @@ def compact_files(
Fragment IDs to exclude from compaction planning. Excluded
fragments remain unchanged and act as boundaries, so fragments
on opposite sides are not combined into the same compaction task.
Duplicate and unknown IDs are ignored.
To preserve row order, only candidates after the last present
excluded fragment are eligible. Duplicate and unknown IDs are
ignored.
data_storage_version: str, optional
Output data file version, such as "2.2", "stable", or "next".
Uses the compaction config target when set, otherwise the dataset's
Expand Down
13 changes: 7 additions & 6 deletions python/python/lance/optimize.py
Original file line number Diff line number Diff line change
Expand Up @@ -91,10 +91,10 @@ class CompactionOptions(TypedDict, total=False):
"""
max_source_fragments: Optional[int]
"""
Maximum number of source fragments to compact in a single run. Tasks
are included until adding the next task would exceed this limit,
allowing for incremental compaction (e.g., compact 20 fragments at a
time). Fragments are processed oldest first.
Maximum number of fragments whose identities may change in a single run.
This includes compacted source fragments and following fragments relabeled
to preserve row order. The planner selects candidates from a suffix within
this limit.
(default: None, no limit)
"""
max_source_rows: Optional[int]
Expand All @@ -118,8 +118,9 @@ class CompactionOptions(TypedDict, total=False):
"""
Fragment IDs to exclude from compaction planning. Excluded fragments
remain unchanged and act as boundaries, so fragments on opposite sides
are not combined into the same task. Duplicate and unknown IDs are
ignored. (default: None)
are not combined into the same task. To preserve row order, only candidates
after the last present excluded fragment are eligible. Duplicate and
unknown IDs are ignored. (default: None)
"""
data_storage_version: Optional[str]
"""
Expand Down
15 changes: 15 additions & 0 deletions rust/lance-table/src/transaction/index_maintenance.rs
Original file line number Diff line number Diff line change
Expand Up @@ -476,6 +476,21 @@ impl Transaction {
groups: &[RewriteGroup],
) {
for group in groups {
// Order-preserving partial compaction represents untouched trailing
// fragments as exact metadata copies with a new ID. Their overlays
// are still live and must not invalidate index coverage.
if let ([old_fragment], [new_fragment]) = (
group.old_fragments.as_slice(),
group.new_fragments.as_slice(),
) {
let mut relabeled_fragment = new_fragment.clone();
relabeled_fragment.id = old_fragment.id;
relabeled_fragment.deletion_file = old_fragment.deletion_file.clone();
if relabeled_fragment == *old_fragment {
continue;
}
}

// field id -> newest overlay committed_version supplying that field
let mut overlaid_field_versions: HashMap<i32, u64> = HashMap::new();
for old_frag in &group.old_fragments {
Expand Down
108 changes: 108 additions & 0 deletions rust/lance-table/src/transaction/manifest_build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3202,6 +3202,114 @@ mod tests {
assert_eq!(last_updated_at_versions(&new_manifest, 2), vec![2; 42]);
}

#[test]
fn test_rewrite_build_manifest_keeps_fragment_ids_sorted() {
let manifest = sample_manifest_with_fragments(0..4);
let transaction = Transaction::new(
manifest.version,
Operation::Rewrite {
groups: vec![RewriteGroup {
old_fragments: vec![Fragment::new(1), Fragment::new(2)],
new_fragments: vec![Fragment::new(0)],
}],
rewritten_indices: vec![],
frag_reuse_index: None,
},
None,
);

let (rewritten, _) = transaction
.build_manifest(Some(&manifest), vec![], "txn", &default_build_config())
.unwrap();

let fragment_ids = rewritten
.fragments
.iter()
.map(|fragment| fragment.id)
.collect::<Vec<_>>();
assert_eq!(fragment_ids, vec![0, 3, 4]);
}

#[test]
fn test_rewrite_after_row_adding_update_keeps_fragment_ids_sorted() {
let manifest = sample_manifest_with_fragments(0..3);
let update = Transaction::new(
manifest.version,
Operation::Update {
removed_fragment_ids: vec![],
updated_fragments: vec![],
new_fragments: vec![Fragment::new(0)],
fields_modified: vec![],
compacted_sstables: vec![],
fields_for_preserving_frag_bitmap: vec![],
update_mode: Some(UpdateMode::RewriteRows),
inserted_rows_filter: None,
updated_fragment_offsets: None,
},
None,
);
let (updated, _) = update
.build_manifest(
Some(&manifest),
vec![],
"update-txn",
&default_build_config(),
)
.unwrap();

let rewrite = Transaction::new(
manifest.version,
Operation::Rewrite {
groups: vec![RewriteGroup {
old_fragments: vec![Fragment::new(0)],
new_fragments: vec![Fragment::new(0)],
}],
rewritten_indices: vec![],
frag_reuse_index: None,
},
None,
);
let (rewritten, _) = rewrite
.build_manifest(
Some(&updated),
vec![],
"rewrite-txn",
&default_build_config(),
)
.unwrap();

let fragment_ids = rewritten
.fragments
.iter()
.map(|fragment| fragment.id)
.collect::<Vec<_>>();
assert_eq!(fragment_ids, vec![1, 2, 3, 4]);
}

#[test]
fn test_merge_build_manifest_keeps_fragment_ids_sorted() {
let manifest = sample_manifest_with_fragments(0..3);
let fragments = manifest.fragments.iter().rev().cloned().collect::<Vec<_>>();
let operation = Operation::Merge {
fragments,
schema: manifest.schema.clone(),
preserves_nullability: true,
};
validate_operation(Some(&manifest), &operation).unwrap();
let transaction = Transaction::new(manifest.version, operation, None);

let (merged, _) = transaction
.build_manifest(Some(&manifest), vec![], "txn", &default_build_config())
.unwrap();

let fragment_ids = merged
.fragments
.iter()
.map(|fragment| fragment.id)
.collect::<Vec<_>>();
assert_eq!(fragment_ids, vec![0, 1, 2]);
}

#[test]
fn test_remove_tombstoned_data_files() {
// Create a fragment with mixed data files: some normal, some fully tombstoned
Expand Down
12 changes: 7 additions & 5 deletions rust/lance/src/dataset/index/frag_reuse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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));

Expand Down Expand Up @@ -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
Expand Down
Loading
Loading