From cd8b9bb6e7d61b52c951b88d28b4b577f97a54fc Mon Sep 17 00:00:00 2001 From: zhangyue19921010 Date: Mon, 13 Apr 2026 21:14:42 +0800 Subject: [PATCH 1/6] make sure test passed --- java/lance-jni/Cargo.lock | 2 + rust/lance/src/dataset/optimize.rs | 113 ++++++++++++++++++++++------- 2 files changed, 89 insertions(+), 26 deletions(-) diff --git a/java/lance-jni/Cargo.lock b/java/lance-jni/Cargo.lock index 2e09f58b0e6..649b1fbbeb3 100644 --- a/java/lance-jni/Cargo.lock +++ b/java/lance-jni/Cargo.lock @@ -3409,6 +3409,7 @@ dependencies = [ "arrow-arith", "arrow-array", "arrow-buffer", + "arrow-cast", "arrow-ipc", "arrow-ord", "arrow-row", @@ -3541,6 +3542,7 @@ dependencies = [ "arrow", "arrow-array", "arrow-buffer", + "arrow-cast", "arrow-ord", "arrow-schema", "arrow-select", diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index c70eb93bbcd..e7ed11d95f3 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -704,6 +704,7 @@ impl CompactionPlanner for DefaultCompactionPlanner { .flat_map(|bin| bin.split_for_size(self.options.target_rows_per_fragment)) .map(|bin| TaskData { fragments: bin.fragments, + has_indexed_fragments: Some(!bin.indices.is_empty()), }) .collect(); @@ -923,6 +924,9 @@ async fn prepare_reader( pub struct TaskData { /// The fragments to compact. pub fragments: Vec, + /// Whether any fragment in this task is covered by an index. + #[serde(default)] + pub has_indexed_fragments: Option, } /// A standalone task that can be serialized and sent to another machine for @@ -1076,6 +1080,8 @@ pub struct RewriteResult { /// /// - `None` when configured with stable row IDs because the row ID /// sequences are rechunked directly. + /// - `None` when the rewritten task did not affect any indexed fragments, + /// and so no row address map is needed. /// - `Some` then these addresses are either (1) written to storage for /// deferred index remap post-processing, or (2) used with reserved /// fragment IDs to build old-to-new mappings. @@ -1152,8 +1158,10 @@ async fn rewrite_files( .iter() .map(|f| f.physical_rows.unwrap() as u64) .sum::(); - // If we aren't using stable row ids, then we need to remap indices. - let needs_remapping = !dataset.manifest.uses_stable_row_ids(); + // If we aren't using stable row ids, then we only need to capture row + // addresses for tasks that actually affect at least one index. + let needs_row_addr_capture = + !dataset.manifest.uses_stable_row_ids() && task.has_indexed_fragments.unwrap_or(true); let mut new_fragments: Vec; let task_id = uuid::Uuid::new_v4(); log::info!( @@ -1178,7 +1186,7 @@ async fn rewrite_files( &fragments, options.batch_size, true, - needs_remapping, + needs_row_addr_capture, ) .await?; row_ids_rx = rx_initial; @@ -1229,7 +1237,7 @@ async fn rewrite_files( )); } - if needs_remapping { + if needs_row_addr_capture { let (tx, rx) = std::sync::mpsc::channel(); let mut addrs = RoaringTreemap::new(); for frag in &fragments { @@ -1490,12 +1498,13 @@ pub async fn commit_compaction( let mut completed_tasks = completed_tasks; - // Single reserve_fragment_ids for all address-style tasks - let has_address_style = completed_tasks.iter().any(|t| t.row_addrs.is_some()); - if has_address_style { + // Rewritten fragments need final ids before we can commit them. Address-style + // remapping also depends on these ids to transpose old row addresses into the + // rewritten fragments. + let has_new_fragments = completed_tasks.iter().any(|t| !t.new_fragments.is_empty()); + if has_new_fragments { let frags: Vec<&mut Fragment> = completed_tasks .iter_mut() - .filter(|t| t.row_addrs.is_some()) .flat_map(|t| t.new_fragments.iter_mut()) .collect(); reserve_fragment_ids(dataset, frags.into_iter()).await?; @@ -1526,12 +1535,9 @@ pub async fn commit_compaction( ); row_id_map.extend(transposed); } - } else if options.defer_index_remap { - let changed_row_addrs = task.row_addrs.ok_or_else(|| { - Error::internal( - "defer_index_remap requires row_addrs but none were provided".to_string(), - ) - })?; + } else if options.defer_index_remap + && let Some(changed_row_addrs) = task.row_addrs + { frag_reuse_groups.push(FragReuseGroup { changed_row_addrs, old_frags: task.original_fragments.iter().map(|f| f.into()).collect(), @@ -1565,21 +1571,11 @@ pub async fn commit_compaction( new_index_files: rewritten.files, }) .collect() - } else if !options.defer_index_remap && !has_address_style { - // We need to reserve fragment ids here so that the fragment bitmap - // can be updated for each index. Only needed for stable row IDs - // since address-style IDs were already reserved above. - let new_fragments = rewrite_groups - .iter_mut() - .flat_map(|group| group.new_fragments.iter_mut()) - .collect::>(); - reserve_fragment_ids(dataset, new_fragments.into_iter()).await?; - Vec::new() } else { Vec::new() }; - let frag_reuse_index = if options.defer_index_remap { + let frag_reuse_index = if options.defer_index_remap && !frag_reuse_groups.is_empty() { Some(build_new_frag_reuse_index(dataset, frag_reuse_groups, new_fragment_bitmap).await?) } else { None @@ -2519,6 +2515,69 @@ mod tests { .expect("loading large frag reuse index details must not fail"); } + #[tokio::test] + async fn test_compact_without_indices_skips_row_addrs_capture() { + let test_dir = TempStrDir::default(); + let mut data_gen = + BatchGenerator::new().col(Box::new(IncrementingInt32::new().named("i".to_owned()))); + let mut dataset = Dataset::write( + data_gen.batch(6_000), + &test_dir, + Some(WriteParams { + max_rows_per_file: 1_000, + ..Default::default() + }), + ) + .await + .unwrap(); + + let options = CompactionOptions { + target_rows_per_fragment: 2_000, + defer_index_remap: true, + ..Default::default() + }; + assert_eq!(dataset.get_fragments().len(), 6); + assert_eq!(dataset.count_rows(None).await.unwrap(), 6_000); + let plan = plan_compaction(&dataset, &options).await.unwrap(); + assert_eq!(plan.tasks().len(), 3); + assert!( + plan.tasks() + .iter() + .all(|task| task.has_indexed_fragments == Some(false)) + ); + + let mut rewrite_results = Vec::with_capacity(plan.tasks().len()); + for task in plan.tasks() { + let rewrite_result = rewrite_files(Cow::Borrowed(&dataset), task.clone(), &options) + .await + .unwrap(); + assert!(rewrite_result.row_addrs.is_none()); + rewrite_results.push(rewrite_result); + } + + let metrics = commit_compaction( + &mut dataset, + rewrite_results, + Arc::new(IgnoreRemap::default()), + &options, + ) + .await + .unwrap(); + + dataset.validate().await.unwrap(); + assert_eq!(metrics.fragments_removed, 6); + assert_eq!(metrics.fragments_added, 3); + assert_eq!(dataset.get_fragments().len(), 3); + assert_eq!(dataset.count_rows(None).await.unwrap(), 6_000); + assert!( + dataset + .load_index_by_name(FRAG_REUSE_INDEX_NAME) + .await + .unwrap() + .is_none() + ); + } + #[tokio::test] async fn test_defer_index_remap() { let mut data_gen = BatchGenerator::new() @@ -2642,7 +2701,9 @@ mod tests { // Build expected values by transposing using the immediate results for immediate_result in &immediate_results { - let row_addrs_bytes = immediate_result.row_addrs.as_ref().unwrap(); + let Some(row_addrs_bytes) = immediate_result.row_addrs.as_ref() else { + continue; + }; let row_addrs = RoaringTreemap::deserialize_from(&mut Cursor::new(row_addrs_bytes)).unwrap(); let transposed = transpose_row_addrs( From 7e3d35fb70cda701e567e318c326c6047dc17f84 Mon Sep 17 00:00:00 2001 From: zhangyue19921010 Date: Mon, 13 Apr 2026 21:20:42 +0800 Subject: [PATCH 2/6] fix(compaction): skip collect rowids for unindexed fragment --- rust/lance/src/dataset/optimize.rs | 33 ++++++++++++++++++++++-------- 1 file changed, 25 insertions(+), 8 deletions(-) diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index e7ed11d95f3..d41c42f4635 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -2630,6 +2630,16 @@ mod tests { ) .await .unwrap(); + dataset2 + .create_index( + &["i"], + IndexType::Scalar, + Some("scalar".into()), + &ScalarIndexParams::default(), + false, + ) + .await + .unwrap(); // Verify the initial state - no fragment reuse index should exist let initial_indices = dataset.load_indices().await.unwrap(); @@ -2670,19 +2680,26 @@ mod tests { .await .unwrap(); - // Both should produce row_addrs (address-style row IDs) - assert!(deferred_result.row_addrs.is_some()); - assert!(!deferred_result.row_addrs.as_ref().unwrap().is_empty()); - assert!(!deferred_result.row_addrs.as_ref().unwrap().is_empty()); assert!(!deferred_result.original_fragments.is_empty()); assert!(!deferred_result.new_fragments.is_empty()); - - assert!(immediate_result.row_addrs.is_some()); assert!(!immediate_result.original_fragments.is_empty()); assert!(!immediate_result.new_fragments.is_empty()); + assert_eq!( + task.has_indexed_fragments, task2.has_indexed_fragments, + "deferred and immediate plans should agree on indexed task coverage" + ); - // Both should capture the same row addresses - assert_eq!(deferred_result.row_addrs, immediate_result.row_addrs); + if task.has_indexed_fragments == Some(true) { + assert!(deferred_result.row_addrs.is_some()); + assert!(!deferred_result.row_addrs.as_ref().unwrap().is_empty()); + assert!(immediate_result.row_addrs.is_some()); + assert!(!immediate_result.row_addrs.as_ref().unwrap().is_empty()); + // Both should capture the same row addresses + assert_eq!(deferred_result.row_addrs, immediate_result.row_addrs); + } else { + assert!(deferred_result.row_addrs.is_none()); + assert!(immediate_result.row_addrs.is_none()); + } deferred_results.push(deferred_result); immediate_results.push(immediate_result); From c13db6a41ca958b284b5997cee8e1054b086fd60 Mon Sep 17 00:00:00 2001 From: zhangyue19921010 Date: Mon, 13 Apr 2026 21:26:19 +0800 Subject: [PATCH 3/6] fix(compaction): skip collect rowids for unindexed fragment --- java/lance-jni/Cargo.lock | 2 -- 1 file changed, 2 deletions(-) diff --git a/java/lance-jni/Cargo.lock b/java/lance-jni/Cargo.lock index 649b1fbbeb3..2e09f58b0e6 100644 --- a/java/lance-jni/Cargo.lock +++ b/java/lance-jni/Cargo.lock @@ -3409,7 +3409,6 @@ dependencies = [ "arrow-arith", "arrow-array", "arrow-buffer", - "arrow-cast", "arrow-ipc", "arrow-ord", "arrow-row", @@ -3542,7 +3541,6 @@ dependencies = [ "arrow", "arrow-array", "arrow-buffer", - "arrow-cast", "arrow-ord", "arrow-schema", "arrow-select", From fde798f491bff8c429bdeeceb62a714e4f7ea5d0 Mon Sep 17 00:00:00 2001 From: zhangyue19921010 Date: Mon, 13 Apr 2026 23:33:28 +0800 Subject: [PATCH 4/6] fix ut --- java/lance-jni/Cargo.lock | 2 + java/lance-jni/src/optimize.rs | 20 +++++- .../java/org/lance/compaction/TaskData.java | 13 ++++ rust/lance/src/dataset/optimize.rs | 65 ++++++++++++++++--- 4 files changed, 89 insertions(+), 11 deletions(-) diff --git a/java/lance-jni/Cargo.lock b/java/lance-jni/Cargo.lock index 2e09f58b0e6..649b1fbbeb3 100644 --- a/java/lance-jni/Cargo.lock +++ b/java/lance-jni/Cargo.lock @@ -3409,6 +3409,7 @@ dependencies = [ "arrow-arith", "arrow-array", "arrow-buffer", + "arrow-cast", "arrow-ipc", "arrow-ord", "arrow-row", @@ -3541,6 +3542,7 @@ dependencies = [ "arrow", "arrow-array", "arrow-buffer", + "arrow-cast", "arrow-ord", "arrow-schema", "arrow-select", diff --git a/java/lance-jni/src/optimize.rs b/java/lance-jni/src/optimize.rs index 0ce92baeec8..cdd129f9f82 100644 --- a/java/lance-jni/src/optimize.rs +++ b/java/lance-jni/src/optimize.rs @@ -302,7 +302,7 @@ fn inner_execute_task<'local>( } const TASK_DATA_CLASS: &str = "org/lance/compaction/TaskData"; -const TASK_DATA_CONSTRUCTOR_SIG: &str = "(Ljava/util/List;)V"; +const TASK_DATA_CONSTRUCTOR_SIG: &str = "(Ljava/util/List;Ljava/lang/Boolean;)V"; const COMPACTION_METRICS_CLASS: &str = "org/lance/compaction/CompactionMetrics"; const COMPACTION_METRICS_CONSTRUCTOR_SIG: &str = "(JJJJ)V"; const COMPACTION_PLAN_CLASS: &str = "org/lance/compaction/CompactionPlan"; @@ -317,10 +317,14 @@ const COMPACTION_OPTIONS_CONSTRUCTOR_SIG: &str = "(Ljava/util/Optional;Ljava/uti impl IntoJava for &TaskData { fn into_java<'a>(self, env: &mut JNIEnv<'a>) -> Result> { let fragments = export_vec(env, &self.fragments)?; + let has_indexed_fragments = to_java_boolean_obj(env, self.has_indexed_fragments)?; Ok(env.new_object( TASK_DATA_CLASS, TASK_DATA_CONSTRUCTOR_SIG, - &[JValueGen::Object(&fragments)], + &[ + JValueGen::Object(&fragments), + JValueGen::Object(&has_indexed_fragments), + ], )?) } } @@ -462,8 +466,20 @@ impl FromJObjectWithEnv for JObject<'_> { let task_data = import_vec_from_method(env, self, "getFragments", |env, fragment| { fragment.extract_object(env) })?; + let has_indexed_fragments_obj = env + .call_method(self, "getHasIndexedFragments", "()Ljava/lang/Boolean;", &[])? + .l()?; + let has_indexed_fragments = if has_indexed_fragments_obj.is_null() { + None + } else { + Some( + env.call_method(&has_indexed_fragments_obj, "booleanValue", "()Z", &[])? + .z()?, + ) + }; Ok(TaskData { fragments: task_data, + has_indexed_fragments, }) } } diff --git a/java/src/main/java/org/lance/compaction/TaskData.java b/java/src/main/java/org/lance/compaction/TaskData.java index 4ec5958215c..10900b16f1d 100644 --- a/java/src/main/java/org/lance/compaction/TaskData.java +++ b/java/src/main/java/org/lance/compaction/TaskData.java @@ -15,18 +15,31 @@ import org.lance.FragmentMetadata; +import javax.annotation.Nullable; + import java.io.Serializable; import java.util.List; /** Data of compaction task. */ public class TaskData implements Serializable { private final List fragments; + @Nullable private final Boolean hasIndexedFragments; public TaskData(List fragments) { + this(fragments, null); + } + + public TaskData(List fragments, @Nullable Boolean hasIndexedFragments) { this.fragments = fragments; + this.hasIndexedFragments = hasIndexedFragments; } public List getFragments() { return fragments; } + + @Nullable + public Boolean getHasIndexedFragments() { + return hasIndexedFragments; + } } diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index d41c42f4635..3ab2fc90f82 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -106,7 +106,7 @@ use lance_core::Error; use lance_core::datatypes::BlobHandling; use lance_core::utils::tokio::get_num_compute_intensive_cpus; use lance_core::utils::tracing::{DATASET_COMPACTING_EVENT, TRACE_DATASET_EVENTS}; -use lance_index::frag_reuse::FragReuseGroup; +use lance_index::{frag_reuse::FragReuseGroup, is_system_index}; use lance_table::format::{Fragment, RowIdMeta}; use roaring::{RoaringBitmap, RoaringTreemap}; use serde::{Deserialize, Serialize}; @@ -1026,8 +1026,7 @@ impl CandidateBin { pos_range: self.pos_range.start..(self.pos_range.start + bin_len), candidacy: self.candidacy.drain(0..bin_len).collect(), row_counts: self.row_counts.drain(0..bin_len).collect(), - // By the time we are splitting for size we are done considering indices - indices: Vec::new(), + indices: self.indices.clone(), }); self.pos_range.start += bin_len; } else { @@ -1123,6 +1122,32 @@ async fn reserve_fragment_ids( Ok(()) } +async fn needs_row_addr_capture( + dataset: &Dataset, + task: &TaskData, + options: &CompactionOptions, +) -> Result { + if dataset.manifest.uses_stable_row_ids() { + return Ok(false); + } + + let has_indexed_fragments = task.has_indexed_fragments.unwrap_or(true); + if has_indexed_fragments { + return Ok(true); + } + + if !options.defer_index_remap { + return Ok(false); + } + + let has_user_indices = dataset + .load_indices() + .await? + .iter() + .any(|index| !is_system_index(index)); + Ok(has_user_indices) +} + /// Rewrite the files in a single task. /// /// This assumes that the dataset is the correct read version to be compacted. @@ -1158,10 +1183,7 @@ async fn rewrite_files( .iter() .map(|f| f.physical_rows.unwrap() as u64) .sum::(); - // If we aren't using stable row ids, then we only need to capture row - // addresses for tasks that actually affect at least one index. - let needs_row_addr_capture = - !dataset.manifest.uses_stable_row_ids() && task.has_indexed_fragments.unwrap_or(true); + let needs_row_addr_capture = needs_row_addr_capture(dataset.as_ref(), &task, options).await?; let mut new_fragments: Vec; let task_id = uuid::Uuid::new_v4(); log::info!( @@ -1551,7 +1573,7 @@ pub async fn commit_compaction( rewrite_groups.push(rewrite_group); } - let rewritten_indices = if needs_remapping { + let rewritten_indices = if needs_remapping && !row_id_map.is_empty() { let index_remapper = remap_options.create_remapper(dataset)?; let affected_ids = rewrite_groups .iter() @@ -2266,6 +2288,19 @@ mod tests { } } + async fn create_scalar_index(dataset: &mut Dataset) { + dataset + .create_index( + &["i"], + IndexType::Scalar, + Some("scalar".into()), + &ScalarIndexParams::default(), + false, + ) + .await + .unwrap(); + } + #[rstest::rstest] #[tokio::test] async fn test_compact_distributed( @@ -2490,6 +2525,7 @@ mod tests { // Delete a few rows from each fragment so compaction has something to do. dataset.delete("i % 1000 = 0").await.unwrap(); + create_scalar_index(&mut dataset).await; compact_files( &mut dataset, @@ -2858,6 +2894,7 @@ mod tests { ) .await .unwrap(); + create_scalar_index(&mut dataset).await; let options = CompactionOptions { target_rows_per_fragment: 2_000, @@ -3230,6 +3267,7 @@ mod tests { .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000)) .await .unwrap(); + create_scalar_index(&mut dataset).await; let options = CompactionOptions { target_rows_per_fragment: 2_000, @@ -3270,6 +3308,10 @@ mod tests { .unwrap(); assert_eq!(frag_reuse_details.versions.len(), 1); + remapping::remap_column_index(&mut dataset_clone, &["i"], Some("scalar".into())) + .await + .unwrap(); + // First commit the remaining 2 compaction tasks. let rewrite_result2 = rewrite_files(Cow::Borrowed(&dataset), tasks[1].clone(), &options) .await @@ -3368,6 +3410,7 @@ mod tests { .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000)) .await .unwrap(); + create_scalar_index(&mut dataset).await; let options = CompactionOptions { target_rows_per_fragment: 2_000, @@ -3408,7 +3451,10 @@ mod tests { assert_eq!(frag_reuse_details.versions.len(), 1); // First commit the frag_reuse_index cleanup - // Because there is no index, it should remove the first version. + // The scalar index is remapped first so cleanup can remove the first version. + remapping::remap_column_index(&mut dataset, &["i"], Some("scalar".into())) + .await + .unwrap(); cleanup_frag_reuse_index(&mut dataset).await.unwrap(); // Load and verify the fragment reuse index content @@ -3478,6 +3524,7 @@ mod tests { .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000)) .await .unwrap(); + create_scalar_index(&mut dataset).await; let options = CompactionOptions { target_rows_per_fragment: 2_000, From 354e08e087c6af206a69a664650c1cd699eb833f Mon Sep 17 00:00:00 2001 From: zhangyue19921010 Date: Mon, 13 Apr 2026 23:34:26 +0800 Subject: [PATCH 5/6] fix ut --- java/lance-jni/Cargo.lock | 2 -- 1 file changed, 2 deletions(-) diff --git a/java/lance-jni/Cargo.lock b/java/lance-jni/Cargo.lock index 649b1fbbeb3..2e09f58b0e6 100644 --- a/java/lance-jni/Cargo.lock +++ b/java/lance-jni/Cargo.lock @@ -3409,7 +3409,6 @@ dependencies = [ "arrow-arith", "arrow-array", "arrow-buffer", - "arrow-cast", "arrow-ipc", "arrow-ord", "arrow-row", @@ -3542,7 +3541,6 @@ dependencies = [ "arrow", "arrow-array", "arrow-buffer", - "arrow-cast", "arrow-ord", "arrow-schema", "arrow-select", From 8b9ef563540c54bd027237038efa7af424fcb98d Mon Sep 17 00:00:00 2001 From: zhangyue19921010 Date: Wed, 20 May 2026 19:54:57 +0800 Subject: [PATCH 6/6] code review --- rust/lance/src/dataset/optimize.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index 9756d4f141f..bb97c684551 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -1044,7 +1044,7 @@ impl CandidateBin { async fn load_index_fragmaps(dataset: &Dataset) -> Result> { let indices = dataset.load_indices().await?; let mut index_fragmaps = Vec::with_capacity(indices.len()); - for index in indices.iter() { + for index in indices.iter().filter(|index| !is_system_index(index)) { if let Some(fragment_bitmap) = index.fragment_bitmap.as_ref() { index_fragmaps.push(fragment_bitmap.clone()); } else {