Skip to content

Commit 382d8d3

Browse files
authored
feat: bounding source fragments for compaction execution (#6232)
Adding a `CompationOption` to allow bounding the number of source fragments to be included in a compaction execution. This effectively allows incremental compaction, for example compact the first 100k fragments, then next, etc; which is especially useful for heavily fragmented datasets.
1 parent 43c7f0b commit 382d8d3

9 files changed

Lines changed: 218 additions & 14 deletions

File tree

java/lance-jni/src/blocking_dataset.rs

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2633,6 +2633,14 @@ fn convert_java_compaction_options_to_rust(
26332633
&[],
26342634
)?
26352635
.l()?;
2636+
let max_source_fragments = env
2637+
.call_method(
2638+
&java_options,
2639+
"getMaxSourceFragments",
2640+
"()Ljava/util/Optional;",
2641+
&[],
2642+
)?
2643+
.l()?;
26362644

26372645
build_compaction_options(
26382646
env,
@@ -2646,6 +2654,7 @@ fn convert_java_compaction_options_to_rust(
26462654
&defer_index_remap,
26472655
&compaction_mode,
26482656
&binary_copy_read_batch_bytes,
2657+
&max_source_fragments,
26492658
config,
26502659
)
26512660
}

java/lance-jni/src/optimize.rs

Lines changed: 19 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,7 @@ pub extern "system" fn Java_org_lance_compaction_Compaction_nativePlanCompaction
4545
defer_index_remap: JObject, // Optional<Boolean>
4646
compaction_mode: JObject, // Optional<String>
4747
binary_copy_read_batch_bytes: JObject, // Optional<Long>
48+
max_source_fragments: JObject, // Optional<Long>
4849
) -> JObject<'local> {
4950
ok_or_throw_with_return!(
5051
env,
@@ -60,7 +61,8 @@ pub extern "system" fn Java_org_lance_compaction_Compaction_nativePlanCompaction
6061
batch_size,
6162
defer_index_remap,
6263
compaction_mode,
63-
binary_copy_read_batch_bytes
64+
binary_copy_read_batch_bytes,
65+
max_source_fragments
6466
),
6567
JObject::null()
6668
)
@@ -80,6 +82,7 @@ fn inner_plan_compaction<'local>(
8082
defer_index_remap: JObject, // Optional<Boolean>
8183
compaction_mode: JObject, // Optional<String>
8284
binary_copy_read_batch_bytes: JObject, // Optional<Long>
85+
max_source_fragments: JObject, // Optional<Long>
8386
) -> Result<JObject<'local>> {
8487
let config = {
8588
let dataset =
@@ -98,6 +101,7 @@ fn inner_plan_compaction<'local>(
98101
&defer_index_remap,
99102
&compaction_mode,
100103
&binary_copy_read_batch_bytes,
104+
&max_source_fragments,
101105
&config,
102106
)?;
103107

