diff --git a/rust/lance-index/src/scalar/zonemap.rs b/rust/lance-index/src/scalar/zonemap.rs index 6205444b678..8c1a195f007 100644 --- a/rust/lance-index/src/scalar/zonemap.rs +++ b/rust/lance-index/src/scalar/zonemap.rs @@ -35,7 +35,7 @@ use arrow_array::{ use arrow_schema::{DataType, Field}; use datafusion::execution::SendableRecordBatchStream; use datafusion_common::ScalarValue; -use lance_select::RowAddrTreeMap; +use lance_select::{RowAddrTreeMap, RowSetOps}; use std::{collections::HashMap, sync::Arc}; use super::{AnyQuery, IndexStore, MetricsCollector, ScalarIndex, SearchResult}; @@ -706,14 +706,26 @@ impl ScalarIndex for ZoneMapIndex { metrics: &dyn MetricsCollector, ) -> Result { let query = query.as_any().downcast_ref::().unwrap(); - if let SargableQuery::IsNull() = query + let result = if let SargableQuery::IsNull() = query && let Some(null_rows) = &self.null_rows { - return Ok(SearchResult::exact(null_rows.clone())); - } + SearchResult::exact(null_rows.clone()) + } else { + search_zones(&self.zones, metrics, |zone| { + self.evaluate_zone_against_query(zone, query) + })? + }; - search_zones(&self.zones, metrics, |zone| { - self.evaluate_zone_against_query(zone, query) + let Some(remapper) = &self.fri else { + return Ok(result); + }; + let selected = remapper.remap_row_addrs_tree_map(result.row_addrs().selected_rows()); + let nulls = remapper.remap_row_addrs_tree_map(result.row_addrs().null_rows()); + + Ok(match result { + SearchResult::Exact(_) => SearchResult::exact(selected).with_nulls(nulls), + SearchResult::AtMost(_) => SearchResult::at_most(selected).with_nulls(nulls), + SearchResult::AtLeast(_) => SearchResult::at_least(selected).with_nulls(nulls), }) } @@ -868,6 +880,104 @@ impl ZoneMapIndex { } } +fn remap_zone( + zone: &ZoneMapStatistics, + remapper: &dyn RowIdRemapper, + remapped_null_rows: Option<&RowAddrTreeMap>, + is_nested: bool, +) -> Result> { + let zone_start = (zone.bound.fragment_id << 32).saturating_add(zone.bound.start); + let mut remapped = (0..zone.bound.length as u64) + .filter_map(|offset| remapper.remap_row_id(zone_start.saturating_add(offset))) + .collect::>(); + remapped.sort_unstable(); + remapped.dedup(); + + let make_zone = |start: u64, end: u64| -> Result { + let length = (end - start + 1) as usize; + let null_count = if let Some(null_rows) = remapped_null_rows { + u32::try_from( + (start..=end) + .filter(|row_id| null_rows.contains(*row_id)) + .count(), + ) + .map_err(|_| { + Error::invalid_input(format!( + "remapped ZoneMap zone has more null rows than can be represented: \ + fragment_id={}, start={}, length={}", + start >> 32, + start & u64::from(u32::MAX), + length + )) + })? + } else if length == zone.bound.length { + zone.null_count + } else if zone.null_count == 0 { + 0 + } else if zone.null_count as usize == zone.bound.length { + u32::try_from(length).map_err(|_| { + Error::invalid_input(format!( + "remapped all-null ZoneMap zone length cannot be represented: \ + fragment_id={}, start={}, length={}", + start >> 32, + start & u64::from(u32::MAX), + length + )) + })? + } else if length > 1 || (!is_nested && !ZoneMapIndex::zone_has_missing_extrema(zone)) { + // Without exact null positions, one null conservatively preserves both + // null and non-null candidates when the run has multiple rows. A scalar + // singleton is also safe when its extrema independently prove that it + // may contain a comparable value. Otherwise null_count == length would + // fabricate an all-null zone and could prune live non-null rows. + 1 + } else { + return Err(Error::not_supported(format!( + "cannot safely remap a mixed-null ZoneMap zone without an exact null bitmap: \ + fragment_id={}, start={}, original_length={}, null_count={}, remapped_length={}, \ + nested={}, missing_extrema={}", + zone.bound.fragment_id, + zone.bound.start, + zone.bound.length, + zone.null_count, + length, + is_nested, + ZoneMapIndex::zone_has_missing_extrema(zone) + ))); + }; + + Ok(ZoneMapStatistics { + min: zone.min.clone(), + max: zone.max.clone(), + null_count, + nan_count: zone.nan_count, + bound: ZoneBound { + fragment_id: start >> 32, + start: start & u64::from(u32::MAX), + length, + }, + }) + }; + let mut zones = Vec::new(); + let mut run_start = None; + let mut previous = 0u64; + for row_id in remapped { + if run_start.is_none() { + run_start = Some(row_id); + } else if row_id != previous.saturating_add(1) || row_id >> 32 != previous >> 32 { + if let Some(start) = run_start { + zones.push(make_zone(start, previous)?); + } + run_start = Some(row_id); + } + previous = row_id; + } + if let Some(start) = run_start { + zones.push(make_zone(start, previous)?); + } + Ok(zones) +} + /// Merge caller-selected ZoneMap segments into one self-contained segment. pub async fn merge_zonemap_indices( source_indices: &[&ZoneMapIndex], @@ -897,19 +1007,31 @@ pub async fn merge_zonemap_indices( data_type, source.data_type ))); } - zones.extend( - source - .zones - .iter() - .filter(|zone| { - u32::try_from(zone.bound.fragment_id) - .is_ok_and(|fragment_id| fragment_filter.contains(fragment_id)) - }) - .cloned(), - ); - match &source.null_rows { - Some(null_rows) => { - let mut filtered = null_rows.clone(); + let remapped_null_rows = source.null_rows.as_ref().map(|null_rows| { + source.fri.as_deref().map_or_else( + || null_rows.clone(), + |remapper| remapper.remap_row_addrs_tree_map(null_rows), + ) + }); + for zone in &source.zones { + let source_zones = source.fri.as_deref().map_or_else( + || Ok(vec![zone.clone()]), + |remapper| { + remap_zone( + zone, + remapper, + remapped_null_rows.as_ref(), + source.data_type.is_nested(), + ) + }, + )?; + zones.extend(source_zones.into_iter().filter(|zone| { + u32::try_from(zone.bound.fragment_id) + .is_ok_and(|fragment_id| fragment_filter.contains(fragment_id)) + })); + } + match remapped_null_rows { + Some(mut filtered) => { filtered.retain_fragments(fragment_filter.iter()); merged_null_rows |= &filtered; } @@ -1804,11 +1926,11 @@ mod tests { use lance_select::RowAddrTreeMap; use crate::scalar::{ - SargableQuery, ScalarIndex, SearchResult, + RowIdRemapper, SargableQuery, ScalarIndex, SearchResult, lance_format::LanceIndexStore, zonemap::{ ZONEMAP_FILENAME, ZONEMAP_SIZE_META_KEY, ZoneMapIndex, ZoneMapIndexBuilderParams, - merge_zonemap_indices, + merge_zonemap_indices, remap_zone, }, }; @@ -1816,7 +1938,7 @@ mod tests { use crate::Index; // Import Index trait to access calculate_included_frags use crate::metrics::NoOpMetricsCollector; use roaring::RoaringBitmap; // Import RoaringBitmap for the test - use std::collections::Bound; + use std::collections::{Bound, HashMap}; // Adds a _rowaddr column emulating each batch as a new fragment fn add_row_addr(stream: SendableRecordBatchStream) -> SendableRecordBatchStream { @@ -1887,6 +2009,198 @@ mod tests { .expect("Failed to load ZoneMapIndex") } + #[derive(Debug)] + struct TestRemapper { + mappings: HashMap, + } + + impl TestRemapper { + fn new(mappings: impl IntoIterator) -> Self { + Self { + mappings: mappings.into_iter().collect(), + } + } + } + + impl RowIdRemapper for TestRemapper { + fn remap_row_id(&self, row_id: u64) -> Option { + self.mappings.get(&row_id).copied() + } + + fn remap_row_addrs_tree_map(&self, _row_addrs: &RowAddrTreeMap) -> RowAddrTreeMap { + unreachable!() + } + + fn remap_row_ids_roaring_tree_map( + &self, + _row_ids: &roaring::RoaringTreemap, + ) -> roaring::RoaringTreemap { + unreachable!() + } + + fn remap_row_ids_record_batch( + &self, + _batch: RecordBatch, + _row_id_idx: usize, + ) -> lance_core::Result { + unreachable!() + } + } + + #[test] + fn test_remap_zone_splits_discontiguous_runs() { + let old_start = (2_u64 << 32) + 10; + let new_start = 3_u64 << 32; + let zone = ZoneMapStatistics { + min: ScalarValue::Int64(Some(10)), + max: ScalarValue::Int64(Some(40)), + null_count: 1, + nan_count: 0, + bound: ZoneBound { + fragment_id: 2, + start: 10, + length: 4, + }, + }; + let mut first_run = zone.clone(); + first_run.bound = ZoneBound { + fragment_id: 3, + start: 1, + length: 2, + }; + let mut second_run = zone.clone(); + second_run.bound = ZoneBound { + fragment_id: 3, + start: 7, + length: 2, + }; + second_run.null_count = 0; + + let remapper = TestRemapper::new([ + (old_start, new_start + 1), + (old_start + 1, new_start + 2), + (old_start + 2, new_start + 7), + (old_start + 3, new_start + 8), + ]); + let mut remapped_null_rows = RowAddrTreeMap::new(); + remapped_null_rows.insert(new_start + 1); + + assert_eq!( + remap_zone(&zone, &remapper, Some(&remapped_null_rows), false).unwrap(), + vec![first_run, second_run] + ); + } + + #[tokio::test] + async fn test_remap_nested_mixed_null_zone_keeps_non_null_candidates() { + let index = train_and_load_fsl( + vec![ + vec![Some(0.0), Some(0.0)], + vec![Some(0.0), Some(0.0)], + vec![Some(3.0), Some(4.0)], + vec![Some(5.0), Some(6.0)], + ], + 2, + ) + .await; + let mut mixed_null_zone = index.zones[0].clone(); + mixed_null_zone.null_count = 2; + + let new_start = 3_u64 << 32; + let remapper = TestRemapper::new([(2, new_start), (3, new_start + 1)]); + let remapped_null_rows = RowAddrTreeMap::new(); + let zones = + remap_zone(&mixed_null_zone, &remapper, Some(&remapped_null_rows), true).unwrap(); + + assert_eq!(zones.len(), 1); + assert_eq!(zones[0].bound.length, 2); + assert_eq!(zones[0].null_count, 0); + assert!( + index + .evaluate_zone_against_query( + &zones[0], + &SargableQuery::Equals(fsl_scalar(vec![Some(3.0), Some(4.0)])), + ) + .unwrap(), + "the two surviving rows are non-null candidates" + ); + } + + #[test] + fn test_remap_nested_mixed_null_singleton_without_bitmap_is_rejected() { + let zone = ZoneMapStatistics { + min: ScalarValue::Null, + max: ScalarValue::Null, + null_count: 2, + nan_count: 0, + bound: ZoneBound { + fragment_id: 0, + start: 0, + length: 4, + }, + }; + let remapper = TestRemapper::new([(2, 0)]); + + let error = remap_zone(&zone, &remapper, None, true).unwrap_err(); + assert!(matches!(error, lance_core::Error::NotSupported { .. })); + assert!( + error + .to_string() + .contains("cannot safely remap a mixed-null ZoneMap zone") + ); + } + + #[tokio::test] + async fn test_merge_legacy_decimal_mixed_null_singleton_is_rejected() { + let mut source = train_and_load::(vec![vec![None, Some(200)]]).await; + let source_mut = Arc::get_mut(&mut source).unwrap(); + let zone = &mut source_mut.zones[0]; + zone.min = ScalarValue::Decimal128(None, 38, 10); + zone.max = ScalarValue::Decimal128(None, 38, 10); + source_mut.null_rows = None; + source_mut.fri = Some(Arc::new(TestRemapper::new([(1, 3_u64 << 32)]))); + + let candidate = ScalarValue::Decimal128(Some(200), 38, 10); + let queries = [ + SargableQuery::Equals(candidate.clone()), + SargableQuery::Range( + Bound::Included(candidate.clone()), + Bound::Included(candidate.clone()), + ), + SargableQuery::IsIn(vec![candidate]), + ]; + for query in queries { + assert!( + source + .evaluate_zone_against_query(&source.zones[0], &query) + .unwrap(), + "legacy Decimal source must retain the non-null candidate for {query:?}" + ); + } + + let dest_tmpdir = TempObjDir::default(); + let dest_store = Arc::new(LanceIndexStore::new( + Arc::new(ObjectStore::local()), + dest_tmpdir.clone(), + Arc::new(LanceCache::no_cache()), + )); + let error = merge_zonemap_indices( + &[source.as_ref()], + dest_store.as_ref(), + &RoaringBitmap::from_iter([3]), + ) + .await + .err() + .expect("ambiguous legacy Decimal singleton must reject consolidation"); + + assert!(matches!(error, lance_core::Error::NotSupported { .. })); + assert!( + error + .to_string() + .contains("cannot safely remap a mixed-null ZoneMap zone") + ); + } + #[tokio::test] async fn test_value_range_spans_fragments() { // Two fragments, multiple zones each; global min/max straddle both. diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index c0f19a5aea1..4c8b4675871 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -5635,6 +5635,117 @@ mod tests { ); } + #[tokio::test] + async fn test_read_zonemap_index_with_defer_index_remap() { + let batch = arrow_array::record_batch!( + ("id", Int32, (0..12).collect::>()), + ( + "value", + Int64, + [ + Some(0), + None, + Some(20), + Some(30), + Some(40), + None, + Some(60), + Some(70), + Some(80), + None, + Some(100), + Some(110) + ] + ) + ) + .unwrap(); + let reader = RecordBatchIterator::new([Ok(batch.clone())], batch.schema()); + let mut dataset = Dataset::write( + reader, + "memory://", + Some(WriteParams { + max_rows_per_file: 4, + max_rows_per_group: 4, + enable_stable_row_ids: false, + ..Default::default() + }), + ) + .await + .unwrap(); + assert_eq!(dataset.get_fragments().len(), 3); + + dataset + .create_index( + &["value"], + IndexType::ZoneMap, + Some("value_idx".into()), + &ScalarIndexParams::for_builtin(BuiltinIndexType::ZoneMap), + false, + ) + .await + .unwrap(); + + let metrics = compact_files( + &mut dataset, + CompactionOptions { + target_rows_per_fragment: 512, + defer_index_remap: true, + ..Default::default() + }, + None, + ) + .await + .unwrap(); + assert_eq!(metrics.fragments_removed, 3); + assert_eq!(metrics.fragments_added, 1); + + async fn scan_ids(dataset: &Dataset, filter: &str, use_scalar_index: bool) -> Vec { + let mut scanner = dataset.scan(); + scanner.filter(filter).unwrap(); + scanner.project(&["id"]).unwrap(); + scanner.use_scalar_index(use_scalar_index); + scanner.try_into_batch().await.unwrap()["id"] + .as_primitive::() + .values() + .to_vec() + } + + for (filter, expected) in [ + ("value IS NULL", vec![1, 5, 9]), + ("value = 20", vec![2]), + ("value > 90", vec![10, 11]), + ] { + assert_eq!(scan_ids(&dataset, filter, false).await, expected); + assert_eq!(scan_ids(&dataset, filter, true).await, expected); + } + + let merged = dataset + .merge_existing_index_segments(dataset.load_indices_by_name("value_idx").await.unwrap()) + .await + .unwrap(); + dataset + .commit_existing_index_segments("value_idx", "value", vec![merged]) + .await + .unwrap(); + + for (filter, expected) in [ + ("value IS NULL", vec![1, 5, 9]), + ("value = 20", vec![2]), + ("value > 90", vec![10, 11]), + ] { + assert_eq!(scan_ids(&dataset, filter, false).await, expected); + assert_eq!(scan_ids(&dataset, filter, true).await, expected); + } + + let mut scanner = dataset.scan(); + scanner.filter("value IS NULL").unwrap(); + let plan = scanner.explain_plan(false).await.unwrap(); + assert!( + plan.contains("ScalarIndexQuery: query=[value IS NULL]@value_idx(ZoneMap)"), + "Expected ZoneMap index query in plan: {plan}" + ); + } + #[tokio::test] async fn test_read_btree_index_with_defer_index_remap() { // Create a dataset with an incremental ID column