Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
6d3f3eb
feat: horizontal compaction via pluggable executor/committer
LuciferYang Sep 29, 2026
0cb9d48
refactor(compaction): make the executor/committer seam load-bearing
LuciferYang Sep 29, 2026
051f4ad
docs: avoid a public-to-private intra-doc link on RewriteExecutor
LuciferYang Sep 29, 2026
5d40322
Merge remote-tracking branch 'upstream/main' into feat/horizontal-com…
LuciferYang Sep 29, 2026
c873253
fix(compaction): bind horizontal rewrite results to their read version
LuciferYang Sep 29, 2026
b15e1e9
test(compaction): cover the horizontal rewrite read-version conflict
LuciferYang Sep 29, 2026
453db7b
fix(compaction): conflict horizontal rewrites with concurrent drops, …
LuciferYang Oct 2, 2026
f8b46ba
fix(compaction): count only files holding schema columns in column_la…
LuciferYang Oct 2, 2026
9e7f6e7
Merge remote-tracking branch 'upstream/main' into feat/horizontal-com…
LuciferYang Oct 2, 2026
30888e8
fix(compaction): pass the new write_fragments_direct argument in the …
LuciferYang Oct 2, 2026
668d999
feat(transaction): add data_change to DataReplacement and allow sever…
LuciferYang Oct 2, 2026
6c15726
feat(compaction): repack columns as a compact_files task committed vi…
LuciferYang Oct 2, 2026
f800078
refactor(compaction): keep scope a run option, warn once about stale …
LuciferYang Oct 2, 2026
e0b5600
fix(compaction): address review of the repack planner
LuciferYang Oct 2, 2026
b047330
fix(compaction): keep non-movable columns out of the merged repack file
LuciferYang Oct 2, 2026
e09c16c
fix(compaction): only merge files a repack empties
LuciferYang Oct 2, 2026
8c3a2f1
fix(compaction): count struct headers when deciding which files a rep…
LuciferYang Oct 2, 2026
513a52b
fix(compaction): apply the empties rule to group merges and leave kep…
LuciferYang Oct 2, 2026
64c13f8
fix(compaction): fall back to merging every other file when the kept …
LuciferYang Oct 2, 2026
715801d
fix(compaction): prefer a repack that leaves the largest file in place
LuciferYang Oct 2, 2026
b7f5e6b
Merge remote-tracking branch 'upstream/main' into feat/horizontal-com…
LuciferYang Oct 2, 2026
3cfb979
fix(compaction): plan nothing when no file holds a column
LuciferYang Oct 2, 2026
bf6f0a1
fix(compaction): box the compaction futures; fix the distributed repa…
LuciferYang Oct 2, 2026
0ccd47b
Merge remote-tracking branch 'upstream/main' into feat/horizontal-com…
LuciferYang Oct 3, 2026
fa4241f
feat(compaction): expose column layout stats in Python and Java; docu…
LuciferYang Oct 3, 2026
dd58904
feat: report file sizes, fields per file and dead slot ratio in colum…
LuciferYang Oct 3, 2026
0f6c23c
fix(java): release per-fragment local refs in getColumnLayoutStatistics
LuciferYang Oct 3, 2026
617ba4e
Merge remote-tracking branch 'upstream/main' into feat/horizontal-com…
LuciferYang Oct 4, 2026
104231c
chore: document DataReplacement group order in the proto
LuciferYang Oct 4, 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
40 changes: 38 additions & 2 deletions docs/src/guide/read_and_write.md
Original file line number Diff line number Diff line change
Expand Up @@ -497,8 +497,9 @@ ordering of the data will be preserved.

!!! note

Compaction creates a new version of the table. It does not delete the old
version of the table and the files referenced by it.
Compaction creates a new version of the table, or two when it both
rewrites fragments and repacks columns (see below). It does not delete the
old versions of the table and the files referenced by them.

```python
import lance
Expand All @@ -517,6 +518,41 @@ of this, it's recommended to rewrite files before re-building indices.

<!-- TODO: remove this last comment once stable row ids are default. -->

#### Repack columns

Each `add_columns` backfill gives every fragment one more data file. A large
fragment with few deletions is never rewritten, so those files pile up and slow
down scans and random reads. Compaction can repack a fragment's columns into
fewer data files instead. A repack moves no rows, so fragment ids, row
addresses, deletions, data overlays and index coverage stay as they are.

Repacking is off by default. Set either option, as a `compact_files` argument or
as table config:

- `max_data_files_per_fragment` (`lance.compaction.max_data_files_per_fragment`):
repack a fragment that holds its columns in more files than this.
- `column_groups` (`lance.compaction.column_groups`, groups separated by `;` and
columns by `,`): keep these top-level columns together in a file of their
own. Fragments that compaction rewrites are written this way too.

`scope` limits a run to one kind of work: `"rewrite_fragments"`,
`"repack_columns"`, or `"all"` (the default).

```python
dataset.optimize.compact_files(max_data_files_per_fragment=2)

# Only repack, and keep the embedding in a file of its own.
dataset.optimize.compact_files(column_groups=[["embedding"]], scope="repack_columns")