@@ -125,6 +129,7 @@ pub extern "system" fn Java_org_lance_compaction_Compaction_nativeCommitCompacti
125129
defer_index_remap: JObject, // Optional<Boolean>
126130
compaction_mode: JObject, // Optional<String>
127131
binary_copy_read_batch_bytes: JObject, // Optional<Long>
132+
max_source_fragments: JObject, // Optional<Long>
128133
) -> JObject<'local> {
129134
ok_or_throw_with_return!(
130135
env,
@@ -142,6 +147,7 @@ pub extern "system" fn Java_org_lance_compaction_Compaction_nativeCommitCompacti
142147
defer_index_remap,
143148
compaction_mode,
144149
binary_copy_read_batch_bytes,
150+
max_source_fragments,
145151
),
146152
JObject::null()
147153
)
@@ -162,6 +168,7 @@ fn inner_commit_compaction<'local>(
162168
defer_index_remap: JObject, // Optional<Boolean>
163169
compaction_mode: JObject, // Optional<String>
164170
binary_copy_read_batch_bytes: JObject, // Optional<Long>
171+
max_source_fragments: JObject, // Optional<Long>
165172
) -> Result<JObject<'local>> {
166173
let config = {
167174
let dataset =
@@ -180,6 +187,7 @@ fn inner_commit_compaction<'local>(
180187
&defer_index_remap,
181188
&compaction_mode,
182189
&binary_copy_read_batch_bytes,
190+
&max_source_fragments,
183191
&config,
184192
)?;
185193
let completed_tasks = import_vec_to_rust(env, &rewrite_results, |env, rewrite_result| {
@@ -216,6 +224,7 @@ pub extern "system" fn Java_org_lance_compaction_CompactionTask_nativeExecute<'l
216224
defer_index_remap: JObject, // Optional<Boolean>
217225
compaction_mode: JObject, // Optional<String>
218226
binary_copy_read_batch_bytes: JObject, // Optional<Long>
227+
max_source_fragments: JObject, // Optional<Long>
219228
) -> JObject<'local> {
220229
ok_or_throw_with_return!(
221230
env,
@@ -233,7 +242,8 @@ pub extern "system" fn Java_org_lance_compaction_CompactionTask_nativeExecute<'l
233242
batch_size,
234243
defer_index_remap,
235244
compaction_mode,
236-
binary_copy_read_batch_bytes
245+
binary_copy_read_batch_bytes,
246+
max_source_fragments
237247
),
238248
JObject::null()
239249
)
@@ -255,6 +265,7 @@ fn inner_execute_task<'local>(
255265
defer_index_remap: JObject, // Optional<Boolean>
256266
compaction_mode: JObject, // Optional<String>
257267
binary_copy_read_batch_bytes: JObject, // Optional<Long>
268+
max_source_fragments: JObject, // Optional<Long>
258269
) -> Result<JObject<'local>> {
259270
let task_data: TaskData = task_data.extract_object(env)?;
260271
let config = {
@@ -274,6 +285,7 @@ fn inner_execute_task<'local>(
274285
&defer_index_remap,
275286
&compaction_mode,
276287
&binary_copy_read_batch_bytes,
288+
&max_source_fragments,
277289
&config,
278290
)?;
279291
let compaction_task = CompactionTask {
@@ -300,7 +312,7 @@ const REWRITE_RESULT_CLASS: &str = "org/lance/compaction/RewriteResult";
300312
const REWRITE_RESULT_CONSTRUCTOR_SIG: &str =
301313
"(Lorg/lance/compaction/CompactionMetrics;Ljava/util/List;Ljava/util/List;J[B)V";
302314
const COMPACTION_OPTIONS_CLASS: &str = "org/lance/compaction/CompactionOptions";
303-
const COMPACTION_OPTIONS_CONSTRUCTOR_SIG: &str = "(Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;)V";
315+
const COMPACTION_OPTIONS_CONSTRUCTOR_SIG: &str = "(Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;Ljava/util/Optional;)V";
304316

305317
impl IntoJava for &TaskData {
306318
fn into_java<'a>(self, env: &mut JNIEnv<'a>) -> Result<JObject<'a>> {
@@ -362,6 +374,9 @@ impl IntoJava for &CompactionOptions {
362374
let binary_copy_read_batch_bytes =
363375
to_java_long_obj(env, self.binary_copy_read_batch_bytes.map(|v| v as i64))?;
364376
let binary_copy_read_batch_bytes_opt = to_java_optional(env, binary_copy_read_batch_bytes)?;
377+
let max_source_fragments =
378+
to_java_long_obj(env, self.max_source_fragments.map(|v| v as i64))?;
379+
let max_source_fragments_opt = to_java_optional(env, max_source_fragments)?;
365380

366381
Ok(env.new_object(
367382
COMPACTION_OPTIONS_CLASS,
@@ -377,6 +392,7 @@ impl IntoJava for &CompactionOptions {
377392
JValueGen::Object(&defer_index_remap_opt),
378393
JValueGen::Object(&compaction_mode_opt),
379394
JValueGen::Object(&binary_copy_read_batch_bytes_opt),
395+
JValueGen::Object(&max_source_fragments_opt),
380396
],
381397
)?)
382398
}

java/lance-jni/src/utils.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -141,6 +141,7 @@ pub fn build_compaction_options(
141141
defer_index_remap: &JObject, // Optional<Boolean>
142142
compaction_mode: &JObject, // Optional<String>
143143
binary_copy_read_batch_bytes: &JObject, // Optional<Long>
144+
max_source_fragments: &JObject, // Optional<Long>
144145
config: &std::collections::HashMap<String, String>,
145146
) -> Result<CompactionOptions> {
146147
let mut compaction_options = CompactionOptions::from_dataset_config(config)?;
@@ -181,6 +182,9 @@ pub fn build_compaction_options(
181182
compaction_options.binary_copy_read_batch_bytes =
182183
Some(binary_copy_read_batch_bytes_val as usize);
183184
}
185+
if let Some(max_source_fragments_val) = env.get_long_opt(max_source_fragments)? {
186+
compaction_options.max_source_fragments = Some(max_source_fragments_val as usize);
187+
}
184188

185189
Ok(compaction_options)
186190
}

java/src/main/java/org/lance/compaction/Compaction.java

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,8 @@ public static CompactionPlan planCompaction(
4343
compactionOptions.getBatchSize(),
4444
compactionOptions.getDeferIndexRemap(),
4545
compactionOptions.getCompactionMode(),
46-
compactionOptions.getBinaryCopyReadBatchBytes());
46+
compactionOptions.getBinaryCopyReadBatchBytes(),
47+
compactionOptions.getMaxSourceFragments());
4748
}
4849

