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 790ce9a04f5..bb97c684551 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}; @@ -705,6 +705,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(); @@ -924,6 +925,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 @@ -1023,8 +1027,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 { @@ -1041,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 { @@ -1077,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. @@ -1118,6 +1123,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. @@ -1153,8 +1184,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 need to remap indices. - let needs_remapping = !dataset.manifest.uses_stable_row_ids(); + 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!( @@ -1179,7 +1209,7 @@ async fn rewrite_files( &fragments, options.batch_size, true, - needs_remapping, + needs_row_addr_capture, ) .await?; row_ids_rx = rx_initial; @@ -1230,7 +1260,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 { @@ -1510,12 +1540,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?; @@ -1546,12 +1577,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,7 +1593,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() @@ -1585,21 +1613,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 @@ -2297,6 +2315,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( @@ -2521,6 +2552,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, @@ -2546,6 +2578,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() @@ -2598,6 +2693,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(); @@ -2638,19 +2743,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); @@ -2669,7 +2781,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( @@ -2807,6 +2921,7 @@ mod tests { ) .await .unwrap(); + create_scalar_index(&mut dataset).await; let options = CompactionOptions { target_rows_per_fragment: 2_000, @@ -3179,6 +3294,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, @@ -3219,6 +3335,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 @@ -3317,6 +3437,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, @@ -3357,7 +3478,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 @@ -3427,6 +3551,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,