# Per fragment: live data files, file sizes, fields per file, dead field slots.
dataset.stats.column_layout_stats()
```

A repack only merges files it can empty. Files holding a blob column, a column
only partly present in the fragment, or spilled row lineage stay as they are,
so a fragment can stay above the limit. Fragments with legacy (V1) data files
are not repacked.

### Cleanup old versions

Lance is an immutable format — every write creates a new version. The new version
Expand Down
113 changes: 108 additions & 5 deletions java/lance-jni/src/blocking_dataset.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,9 @@ use crate::namespace::{
use crate::session::{handle_from_session, session_from_handle};
use crate::traits::{FromJObjectWithEnv, FromJString, export_vec, import_vec, import_vec_to_rust};
use crate::utils::{
build_compaction_options, extract_base_store_params, extract_storage_options,
extract_write_params, get_scalar_index_params, get_vector_index_params, to_java_map,
to_rust_map,
apply_repack_options, build_compaction_options, extract_base_store_params,
extract_storage_options, extract_write_params, get_scalar_index_params,
get_vector_index_params, to_java_map, to_rust_map,
};
use crate::{block_on, traits::IntoJava};
use arrow::array::RecordBatchReader;
Expand Down Expand Up @@ -1964,6 +1964,86 @@ fn inner_get_fragment_statistics<'local>(
)?)
}

#[unsafe(no_mangle)]
pub extern "system" fn Java_org_lance_Dataset_nativeGetColumnLayoutStatistics<'a>(
mut env: JNIEnv<'a>,
jdataset: JObject,
) -> JObject<'a> {
ok_or_throw!(env, inner_get_column_layout_statistics(&mut env, jdataset))
}

/// `Dataset::column_layout_stats` as parallel Java arrays.
fn inner_get_column_layout_statistics<'local>(
env: &mut JNIEnv<'local>,
jdataset: JObject,
) -> Result<JObject<'local>> {
let stats = {
let dataset =
unsafe { env.get_rust_field::<_, _, BlockingDataset>(jdataset, NATIVE_DATASET) }?;
dataset.inner.column_layout_stats()
};
let to_int = |value: usize| {
i32::try_from(value)
.map_err(|_| Error::runtime_error(format!("{value} does not fit in a Java int")))
};
let len = to_int(stats.len())?;
let fragment_ids: Vec<i64> = stats.iter().map(|s| s.fragment_id as i64).collect();
let live_file_counts = stats
.iter()
.map(|s| to_int(s.live_file_count))
.collect::<Result<Vec<_>>>()?;
let ratios: Vec<f64> = stats.iter().map(|s| s.tombstoned_field_ratio).collect();
let overlay_counts = stats
.iter()
.map(|s| to_int(s.overlay_count))
.collect::<Result<Vec<_>>>()?;

let jfragment_ids = env.new_long_array(len)?;
let jlive_file_counts = env.new_int_array(len)?;
let jratios = env.new_double_array(len)?;
let joverlay_counts = env.new_int_array(len)?;
env.set_long_array_region(&jfragment_ids, 0, &fragment_ids)?;
env.set_int_array_region(&jlive_file_counts, 0, &live_file_counts)?;
env.set_double_array_region(&jratios, 0, &ratios)?;
env.set_int_array_region(&joverlay_counts, 0, &overlay_counts)?;

let jfile_sizes = env.new_object_array(len, "[J", JObject::null())?;
let jfields_per_file = env.new_object_array(len, "[I", JObject::null())?;
for (index, s) in stats.iter().enumerate() {
let sizes: Vec<i64> = s
.file_sizes
.iter()
.map(|size| size.map_or(-1, |size| size as i64))
.collect();
let fields = s
.fields_per_file
.iter()
.map(|count| to_int(*count))
.collect::<Result<Vec<_>>>()?;
let jsizes = env.new_long_array(to_int(sizes.len())?)?;
env.set_long_array_region(&jsizes, 0, &sizes)?;
env.set_object_array_element(&jfile_sizes, index as i32, &jsizes)?;
env.delete_local_ref(jsizes)?;
let jfields = env.new_int_array(to_int(fields.len())?)?;
env.set_int_array_region(&jfields, 0, &fields)?;
env.set_object_array_element(&jfields_per_file, index as i32, &jfields)?;
env.delete_local_ref(jfields)?;
}

Ok(env.new_object(
"org/lance/ColumnLayoutStatistics",
"([J[I[[J[[I[D[I)V",
&[
JValue::Object(&jfragment_ids),
JValue::Object(&jlive_file_counts),
JValue::Object(&jfile_sizes),
JValue::Object(&jfields_per_file),
JValue::Object(&jratios),
JValue::Object(&joverlay_counts),
],
)?)
}

#[unsafe(no_mangle)]
pub extern "system" fn Java_org_lance_Dataset_getFragmentNative<'a>(
mut env: JNIEnv<'a>,
Expand Down Expand Up @@ -3628,7 +3708,22 @@ fn convert_java_compaction_options_to_rust(
)?
.l()?;

build_compaction_options(
let max_data_files_per_fragment = env
.call_method(
&java_options,
"getMaxDataFilesPerFragment",
"()Ljava/util/Optional;",
&[],
)?
.l()?;
let column_groups = env
.call_method(&java_options, "getColumnGroups", "()Ljava/util/List;", &[])?
.l()?;
let scope = env
.call_method(&java_options, "getScope", "()Ljava/util/Optional;", &[])?
.l()?;

let mut options = build_compaction_options(
env,
&target_rows_per_fragment,
&max_rows_per_group,
Expand All @@ -3646,7 +3741,15 @@ fn convert_java_compaction_options_to_rust(
&excluded_fragment_ids,
&data_storage_version,
config,
)
)?;
apply_repack_options(
env,
&mut options,
&max_data_files_per_fragment,
&column_groups,
&scope,
)?;
Ok(options)
}

#[unsafe(no_mangle)]
Expand Down
Loading
Loading