4950
public static CompactionMetrics commitCompaction(
@@ -63,7 +64,8 @@ public static CompactionMetrics commitCompaction(
6364
compactionOptions.getBatchSize(),
6465
compactionOptions.getDeferIndexRemap(),
6566
compactionOptions.getCompactionMode(),
66-
compactionOptions.getBinaryCopyReadBatchBytes());
67+
compactionOptions.getBinaryCopyReadBatchBytes(),
68+
compactionOptions.getMaxSourceFragments());
6769
}
6870

6971
public static native CompactionMetrics nativeCommitCompaction(
@@ -78,7 +80,8 @@ public static native CompactionMetrics nativeCommitCompaction(
7880
Optional<Long> batchSize,
7981
Optional<Boolean> deferIndexRemap,
8082
Optional<String> compactionMode,
81-
Optional<Long> binaryCopyReadBatchBytes);
83+
Optional<Long> binaryCopyReadBatchBytes,
84+
Optional<Long> maxSourceFragments);
8285

8386
private static native CompactionPlan nativePlanCompaction(
8487
Dataset dataset,
@@ -91,5 +94,6 @@ private static native CompactionPlan nativePlanCompaction(
9194
Optional<Long> batchSize,
9295
Optional<Boolean> deferIndexRemap,
9396
Optional<String> compactionMode,
94-
Optional<Long> binaryCopyReadBatchBytes);
97+
Optional<Long> binaryCopyReadBatchBytes,
98+
Optional<Long> maxSourceFragments);
9599
}

java/src/main/java/org/lance/compaction/CompactionOptions.java

Lines changed: 24 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@ public class CompactionOptions implements Serializable {
3939
private Optional<Boolean> deferIndexRemap;
4040
private Optional<CompactionMode> compactionMode;
4141
private Optional<Long> binaryCopyReadBatchBytes;
42+
private Optional<Long> maxSourceFragments;
4243

4344
private CompactionOptions(
4445
Optional<Long> targetRowsPerFragment,
@@ -50,7 +51,8 @@ private CompactionOptions(
5051
Optional<Long> batchSize,
5152
Optional<Boolean> deferIndexRemap,
5253
Optional<CompactionMode> compactionMode,
53-
Optional<Long> binaryCopyReadBatchBytes) {
54+
Optional<Long> binaryCopyReadBatchBytes,
55+
Optional<Long> maxSourceFragments) {
5456
this.targetRowsPerFragment = targetRowsPerFragment;
5557
this.maxRowsPerGroup = maxRowsPerGroup;
5658
this.maxBytesPerFile = maxBytesPerFile;
@@ -61,6 +63,7 @@ private CompactionOptions(
6163
this.deferIndexRemap = deferIndexRemap;
6264
this.compactionMode = compactionMode;
6365
this.binaryCopyReadBatchBytes = binaryCopyReadBatchBytes;
66+
this.maxSourceFragments = maxSourceFragments;
6467
}
6568

6669
public Optional<Boolean> getDeferIndexRemap() {
@@ -76,6 +79,10 @@ public Optional<Long> getBinaryCopyReadBatchBytes() {
7679
return binaryCopyReadBatchBytes;
7780
}
7881

82+
public Optional<Long> getMaxSourceFragments() {
83+
return maxSourceFragments;
84+
}
85+
7986
public Optional<Boolean> getMaterializeDeletions() {
8087
return materializeDeletions;
8188
}
@@ -121,6 +128,7 @@ public String toString() {
121128
.add("deferIndexRemap", deferIndexRemap.orElse(null))
122129
.add("compactionMode", compactionMode.orElse(null))
123130
.add("binaryCopyReadBatchBytes", binaryCopyReadBatchBytes.orElse(null))
131+
.add("maxSourceFragments", maxSourceFragments.orElse(null))
124132
.toString();
125133
}
126134

@@ -135,6 +143,7 @@ private void writeObject(ObjectOutputStream output) throws IOException {
135143
output.writeObject(deferIndexRemap.orElse(null));
136144
output.writeObject(compactionMode.map(CompactionMode::getValue).orElse(null));
137145
output.writeObject(binaryCopyReadBatchBytes.orElse(null));
146+
output.writeObject(maxSourceFragments.orElse(null));
138147
}
139148

140149
private void readObject(ObjectInputStream input) throws IOException, ClassNotFoundException {
@@ -157,6 +166,7 @@ private void readObject(ObjectInputStream input) throws IOException, ClassNotFou
157166
}
158167
}
159168
this.binaryCopyReadBatchBytes = Optional.ofNullable((Long) input.readObject());
169+
this.maxSourceFragments = Optional.ofNullable((Long) input.readObject());
160170
}
161171

162172
/** Builder for CompactionOptions. */
@@ -171,6 +181,7 @@ public static class Builder {
171181
private Optional<Boolean> deferIndexRemap = Optional.empty();
172182
private Optional<CompactionMode> compactionMode = Optional.empty();
173183
private Optional<Long> binaryCopyReadBatchBytes = Optional.empty();
184+
private Optional<Long> maxSourceFragments = Optional.empty();
174185

175186
private Builder() {}
176187

@@ -224,6 +235,16 @@ public Builder withBinaryCopyReadBatchBytes(long binaryCopyReadBatchBytes) {
224235
return this;
225236
}
226237

238+
/**
239+
* Maximum number of source fragments to compact in a single run. Tasks are included until
240+
* adding the next task would exceed this limit, allowing for incremental compaction. Fragments
241+
* are processed oldest first.
242+
*/
243+
public Builder withMaxSourceFragments(long maxSourceFragments) {
244+
this.maxSourceFragments = Optional.of(maxSourceFragments);
245+
return this;
246+
}
247+
227248
public CompactionOptions build() {
228249
return new CompactionOptions(
229250
targetRowsPerFragment,
@@ -235,7 +256,8 @@ public CompactionOptions build() {
235256
batchSize,
236257
deferIndexRemap,
237258
compactionMode,
238-
binaryCopyReadBatchBytes);
259+
binaryCopyReadBatchBytes,
260+
maxSourceFragments);
239261
}
240262
}
241263
}

java/src/main/java/org/lance/compaction/CompactionTask.java

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -55,7 +55,8 @@ public RewriteResult execute(Dataset dataset) {
5555
compactionOptions.getBatchSize(),
5656
compactionOptions.getDeferIndexRemap(),
5757
compactionOptions.getCompactionMode(),
58-
compactionOptions.getBinaryCopyReadBatchBytes());
58+
compactionOptions.getBinaryCopyReadBatchBytes(),
59+
compactionOptions.getMaxSourceFragments());
5960
}
6061

6162
private native RewriteResult nativeExecute(
@@ -71,7 +72,8 @@ private native RewriteResult nativeExecute(
7172
Optional<Long> batchSize,
7273
Optional<Boolean> deferIndexRemap,
7374
Optional<String> compactionMode,
74-
Optional<Long> binaryCopyReadBatchBytes);
75+
Optional<Long> binaryCopyReadBatchBytes,
76+
Optional<Long> maxSourceFragments);
7577

7678
public CompactionOptions getCompactionOptions() {
7779
return compactionOptions;

python/python/lance/optimize.py

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -81,3 +81,11 @@ class CompactionOptions(TypedDict):
8181
"""
8282
Whether to defer index remapping during compaction (default: False).
8383
"""
84+
max_source_fragments: Optional[int]
85+
"""
86+
Maximum number of source fragments to compact in a single run. Tasks
87+
are included until adding the next task would exceed this limit,
88+
allowing for incremental compaction (e.g., compact 20 fragments at a
89+
time). Fragments are processed oldest first.
90+
(default: None, no limit)
91+
"""

python/src/dataset/optimize.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,9 @@ fn parse_compaction_options(
7070
"binary_copy_read_batch_bytes" => {
7171
opts.binary_copy_read_batch_bytes = value.extract()?;
7272
}
73+
"max_source_fragments" => {
74+
opts.max_source_fragments = value.extract()?;
75+
}
7376
_ => {
7477
return Err(PyValueError::new_err(format!(
7578
"Invalid compaction option: {}",

0 commit comments

Comments
 (0)