1use lance_core::utils::row_addr_remap::{GroupInput, RowAddrRemap};
84use std::borrow::Cow;
85use std::collections::HashMap;
86use std::io::Cursor;
87use std::ops::{AddAssign, Range};
88use std::sync::Arc;
89
90use super::fragment::FileFragment;
91use super::index::{DatasetIndexRemapperOptions, load_indices_for_remapping};
92use super::rowids::load_row_id_sequences;
93use super::transaction::{
94 Operation, RewriteGroup, RewrittenIndex, Transaction, TransactionBuilder,
95};
96use super::utils::make_rowid_capture_stream;
97use super::{WriteMode, WriteParams, cleanup_data_fragments, write_fragments_internal};
98use crate::Dataset;
99use crate::Result;
100use crate::dataset::utils::CapturedRowIds;
101use crate::index::DatasetIndexExt;
102use crate::io::commit::{commit_transaction, migrate_fragments};
103use arrow::array::AsArray;
104use arrow::datatypes::{UInt8Type, UInt32Type, UInt64Type};
105use arrow_array::Array;
106use arrow_array::RecordBatch;
107use arrow_array::StructArray;
108use arrow_array::builder::{LargeBinaryBuilder, PrimitiveBuilder, StringBuilder};
109use arrow_buffer::NullBuffer;
110use datafusion::physical_plan::SendableRecordBatchStream;
111use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
112use futures::{StreamExt, TryStreamExt};
113use lance_core::Error;
114use lance_core::datatypes::{BlobHandling, BlobKind};
115use lance_core::utils::tokio::get_num_compute_intensive_cpus;
116use lance_core::utils::tracing::{DATASET_COMPACTING_EVENT, TRACE_DATASET_EVENTS};
117use lance_index::frag_reuse::{FRAG_REUSE_INDEX_NAME, FragReuseGroup};
118use lance_index::is_system_index;
119use lance_table::format::{Fragment, RowIdMeta};
120use roaring::{RoaringBitmap, RoaringTreemap};
121use serde::{Deserialize, Serialize};
122use tracing::{info, warn};
123
124mod binary_copy;
125pub mod remapping;
126
127use crate::index::frag_reuse::build_new_frag_reuse_index;
128use crate::io::deletion::read_dataset_deletion_file;
129use binary_copy::rewrite_files_binary_copy;
130pub use remapping::{IgnoreRemap, IndexRemapper, IndexRemapperOptions, RemappedIndex};
131
132#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)]
134pub enum CompactionMode {
135 Reencode,
137 TryBinaryCopy,
139 ForceBinaryCopy,
141}
142
143impl TryFrom<&str> for CompactionMode {
144 type Error = Error;
145
146 fn try_from(value: &str) -> std::result::Result<Self, Self::Error> {
147 match value.to_lowercase().as_str() {
148 "reencode" => Ok(Self::Reencode),
149 "try_binary_copy" => Ok(Self::TryBinaryCopy),
150 "force_binary_copy" => Ok(Self::ForceBinaryCopy),
151 _ => Err(Error::invalid_input(format!(
152 "Invalid compaction mode \"{}\". Valid values: \"reencode\", \"try_binary_copy\", \"force_binary_copy\"",
153 value
154 ))),
155 }
156 }
157}
158
159#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
162pub enum IndexRemapMode {
163 Compact,
168 #[default]
174 Direct,
175}
176
177impl TryFrom<&str> for IndexRemapMode {
178 type Error = Error;
179
180 fn try_from(value: &str) -> std::result::Result<Self, Self::Error> {
181 match value.to_lowercase().as_str() {
182 "compact" => Ok(Self::Compact),
183 "direct" => Ok(Self::Direct),
184 _ => Err(Error::invalid_input(format!(
185 "Invalid index remap mode \"{}\". Valid values: \"compact\", \"direct\"",
186 value
187 ))),
188 }
189 }
190}
191
192#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
194pub struct CompactionOptions {
195 pub target_rows_per_fragment: usize,
201 pub max_rows_per_group: usize,
206 pub max_bytes_per_file: Option<usize>,
213 pub materialize_deletions: bool,
216 pub materialize_deletions_threshold: f32,
221 pub num_threads: Option<usize>,
225 pub batch_size: Option<usize>,
229 pub io_buffer_size: Option<u64>,
236 pub defer_index_remap: bool,
240 #[serde(default)]
243 pub index_remap_mode: IndexRemapMode,
244 pub compaction_mode: Option<CompactionMode>,
249 #[deprecated(note = "Use `compaction_mode` instead")]
251 pub enable_binary_copy: bool,
252 #[deprecated(note = "Use `compaction_mode` instead")]
254 pub enable_binary_copy_force: bool,
255 pub binary_copy_read_batch_bytes: Option<usize>,
259 pub max_source_fragments: Option<usize>,
265 pub max_overlays_per_fragment: Option<usize>,
274 #[serde(skip)]
280 pub transaction_properties: Option<Arc<HashMap<String, String>>>,
281}
282
283#[allow(deprecated)]
284impl Default for CompactionOptions {
285 fn default() -> Self {
286 Self {
287 target_rows_per_fragment: 1024 * 1024,
289 max_rows_per_group: 1024,
290 materialize_deletions: true,
291 materialize_deletions_threshold: 0.1,
292 num_threads: None,
293 max_bytes_per_file: None,
294 batch_size: None,
295 io_buffer_size: None,
296 defer_index_remap: false,
297 index_remap_mode: IndexRemapMode::Direct,
298 compaction_mode: None,
299 enable_binary_copy: false,
300 enable_binary_copy_force: false,
301 binary_copy_read_batch_bytes: Some(16 * 1024 * 1024),
302 max_source_fragments: None,
303 max_overlays_per_fragment: Some(10),
304 transaction_properties: None,
305 }
306 }
307}
308
309pub const COMPACTION_CONFIG_PREFIX: &str = "lance.compaction.";
311
312#[allow(deprecated)]
313impl CompactionOptions {
314 pub fn from_dataset_config(config: &HashMap<String, String>) -> Result<Self> {
332 let mut opts = Self::default();
333 opts.apply_dataset_config(config)?;
334 Ok(opts)
335 }
336
337 pub fn apply_dataset_config(&mut self, config: &HashMap<String, String>) -> Result<()> {
342 for (key, value) in config {
343 let Some(field) = key.strip_prefix(COMPACTION_CONFIG_PREFIX) else {
344 continue;
345 };
346 match field {
347 "target_rows_per_fragment" => {
348 self.target_rows_per_fragment = value.parse().map_err(|_| {
349 Error::invalid_input(format!(
350 "Invalid value for {}: '{}' (expected a non-negative integer)",
351 key, value
352 ))
353 })?;
354 }
355 "max_rows_per_group" => {
356 self.max_rows_per_group = value.parse().map_err(|_| {
357 Error::invalid_input(format!(
358 "Invalid value for {}: '{}' (expected a non-negative integer)",
359 key, value
360 ))
361 })?;
362 }
363 "max_bytes_per_file" => {
364 self.max_bytes_per_file = Some(value.parse().map_err(|_| {
365 Error::invalid_input(format!(
366 "Invalid value for {}: '{}' (expected a non-negative integer)",
367 key, value
368 ))
369 })?);
370 }
371 "materialize_deletions" => {
372 self.materialize_deletions = match value.to_lowercase().as_str() {
373 "true" => true,
374 "false" => false,
375 _ => {
376 return Err(Error::invalid_input(format!(
377 "Invalid value for {}: '{}' (expected 'true' or 'false')",
378 key, value
379 )));
380 }
381 };
382 }
383 "materialize_deletions_threshold" => {
384 self.materialize_deletions_threshold = value.parse().map_err(|_| {
385 Error::invalid_input(format!(
386 "Invalid value for {}: '{}' (expected a float between 0.0 and 1.0)",
387 key, value
388 ))
389 })?;
390 }
391 "defer_index_remap" => {
392 self.defer_index_remap = match value.to_lowercase().as_str() {
393 "true" => true,
394 "false" => false,
395 _ => {
396 return Err(Error::invalid_input(format!(
397 "Invalid value for {}: '{}' (expected 'true' or 'false')",
398 key, value
399 )));
400 }
401 };
402 }
403 "index_remap_mode" => {
404 self.index_remap_mode = IndexRemapMode::try_from(value.as_str())?;
405 }
406 "batch_size" => {
407 self.batch_size = Some(value.parse().map_err(|_| {
408 Error::invalid_input(format!(
409 "Invalid value for {}: '{}' (expected a non-negative integer)",
410 key, value
411 ))
412 })?);
413 }
414 "io_buffer_size" => {
415 self.io_buffer_size = Some(value.parse().map_err(|_| {
416 Error::invalid_input(format!(
417 "Invalid value for {}: '{}' (expected a non-negative integer)",
418 key, value
419 ))
420 })?);
421 }
422 "compaction_mode" => {
423 self.compaction_mode = Some(CompactionMode::try_from(value.as_str())?);
424 }
425 "binary_copy_read_batch_bytes" => {
426 self.binary_copy_read_batch_bytes = Some(value.parse().map_err(|_| {
427 Error::invalid_input(format!(
428 "Invalid value for {}: '{}' (expected a non-negative integer)",
429 key, value
430 ))
431 })?);
432 }
433 "max_source_fragments" => {
434 self.max_source_fragments = Some(value.parse().map_err(|_| {
435 Error::invalid_input(format!(
436 "Invalid value for {}: '{}' (expected a non-negative integer)",
437 key, value
438 ))
439 })?);
440 }
441 "max_overlays_per_fragment" => {
442 self.max_overlays_per_fragment = match value.to_ascii_lowercase().as_str() {
445 "none" => None,
446 _ => Some(value.parse().map_err(|_| {
447 Error::invalid_input(format!(
448 "Invalid value for {}: '{}' (expected a non-negative integer or 'none')",
449 key, value
450 ))
451 })?),
452 };
453 }
454 _ => {
455 warn!("Ignoring unknown compaction config key: {}", key);
456 }
457 }
458 }
459 Ok(())
460 }
461
462 pub fn validate(&mut self) {
463 if self.materialize_deletions && self.materialize_deletions_threshold >= 1.0 {
465 self.materialize_deletions = false;
466 }
467 }
468
469 pub fn compaction_mode(&self) -> CompactionMode {
473 if let Some(mode) = self.compaction_mode {
474 return mode;
475 }
476 match (self.enable_binary_copy, self.enable_binary_copy_force) {
478 (true, true) => CompactionMode::ForceBinaryCopy,
479 (true, false) => CompactionMode::TryBinaryCopy,
480 _ => CompactionMode::Reencode,
481 }
482 }
483
484 pub fn transaction_properties(mut self, properties: HashMap<String, String>) -> Self {
486 self.transaction_properties = Some(Arc::new(properties));
487 self
488 }
489}
490
491async fn can_use_binary_copy(
503 dataset: &Dataset,
504 options: &CompactionOptions,
505 fragments: &[Fragment],
506) -> bool {
507 can_use_binary_copy_impl(dataset, options, fragments)
508 .await
509 .unwrap_or_else(|err| {
510 log::warn!("Binary copy disabled due to error: {}", err);
511 false
512 })
513}
514
515async fn can_use_binary_copy_impl(
516 dataset: &Dataset,
517 options: &CompactionOptions,
518 fragments: &[Fragment],
519) -> Result<bool> {
520 use lance_file::reader::FileReader as LFReader;
521 use lance_file::version::{ConcreteFileVersion, LanceFileVersion};
522 use lance_io::scheduler::{ScanScheduler, SchedulerConfig};
523
524 if matches!(options.compaction_mode(), CompactionMode::Reencode) {
525 log::debug!("Binary copy disabled: compaction mode is Reencode");
526 return Ok(false);
527 }
528
529 let has_blob_columns = dataset
530 .schema()
531 .fields_pre_order()
532 .any(|field| field.is_blob());
533 if has_blob_columns {
534 log::debug!("Binary copy disabled: dataset contains blob columns");
535 return Ok(false);
536 }
537
538 let storage_ok = dataset
539 .manifest
540 .data_storage_format
541 .lance_file_version()
542 .map(|v| !matches!(v.resolve(), LanceFileVersion::Legacy))
543 .unwrap_or(false);
544 if !storage_ok {
545 log::debug!("Binary copy disabled: dataset uses legacy storage format");
546 return Ok(false);
547 }
548
549 if fragments.is_empty() {
550 log::debug!("Binary copy disabled: no fragments to compact");
551 return Ok(false);
552 }
553
554 let storage_file_version = dataset
555 .manifest
556 .data_storage_format
557 .lance_file_version()?
558 .resolve();
559
560 if fragments[0].files.is_empty() {
561 log::debug!(
562 "Binary copy disabled: fragment {} has no data files",
563 fragments[0].id
564 );
565 return Ok(false);
566 }
567 let ref_fields = &fragments[0].files[0].fields;
568 let ref_cols = &fragments[0].files[0].column_indices;
569 let mut is_same_version = true;
570
571 for fragment in fragments {
572 if fragment.deletion_file.is_some() {
573 log::debug!(
574 "Binary copy disabled: fragment {} has a deletion file",
575 fragment.id
576 );
577 return Ok(false);
578 }
579
580 for data_file in &fragment.files {
581 let version_ok = data_file
582 .file_version()
583 .is_ok_and(|v| v == ConcreteFileVersion::from(storage_file_version));
584
585 if !version_ok {
586 is_same_version = false;
587 }
588 if data_file.fields != *ref_fields || data_file.column_indices != *ref_cols {
589 return Ok(false);
590 }
591
592 let object_store = match data_file.base_id {
594 Some(base_id) => dataset.object_store(Some(base_id)).await?,
595 None => dataset.object_store.clone(),
596 };
597 let full_path = dataset
598 .data_file_dir(data_file)?
599 .clone()
600 .join(data_file.path.as_str());
601 let scan_scheduler = ScanScheduler::new(
602 object_store.clone(),
603 SchedulerConfig::max_bandwidth(&object_store),
604 );
605 let file_scheduler = scan_scheduler
606 .open_file_with_priority(&full_path, 0, &data_file.file_size_bytes)
607 .await?;
608 let file_meta = LFReader::read_all_metadata(&file_scheduler).await?;
609 if file_meta.file_buffers.len() > 1 {
615 log::debug!(
616 "Binary copy disabled: data file has extra global buffers (len={})",
617 file_meta.file_buffers.len()
618 );
619 return Ok(false);
620 }
621 }
622 }
623
624 if !is_same_version {
625 log::debug!("Binary copy disabled: data files use different file versions");
626 return Ok(false);
627 }
628
629 Ok(true)
630}
631
632#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
634pub struct CompactionMetrics {
635 pub fragments_removed: usize,
637 pub fragments_added: usize,
639 pub files_removed: usize,
641 pub files_added: usize,
644}
645
646impl AddAssign for CompactionMetrics {
647 fn add_assign(&mut self, rhs: Self) {
648 self.fragments_removed += rhs.fragments_removed;
649 self.fragments_added += rhs.fragments_added;
650 self.files_removed += rhs.files_removed;
651 self.files_added += rhs.files_added;
652 }
653}
654
655#[async_trait::async_trait]
660pub trait CompactionPlanner: Send + Sync {
661 async fn plan(&self, dataset: &Dataset) -> Result<CompactionPlan>;
670}
671
672#[derive(Debug, Clone, Default)]
681pub struct DefaultCompactionPlanner {
682 options: CompactionOptions,
683}
684
685impl DefaultCompactionPlanner {
686 pub fn new(mut options: CompactionOptions) -> Self {
687 options.validate();
688 Self { options }
689 }
690}
691
692#[async_trait::async_trait]
693impl CompactionPlanner for DefaultCompactionPlanner {
694 async fn plan(&self, dataset: &Dataset) -> Result<CompactionPlan> {
695 if self.options.defer_index_remap && dataset.manifest.uses_stable_row_ids() {
696 return Err(Error::invalid_input(
697 "defer_index_remap=true is not supported on datasets with stable row IDs: \
698 stable row IDs do not require index remapping during compaction, so there \
699 is nothing to defer."
700 .to_string(),
701 ));
702 }
703
704 let fragments = dataset.get_fragments();
707
708 debug_assert!(
709 fragments.windows(2).all(|w| w[0].id() < w[1].id()),
710 "fragments in manifest are not sorted"
711 );
712 let mut fragment_metrics = futures::stream::iter(fragments)
713 .map(|fragment| async move {
714 match collect_metrics(&fragment).await {
715 Ok(metrics) => Ok((fragment.metadata, metrics)),
716 Err(e) => Err(e),
717 }
718 })
719 .buffered(dataset.object_store.as_ref().io_parallelism());
720
721 let index_fragmaps = load_index_fragmaps(dataset).await?;
722 let indices_containing_frag = |frag_id: u32| {
723 index_fragmaps
724 .iter()
725 .enumerate()
726 .filter(|(_, bitmap)| bitmap.contains(frag_id))
727 .map(|(pos, _)| pos)
728 .collect::<Vec<_>>()
729 };
730
731 let mut candidate_bins: Vec<CandidateBin> = Vec::new();
732 let mut current_bin: Option<CandidateBin> = None;
733 let mut i = 0;
734
735 while let Some(res) = fragment_metrics.next().await {
736 let (fragment, metrics) = res?;
737
738 let over_overlay_limit = self
739 .options
740 .max_overlays_per_fragment
741 .is_some_and(|max| fragment.overlays.len() > max);
742
743 let candidacy = if over_overlay_limit {
744 Some(CompactionCandidacy::CompactItself)
747 } else if self.options.materialize_deletions
748 && metrics.deletion_percentage() > self.options.materialize_deletions_threshold
749 {
750 Some(CompactionCandidacy::CompactItself)
751 } else if metrics.physical_rows < self.options.target_rows_per_fragment {
752 Some(CompactionCandidacy::CompactWithNeighbors)
755 } else {
756 None
758 };
759
760 let indices = indices_containing_frag(fragment.id as u32);
761
762 match (candidacy, &mut current_bin) {
763 (None, None) => {} (Some(candidacy), None) => {
765 current_bin = Some(CandidateBin {
767 fragments: vec![fragment],
768 pos_range: i..(i + 1),
769 candidacy: vec![candidacy],
770 row_counts: vec![metrics.num_rows()],
771 indices,
772 });
773 }
774 (Some(candidacy), Some(bin)) => {
775 if bin.indices == indices {
778 bin.fragments.push(fragment);
780 bin.pos_range.end += 1;
781 bin.candidacy.push(candidacy);
782 bin.row_counts.push(metrics.num_rows());
783 } else {
784 candidate_bins.push(current_bin.take().unwrap());
786 current_bin = Some(CandidateBin {
787 fragments: vec![fragment],
788 pos_range: i..(i + 1),
789 candidacy: vec![candidacy],
790 row_counts: vec![metrics.num_rows()],
791 indices,
792 });
793 }
794 }
795 (None, Some(_)) => {
796 candidate_bins.push(current_bin.take().unwrap());
798 }
799 }
800
801 i += 1;
802 }
803
804 if let Some(bin) = current_bin {
806 candidate_bins.push(bin);
807 }
808
809 let all_tasks: Vec<TaskData> = candidate_bins
810 .into_iter()
811 .filter(|bin| !bin.is_noop())
812 .flat_map(|bin| bin.split_for_size(self.options.target_rows_per_fragment))
813 .map(|bin| TaskData {
814 fragments: bin.fragments,
815 })
816 .collect();
817
818 let tasks = if let Some(max_frags) = self.options.max_source_fragments {
819 let mut total_frags = 0;
820 all_tasks
821 .into_iter()
822 .take_while(|task| {
823 total_frags += task.fragments.len();
824 total_frags <= max_frags
825 })
826 .collect()
827 } else {
828 all_tasks
829 };
830
831 let mut compaction_plan =
832 CompactionPlan::new(dataset.manifest.version, self.options.clone());
833 compaction_plan.extend_tasks(tasks);
834
835 Ok(compaction_plan)
836 }
837}
838
839pub async fn compact_files(
850 dataset: &mut Dataset,
851 options: CompactionOptions,
852 remap_options: Option<Arc<dyn IndexRemapperOptions>>, ) -> Result<CompactionMetrics> {
854 info!(target: TRACE_DATASET_EVENTS, event=DATASET_COMPACTING_EVENT, uri = &dataset.uri);
855 let planner = DefaultCompactionPlanner::new(options);
856 compact_files_with_planner(dataset, remap_options, &planner).await
857}
858
859pub async fn compact_files_with_planner(
860 dataset: &mut Dataset,
861 remap_options: Option<Arc<dyn IndexRemapperOptions>>, planner: &dyn CompactionPlanner,
863) -> Result<CompactionMetrics> {
864 let compaction_plan: CompactionPlan = planner.plan(dataset).await?;
865
866 if compaction_plan.tasks().is_empty() {
868 return Ok(CompactionMetrics::default());
869 }
870
871 let dataset_ref = &dataset.clone();
872
873 let result_stream = futures::stream::iter(compaction_plan.tasks)
874 .map(|task| rewrite_files(Cow::Borrowed(dataset_ref), task, &compaction_plan.options))
875 .buffer_unordered(
876 compaction_plan
877 .options
878 .num_threads
879 .unwrap_or_else(get_num_compute_intensive_cpus),
880 );
881
882 let completed_tasks: Vec<RewriteResult> = result_stream.try_collect().await?;
883 let remap_options = remap_options.unwrap_or(Arc::new(DatasetIndexRemapperOptions::default()));
884 let metrics = commit_compaction(
885 dataset,
886 completed_tasks,
887 remap_options,
888 &compaction_plan.options,
889 )
890 .await?;
891
892 Ok(metrics)
893}
894
895#[derive(Debug)]
897struct FragmentMetrics {
898 pub physical_rows: usize,
900 pub num_deletions: usize,
902}
903
904impl FragmentMetrics {
905 fn deletion_percentage(&self) -> f32 {
907 if self.physical_rows > 0 {
908 self.num_deletions as f32 / self.physical_rows as f32
909 } else {
910 0.0
911 }
912 }
913
914 fn num_rows(&self) -> usize {
916 self.physical_rows - self.num_deletions
917 }
918}
919
920async fn collect_metrics(fragment: &FileFragment) -> Result<FragmentMetrics> {
921 let physical_rows = fragment.physical_rows();
922 let num_deletions = fragment.count_deletions();
923 let (physical_rows, num_deletions) =
924 futures::future::try_join(physical_rows, num_deletions).await?;
925 Ok(FragmentMetrics {
926 physical_rows,
927 num_deletions,
928 })
929}
930
931#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
935pub struct CompactionPlan {
936 pub tasks: Vec<TaskData>,
937 pub read_version: u64,
938 pub options: CompactionOptions,
939}
940
941impl CompactionPlan {
942 pub fn compaction_tasks(&self) -> impl Iterator<Item = CompactionTask> + '_ {
944 let read_version = self.read_version;
945 let options = self.options.clone();
946 self.tasks.iter().map(move |task| CompactionTask {
947 task: task.clone(),
948 read_version,
949 options: options.clone(),
950 })
951 }
952
953 pub fn num_tasks(&self) -> usize {
955 self.tasks.len()
956 }
957
958 pub fn read_version(&self) -> u64 {
960 self.read_version
961 }
962
963 pub fn options(&self) -> &CompactionOptions {
965 &self.options
966 }
967}
968
969enum RowClass {
975 Null,
976 External,
977 DataBlob,
978}
979
980struct BlobV2Descriptor<'a> {
982 kind_col: &'a arrow::array::UInt8Array,
983 position_col: &'a arrow::array::UInt64Array,
984 size_col: &'a arrow::array::UInt64Array,
985 blob_uri_col: &'a arrow::array::StringArray,
986 blob_id_col: &'a arrow::array::UInt32Array,
987}
988
989impl<'a> BlobV2Descriptor<'a> {
990 fn try_from_struct(struct_arr: &'a StructArray, column_name: &str) -> Result<Self> {
992 let kind_col = struct_arr
993 .column_by_name("kind")
994 .ok_or_else(|| {
995 Error::internal(format!(
996 "Blob v2 descriptor for column '{}' missing `kind` field",
997 column_name
998 ))
999 })?
1000 .as_primitive::<UInt8Type>();
1001 let position_col = struct_arr
1002 .column_by_name("position")
1003 .ok_or_else(|| {
1004 Error::internal(format!(
1005 "Blob v2 descriptor for column '{}' missing `position` field",
1006 column_name
1007 ))
1008 })?
1009 .as_primitive::<UInt64Type>();
1010 let size_col = struct_arr
1011 .column_by_name("size")
1012 .ok_or_else(|| {
1013 Error::internal(format!(
1014 "Blob v2 descriptor for column '{}' missing `size` field",
1015 column_name
1016 ))
1017 })?
1018 .as_primitive::<UInt64Type>();
1019 let blob_uri_col = struct_arr
1020 .column_by_name("blob_uri")
1021 .ok_or_else(|| {
1022 Error::internal(format!(
1023 "Blob v2 descriptor for column '{}' missing `blob_uri` field",
1024 column_name
1025 ))
1026 })?
1027 .as_string::<i32>();
1028 let blob_id_col = struct_arr
1029 .column_by_name("blob_id")
1030 .ok_or_else(|| {
1031 Error::internal(format!(
1032 "Blob v2 descriptor for column '{}' missing `blob_id` field",
1033 column_name
1034 ))
1035 })?
1036 .as_primitive::<UInt32Type>();
1037 Ok(Self {
1038 kind_col,
1039 position_col,
1040 size_col,
1041 blob_uri_col,
1042 blob_id_col,
1043 })
1044 }
1045}
1046
1047struct RowClassification {
1049 row_classes: Vec<RowClass>,
1050 blob_read_addrs: Vec<u64>,
1051}
1052
1053fn classify_rows(
1055 struct_arr: &StructArray,
1056 descriptor: &BlobV2Descriptor<'_>,
1057 row_addrs: &arrow::array::UInt64Array,
1058 column_name: &str,
1059) -> Result<RowClassification> {
1060 let num_rows = struct_arr.len();
1061 let mut row_classes = Vec::with_capacity(num_rows);
1062 let mut blob_read_addrs = Vec::with_capacity(num_rows);
1063
1064 for i in 0..num_rows {
1065 if struct_arr.is_null(i) || descriptor.kind_col.is_null(i) {
1066 row_classes.push(RowClass::Null);
1067 } else {
1068 let kind = BlobKind::try_from(descriptor.kind_col.value(i)).map_err(|e| {
1069 Error::internal(format!(
1070 "Blob v2 column '{}' has invalid kind at row {}: {e}",
1071 column_name, i
1072 ))
1073 })?;
1074 if kind == BlobKind::External {
1075 row_classes.push(RowClass::External);
1076 } else {
1077 row_classes.push(RowClass::DataBlob);
1078 blob_read_addrs.push(row_addrs.value(i));
1079 }
1080 }
1081 }
1082
1083 Ok(RowClassification {
1084 row_classes,
1085 blob_read_addrs,
1086 })
1087}
1088
1089async fn build_user_view_struct(
1094 dataset: &Arc<Dataset>,
1095 descriptor: &BlobV2Descriptor<'_>,
1096 classification: &RowClassification,
1097 column_name: &str,
1098 num_rows: usize,
1099 null_buffer: Option<NullBuffer>,
1100) -> Result<StructArray> {
1101 let blob_files = if classification.blob_read_addrs.is_empty() {
1102 Vec::new()
1103 } else {
1104 super::blob::take_blobs_by_addresses(dataset, &classification.blob_read_addrs, column_name)
1105 .await?
1106 };
1107
1108 let mut data_builder = LargeBinaryBuilder::with_capacity(num_rows, 0);
1109 let mut uri_builder = StringBuilder::with_capacity(num_rows, 0);
1110 let mut out_position_builder = PrimitiveBuilder::<UInt64Type>::with_capacity(num_rows);
1111 let mut out_size_builder = PrimitiveBuilder::<UInt64Type>::with_capacity(num_rows);
1112
1113 let mut blob_file_idx = 0;
1114 #[allow(clippy::needless_range_loop)]
1115 for i in 0..num_rows {
1116 match classification.row_classes[i] {
1117 RowClass::Null => {
1118 data_builder.append_null();
1119 uri_builder.append_null();
1120 out_position_builder.append_null();
1121 out_size_builder.append_null();
1122 }
1123 RowClass::External => {
1124 data_builder.append_null();
1125 let base_id = descriptor.blob_id_col.value(i);
1126 let uri_val = descriptor.blob_uri_col.value(i);
1127 if base_id == 0 {
1128 uri_builder.append_value(uri_val);
1129 } else {
1130 let base = dataset.manifest().base_paths.get(&base_id).ok_or_else(|| {
1131 Error::internal(format!(
1132 "External blob in column '{}' references unknown base_id {}",
1133 column_name, base_id
1134 ))
1135 })?;
1136 let absolute_uri = format!("{}/{}", base.path.trim_end_matches('/'), uri_val);
1137 uri_builder.append_value(&absolute_uri);
1138 }
1139 if descriptor.position_col.is_null(i) {
1140 out_position_builder.append_null();
1141 } else {
1142 out_position_builder.append_value(descriptor.position_col.value(i));
1143 }
1144 if descriptor.size_col.is_null(i) {
1145 out_size_builder.append_null();
1146 } else {
1147 out_size_builder.append_value(descriptor.size_col.value(i));
1148 }
1149 }
1150 RowClass::DataBlob => {
1151 let blob_file = blob_files[blob_file_idx].as_ref().ok_or_else(|| {
1152 Error::internal(format!(
1153 "Non-null blob row {} in column '{}' resolved to null",
1154 i, column_name
1155 ))
1156 })?;
1157 let data = blob_file.read().await?;
1158 blob_file_idx += 1;
1159 data_builder.append_value(data.as_ref());
1160 uri_builder.append_null();
1161 out_position_builder.append_null();
1162 out_size_builder.append_null();
1163 }
1164 }
1165 }
1166
1167 Ok(StructArray::try_new(
1168 lance_core::datatypes::BLOB_V2_USER_FIELDS.clone(),
1169 vec![
1170 Arc::new(data_builder.finish()),
1171 Arc::new(uri_builder.finish()),
1172 Arc::new(out_position_builder.finish()),
1173 Arc::new(out_size_builder.finish()),
1174 ],
1175 null_buffer,
1176 )?)
1177}
1178
1179pub(crate) async fn transform_blob_v2_batch(
1180 dataset: &Arc<Dataset>,
1181 schema: &lance_core::datatypes::Schema,
1182 batch: RecordBatch,
1183 keep_row_addr: bool,
1184) -> Result<RecordBatch> {
1185 let row_addr_idx = batch
1186 .schema()
1187 .column_with_name(lance_core::ROW_ADDR)
1188 .ok_or_else(|| {
1189 Error::internal(format!(
1190 "_rowaddr column missing from batch for blob v2 compaction, columns: {:?}",
1191 batch
1192 .schema()
1193 .fields()
1194 .iter()
1195 .map(|f| f.name())
1196 .collect::<Vec<_>>()
1197 ))
1198 })?
1199 .0;
1200 let row_addrs = batch.column(row_addr_idx).as_primitive::<UInt64Type>();
1201
1202 let mut new_columns: Vec<Arc<dyn Array>> = Vec::new();
1203 let mut new_fields: Vec<Arc<arrow_schema::Field>> = Vec::new();
1204
1205 let batch_schema = batch.schema();
1206 for (col_idx, field) in batch_schema.fields().iter().enumerate() {
1207 if field.name() == lance_core::ROW_ADDR && !keep_row_addr {
1208 continue;
1209 }
1210
1211 let lance_field = schema.field(field.name());
1212 let is_blob_v2 = lance_field.is_some_and(|f| f.is_blob_v2());
1213
1214 if !is_blob_v2 {
1215 new_columns.push(batch.column(col_idx).clone());
1216 new_fields.push(field.clone());
1217 continue;
1218 }
1219
1220 let struct_arr = batch
1221 .column(col_idx)
1222 .as_any()
1223 .downcast_ref::<StructArray>()
1224 .ok_or_else(|| {
1225 Error::internal(format!(
1226 "Blob v2 column '{}' expected StructArray, got {:?}",
1227 field.name(),
1228 batch.column(col_idx).data_type()
1229 ))
1230 })?;
1231
1232 if struct_arr.column_by_name("kind").is_none() {
1236 new_columns.push(batch.column(col_idx).clone());
1237 new_fields.push(field.clone());
1238 continue;
1239 }
1240
1241 let column_name = field.name();
1242 let descriptor = BlobV2Descriptor::try_from_struct(struct_arr, column_name)?;
1243 let classification = classify_rows(struct_arr, &descriptor, row_addrs, column_name)?;
1244 let num_rows = struct_arr.len();
1245
1246 let new_struct = build_user_view_struct(
1247 dataset,
1248 &descriptor,
1249 &classification,
1250 column_name,
1251 num_rows,
1252 struct_arr.nulls().cloned(),
1253 )
1254 .await?;
1255
1256 new_columns.push(Arc::new(new_struct));
1257 let logical_field = arrow_schema::Field::from(lance_field.ok_or_else(|| {
1258 Error::internal(format!(
1259 "Blob v2 column '{}' missing from dataset schema during compaction",
1260 field.name()
1261 ))
1262 })?);
1263 new_fields.push(Arc::new(
1264 arrow_schema::Field::new(
1265 field.name(),
1266 lance_core::datatypes::BLOB_V2_USER_TYPE.clone(),
1267 field.is_nullable(),
1268 )
1269 .with_metadata(logical_field.metadata().clone()),
1270 ));
1271 }
1272
1273 let new_schema = Arc::new(arrow_schema::Schema::new_with_metadata(
1274 new_fields
1275 .iter()
1276 .map(|f| f.as_ref().clone())
1277 .collect::<Vec<_>>(),
1278 batch_schema.metadata().clone(),
1279 ));
1280
1281 Ok(RecordBatch::try_new(new_schema, new_columns)?)
1282}
1283
1284async fn prepare_reader(
1306 dataset: &Dataset,
1307 fragments: &[Fragment],
1308 batch_size: Option<usize>,
1309 io_buffer_size: Option<u64>,
1310 with_frags: bool,
1311 capture_row_ids: bool,
1312) -> Result<(
1313 SendableRecordBatchStream,
1314 Option<std::sync::mpsc::Receiver<CapturedRowIds>>,
1315 bool,
1316)> {
1317 let mut scanner = dataset.scan();
1318 let has_legacy_blob_columns = dataset
1319 .schema()
1320 .fields_pre_order()
1321 .any(|field| field.is_blob() && !field.is_blob_v2());
1322 if has_legacy_blob_columns {
1323 scanner.blob_handling(BlobHandling::AllBinary);
1324 }
1325 let has_blob_v2_columns = dataset
1326 .schema()
1327 .fields_pre_order()
1328 .any(|field| field.is_blob_v2());
1329 if has_blob_v2_columns {
1330 scanner.with_row_address();
1331 }
1332 if let Some(bs) = batch_size {
1333 scanner.batch_size(bs);
1334 }
1335 if let Some(io_buffer_size) = io_buffer_size {
1336 scanner.io_buffer_size(io_buffer_size);
1337 }
1338 if with_frags {
1339 scanner
1340 .with_fragments(fragments.to_vec())
1341 .scan_in_order(true);
1342 }
1343 if capture_row_ids {
1344 scanner.with_row_id();
1345 let data = SendableRecordBatchStream::from(scanner.try_into_stream().await?);
1346 let (data_no_row_ids, rx) =
1347 make_rowid_capture_stream(data, dataset.manifest.uses_stable_row_ids())?;
1348 Ok((data_no_row_ids, Some(rx), has_blob_v2_columns))
1349 } else {
1350 Ok((
1351 SendableRecordBatchStream::from(scanner.try_into_stream().await?),
1352 None,
1353 has_blob_v2_columns,
1354 ))
1355 }
1356}
1357
1358#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1362pub struct TaskData {
1363 pub fragments: Vec<Fragment>,
1365}
1366
1367#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1370pub struct CompactionTask {
1371 pub task: TaskData,
1372 pub read_version: u64,
1373 pub options: CompactionOptions,
1374}
1375
1376impl CompactionTask {
1377 pub async fn execute(&self, dataset: &Dataset) -> Result<RewriteResult> {
1386 let dataset = if dataset.manifest.version == self.read_version {
1387 Cow::Borrowed(dataset)
1388 } else {
1389 Cow::Owned(dataset.checkout_version(self.read_version).await?)
1390 };
1391 rewrite_files(dataset, self.task.clone(), &self.options).await
1392 }
1393}
1394
1395impl CompactionPlan {
1396 fn new(read_version: u64, options: CompactionOptions) -> Self {
1397 Self {
1398 tasks: Vec::new(),
1399 read_version,
1400 options,
1401 }
1402 }
1403
1404 fn extend_tasks(&mut self, tasks: impl IntoIterator<Item = TaskData>) {
1405 self.tasks.extend(tasks);
1406 }
1407
1408 fn tasks(&self) -> &[TaskData] {
1409 &self.tasks
1410 }
1411}
1412
1413#[derive(Debug, Clone)]
1414enum CompactionCandidacy {
1415 CompactWithNeighbors,
1417 CompactItself,
1419}
1420
1421struct CandidateBin {
1423 pub fragments: Vec<Fragment>,
1424 pub pos_range: Range<usize>,
1425 pub candidacy: Vec<CompactionCandidacy>,
1426 pub row_counts: Vec<usize>,
1427 pub indices: Vec<usize>,
1428}
1429
1430impl CandidateBin {
1431 fn is_noop(&self) -> bool {
1433 if self.fragments.is_empty() {
1434 return true;
1435 }
1436 if self.fragments.len() == 1 {
1438 matches!(self.candidacy[0], CompactionCandidacy::CompactWithNeighbors)
1439 } else {
1440 false
1441 }
1442 }
1443
1444 fn split_for_size(self, min_num_rows: usize) -> Vec<Self> {
1446 let total_rows = self.row_counts.iter().sum::<usize>();
1447 let mut remaining_rows = total_rows;
1448 let mut current_rows = 0;
1449 let mut current_len = 0;
1450 let mut split_lengths = Vec::new();
1451
1452 for row_count in &self.row_counts {
1453 current_rows += *row_count;
1454 current_len += 1;
1455 remaining_rows -= *row_count;
1456
1457 if current_rows >= min_num_rows && remaining_rows > 0 && remaining_rows >= min_num_rows
1460 {
1461 split_lengths.push(current_len);
1462 current_rows = 0;
1463 current_len = 0;
1464 }
1465 }
1466
1467 if split_lengths.is_empty() {
1468 return vec![self];
1469 }
1470
1471 let mut bins = Vec::with_capacity(split_lengths.len() + 1);
1472 let mut fragments = self.fragments.into_iter();
1473 let mut candidacy = self.candidacy.into_iter();
1474 let mut row_counts = self.row_counts.into_iter();
1475 let mut pos_start = self.pos_range.start;
1476
1477 for bin_len in split_lengths {
1478 bins.push(Self {
1479 fragments: fragments.by_ref().take(bin_len).collect(),
1480 pos_range: pos_start..(pos_start + bin_len),
1481 candidacy: candidacy.by_ref().take(bin_len).collect(),
1482 row_counts: row_counts.by_ref().take(bin_len).collect(),
1483 indices: Vec::new(),
1485 });
1486 pos_start += bin_len;
1487 }
1488
1489 bins.push(Self {
1490 fragments: fragments.collect(),
1491 pos_range: pos_start..self.pos_range.end,
1492 candidacy: candidacy.collect(),
1493 row_counts: row_counts.collect(),
1494 indices: self.indices,
1495 });
1496
1497 bins
1498 }
1499}
1500
1501async fn load_index_fragmaps(dataset: &Dataset) -> Result<Vec<RoaringBitmap>> {
1502 let indices = dataset.load_indices().await?;
1503 let mut index_fragmaps = Vec::with_capacity(indices.len());
1504 for index in indices.iter().filter(|idx| !is_system_index(idx)) {
1509 if let Some(fragment_bitmap) = index.fragment_bitmap.as_ref() {
1510 index_fragmaps.push(fragment_bitmap.clone());
1511 } else {
1512 let dataset_at_index = dataset.checkout_version(index.dataset_version).await?;
1513 let frags = 0..dataset_at_index
1516 .manifest
1517 .max_fragment_id
1518 .map_or(0, |m| m + 1);
1519 index_fragmaps.push(RoaringBitmap::from_sorted_iter(frags).unwrap());
1520 }
1521 }
1522 Ok(index_fragmaps)
1523}
1524
1525pub async fn plan_compaction(
1526 dataset: &Dataset,
1527 options: &CompactionOptions,
1528) -> Result<CompactionPlan> {
1529 let planner = DefaultCompactionPlanner::new(options.clone());
1530 planner.plan(dataset).await
1531}
1532
1533#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
1537pub struct RewriteResult {
1538 pub metrics: CompactionMetrics,
1539 pub new_fragments: Vec<Fragment>,
1540 pub read_version: u64,
1542 pub original_fragments: Vec<Fragment>,
1544 pub row_addrs: Option<Vec<u8>>,
1553}
1554
1555async fn reserve_fragment_ids(
1556 dataset: &Dataset,
1557 fragments: impl ExactSizeIterator<Item = &mut Fragment>,
1558) -> Result<()> {
1559 let transaction = Transaction::new(
1560 dataset.manifest.version,
1561 Operation::ReserveFragments {
1562 num_fragments: fragments.len() as u32,
1563 },
1564 None,
1565 );
1566
1567 let (manifest, _) = commit_transaction(
1568 dataset,
1569 dataset.object_store.as_ref(),
1570 dataset.commit_handler.as_ref(),
1571 &transaction,
1572 &Default::default(),
1573 &Default::default(),
1574 dataset.manifest_location.naming_scheme,
1575 None,
1576 )
1577 .await?;
1578
1579 let new_max_exclusive = manifest.max_fragment_id.unwrap_or(0) + 1;
1581 let reserved_ids = (new_max_exclusive - fragments.len() as u32)..(new_max_exclusive);
1582
1583 for (fragment, new_id) in fragments.zip(reserved_ids) {
1584 fragment.id = new_id as u64;
1585 }
1586
1587 Ok(())
1588}
1589
1590async fn rewrite_files(
1594 dataset: Cow<'_, Dataset>,
1595 task: TaskData,
1596 options: &CompactionOptions,
1597) -> Result<RewriteResult> {
1598 let mut metrics = CompactionMetrics::default();
1599
1600 if task.fragments.is_empty() {
1601 return Ok(RewriteResult {
1602 metrics,
1603 new_fragments: Vec::new(),
1604 read_version: dataset.manifest.version,
1605 original_fragments: task.fragments,
1606 row_addrs: None,
1607 });
1608 }
1609
1610 let previous_writer_version = &dataset.manifest.writer_version;
1611 let recompute_stats = previous_writer_version.is_none();
1616
1617 let fragments = migrate_fragments(dataset.as_ref(), &task.fragments, recompute_stats).await?;
1621 let num_rows = fragments
1622 .iter()
1623 .map(|f| f.physical_rows.unwrap() as u64)
1624 .sum::<u64>();
1625 let capture_row_addrs = !dataset.manifest.uses_stable_row_ids()
1628 && (options.defer_index_remap
1629 || load_indices_for_remapping(dataset.as_ref())
1630 .await?
1631 .is_some());
1632 let mut new_fragments: Vec<Fragment>;
1633 let task_id = uuid::Uuid::new_v4();
1634 log::info!(
1635 "Compaction task {}: Begin compacting {} rows across {} fragments",
1636 task_id,
1637 num_rows,
1638 fragments.len()
1639 );
1640 let mode = options.compaction_mode();
1641 let can_binary_copy = can_use_binary_copy(dataset.as_ref(), options, &fragments).await;
1642 if !can_binary_copy && matches!(mode, CompactionMode::ForceBinaryCopy) {
1643 return Err(Error::not_supported_source(
1644 format!("compaction task {}: binary copy is not supported", task_id).into(),
1645 ));
1646 }
1647 let mut row_ids_rx: Option<std::sync::mpsc::Receiver<CapturedRowIds>> = None;
1648 let mut reader: Option<SendableRecordBatchStream> = None;
1649
1650 if !can_binary_copy {
1651 let (prepared_reader, rx_initial, has_blob_v2_columns) = prepare_reader(
1652 dataset.as_ref(),
1653 &fragments,
1654 options.batch_size,
1655 options.io_buffer_size,
1656 true,
1657 capture_row_addrs,
1658 )
1659 .await?;
1660 row_ids_rx = rx_initial;
1661
1662 let mut rows_read = 0;
1663 let schema = prepared_reader.schema();
1664 let reader_with_progress = prepared_reader.inspect_ok(move |batch| {
1665 rows_read += batch.num_rows();
1666 log::info!(
1667 "Compaction task {}: Read progress {}/{}",
1668 task_id,
1669 rows_read,
1670 num_rows,
1671 );
1672 });
1673
1674 if has_blob_v2_columns {
1675 let dataset_arc = Arc::new(dataset.as_ref().clone());
1676 let dataset_schema = dataset.schema().clone();
1677 let transformed = reader_with_progress.then(move |batch_result| {
1678 let dataset = dataset_arc.clone();
1679 let schema = dataset_schema.clone();
1680 async move {
1681 let batch = batch_result?;
1682 transform_blob_v2_batch(&dataset, &schema, batch, false)
1683 .await
1684 .map_err(|e| datafusion::error::DataFusionError::External(Box::new(e)))
1685 }
1686 });
1687 let transformed_schema = {
1688 let mut fields: Vec<Arc<arrow_schema::Field>> = Vec::new();
1689 for field in schema.fields().iter() {
1690 if field.name() == lance_core::ROW_ADDR {
1691 continue;
1692 }
1693 let lance_field = dataset.schema().field(field.name());
1694 if let Some(lance_field) = lance_field.filter(|f| f.is_blob_v2()) {
1695 let logical_field = arrow_schema::Field::from(lance_field);
1696 fields.push(Arc::new(
1697 arrow_schema::Field::new(
1698 field.name(),
1699 lance_core::datatypes::BLOB_V2_USER_TYPE.clone(),
1700 field.is_nullable(),
1701 )
1702 .with_metadata(logical_field.metadata().clone()),
1703 ));
1704 } else {
1705 fields.push(field.clone());
1706 }
1707 }
1708 Arc::new(arrow_schema::Schema::new_with_metadata(
1709 fields
1710 .iter()
1711 .map(|f| f.as_ref().clone())
1712 .collect::<Vec<_>>(),
1713 schema.metadata().clone(),
1714 ))
1715 };
1716 reader = Some(Box::pin(RecordBatchStreamAdapter::new(
1717 transformed_schema,
1718 transformed,
1719 )));
1720 } else {
1721 reader = Some(Box::pin(RecordBatchStreamAdapter::new(
1722 schema,
1723 reader_with_progress,
1724 )));
1725 }
1726 }
1727
1728 let mut params = WriteParams {
1729 max_rows_per_file: options.target_rows_per_fragment,
1730 max_rows_per_group: options.max_rows_per_group,
1731 mode: WriteMode::Append,
1732 allow_external_blob_outside_bases: true,
1736 ..Default::default()
1737 };
1738 if let Some(max_bytes_per_file) = options.max_bytes_per_file {
1739 params.max_bytes_per_file = max_bytes_per_file;
1740 }
1741
1742 if dataset.manifest.uses_stable_row_ids() {
1743 params.enable_stable_row_ids = true;
1744 }
1745
1746 if can_binary_copy {
1747 new_fragments = rewrite_files_binary_copy(
1748 dataset.as_ref(),
1749 &fragments,
1750 ¶ms,
1751 options.binary_copy_read_batch_bytes,
1752 )
1753 .await?;
1754
1755 if new_fragments.is_empty() && matches!(mode, CompactionMode::ForceBinaryCopy) {
1756 return Err(Error::not_supported_source(
1757 format!("compaction task {}: binary copy is not supported", task_id).into(),
1758 ));
1759 }
1760
1761 if capture_row_addrs {
1762 let (tx, rx) = std::sync::mpsc::channel();
1763 let mut addrs = RoaringTreemap::new();
1764 for frag in &fragments {
1765 let frag_id = frag.id as u32;
1766 let count = u64::try_from(frag.physical_rows.unwrap_or(0)).map_err(|_| {
1767 Error::internal(format!(
1768 "Fragment {} has too many physical rows to represent as row addresses",
1769 frag.id
1770 ))
1771 })?;
1772 let start = u64::from(lance_core::utils::address::RowAddress::first_row(frag_id));
1773 addrs.insert_range(start..start + count);
1774 }
1775 let captured = CapturedRowIds::AddressStyle(addrs);
1776 let _ = tx.send(captured);
1777 row_ids_rx = Some(rx);
1778 }
1779 } else {
1780 let (frags, _) = write_fragments_internal(
1781 Some(dataset.as_ref()),
1782 dataset.object_store.clone(),
1783 &dataset.base,
1784 dataset.schema().clone(),
1785 reader.expect("reader must be prepared for non-binary-copy path"),
1786 params,
1787 None,
1788 )
1789 .await?;
1790 new_fragments = frags;
1791 }
1792
1793 log::info!("Compaction task {}: file written", task_id);
1794
1795 let row_addrs_result: Result<Option<Vec<u8>>> = async {
1798 if let Some(row_ids_rx) = row_ids_rx {
1799 let captured_ids = row_ids_rx
1800 .try_recv()
1801 .map_err(|err| Error::internal(format!("Failed to receive row ids: {}", err)))?;
1802 let row_addrs = captured_ids.row_addrs(None).into_owned();
1803 let mut serialized = Vec::with_capacity(row_addrs.serialized_size());
1804 row_addrs.serialize_into(&mut serialized)?;
1805 Ok(Some(serialized))
1806 } else {
1807 if dataset.manifest.uses_stable_row_ids() {
1808 log::info!("Compaction task {}: rechunking stable row ids", task_id);
1809 rechunk_stable_row_ids(dataset.as_ref(), &mut new_fragments, &fragments).await?;
1810 recalc_versions_for_rewritten_fragments(
1811 dataset.as_ref(),
1812 &mut new_fragments,
1813 &fragments,
1814 )
1815 .await?;
1816 }
1817 Ok(None)
1818 }
1819 }
1820 .await;
1821
1822 let row_addrs = match row_addrs_result {
1823 Ok(v) => v,
1824 Err(e) => {
1825 cleanup_data_fragments(&dataset.object_store, &dataset.base, None, &new_fragments)
1826 .await;
1827 return Err(e);
1828 }
1829 };
1830
1831 metrics.files_removed = task
1832 .fragments
1833 .iter()
1834 .map(|f| f.files.len() + f.deletion_file.is_some() as usize)
1835 .sum();
1836 metrics.fragments_removed = task.fragments.len();
1837 metrics.fragments_added = new_fragments.len();
1838 metrics.files_added = new_fragments
1839 .iter()
1840 .map(|f| f.files.len() + f.deletion_file.is_some() as usize)
1841 .sum();
1842
1843 log::info!("Compaction task {}: completed", task_id);
1844
1845 Ok(RewriteResult {
1846 metrics,
1847 new_fragments,
1848 read_version: dataset.manifest.version,
1849 original_fragments: fragments,
1850 row_addrs,
1851 })
1852}
1853
1854async fn rechunk_stable_row_ids(
1855 dataset: &Dataset,
1856 new_fragments: &mut [Fragment],
1857 old_fragments: &[Fragment],
1858) -> Result<()> {
1859 let mut old_sequences = load_row_id_sequences(dataset, old_fragments)
1860 .try_collect::<Vec<_>>()
1861 .await?;
1862 old_sequences.sort_by_key(|(frag_id, _)| {
1864 old_fragments
1865 .iter()
1866 .position(|frag| frag.id as u32 == *frag_id)
1867 .expect("Fragment not found")
1868 });
1869
1870 futures::stream::iter(old_sequences.iter_mut().zip(old_fragments.iter()))
1872 .map(Ok)
1873 .try_for_each(|((_, seq), frag)| async move {
1874 if let Some(deletion_file) = &frag.deletion_file {
1875 let deletions = read_dataset_deletion_file(dataset, frag.id, deletion_file).await?;
1876
1877 let mut new_seq = seq.as_ref().clone();
1878 new_seq.mask(deletions.to_sorted_iter())?;
1879 *seq = Arc::new(new_seq);
1880 }
1881 Ok::<(), crate::Error>(())
1882 })
1883 .await?;
1884
1885 debug_assert_eq!(
1886 { old_sequences.iter().map(|(_, seq)| seq.len()).sum::<u64>() },
1887 {
1888 new_fragments
1889 .iter()
1890 .map(|frag| frag.physical_rows.unwrap() as u64)
1891 .sum::<u64>()
1892 },
1893 "{:?}",
1894 old_sequences
1895 );
1896
1897 let new_sequences = lance_table::rowids::rechunk_sequences(
1898 old_sequences
1899 .into_iter()
1900 .map(|(_, seq)| seq.as_ref().clone()),
1901 new_fragments
1902 .iter()
1903 .map(|frag| frag.physical_rows.unwrap() as u64),
1904 false,
1905 )?;
1906
1907 for (fragment, sequence) in new_fragments.iter_mut().zip(new_sequences) {
1908 let serialized = lance_table::rowids::write_row_ids(&sequence);
1910 fragment.row_id_meta = Some(RowIdMeta::Inline(serialized));
1911 }
1912
1913 Ok(())
1914}
1915
1916async fn recalc_versions_for_rewritten_fragments(
1918 dataset: &Dataset,
1919 new_fragments: &mut [Fragment],
1920 old_fragments: &[Fragment],
1921) -> Result<()> {
1922 let mut old_last_updated_sequences: Vec<lance_table::format::RowDatasetVersionSequence> =
1924 Vec::with_capacity(old_fragments.len());
1925 let mut old_created_at_sequences: Vec<lance_table::format::RowDatasetVersionSequence> =
1927 Vec::with_capacity(old_fragments.len());
1928
1929 for frag in old_fragments.iter() {
1930 let row_count = if let Some(row_id_meta) = &frag.row_id_meta {
1931 match row_id_meta {
1932 RowIdMeta::Inline(data) => lance_table::rowids::read_row_ids(data)?.len(),
1933 RowIdMeta::External(_file) => frag.physical_rows.unwrap_or(0) as u64,
1934 }
1935 } else {
1936 frag.physical_rows.unwrap_or(0) as u64
1937 };
1938
1939 let mut created_at_seq = if let Some(version_meta) = &frag.created_at_version_meta {
1941 version_meta.load_sequence().map_err(|e| {
1942 Error::internal(format!("Failed to load created_at version sequence: {}", e))
1943 })?
1944 } else {
1945 lance_table::format::RowDatasetVersionSequence::from_uniform_row_count(row_count, 1)
1947 };
1948
1949 let mut last_updated_seq = if let Some(version_meta) = &frag.last_updated_at_version_meta {
1951 version_meta.load_sequence().map_err(|e| {
1952 Error::internal(format!(
1953 "Failed to load last_updated_at version sequence: {}",
1954 e
1955 ))
1956 })?
1957 } else {
1958 created_at_seq.clone()
1959 };
1960
1961 if let Some(deletion_file) = &frag.deletion_file {
1963 let deletions = read_dataset_deletion_file(dataset, frag.id, deletion_file).await?;
1964 last_updated_seq.mask(deletions.to_sorted_iter())?;
1965 created_at_seq.mask(deletions.to_sorted_iter())?;
1966 }
1967
1968 old_last_updated_sequences.push(last_updated_seq);
1969 old_created_at_sequences.push(created_at_seq);
1970 }
1971
1972 let old_total: u64 = old_last_updated_sequences.iter().map(|s| s.len()).sum();
1974 let new_total: u64 = new_fragments
1975 .iter()
1976 .map(|f| f.physical_rows.unwrap_or(0) as u64)
1977 .sum();
1978 debug_assert_eq!(old_total, new_total);
1979
1980 let chunk_sizes: Vec<u64> = new_fragments
1982 .iter()
1983 .map(|f| f.physical_rows.unwrap_or(0) as u64)
1984 .collect();
1985
1986 let new_last_updated_sequences = lance_table::rowids::version::rechunk_version_sequences(
1987 old_last_updated_sequences,
1988 chunk_sizes.clone(),
1989 false,
1990 )?;
1991
1992 let new_created_at_sequences = lance_table::rowids::version::rechunk_version_sequences(
1993 old_created_at_sequences,
1994 chunk_sizes,
1995 false,
1996 )?;
1997
1998 for ((fragment, last_updated_seq), created_at_seq) in new_fragments
2000 .iter_mut()
2001 .zip(new_last_updated_sequences)
2002 .zip(new_created_at_sequences)
2003 {
2004 fragment.last_updated_at_version_meta = Some(
2005 lance_table::format::RowDatasetVersionMeta::from_sequence(&last_updated_seq).unwrap(),
2006 );
2007 fragment.created_at_version_meta = Some(
2008 lance_table::format::RowDatasetVersionMeta::from_sequence(&created_at_seq).unwrap(),
2009 );
2010 }
2011
2012 Ok(())
2013}
2014
2015pub async fn commit_compaction(
2022 dataset: &mut Dataset,
2023 completed_tasks: Vec<RewriteResult>,
2024 remap_options: Arc<dyn IndexRemapperOptions>,
2025 options: &CompactionOptions,
2026) -> Result<CompactionMetrics> {
2027 if completed_tasks.is_empty() {
2028 return Ok(CompactionMetrics::default());
2029 }
2030
2031 let has_address_style = completed_tasks.iter().any(|t| t.row_addrs.is_some());
2032 let needs_remapping =
2034 !dataset.manifest.uses_stable_row_ids() && !options.defer_index_remap && has_address_style;
2035
2036 let index_remapper = if needs_remapping {
2038 remap_options.create_remapper(dataset).await?
2039 } else {
2040 None
2041 };
2042
2043 let tasks_read_version = completed_tasks
2057 .iter()
2058 .map(|t| t.read_version)
2059 .min()
2060 .unwrap_or(dataset.manifest.version);
2061
2062 let mut completed_tasks = completed_tasks;
2063
2064 if has_address_style {
2066 let frags: Vec<&mut Fragment> = completed_tasks
2067 .iter_mut()
2068 .filter(|t| t.row_addrs.is_some())
2069 .flat_map(|t| t.new_fragments.iter_mut())
2070 .collect();
2071 reserve_fragment_ids(dataset, frags.into_iter()).await?;
2072 }
2073
2074 let mut rewrite_groups = Vec::with_capacity(completed_tasks.len());
2075 let mut metrics = CompactionMetrics::default();
2076
2077 let mut remap_group_inputs: Vec<GroupInput> = Vec::new();
2078 let mut direct_row_id_map: HashMap<u64, Option<u64>> = HashMap::default();
2079 let mut frag_reuse_groups: Vec<FragReuseGroup> = Vec::new();
2080 let mut new_fragment_bitmap: RoaringBitmap = RoaringBitmap::new();
2081
2082 let indexed_frags: RoaringBitmap = if options.defer_index_remap {
2090 let mut covered = RoaringBitmap::new();
2091 for bm in load_index_fragmaps(dataset).await? {
2092 covered |= bm;
2093 }
2094 if let Some(bm) = dataset
2095 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
2096 .await?
2097 .and_then(|fri| fri.fragment_bitmap)
2098 {
2099 covered |= bm;
2100 }
2101 covered
2102 } else {
2103 RoaringBitmap::new()
2104 };
2105 let mut any_group_indexed = false;
2106
2107 for task in completed_tasks {
2108 metrics += task.metrics;
2109 let rewrite_group = RewriteGroup {
2110 old_fragments: task.original_fragments.clone(),
2111 new_fragments: task.new_fragments.clone(),
2112 };
2113
2114 if index_remapper.is_some() {
2115 if let Some(row_addrs_bytes) = task.row_addrs {
2116 let row_addrs =
2117 RoaringTreemap::deserialize_from(&mut Cursor::new(&row_addrs_bytes))?;
2118 match options.index_remap_mode {
2119 IndexRemapMode::Direct => {
2120 let transposed = remapping::transpose_row_addrs(
2121 row_addrs,
2122 &task.original_fragments,
2123 &task.new_fragments,
2124 );
2125 direct_row_id_map.extend(transposed);
2126 }
2127 IndexRemapMode::Compact => {
2128 let new_frags = task
2129 .new_fragments
2130 .iter()
2131 .map(|f| {
2132 let physical_rows = f.physical_rows.ok_or_else(|| {
2133 Error::invalid_input(format!(
2134 "compacted fragment {} is missing physical_rows",
2135 f.id
2136 ))
2137 })?;
2138 Ok((f.id as u32, physical_rows as u32))
2139 })
2140 .collect::<Result<Vec<_>>>()?;
2141
2142 remap_group_inputs.push(GroupInput {
2143 rewritten_old_row_addrs: row_addrs,
2144 old_frag_ids: task
2145 .original_fragments
2146 .iter()
2147 .map(|f| f.id as u32)
2148 .collect(),
2149 new_frags,
2150 });
2151 }
2152 }
2153 }
2154 } else if options.defer_index_remap {
2155 if task
2157 .original_fragments
2158 .iter()
2159 .any(|f| indexed_frags.contains(f.id as u32))
2160 {
2161 any_group_indexed = true;
2162 }
2163 let changed_row_addrs = task.row_addrs.ok_or_else(|| {
2164 Error::internal(
2165 "defer_index_remap requires row_addrs but none were provided".to_string(),
2166 )
2167 })?;
2168 frag_reuse_groups.push(FragReuseGroup {
2169 changed_row_addrs,
2170 old_frags: task.original_fragments.iter().map(|f| f.into()).collect(),
2171 new_frags: task.new_fragments.iter().map(|f| f.into()).collect(),
2172 });
2173
2174 task.new_fragments.iter().for_each(|frag| {
2175 new_fragment_bitmap.insert(frag.id as u32);
2176 });
2177 }
2178 rewrite_groups.push(rewrite_group);
2179 }
2180
2181 let rewritten_indices = if let Some(index_remapper) = index_remapper {
2182 let affected_ids = rewrite_groups
2183 .iter()
2184 .flat_map(|group| group.old_fragments.iter().map(|frag| frag.id))
2185 .collect::<Vec<_>>();
2186
2187 let remap = match options.index_remap_mode {
2188 IndexRemapMode::Direct => RowAddrRemap::direct(direct_row_id_map),
2189 IndexRemapMode::Compact => RowAddrRemap::compact(remap_group_inputs)?,
2190 };
2191 let remapped_indices = index_remapper.remap_indices(remap, &affected_ids).await?;
2192 remapped_indices
2193 .into_iter()
2194 .map(|rewritten| RewrittenIndex {
2195 old_id: rewritten.old_id,
2196 new_id: rewritten.new_id,
2197 new_index_details: rewritten.index_details,
2198 new_index_version: rewritten.index_version,
2199 new_index_files: rewritten.files,
2200 })
2201 .collect()
2202 } else if !options.defer_index_remap && !has_address_style {
2203 let new_fragments = rewrite_groups
2207 .iter_mut()
2208 .flat_map(|group| group.new_fragments.iter_mut())
2209 .collect::<Vec<_>>();
2210 reserve_fragment_ids(dataset, new_fragments.into_iter()).await?;
2211 Vec::new()
2212 } else {
2213 Vec::new()
2214 };
2215
2216 let frag_reuse_index = if options.defer_index_remap && any_group_indexed {
2218 Some(build_new_frag_reuse_index(dataset, frag_reuse_groups, new_fragment_bitmap).await?)
2219 } else {
2220 if options.defer_index_remap {
2221 log::debug!(
2222 "skipping fragment-reuse index: no rewritten fragments were covered by an index"
2223 );
2224 }
2225 None
2226 };
2227
2228 let all_new_fragments: Vec<Fragment> = rewrite_groups
2231 .iter()
2232 .flat_map(|g| g.new_fragments.iter().cloned())
2233 .collect();
2234
2235 let transaction = TransactionBuilder::new(
2236 tasks_read_version,
2244 Operation::Rewrite {
2245 groups: rewrite_groups,
2246 rewritten_indices,
2247 frag_reuse_index,
2248 },
2249 )
2250 .transaction_properties(options.transaction_properties.clone())
2251 .build();
2252
2253 if let Err(e) = dataset
2254 .apply_commit(transaction, &Default::default(), &Default::default())
2255 .await
2256 {
2257 cleanup_data_fragments(
2258 &dataset.object_store,
2259 &dataset.base,
2260 None,
2261 &all_new_fragments,
2262 )
2263 .await;
2264 return Err(e);
2265 }
2266
2267 Ok(metrics)
2268}
2269
2270#[cfg(test)]
2271mod tests {
2272
2273 mod binary_copy;
2274 use self::remapping::RemappedIndex;
2275 use super::*;
2276 use crate::dataset::WriteDestination;
2277 use crate::dataset::index::frag_reuse::cleanup_frag_reuse_index;
2278 use crate::dataset::optimize::remapping::{transpose_row_addrs, transpose_row_ids_from_digest};
2279 use crate::index::frag_reuse::{load_frag_reuse_index_details, open_frag_reuse_index};
2280 use crate::index::vector::{StageParams, VectorIndexParams};
2281 use crate::utils::test::{DatagenExt, FragmentCount, FragmentRowCount};
2282 use arrow_array::types::{Float32Type, Float64Type, Int32Type, Int64Type};
2283 use arrow_array::{
2284 ArrayRef, Float32Array, Int32Array, Int64Array, LargeBinaryArray, LargeStringArray,
2285 PrimitiveArray, RecordBatch, RecordBatchIterator,
2286 };
2287 use arrow_schema::{DataType, Field, Schema};
2288 use arrow_select::concat::concat_batches;
2289 use async_trait::async_trait;
2290 use lance_arrow::BLOB_META_KEY;
2291 use lance_core::Error;
2292 use lance_core::ROW_ID;
2293 use lance_core::utils::address::RowAddress;
2294 use lance_core::utils::tempfile::TempStrDir;
2295 use lance_datagen::Dimension;
2296 use lance_file::version::LanceFileVersion;
2297 use lance_index::frag_reuse::FRAG_REUSE_INDEX_NAME;
2298 use lance_index::frag_reuse::FragReuseIndexHandle;
2299 use lance_index::scalar::{
2300 BuiltinIndexType, FullTextSearchQuery, InvertedIndexParams, ScalarIndexParams,
2301 };
2302 use lance_index::vector::ivf::IvfBuildParams;
2303 use lance_index::vector::pq::PQBuildParams;
2304 use lance_index::{Index, IndexType};
2305 use lance_linalg::distance::{DistanceType, MetricType};
2306 use lance_table::io::manifest::read_manifest_indexes;
2307 use lance_testing::datagen::{BatchGenerator, IncrementingInt32, RandomVector};
2308 use rstest::rstest;
2309 use std::collections::HashSet;
2310 use std::io::Cursor;
2311 use std::sync::Arc;
2312 use uuid::Uuid;
2313
2314 #[test]
2315 fn test_candidate_bin() {
2316 let empty_bin = CandidateBin {
2317 fragments: vec![],
2318 pos_range: 0..0,
2319 candidacy: vec![],
2320 row_counts: vec![],
2321 indices: vec![],
2322 };
2323 assert!(empty_bin.is_noop());
2324
2325 let fragment = Fragment {
2326 id: 0,
2327 files: vec![],
2328 overlays: vec![],
2329 deletion_file: None,
2330 row_id_meta: None,
2331 physical_rows: Some(0),
2332 last_updated_at_version_meta: None,
2333 created_at_version_meta: None,
2334 };
2335 let single_bin = CandidateBin {
2336 fragments: vec![fragment.clone()],
2337 pos_range: 0..1,
2338 candidacy: vec![CompactionCandidacy::CompactWithNeighbors],
2339 row_counts: vec![100],
2340 indices: vec![],
2341 };
2342 assert!(single_bin.is_noop());
2343
2344 let single_bin = CandidateBin {
2345 fragments: vec![fragment.clone()],
2346 pos_range: 0..1,
2347 candidacy: vec![CompactionCandidacy::CompactItself],
2348 row_counts: vec![100],
2349 indices: vec![],
2350 };
2351 assert!(!single_bin.is_noop());
2353
2354 let big_bin = CandidateBin {
2355 fragments: std::iter::repeat_n(fragment, 8).collect(),
2356 pos_range: 0..8,
2357 candidacy: std::iter::repeat_n(CompactionCandidacy::CompactItself, 8).collect(),
2358 row_counts: vec![100, 400, 200, 200, 400, 300, 300, 100],
2359 indices: vec![],
2360 };
2363 assert!(!big_bin.is_noop());
2364 let split = big_bin.split_for_size(500);
2365 assert_eq!(split.len(), 3);
2366 assert_eq!(split[0].pos_range, 0..2);
2367 assert_eq!(split[1].pos_range, 2..5);
2368 assert_eq!(split[2].pos_range, 5..8);
2369
2370 let zero_min_split_bin = CandidateBin {
2371 fragments: std::iter::repeat_n(
2372 Fragment {
2373 id: 0,
2374 files: vec![],
2375 overlays: vec![],
2376 deletion_file: None,
2377 row_id_meta: None,
2378 physical_rows: Some(0),
2379 last_updated_at_version_meta: None,
2380 created_at_version_meta: None,
2381 },
2382 3,
2383 )
2384 .collect(),
2385 pos_range: 0..3,
2386 candidacy: std::iter::repeat_n(CompactionCandidacy::CompactItself, 3).collect(),
2387 row_counts: vec![100, 200, 300],
2388 indices: vec![],
2389 };
2390 let split = zero_min_split_bin.split_for_size(0);
2391 assert_eq!(split.len(), 3);
2392 assert!(split.iter().all(|bin| !bin.fragments.is_empty()));
2393 assert_eq!(split[0].pos_range, 0..1);
2394 assert_eq!(split[1].pos_range, 1..2);
2395 assert_eq!(split[2].pos_range, 2..3);
2396 }
2397
2398 fn sample_data() -> RecordBatch {
2399 let schema = Schema::new(vec![Field::new("a", DataType::Int64, false)]);
2400
2401 RecordBatch::try_new(
2402 Arc::new(schema),
2403 vec![Arc::new(Int64Array::from_iter_values(0..10_000))],
2404 )
2405 .unwrap()
2406 }
2407
2408 async fn create_scalar_index(dataset: &mut Dataset, col: &str, replace: bool) {
2410 dataset
2411 .create_index(
2412 &[col],
2413 IndexType::Scalar,
2414 Some("scalar".into()),
2415 &ScalarIndexParams::default(),
2416 replace,
2417 )
2418 .await
2419 .unwrap();
2420 }
2421
2422 #[derive(Debug, Default, Clone, PartialEq)]
2423 struct MockIndexRemapperExpectation {
2424 expected: HashMap<u64, Option<u64>>,
2425 answer: Vec<RemappedIndex>,
2426 }
2427
2428 #[derive(Debug, Default, Clone, PartialEq)]
2429 struct MockIndexRemapper {
2430 expectations: Vec<MockIndexRemapperExpectation>,
2431 }
2432
2433 impl MockIndexRemapper {
2434 fn stringify_map(map: &HashMap<u64, Option<u64>>) -> String {
2435 let mut sorted_keys = map.keys().collect::<Vec<_>>();
2436 sorted_keys.sort();
2437 let mut first_keys = sorted_keys
2438 .into_iter()
2439 .take(10)
2440 .map(|key| {
2441 format!(
2442 "{}:{:?}",
2443 RowAddress::from(*key),
2444 map[key].map(RowAddress::from)
2445 )
2446 })
2447 .collect::<Vec<_>>()
2448 .join(",");
2449 if map.len() > 10 {
2450 first_keys.push_str(", ...");
2451 }
2452 let mut result_str = format!("(len={})", map.len());
2453 result_str.push_str(&first_keys);
2454 result_str
2455 }
2456
2457 fn in_any_order(expectations: &[Self]) -> Self {
2458 let expectations = expectations
2459 .iter()
2460 .flat_map(|item| item.expectations.clone())
2461 .collect::<Vec<_>>();
2462 Self { expectations }
2463 }
2464 }
2465
2466 #[async_trait]
2467 impl IndexRemapper for MockIndexRemapper {
2468 async fn remap_indices(
2469 &self,
2470 index_map: RowAddrRemap,
2471 _: &[u64],
2472 ) -> Result<Vec<RemappedIndex>> {
2473 for expectation in &self.expectations {
2474 let matches = match &index_map {
2475 RowAddrRemap::Direct(map) => map == &expectation.expected,
2476 RowAddrRemap::Compact(_) => {
2477 let expected_frags: RoaringBitmap = expectation
2478 .expected
2479 .keys()
2480 .map(|addr| (addr >> 32) as u32)
2481 .collect();
2482 index_map.affected_fragments() == expected_frags
2483 && expectation
2484 .expected
2485 .iter()
2486 .all(|(k, v)| index_map.get(*k) == Some(*v))
2487 }
2488 };
2489 if matches {
2490 return Ok(expectation.answer.clone());
2491 }
2492 }
2493 panic!(
2494 "Unexpected index map; expected one of:\n {}",
2495 self.expectations
2496 .iter()
2497 .map(|expectation| Self::stringify_map(&expectation.expected))
2498 .collect::<Vec<_>>()
2499 .join("\n ")
2500 );
2501 }
2502 }
2503
2504 #[async_trait]
2505 impl IndexRemapperOptions for MockIndexRemapper {
2506 async fn create_remapper(&self, _: &Dataset) -> Result<Option<Box<dyn IndexRemapper>>> {
2507 Ok(Some(Box::new(self.clone())))
2508 }
2509 }
2510
2511 #[rstest]
2512 #[tokio::test]
2513 async fn test_compact_empty(
2514 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
2515 data_storage_version: LanceFileVersion,
2516 ) {
2517 let test_dir = TempStrDir::default();
2518 let test_uri = &test_dir;
2519
2520 let schema = Schema::new(vec![Field::new("a", DataType::Int64, false)]);
2522
2523 let reader = RecordBatchIterator::new(vec![].into_iter().map(Ok), Arc::new(schema));
2524 let mut dataset = Dataset::write(
2525 reader,
2526 test_uri,
2527 Some(WriteParams {
2528 data_storage_version: Some(data_storage_version),
2529 ..Default::default()
2530 }),
2531 )
2532 .await
2533 .unwrap();
2534
2535 let plan = plan_compaction(&dataset, &CompactionOptions::default())
2536 .await
2537 .unwrap();
2538 assert_eq!(plan.tasks().len(), 0);
2539
2540 let metrics = compact_files(&mut dataset, CompactionOptions::default(), None)
2541 .await
2542 .unwrap();
2543
2544 assert_eq!(metrics, CompactionMetrics::default());
2545 assert_eq!(dataset.manifest.version, 1);
2546 }
2547
2548 #[rstest]
2549 #[tokio::test]
2550 async fn test_compact_all_good(
2551 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
2552 data_storage_version: LanceFileVersion,
2553 ) {
2554 let test_dir = TempStrDir::default();
2556 let test_uri = &test_dir;
2557
2558 let data = sample_data();
2559 let reader = RecordBatchIterator::new(vec![Ok(data.clone())], data.schema());
2560 let write_params = WriteParams {
2562 max_rows_per_file: 10_000,
2563 data_storage_version: Some(data_storage_version),
2564 ..Default::default()
2565 };
2566 let dataset = Dataset::write(reader, test_uri, Some(write_params))
2567 .await
2568 .unwrap();
2569
2570 let plan = plan_compaction(&dataset, &CompactionOptions::default())
2572 .await
2573 .unwrap();
2574 assert_eq!(plan.tasks().len(), 0);
2575
2576 let reader = RecordBatchIterator::new(vec![Ok(data.clone())], data.schema());
2578 let write_params = WriteParams {
2579 max_rows_per_file: 3_000,
2580 max_rows_per_group: 1_000,
2581 data_storage_version: Some(data_storage_version),
2582 mode: WriteMode::Overwrite,
2583 ..Default::default()
2584 };
2585 let dataset = Dataset::write(reader, test_uri, Some(write_params))
2586 .await
2587 .unwrap();
2588
2589 let options = CompactionOptions {
2590 target_rows_per_fragment: 3_000,
2591 ..Default::default()
2592 };
2593 let plan = plan_compaction(&dataset, &options).await.unwrap();
2594 assert_eq!(plan.tasks().len(), 0);
2595 }
2596
2597 #[tokio::test]
2598 async fn test_compact_blob_columns() {
2599 let test_dir = TempStrDir::default();
2600 let schema = Arc::new(Schema::new(vec![
2601 Field::new("id", DataType::Int32, false),
2602 Field::new("blob", DataType::LargeBinary, false)
2603 .with_metadata([(BLOB_META_KEY.to_string(), "true".to_string())].into()),
2604 ]));
2605 let expected_payload: Vec<Vec<u8>> =
2606 vec![vec![1, 2, 3], vec![4, 5, 6], vec![7, 8, 9, 10], vec![11]];
2607 let id_column: ArrayRef = Arc::new(Int32Array::from_iter_values(
2608 0..expected_payload.len() as i32,
2609 ));
2610 let blob_array: ArrayRef = Arc::new(LargeBinaryArray::from_iter(
2611 expected_payload.iter().map(|value| Some(value.as_slice())),
2612 ));
2613 let batch = RecordBatch::try_new(schema.clone(), vec![id_column, blob_array]).unwrap();
2614 let reader = RecordBatchIterator::new(vec![Ok(batch)], schema.clone());
2615
2616 let mut dataset = Dataset::write(
2617 reader,
2618 &test_dir,
2619 Some(WriteParams {
2620 max_rows_per_file: 1,
2621 ..Default::default()
2622 }),
2623 )
2624 .await
2625 .unwrap();
2626 dataset.validate().await.unwrap();
2627 assert!(dataset.get_fragments().len() > 1);
2628
2629 compact_files(&mut dataset, CompactionOptions::default(), None)
2630 .await
2631 .unwrap();
2632 dataset.validate().await.unwrap();
2633 assert_eq!(dataset.get_fragments().len(), 1);
2634
2635 let dataset = Arc::new(dataset);
2636 let row_indices: Vec<u64> = (0..expected_payload.len() as u64).collect();
2637 let blobs = dataset
2638 .take_blobs_by_indices(&row_indices, "blob")
2639 .await
2640 .unwrap();
2641 assert_eq!(blobs.len(), expected_payload.len());
2642 for (blob, expected) in blobs.iter().zip(expected_payload.iter()) {
2643 let bytes = blob.as_ref().unwrap().read().await.unwrap();
2644 assert_eq!(bytes.as_ref(), expected.as_slice());
2645 }
2646 }
2647
2648 fn row_addrs(frag_idx: u32, offsets: Range<u32>) -> Range<u64> {
2649 let start = RowAddress::new_from_parts(frag_idx, offsets.start);
2650 let end = RowAddress::new_from_parts(frag_idx, offsets.end);
2651 start.into()..end.into()
2652 }
2653
2654 fn expect_remap(
2657 ranges: &[Vec<(Range<u64>, bool)>],
2658 starting_new_frag_idx: u32,
2659 ) -> MockIndexRemapper {
2660 let mut expected_remap: HashMap<u64, Option<u64>> = HashMap::default();
2661 expected_remap.reserve(ranges.iter().map(|r| r.len()).sum());
2662 for (new_frag_offset, new_frag_ranges) in ranges.iter().enumerate() {
2663 let new_frag_idx = starting_new_frag_idx + new_frag_offset as u32;
2664 let mut row_offset = 0;
2665 for (old_id_range, is_found) in new_frag_ranges.iter() {
2666 for old_id in old_id_range.clone() {
2667 if *is_found {
2668 let new_id = RowAddress::new_from_parts(new_frag_idx, row_offset);
2669 expected_remap.insert(old_id, Some(new_id.into()));
2670 row_offset += 1;
2671 } else {
2672 expected_remap.insert(old_id, None);
2673 }
2674 }
2675 }
2676 }
2677 MockIndexRemapper {
2678 expectations: vec![MockIndexRemapperExpectation {
2679 expected: expected_remap,
2680 answer: vec![],
2681 }],
2682 }
2683 }
2684
2685 #[rstest]
2686 #[tokio::test]
2687 async fn test_compact_many(
2688 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
2689 data_storage_version: LanceFileVersion,
2690 ) {
2691 let test_dir = TempStrDir::default();
2692 let test_uri = &test_dir;
2693
2694 let data = sample_data();
2695
2696 let reader = RecordBatchIterator::new(vec![Ok(data.slice(0, 1200))], data.schema());
2698 let write_params = WriteParams {
2699 max_rows_per_file: 400,
2700 data_storage_version: Some(data_storage_version),
2701 ..Default::default()
2702 };
2703 Dataset::write(reader, test_uri, Some(write_params))
2704 .await
2705 .unwrap();
2706
2707 let reader = RecordBatchIterator::new(vec![Ok(data.slice(1200, 2000))], data.schema());
2709 let write_params = WriteParams {
2710 max_rows_per_file: 1000,
2711 data_storage_version: Some(data_storage_version),
2712 mode: WriteMode::Append,
2713 ..Default::default()
2714 };
2715 let mut dataset = Dataset::write(reader, test_uri, Some(write_params))
2716 .await
2717 .unwrap();
2718
2719 dataset.delete("a = 1300").await.unwrap();
2721
2722 dataset.delete("a >= 2400 AND a < 2600").await.unwrap();
2724
2725 let reader = RecordBatchIterator::new(vec![Ok(data.slice(3200, 600))], data.schema());
2727 let write_params = WriteParams {
2728 max_rows_per_file: 300,
2729 data_storage_version: Some(data_storage_version),
2730 mode: WriteMode::Append,
2731 ..Default::default()
2732 };
2733 let mut dataset = Dataset::write(reader, test_uri, Some(write_params))
2734 .await
2735 .unwrap();
2736
2737 let first_new_frag_idx = 7;
2738 let remap_a = expect_remap(
2742 &[
2743 vec![
2744 (row_addrs(0, 0..400), true),
2746 (row_addrs(1, 0..400), true),
2747 (row_addrs(2, 0..200), true),
2748 ],
2749 vec![(row_addrs(2, 200..400), true)],
2750 vec![
2753 (row_addrs(4, 0..200), true),
2755 (row_addrs(4, 200..400), false),
2756 (row_addrs(4, 400..1000), true),
2757 (row_addrs(5, 0..200), true),
2759 ],
2760 vec![(row_addrs(5, 200..300), true), (row_addrs(6, 0..300), true)],
2761 ],
2762 first_new_frag_idx,
2763 );
2764 let remap_b = expect_remap(
2765 &[
2766 vec![
2768 (row_addrs(4, 0..200), true),
2769 (row_addrs(4, 200..400), false),
2770 (row_addrs(4, 400..1000), true),
2771 (row_addrs(5, 0..200), true),
2772 ],
2773 vec![(row_addrs(5, 200..300), true), (row_addrs(6, 0..300), true)],
2774 vec![
2776 (row_addrs(0, 0..400), true),
2777 (row_addrs(1, 0..400), true),
2778 (row_addrs(2, 0..200), true),
2779 ],
2780 vec![(row_addrs(2, 200..400), true)],
2781 ],
2782 first_new_frag_idx,
2783 );
2784
2785 let options = CompactionOptions {
2787 target_rows_per_fragment: 1000,
2788 ..Default::default()
2789 };
2790 let plan = plan_compaction(&dataset, &options).await.unwrap();
2791 assert_eq!(plan.tasks().len(), 2);
2792 assert_eq!(plan.tasks()[0].fragments.len(), 3);
2793 assert_eq!(plan.tasks()[1].fragments.len(), 3);
2794
2795 assert_eq!(
2796 plan.tasks()[0]
2797 .fragments
2798 .iter()
2799 .map(|f| f.id)
2800 .collect::<Vec<_>>(),
2801 vec![0, 1, 2]
2802 );
2803 assert_eq!(
2804 plan.tasks()[1]
2805 .fragments
2806 .iter()
2807 .map(|f| f.id)
2808 .collect::<Vec<_>>(),
2809 vec![4, 5, 6]
2810 );
2811
2812 let mock_remapper = MockIndexRemapper::in_any_order(&[remap_a, remap_b]);
2813
2814 let metrics = compact_files(&mut dataset, options, Some(Arc::new(mock_remapper)))
2816 .await
2817 .unwrap();
2818
2819 assert_eq!(metrics.fragments_removed, 6);
2821 assert_eq!(metrics.fragments_added, 4);
2822 assert_eq!(metrics.files_removed, 7); assert_eq!(metrics.files_added, 4);
2824
2825 let fragment_ids = dataset
2826 .get_fragments()
2827 .iter()
2828 .map(|f| f.id())
2829 .collect::<Vec<_>>();
2830 assert_eq!(fragment_ids, vec![3, 7, 8, 9, 10]);
2831 }
2832
2833 #[rstest]
2834 #[tokio::test]
2835 async fn test_compact_data_files(
2836 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
2837 data_storage_version: LanceFileVersion,
2838 ) {
2839 let test_dir = TempStrDir::default();
2840 let test_uri = &test_dir;
2841
2842 let data = sample_data();
2843
2844 let reader = RecordBatchIterator::new(vec![Ok(data.clone())], data.schema());
2846 let write_params = WriteParams {
2847 max_rows_per_file: 5_000,
2848 max_rows_per_group: 1_000,
2849 data_storage_version: Some(data_storage_version),
2850 ..Default::default()
2851 };
2852 let mut dataset = Dataset::write(reader, test_uri, Some(write_params))
2853 .await
2854 .unwrap();
2855
2856 let schema = Schema::new(vec![
2858 Field::new("a", DataType::Int64, false),
2859 Field::new("x", DataType::Float32, false),
2860 ]);
2861
2862 let data = RecordBatch::try_new(
2863 Arc::new(schema),
2864 vec![
2865 Arc::new(Int64Array::from_iter_values(0..10_000)),
2866 Arc::new(Float32Array::from_iter_values(
2867 (0..10_000).map(|x| x as f32 * std::f32::consts::PI),
2868 )),
2869 ],
2870 )
2871 .unwrap();
2872 let reader = RecordBatchIterator::new(vec![Ok(data.clone())], data.schema());
2873
2874 dataset.merge(reader, "a", "a").await.unwrap();
2875
2876 let expected_remap = expect_remap(
2877 &[vec![
2878 (row_addrs(0, 0..5000), true),
2880 (row_addrs(1, 0..5000), true),
2881 ]],
2882 2,
2883 );
2884
2885 let plan = plan_compaction(
2886 &dataset,
2887 &CompactionOptions {
2888 ..Default::default()
2889 },
2890 )
2891 .await
2892 .unwrap();
2893 assert_eq!(plan.tasks().len(), 1);
2894 assert_eq!(plan.tasks()[0].fragments.len(), 2);
2895
2896 let metrics = compact_files(&mut dataset, plan.options, Some(Arc::new(expected_remap)))
2897 .await
2898 .unwrap();
2899
2900 assert_eq!(metrics.files_removed, 4); assert_eq!(metrics.files_added, 1); assert_eq!(metrics.fragments_removed, 2);
2903 assert_eq!(metrics.fragments_added, 1);
2904
2905 let scanner = dataset.scan();
2907 let batches = scanner
2908 .try_into_stream()
2909 .await
2910 .unwrap()
2911 .try_collect::<Vec<_>>()
2912 .await
2913 .unwrap();
2914 let scanned_data = concat_batches(&batches[0].schema(), &batches).unwrap();
2915
2916 assert_eq!(scanned_data, data);
2917 }
2918
2919 #[rstest]
2920 #[tokio::test]
2921 async fn test_compact_with_io_buffer_size(
2922 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
2923 data_storage_version: LanceFileVersion,
2924 ) {
2925 let test_dir = TempStrDir::default();
2928 let test_uri = &test_dir;
2929
2930 let data = sample_data();
2931
2932 let reader = RecordBatchIterator::new(vec![Ok(data.clone())], data.schema());
2934 let write_params = WriteParams {
2935 max_rows_per_file: 5_000,
2936 max_rows_per_group: 1_000,
2937 data_storage_version: Some(data_storage_version),
2938 ..Default::default()
2939 };
2940 let mut dataset = Dataset::write(reader, test_uri, Some(write_params))
2941 .await
2942 .unwrap();
2943 assert_eq!(dataset.get_fragments().len(), 2);
2944
2945 let options = CompactionOptions {
2946 io_buffer_size: Some(256 * 1024 * 1024),
2948 ..Default::default()
2949 };
2950 let plan = plan_compaction(&dataset, &options).await.unwrap();
2951 assert_eq!(plan.tasks().len(), 1);
2952
2953 let metrics = compact_files(&mut dataset, options, None).await.unwrap();
2954 assert_eq!(metrics.fragments_removed, 2);
2955 assert_eq!(metrics.fragments_added, 1);
2956
2957 let scanner = dataset.scan();
2959 let batches = scanner
2960 .try_into_stream()
2961 .await
2962 .unwrap()
2963 .try_collect::<Vec<_>>()
2964 .await
2965 .unwrap();
2966 let scanned_data = concat_batches(&batches[0].schema(), &batches).unwrap();
2967 assert_eq!(scanned_data.num_rows(), data.num_rows());
2968 }
2969
2970 #[rstest]
2971 #[tokio::test]
2972 async fn test_compact_deletions(
2973 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
2974 data_storage_version: LanceFileVersion,
2975 ) {
2976 let test_dir = TempStrDir::default();
2980 let test_uri = &test_dir;
2981
2982 let data = sample_data();
2983
2984 let reader = RecordBatchIterator::new(vec![Ok(data.slice(0, 1000))], data.schema());
2986 let write_params = WriteParams {
2987 max_rows_per_file: 1000,
2988 data_storage_version: Some(data_storage_version),
2989 ..Default::default()
2990 };
2991 let mut dataset = Dataset::write(reader, test_uri, Some(write_params))
2992 .await
2993 .unwrap();
2994
2995 dataset.delete("a <= 500").await.unwrap();
2996
2997 let mut options = CompactionOptions {
2999 materialize_deletions_threshold: 0.8,
3000 ..Default::default()
3001 };
3002 let plan = plan_compaction(&dataset, &options).await.unwrap();
3003 assert_eq!(plan.tasks().len(), 0);
3004
3005 options.materialize_deletions_threshold = 0.1;
3007 options.materialize_deletions = false;
3008 let plan = plan_compaction(&dataset, &options).await.unwrap();
3009 assert_eq!(plan.tasks().len(), 0);
3010
3011 options.materialize_deletions = true;
3013 let plan = plan_compaction(&dataset, &options).await.unwrap();
3014 assert_eq!(plan.tasks().len(), 1);
3015
3016 let metrics = compact_files(&mut dataset, options, None).await.unwrap();
3017 assert_eq!(metrics.fragments_removed, 1);
3018 assert_eq!(metrics.files_removed, 2);
3019 assert_eq!(metrics.fragments_added, 1);
3020
3021 let fragments = dataset.get_fragments();
3022 assert_eq!(fragments.len(), 1);
3023 assert!(fragments[0].metadata.deletion_file.is_none());
3024 }
3025
3026 #[derive(Debug, Default, Clone, PartialEq, Serialize, Deserialize)]
3027 struct IgnoreRemap {}
3028
3029 #[async_trait]
3030 impl IndexRemapper for IgnoreRemap {
3031 async fn remap_indices(&self, _: RowAddrRemap, _: &[u64]) -> Result<Vec<RemappedIndex>> {
3032 Ok(Vec::new())
3033 }
3034 }
3035
3036 #[async_trait]
3037 impl IndexRemapperOptions for IgnoreRemap {
3038 async fn create_remapper(&self, _: &Dataset) -> Result<Option<Box<dyn IndexRemapper>>> {
3039 Ok(None)
3040 }
3041 }
3042
3043 #[rstest]
3044 #[case::without_index(false)]
3045 #[case::with_index(true)]
3046 #[tokio::test]
3047 async fn test_row_addrs_only_used_with_remappable_index(#[case] has_index: bool) {
3048 let data = sample_data();
3049 let reader = RecordBatchIterator::new(vec![Ok(data.slice(0, 9_000))], data.schema());
3050 let mut dataset = Dataset::write(
3051 reader,
3052 "memory://",
3053 Some(WriteParams {
3054 max_rows_per_file: 3_000,
3055 data_storage_version: Some(LanceFileVersion::Legacy),
3056 ..Default::default()
3057 }),
3058 )
3059 .await
3060 .unwrap();
3061
3062 if has_index {
3063 create_scalar_index(&mut dataset, "a", false).await;
3064 }
3065
3066 let options = CompactionOptions {
3067 target_rows_per_fragment: 9_000,
3068 ..Default::default()
3069 };
3070 let plan = plan_compaction(&dataset, &options).await.unwrap();
3071 assert_eq!(plan.tasks().len(), 1);
3072
3073 let mut result = rewrite_files(Cow::Borrowed(&dataset), plan.tasks()[0].clone(), &options)
3074 .await
3075 .unwrap();
3076 assert_eq!(result.row_addrs.is_some(), has_index);
3077
3078 if has_index {
3079 let row_addrs_bytes = result
3080 .row_addrs
3081 .as_ref()
3082 .expect("indexed compaction should capture row addresses");
3083 let row_addrs =
3084 RoaringTreemap::deserialize_from(&mut Cursor::new(row_addrs_bytes)).unwrap();
3085 assert_eq!(row_addrs.len(), 9_000);
3086 } else {
3087 result.row_addrs = Some(b"not a roaring treemap".to_vec());
3091 commit_compaction(
3092 &mut dataset,
3093 vec![result],
3094 Arc::new(DatasetIndexRemapperOptions::default()),
3095 &options,
3096 )
3097 .await
3098 .unwrap();
3099 assert_eq!(dataset.get_fragments().len(), 1);
3100 }
3101 }
3102
3103 #[rstest::rstest]
3104 #[tokio::test]
3105 async fn test_compact_distributed(
3106 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
3107 data_storage_version: LanceFileVersion,
3108 #[values(false, true)] use_stable_row_id: bool,
3109 ) {
3110 let test_dir = TempStrDir::default();
3114 let test_uri = &test_dir;
3115
3116 let data = sample_data();
3117
3118 let reader = RecordBatchIterator::new(vec![Ok(data.slice(0, 9000))], data.schema());
3120 let write_params = WriteParams {
3121 max_rows_per_file: 1000,
3122 data_storage_version: Some(data_storage_version),
3123 enable_stable_row_ids: use_stable_row_id,
3124 ..Default::default()
3125 };
3126 let mut dataset = Dataset::write(reader, test_uri, Some(write_params))
3127 .await
3128 .unwrap();
3129
3130 let options = CompactionOptions {
3132 target_rows_per_fragment: 3_000,
3133 ..Default::default()
3134 };
3135 let plan = plan_compaction(&dataset, &options).await.unwrap();
3136 assert_eq!(plan.tasks().len(), 3);
3137
3138 let dataset_ref = &dataset;
3139 let mut results = futures::stream::iter(plan.compaction_tasks())
3140 .then(|task| async move { task.execute(dataset_ref).await.unwrap() })
3141 .collect::<Vec<_>>()
3142 .await;
3143
3144 assert_eq!(results.len(), 3);
3145
3146 assert_eq!(
3147 results[0]
3148 .original_fragments
3149 .iter()
3150 .map(|f| f.id)
3151 .collect::<Vec<_>>(),
3152 vec![0, 1, 2]
3153 );
3154 assert_eq!(results[0].metrics.files_removed, 3);
3155 assert_eq!(results[0].metrics.files_added, 1);
3156
3157 commit_compaction(
3159 &mut dataset,
3160 vec![results.pop().unwrap()],
3161 Arc::new(IgnoreRemap::default()),
3162 &options,
3163 )
3164 .await
3165 .unwrap();
3166
3167 assert_eq!(dataset.manifest.version, 3);
3170
3171 commit_compaction(
3173 &mut dataset,
3174 results,
3175 Arc::new(IgnoreRemap::default()),
3176 &options,
3177 )
3178 .await
3179 .unwrap();
3180 assert_eq!(dataset.manifest.version, 5);
3183
3184 assert_eq!(dataset.manifest.uses_stable_row_ids(), use_stable_row_id,);
3185 }
3186
3187 #[tokio::test]
3188 async fn test_stable_row_indices() {
3189 let mut data_gen = BatchGenerator::new()
3191 .col(Box::new(
3192 RandomVector::new().vec_width(16).named("vec".to_owned()),
3193 ))
3194 .col(Box::new(IncrementingInt32::new().named("i".to_owned())));
3195 let mut dataset = Dataset::write(
3196 data_gen.batch(500),
3197 "memory://test/table",
3198 Some(WriteParams {
3199 enable_stable_row_ids: true,
3200 max_rows_per_file: 100, ..Default::default()
3202 }),
3203 )
3204 .await
3205 .unwrap();
3206
3207 dataset.delete("i < 110").await.unwrap();
3211
3212 dataset
3213 .create_index(
3214 &["i"],
3215 IndexType::Scalar,
3216 Some("scalar".into()),
3217 &ScalarIndexParams::default(),
3218 false,
3219 )
3220 .await
3221 .unwrap();
3222 let params = VectorIndexParams::ivf_pq(1, 8, 1, MetricType::L2, 50);
3223 dataset
3224 .create_index(
3225 &["vec"],
3226 IndexType::Vector,
3227 Some("vector".into()),
3228 ¶ms,
3229 false,
3230 )
3231 .await
3232 .unwrap();
3233
3234 async fn index_set(dataset: &Dataset) -> HashSet<Uuid> {
3235 dataset
3236 .load_indices()
3237 .await
3238 .unwrap()
3239 .iter()
3240 .map(|index| index.uuid)
3241 .collect()
3242 }
3243 let indices = index_set(&dataset).await;
3244
3245 async fn vector_query(dataset: &Dataset) -> RecordBatch {
3246 let mut scanner = dataset.scan();
3247
3248 let query = Float32Array::from(vec![0.0f32; 16]);
3249 scanner
3250 .nearest("vec", &query, 10)
3251 .unwrap()
3252 .project(&["i"])
3253 .unwrap();
3254
3255 scanner.try_into_batch().await.unwrap()
3256 }
3257
3258 async fn scalar_query(dataset: &Dataset) -> RecordBatch {
3259 let mut scanner = dataset.scan();
3260
3261 scanner.filter("i = 100").unwrap().project(&["i"]).unwrap();
3262
3263 scanner.try_into_batch().await.unwrap()
3264 }
3265
3266 let before_vec_result = vector_query(&dataset).await;
3267 let before_scalar_result = scalar_query(&dataset).await;
3268
3269 let options = CompactionOptions {
3270 target_rows_per_fragment: 180,
3271 ..Default::default()
3272 };
3273 let _metrics = compact_files(&mut dataset, options, None).await.unwrap();
3274
3275 let current_indices = index_set(&dataset).await;
3278 assert_eq!(indices, current_indices);
3279
3280 let after_vec_result = vector_query(&dataset).await;
3281 assert_eq!(before_vec_result, after_vec_result);
3282
3283 let after_scalar_result = scalar_query(&dataset).await;
3284 assert_eq!(before_scalar_result, after_scalar_result);
3285 }
3286
3287 #[tokio::test]
3293 async fn test_defer_index_remap_large_external_file() {
3294 let test_dir = TempStrDir::default();
3295 let test_uri = &test_dir;
3296
3297 let num_fragments = 150usize;
3300 let rows_per_fragment = 1000usize;
3301 let total_rows = num_fragments * rows_per_fragment;
3302
3303 let schema = Arc::new(Schema::new(vec![Field::new("i", DataType::Int32, false)]));
3304
3305 let mut dataset = Dataset::write(
3306 RecordBatchIterator::new(
3307 vec![Ok(RecordBatch::try_new(
3308 schema.clone(),
3309 vec![Arc::new(Int32Array::from_iter_values(0..total_rows as i32)) as ArrayRef],
3310 )
3311 .unwrap())],
3312 schema.clone(),
3313 ),
3314 test_uri,
3315 Some(WriteParams {
3316 max_rows_per_file: rows_per_fragment,
3317 ..Default::default()
3318 }),
3319 )
3320 .await
3321 .unwrap();
3322
3323 assert_eq!(dataset.get_fragments().len(), num_fragments);
3324
3325 create_scalar_index(&mut dataset, "i", false).await;
3328
3329 dataset.delete("i % 1000 = 0").await.unwrap();
3331
3332 compact_files(
3333 &mut dataset,
3334 CompactionOptions {
3335 defer_index_remap: true,
3336 ..Default::default()
3337 },
3338 None,
3339 )
3340 .await
3341 .unwrap();
3342
3343 let frag_reuse_meta = dataset
3346 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
3347 .await
3348 .unwrap()
3349 .expect("fragment reuse index must exist after compaction");
3350
3351 load_frag_reuse_index_details(&dataset, &frag_reuse_meta)
3352 .await
3353 .expect("loading large frag reuse index details must not fail");
3354 }
3355
3356 #[tokio::test]
3357 async fn test_defer_index_remap_rejected_with_stable_row_ids() {
3358 let test_dir = TempStrDir::default();
3359 let test_uri = &test_dir;
3360
3361 let data = sample_data();
3362 let reader = RecordBatchIterator::new(vec![Ok(data.slice(0, 9000))], data.schema());
3363 let mut dataset = Dataset::write(
3364 reader,
3365 test_uri,
3366 Some(WriteParams {
3367 max_rows_per_file: 1000, enable_stable_row_ids: true,
3369 ..Default::default()
3370 }),
3371 )
3372 .await
3373 .unwrap();
3374 assert!(dataset.manifest.uses_stable_row_ids());
3375
3376 let options = CompactionOptions {
3377 target_rows_per_fragment: 3_000,
3378 defer_index_remap: true,
3379 ..Default::default()
3380 };
3381
3382 let plan_err = plan_compaction(&dataset, &options).await.unwrap_err();
3384 assert!(matches!(plan_err, Error::InvalidInput { .. }));
3385 let msg = plan_err.to_string();
3386 assert!(msg.contains("defer_index_remap"));
3387 assert!(msg.contains("stable row IDs"));
3388
3389 let version_before = dataset.manifest.version;
3392 let compact_err = compact_files(&mut dataset, options, None)
3393 .await
3394 .unwrap_err();
3395 assert!(matches!(compact_err, Error::InvalidInput { .. }));
3396 assert_eq!(dataset.manifest.version, version_before);
3397 }
3398
3399 #[tokio::test]
3400 async fn test_defer_index_remap() {
3401 let mut data_gen = BatchGenerator::new()
3402 .col(Box::new(
3403 RandomVector::new().vec_width(128).named("vec".to_owned()),
3404 ))
3405 .col(Box::new(IncrementingInt32::new().named("i".to_owned())));
3406
3407 let mut dataset = Dataset::write(
3408 data_gen.batch(6_000),
3409 "memory://test/table",
3410 Some(WriteParams {
3411 max_rows_per_file: 1_000, ..Default::default()
3413 }),
3414 )
3415 .await
3416 .unwrap();
3417
3418 let mut data_gen2 = BatchGenerator::new()
3420 .col(Box::new(
3421 RandomVector::new().vec_width(128).named("vec".to_owned()),
3422 ))
3423 .col(Box::new(IncrementingInt32::new().named("i".to_owned())));
3424
3425 let mut dataset2 = Dataset::write(
3426 data_gen2.batch(6_000),
3427 "memory://test/table",
3428 Some(WriteParams {
3429 max_rows_per_file: 1_000, ..Default::default()
3431 }),
3432 )
3433 .await
3434 .unwrap();
3435
3436 dataset.delete("i < 500").await.unwrap();
3438 dataset2.delete("i < 500").await.unwrap();
3439
3440 create_scalar_index(&mut dataset, "i", false).await;
3443 create_scalar_index(&mut dataset2, "i", false).await;
3444
3445 let initial_indices = dataset.load_indices().await.unwrap();
3447 assert_eq!(initial_indices.len(), 1);
3448 assert_eq!(initial_indices[0].name, "scalar");
3449
3450 let original_scalar_uuid = initial_indices[0].uuid;
3452
3453 let options = CompactionOptions {
3455 target_rows_per_fragment: 2_000,
3456 defer_index_remap: true,
3457 ..Default::default()
3458 };
3459 let options2 = CompactionOptions {
3460 target_rows_per_fragment: 2_000,
3461 defer_index_remap: false,
3462 ..Default::default()
3463 };
3464
3465 let plan = plan_compaction(&dataset, &options).await.unwrap();
3466 let plan2 = plan_compaction(&dataset2, &options2).await.unwrap();
3467
3468 let mut expected_all_old_frag_ids = Vec::new();
3469 let mut expected_all_new_frag_ids = Vec::new();
3470 let mut expected_all_new_frag_bitmap = RoaringBitmap::new();
3471 let mut expected_all_row_id_map = HashMap::new();
3472 let mut deferred_results = Vec::new();
3473 let mut immediate_results = Vec::new();
3474
3475 for (task, task2) in plan.tasks().iter().zip(plan2.tasks()) {
3476 let deferred_result = rewrite_files(Cow::Borrowed(&dataset), task.clone(), &options)
3477 .await
3478 .unwrap();
3479 let immediate_result =
3480 rewrite_files(Cow::Borrowed(&dataset2), task2.clone(), &options2)
3481 .await
3482 .unwrap();
3483
3484 assert!(deferred_result.row_addrs.is_some());
3486 assert!(!deferred_result.row_addrs.as_ref().unwrap().is_empty());
3487 assert!(!deferred_result.row_addrs.as_ref().unwrap().is_empty());
3488 assert!(!deferred_result.original_fragments.is_empty());
3489 assert!(!deferred_result.new_fragments.is_empty());
3490
3491 assert!(immediate_result.row_addrs.is_some());
3492 assert!(!immediate_result.original_fragments.is_empty());
3493 assert!(!immediate_result.new_fragments.is_empty());
3494
3495 assert_eq!(deferred_result.row_addrs, immediate_result.row_addrs);
3497
3498 deferred_results.push(deferred_result);
3499 immediate_results.push(immediate_result);
3500 }
3501
3502 {
3504 let frags: Vec<&mut Fragment> = immediate_results
3505 .iter_mut()
3506 .flat_map(|r| r.new_fragments.iter_mut())
3507 .collect();
3508 reserve_fragment_ids(&dataset2, frags.into_iter())
3509 .await
3510 .unwrap();
3511 }
3512
3513 for immediate_result in &immediate_results {
3515 let row_addrs_bytes = immediate_result.row_addrs.as_ref().unwrap();
3516 let row_addrs =
3517 RoaringTreemap::deserialize_from(&mut Cursor::new(row_addrs_bytes)).unwrap();
3518 let transposed = transpose_row_addrs(
3519 row_addrs,
3520 &immediate_result.original_fragments,
3521 &immediate_result.new_fragments,
3522 );
3523 expected_all_row_id_map.extend(transposed);
3524 immediate_result.new_fragments.iter().for_each(|frag| {
3525 expected_all_new_frag_bitmap.insert(frag.id as u32);
3526 });
3527 expected_all_new_frag_ids.extend(
3528 immediate_result
3529 .new_fragments
3530 .iter()
3531 .map(|s| s.id)
3532 .collect::<Vec<_>>(),
3533 );
3534 expected_all_old_frag_ids.extend(
3535 immediate_result
3536 .original_fragments
3537 .iter()
3538 .map(|s| s.id)
3539 .collect::<Vec<_>>(),
3540 );
3541 }
3542
3543 let first_metrics = commit_compaction(
3545 &mut dataset,
3546 deferred_results.clone(),
3547 Arc::new(DatasetIndexRemapperOptions::default()),
3548 &options,
3549 )
3550 .await
3551 .unwrap();
3552
3553 assert!(first_metrics.fragments_removed > 0);
3555 assert!(first_metrics.fragments_added > 0);
3556
3557 let Some(frag_reuse_index_meta) = dataset
3559 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
3560 .await
3561 .unwrap()
3562 else {
3563 panic!("Fragment reuse index must be available");
3564 };
3565
3566 assert_eq!(
3567 frag_reuse_index_meta.fragment_bitmap.clone().unwrap(),
3568 expected_all_new_frag_bitmap
3569 );
3570 let frag_reuse_details = load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta)
3571 .await
3572 .unwrap();
3573 let frag_reuse_index =
3574 open_frag_reuse_index(frag_reuse_index_meta.uuid, frag_reuse_details.as_ref())
3575 .await
3576 .unwrap();
3577 let stats = FragReuseIndexHandle(Arc::new(frag_reuse_index.clone()))
3578 .statistics()
3579 .unwrap();
3580 assert_eq!(
3581 serde_json::to_string(&stats).unwrap(),
3582 dataset
3583 .index_statistics(FRAG_REUSE_INDEX_NAME)
3584 .await
3585 .unwrap()
3586 );
3587
3588 let compaction_version = &frag_reuse_index.details.versions[0];
3590 assert_eq!(frag_reuse_index.details.versions.len(), 1);
3591 assert_eq!(
3592 compaction_version.dataset_version,
3593 frag_reuse_index_meta.dataset_version
3594 );
3595
3596 let mut compacted_all_old_frag_digests = Vec::new();
3598 let mut compacted_all_new_frag_digests = Vec::new();
3599 let mut transposed_map = HashMap::new();
3600 for group in compaction_version.groups.iter() {
3601 let changed_row_addr_bytes = &group.changed_row_addrs;
3602 let mut cursor = Cursor::new(&changed_row_addr_bytes);
3603 let changed_row_addrs = RoaringTreemap::deserialize_from(&mut cursor).unwrap();
3604 compacted_all_old_frag_digests.extend(group.old_frags.clone());
3605 compacted_all_new_frag_digests.extend(group.new_frags.clone());
3606
3607 let group_transposed_map = transpose_row_ids_from_digest(
3608 changed_row_addrs,
3609 &group.old_frags,
3610 &group.new_frags,
3611 );
3612 transposed_map.extend(group_transposed_map);
3613 }
3614 assert_eq!(transposed_map, expected_all_row_id_map);
3615 assert_eq!(
3616 compacted_all_old_frag_digests
3617 .iter()
3618 .map(|f| f.id)
3619 .collect::<Vec<_>>(),
3620 expected_all_old_frag_ids
3621 );
3622 assert_eq!(
3623 compacted_all_new_frag_digests
3624 .iter()
3625 .map(|f| f.id)
3626 .collect::<Vec<_>>(),
3627 expected_all_new_frag_ids
3628 );
3629
3630 let Some(current_scalar_index) = dataset.load_index_by_name("scalar").await.unwrap() else {
3632 panic!("scalar index must be available");
3633 };
3634 assert_eq!(current_scalar_index.uuid, original_scalar_uuid);
3635 }
3636
3637 #[tokio::test]
3638 async fn test_defer_index_remap_skips_fri_when_no_indexed_data() {
3639 let mut data_gen =
3642 BatchGenerator::new().col(Box::new(IncrementingInt32::new().named("i".to_owned())));
3643
3644 let mut dataset = Dataset::write(
3645 data_gen.batch(600),
3646 "memory://test/noindex",
3647 Some(WriteParams {
3648 max_rows_per_file: 100, ..Default::default()
3650 }),
3651 )
3652 .await
3653 .unwrap();
3654
3655 assert!(dataset.load_indices().await.unwrap().is_empty());
3657 let fragments_before = dataset.get_fragments().len();
3658 assert!(fragments_before > 1, "need multiple fragments to compact");
3659
3660 let options = CompactionOptions {
3661 target_rows_per_fragment: 100_000,
3662 defer_index_remap: true,
3663 ..Default::default()
3664 };
3665 compact_files(&mut dataset, options, None).await.unwrap();
3666
3667 assert!(
3669 dataset.get_fragments().len() < fragments_before,
3670 "compaction should have merged fragments"
3671 );
3672 assert!(
3674 dataset
3675 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
3676 .await
3677 .unwrap()
3678 .is_none(),
3679 "deferred compaction with no indexed data must not create an FRI"
3680 );
3681 }
3682
3683 #[tokio::test]
3684 async fn test_defer_index_remap_multiple_compactions() {
3685 let mut data_gen = BatchGenerator::new()
3686 .col(Box::new(
3687 RandomVector::new().vec_width(128).named("vec".to_owned()),
3688 ))
3689 .col(Box::new(IncrementingInt32::new().named("i".to_owned())));
3690
3691 let mut dataset = Dataset::write(
3692 data_gen.batch(6_000),
3693 "memory://test/table",
3694 Some(WriteParams {
3695 max_rows_per_file: 1_000, ..Default::default()
3697 }),
3698 )
3699 .await
3700 .unwrap();
3701
3702 create_scalar_index(&mut dataset, "i", false).await;
3705
3706 let options = CompactionOptions {
3707 target_rows_per_fragment: 2_000,
3708 defer_index_remap: true,
3709 ..Default::default()
3710 };
3711
3712 let mut compact_read_versions = Vec::new();
3713 for i in 0..10 {
3714 dataset
3715 .delete(&format!("i < {}", 500 * (i + 1)))
3716 .await
3717 .unwrap();
3718 let read_version = dataset.manifest.version;
3719 compact_files(&mut dataset, options.clone(), None)
3720 .await
3721 .unwrap();
3722
3723 if dataset.manifest.version > read_version {
3725 compact_read_versions.push(read_version);
3726 }
3727
3728 let Some(frag_reuse_index_meta) = dataset
3730 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
3731 .await
3732 .unwrap()
3733 else {
3734 panic!("Fragment reuse index must be available");
3735 };
3736 let frag_reuse_details =
3737 load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta)
3738 .await
3739 .unwrap();
3740 let frag_reuse_index =
3741 open_frag_reuse_index(frag_reuse_index_meta.uuid, frag_reuse_details.as_ref())
3742 .await
3743 .unwrap();
3744
3745 assert_eq!(
3747 frag_reuse_index
3748 .details
3749 .versions
3750 .iter()
3751 .map(|v| v.dataset_version)
3752 .collect::<Vec<_>>(),
3753 compact_read_versions
3754 );
3755 }
3756 }
3757
3758 #[tokio::test]
3759 async fn test_defer_index_remap_mixed_records_all_groups() {
3760 let mut data_gen =
3763 BatchGenerator::new().col(Box::new(IncrementingInt32::new().named("i".to_owned())));
3764 let mut dataset = Dataset::write(
3765 data_gen.batch(300),
3766 "memory://test/mixed",
3767 Some(WriteParams {
3768 max_rows_per_file: 100, ..Default::default()
3770 }),
3771 )
3772 .await
3773 .unwrap();
3774
3775 create_scalar_index(&mut dataset, "i", false).await;
3777 Dataset::write(
3778 data_gen.batch(300),
3779 WriteDestination::Dataset(Arc::new(dataset.clone())),
3780 Some(WriteParams {
3781 max_rows_per_file: 100, mode: WriteMode::Append,
3783 ..Default::default()
3784 }),
3785 )
3786 .await
3787 .unwrap();
3788 dataset.checkout_latest().await.unwrap();
3789
3790 let indexed: HashSet<u32> = dataset
3792 .load_index_by_name("scalar")
3793 .await
3794 .unwrap()
3795 .unwrap()
3796 .fragment_bitmap
3797 .unwrap()
3798 .iter()
3799 .collect();
3800 let unindexed_frags: Vec<u64> = dataset
3801 .fragments()
3802 .iter()
3803 .map(|f| f.id)
3804 .filter(|id| !indexed.contains(&(*id as u32)))
3805 .collect();
3806 assert!(
3807 !unindexed_frags.is_empty(),
3808 "expected some unindexed fragments"
3809 );
3810
3811 compact_files(
3812 &mut dataset,
3813 CompactionOptions {
3814 target_rows_per_fragment: 100_000,
3815 defer_index_remap: true,
3816 ..Default::default()
3817 },
3818 None,
3819 )
3820 .await
3821 .unwrap();
3822
3823 let fri_meta = dataset
3827 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
3828 .await
3829 .unwrap()
3830 .expect("mixed compaction must write an FRI");
3831 let details = load_frag_reuse_index_details(&dataset, &fri_meta)
3832 .await
3833 .unwrap();
3834 let recorded_old: HashSet<u64> = details
3835 .versions
3836 .iter()
3837 .flat_map(|v| v.old_frag_ids())
3838 .collect();
3839 for f in &unindexed_frags {
3840 assert!(
3841 recorded_old.contains(f),
3842 "unindexed fragment {f} must be recorded in the FRI (all-or-nothing)"
3843 );
3844 }
3845 }
3846
3847 #[tokio::test]
3848 async fn test_deferred_compaction_not_split_by_frag_reuse_index() {
3849 let data = sample_data();
3853 let test_dir = TempStrDir::default();
3854 let test_uri = &test_dir;
3855 let options = CompactionOptions {
3856 defer_index_remap: true,
3857 ..Default::default()
3858 };
3859
3860 let reader = RecordBatchIterator::new(vec![Ok(data.slice(0, 400))], data.schema());
3863 let mut dataset = Dataset::write(
3864 reader,
3865 test_uri,
3866 Some(WriteParams {
3867 max_rows_per_file: 200,
3868 ..Default::default()
3869 }),
3870 )
3871 .await
3872 .unwrap();
3873
3874 create_scalar_index(&mut dataset, "a", false).await;
3878 compact_files(&mut dataset, options.clone(), None)
3879 .await
3880 .unwrap();
3881 assert_eq!(dataset.get_fragments().len(), 1);
3882 assert!(
3883 dataset
3884 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
3885 .await
3886 .unwrap()
3887 .is_some()
3888 );
3889
3890 let reader = RecordBatchIterator::new(vec![Ok(data.slice(400, 400))], data.schema());
3892 let mut dataset = Dataset::write(
3893 reader,
3894 test_uri,
3895 Some(WriteParams {
3896 max_rows_per_file: 200,
3897 mode: WriteMode::Append,
3898 ..Default::default()
3899 }),
3900 )
3901 .await
3902 .unwrap();
3903 assert_eq!(dataset.get_fragments().len(), 3);
3904
3905 create_scalar_index(&mut dataset, "a", true).await;
3909
3910 compact_files(&mut dataset, options, None).await.unwrap();
3911 assert_eq!(
3912 dataset.get_fragments().len(),
3913 1,
3914 "FRI (a system index) must not split the compaction bin; all fragments coalesce"
3915 );
3916 }
3917
3918 #[tokio::test]
3919 async fn test_remap_index_after_compaction() {
3920 let mut data_gen = BatchGenerator::new()
3921 .col(Box::new(
3922 RandomVector::new().vec_width(128).named("vec".to_owned()),
3923 ))
3924 .col(Box::new(IncrementingInt32::new().named("i".to_owned())));
3925
3926 let mut dataset = Dataset::write(
3927 data_gen.batch(6_000),
3928 "memory://test/table",
3929 Some(WriteParams {
3930 max_rows_per_file: 1_000, ..Default::default()
3932 }),
3933 )
3934 .await
3935 .unwrap();
3936
3937 let index_name = Some("scalar".into());
3939 dataset
3940 .create_index(
3941 &["i"],
3942 IndexType::Scalar,
3943 index_name.clone(),
3944 &ScalarIndexParams::default(),
3945 false,
3946 )
3947 .await
3948 .unwrap();
3949
3950 let options = CompactionOptions {
3951 target_rows_per_fragment: 2_000,
3952 defer_index_remap: true,
3953 ..Default::default()
3954 };
3955
3956 let Some(scalar_index) = dataset.load_index_by_name("scalar").await.unwrap() else {
3958 panic!("scalar index must be available");
3959 };
3960
3961 let result = remapping::remap_column_index(&mut dataset, &["i"], index_name.clone()).await;
3962 assert!(matches!(result, Err(Error::NotSupported { .. })));
3963
3964 let plan = plan_compaction(&dataset, &options).await.unwrap();
3965
3966 for task in plan.tasks().iter() {
3969 let rewrite_result = rewrite_files(Cow::Borrowed(&dataset), task.clone(), &options)
3970 .await
3971 .unwrap();
3972
3973 commit_compaction(
3974 &mut dataset,
3975 Vec::from([rewrite_result]),
3976 Arc::new(DatasetIndexRemapperOptions::default()),
3977 &options,
3978 )
3979 .await
3980 .unwrap();
3981 }
3982
3983 let Some(frag_reuse_index_meta) = dataset
3985 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
3986 .await
3987 .unwrap()
3988 else {
3989 panic!("Fragment reuse index must be available");
3990 };
3991 let frag_reuse_details = load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta)
3992 .await
3993 .unwrap();
3994 let frag_reuse_index =
3995 open_frag_reuse_index(frag_reuse_index_meta.uuid, frag_reuse_details.as_ref())
3996 .await
3997 .unwrap();
3998
3999 assert_eq!(frag_reuse_index.details.versions.len(), plan.tasks().len());
4000
4001 let mut all_fragment_bitmap = RoaringBitmap::new();
4003 dataset.fragments().iter().for_each(|f| {
4004 all_fragment_bitmap.insert(f.id as u32);
4005 });
4006 let Some(scalar_index_before_remap) = dataset.load_index_by_name("scalar").await.unwrap()
4007 else {
4008 panic!("scalar index must be available");
4009 };
4010 assert_eq!(
4011 scalar_index_before_remap.fragment_bitmap.unwrap(),
4012 all_fragment_bitmap
4013 );
4014
4015 remapping::remap_column_index(&mut dataset, &["i"], index_name.clone())
4017 .await
4018 .unwrap();
4019
4020 let indices = read_manifest_indexes(
4022 &dataset.object_store,
4023 &dataset.manifest_location,
4024 &dataset.manifest,
4025 )
4026 .await
4027 .unwrap();
4028 let Some(remapped_scalar_index) = indices.into_iter().find(|idx| idx.name == "scalar")
4029 else {
4030 panic!("scalar index must be available");
4031 };
4032 assert_ne!(remapped_scalar_index.uuid, scalar_index.uuid);
4033 assert_eq!(
4034 remapped_scalar_index.fragment_bitmap.unwrap(),
4035 all_fragment_bitmap
4036 );
4037 }
4038
4039 #[tokio::test]
4040 async fn test_concurrent_compaction_reindex_compaction_commit_first() {
4041 let mut data_gen = BatchGenerator::new()
4042 .col(Box::new(
4043 RandomVector::new().vec_width(128).named("vec".to_owned()),
4044 ))
4045 .col(Box::new(IncrementingInt32::new().named("i".to_owned())));
4046
4047 let mut dataset = Dataset::write(
4048 data_gen.batch(6_000),
4049 "memory://test/table",
4050 Some(WriteParams {
4051 max_rows_per_file: 1_000, ..Default::default()
4053 }),
4054 )
4055 .await
4056 .unwrap();
4057
4058 let index_name = Some("scalar".into());
4060 dataset
4061 .create_index(
4062 &["i"],
4063 IndexType::Scalar,
4064 index_name.clone(),
4065 &ScalarIndexParams::default(),
4066 false,
4067 )
4068 .await
4069 .unwrap();
4070
4071 Dataset::write(
4073 data_gen.batch(6_000),
4074 WriteDestination::Dataset(Arc::new(dataset.clone())),
4075 Some(WriteParams {
4076 max_rows_per_file: 1_000, mode: WriteMode::Append,
4078 ..Default::default()
4079 }),
4080 )
4081 .await
4082 .unwrap();
4083
4084 dataset.checkout_latest().await.unwrap();
4085 let mut dataset_clone = dataset.clone();
4086
4087 compact_files(
4089 &mut dataset,
4090 CompactionOptions {
4091 target_rows_per_fragment: 2_000,
4092 defer_index_remap: true,
4093 ..Default::default()
4094 },
4095 None,
4096 )
4097 .await
4098 .unwrap();
4099
4100 dataset_clone
4102 .create_index(
4103 &["i"],
4104 IndexType::Scalar,
4105 index_name.clone(),
4106 &ScalarIndexParams::default(),
4107 true,
4108 )
4109 .await
4110 .unwrap();
4111
4112 dataset.checkout_latest().await.unwrap();
4114
4115 let Some(scalar_index) = dataset.load_index_by_name("scalar").await.unwrap() else {
4116 panic!("scalar index must be available");
4117 };
4118 let index_frags = scalar_index
4119 .fragment_bitmap
4120 .unwrap()
4121 .iter()
4122 .collect::<HashSet<_>>();
4123 assert_eq!(
4124 index_frags,
4125 dataset
4126 .fragments()
4127 .iter()
4128 .map(|f| f.id as u32)
4129 .collect::<HashSet<_>>()
4130 )
4131 }
4132
4133 #[tokio::test]
4134 async fn test_concurrent_compaction_reindex_reindex_commit_first() {
4135 let mut data_gen = BatchGenerator::new()
4136 .col(Box::new(
4137 RandomVector::new().vec_width(128).named("vec".to_owned()),
4138 ))
4139 .col(Box::new(IncrementingInt32::new().named("i".to_owned())));
4140
4141 let mut dataset = Dataset::write(
4142 data_gen.batch(6_000),
4143 "memory://test/table",
4144 Some(WriteParams {
4145 max_rows_per_file: 1_000, ..Default::default()
4147 }),
4148 )
4149 .await
4150 .unwrap();
4151
4152 let index_name = Some("scalar".into());
4154 dataset
4155 .create_index(
4156 &["i"],
4157 IndexType::Scalar,
4158 index_name.clone(),
4159 &ScalarIndexParams::default(),
4160 false,
4161 )
4162 .await
4163 .unwrap();
4164
4165 Dataset::write(
4167 data_gen.batch(6_000),
4168 WriteDestination::Dataset(Arc::new(dataset.clone())),
4169 Some(WriteParams {
4170 max_rows_per_file: 1_000, mode: WriteMode::Append,
4172 ..Default::default()
4173 }),
4174 )
4175 .await
4176 .unwrap();
4177
4178 dataset.checkout_latest().await.unwrap();
4179 let mut dataset_clone = dataset.clone();
4180
4181 dataset
4183 .create_index(
4184 &["i"],
4185 IndexType::Scalar,
4186 index_name.clone(),
4187 &ScalarIndexParams::default(),
4188 true,
4189 )
4190 .await
4191 .unwrap();
4192
4193 compact_files(
4195 &mut dataset_clone,
4196 CompactionOptions {
4197 target_rows_per_fragment: 2_000,
4198 defer_index_remap: true,
4199 ..Default::default()
4200 },
4201 None,
4202 )
4203 .await
4204 .unwrap();
4205
4206 dataset.checkout_latest().await.unwrap();
4208 let Some(scalar_index) = dataset.load_index_by_name("scalar").await.unwrap() else {
4209 panic!("scalar index must be available");
4210 };
4211 let index_frags = scalar_index
4212 .fragment_bitmap
4213 .unwrap()
4214 .iter()
4215 .collect::<HashSet<_>>();
4216 assert_eq!(
4217 index_frags,
4218 dataset
4219 .fragments()
4220 .iter()
4221 .map(|f| f.id as u32)
4222 .collect::<HashSet<_>>()
4223 )
4224 }
4225
4226 #[tokio::test]
4227 async fn test_concurrent_cleanup_and_compaction_rebase_cleanup() {
4228 let mut dataset = lance_datagen::gen_batch()
4229 .col(
4230 "vec",
4231 lance_datagen::array::rand_vec::<Float32Type>(Dimension::from(128)),
4232 )
4233 .col("i", lance_datagen::array::step::<Int32Type>())
4234 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
4235 .await
4236 .unwrap();
4237
4238 create_scalar_index(&mut dataset, "i", false).await;
4240
4241 let options = CompactionOptions {
4242 target_rows_per_fragment: 2_000,
4243 defer_index_remap: true,
4244 ..Default::default()
4245 };
4246
4247 let plan = plan_compaction(&dataset, &options).await.unwrap();
4248 let tasks = plan.tasks();
4249
4250 let rewrite_result = rewrite_files(Cow::Borrowed(&dataset), tasks[0].clone(), &options)
4252 .await
4253 .unwrap();
4254
4255 commit_compaction(
4256 &mut dataset,
4257 Vec::from([rewrite_result]),
4258 Arc::new(DatasetIndexRemapperOptions::default()),
4259 &options,
4260 )
4261 .await
4262 .unwrap();
4263
4264 let mut dataset_clone = dataset.clone();
4265
4266 let Some(frag_reuse_index_meta) = dataset
4268 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
4269 .await
4270 .unwrap()
4271 else {
4272 panic!("Fragment reuse index must be available");
4273 };
4274
4275 let frag_reuse_details = load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta)
4276 .await
4277 .unwrap();
4278 assert_eq!(frag_reuse_details.versions.len(), 1);
4279
4280 let rewrite_result2 = rewrite_files(Cow::Borrowed(&dataset), tasks[1].clone(), &options)
4282 .await
4283 .unwrap();
4284 let rewritten_frags2 = rewrite_result2
4285 .original_fragments
4286 .iter()
4287 .map(|f| f.id)
4288 .collect::<Vec<_>>();
4289 commit_compaction(
4290 &mut dataset,
4291 Vec::from([rewrite_result2]),
4292 Arc::new(DatasetIndexRemapperOptions::default()),
4293 &options,
4294 )
4295 .await
4296 .unwrap();
4297
4298 let frag_reuse_index_meta2 = dataset
4300 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
4301 .await
4302 .unwrap()
4303 .unwrap();
4304 let frag_reuse_details2 = load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta2)
4305 .await
4306 .unwrap();
4307 let new_frags2 = frag_reuse_details2.versions.last().unwrap().new_frag_ids();
4308
4309 let rewrite_result3 = rewrite_files(Cow::Borrowed(&dataset), tasks[2].clone(), &options)
4310 .await
4311 .unwrap();
4312 let rewritten_frags3 = rewrite_result3
4313 .original_fragments
4314 .iter()
4315 .map(|f| f.id)
4316 .collect::<Vec<_>>();
4317 commit_compaction(
4318 &mut dataset,
4319 Vec::from([rewrite_result3]),
4320 Arc::new(DatasetIndexRemapperOptions::default()),
4321 &options,
4322 )
4323 .await
4324 .unwrap();
4325
4326 let frag_reuse_index_meta3 = dataset
4328 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
4329 .await
4330 .unwrap()
4331 .unwrap();
4332 let frag_reuse_details3 = load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta3)
4333 .await
4334 .unwrap();
4335 let new_frags3 = frag_reuse_details3.versions.last().unwrap().new_frag_ids();
4336
4337 remapping::remap_column_index(&mut dataset_clone, &["i"], Some("scalar".into()))
4342 .await
4343 .unwrap();
4344 cleanup_frag_reuse_index(&mut dataset_clone).await.unwrap();
4345
4346 dataset.checkout_latest().await.unwrap();
4348 let Some(frag_reuse_index_meta) = dataset
4349 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
4350 .await
4351 .unwrap()
4352 else {
4353 panic!("Fragment reuse index must be available");
4354 };
4355 let frag_reuse_details = load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta)
4356 .await
4357 .unwrap();
4358 assert_eq!(frag_reuse_details.versions.len(), 2);
4359 assert_eq!(
4360 frag_reuse_details.versions[0].old_frag_ids(),
4361 rewritten_frags2
4362 );
4363 assert_eq!(frag_reuse_details.versions[0].new_frag_ids(), new_frags2);
4364 assert_eq!(
4365 frag_reuse_details.versions[1].old_frag_ids(),
4366 rewritten_frags3
4367 );
4368 assert_eq!(frag_reuse_details.versions[1].new_frag_ids(), new_frags3);
4369 }
4370
4371 #[tokio::test]
4372 async fn test_concurrent_cleanup_and_compaction_rebase_compaction() {
4373 let mut dataset = lance_datagen::gen_batch()
4374 .col(
4375 "vec",
4376 lance_datagen::array::rand_vec::<Float32Type>(Dimension::from(128)),
4377 )
4378 .col("i", lance_datagen::array::step::<Int32Type>())
4379 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
4380 .await
4381 .unwrap();
4382
4383 create_scalar_index(&mut dataset, "i", false).await;
4385
4386 let options = CompactionOptions {
4387 target_rows_per_fragment: 2_000,
4388 defer_index_remap: true,
4389 ..Default::default()
4390 };
4391
4392 let plan = plan_compaction(&dataset, &options).await.unwrap();
4393 let tasks = plan.tasks();
4394
4395 let rewrite_result = rewrite_files(Cow::Borrowed(&dataset), tasks[0].clone(), &options)
4397 .await
4398 .unwrap();
4399
4400 commit_compaction(
4401 &mut dataset,
4402 Vec::from([rewrite_result]),
4403 Arc::new(DatasetIndexRemapperOptions::default()),
4404 &options,
4405 )
4406 .await
4407 .unwrap();
4408
4409 let mut dataset_clone = dataset.clone();
4410
4411 let Some(frag_reuse_index_meta) = dataset
4413 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
4414 .await
4415 .unwrap()
4416 else {
4417 panic!("Fragment reuse index must be available");
4418 };
4419 let frag_reuse_details = load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta)
4420 .await
4421 .unwrap();
4422 assert_eq!(frag_reuse_details.versions.len(), 1);
4423
4424 remapping::remap_column_index(&mut dataset, &["i"], Some("scalar".into()))
4428 .await
4429 .unwrap();
4430 cleanup_frag_reuse_index(&mut dataset).await.unwrap();
4431
4432 dataset.checkout_latest().await.unwrap();
4434 let Some(frag_reuse_index_meta) = dataset
4435 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
4436 .await
4437 .unwrap()
4438 else {
4439 panic!("Fragment reuse index must be available");
4440 };
4441 let frag_reuse_details = load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta)
4442 .await
4443 .unwrap();
4444 assert_eq!(frag_reuse_details.versions.len(), 0);
4445
4446 let rewrite_result2 =
4449 rewrite_files(Cow::Borrowed(&dataset_clone), tasks[1].clone(), &options)
4450 .await
4451 .unwrap();
4452 let rewritten_frags2 = rewrite_result2
4453 .original_fragments
4454 .iter()
4455 .map(|f| f.id)
4456 .collect::<Vec<_>>();
4457 commit_compaction(
4458 &mut dataset_clone,
4459 Vec::from([rewrite_result2]),
4460 Arc::new(DatasetIndexRemapperOptions::default()),
4461 &options,
4462 )
4463 .await
4464 .unwrap();
4465
4466 dataset.checkout_latest().await.unwrap();
4468 let Some(frag_reuse_index_meta) = dataset
4469 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
4470 .await
4471 .unwrap()
4472 else {
4473 panic!("Fragment reuse index must be available");
4474 };
4475 let frag_reuse_details = load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta)
4476 .await
4477 .unwrap();
4478 assert_eq!(frag_reuse_details.versions.len(), 1);
4479 assert_eq!(
4480 frag_reuse_details.versions[0].old_frag_ids(),
4481 rewritten_frags2
4482 );
4483 let new_frags2 = frag_reuse_details.versions[0].new_frag_ids();
4485 assert!(new_frags2.iter().all(|id| *id != 0));
4486 }
4487
4488 #[tokio::test]
4489 async fn test_concurrent_compactions_with_defer_index_remap() {
4490 let mut dataset = lance_datagen::gen_batch()
4491 .col(
4492 "vec",
4493 lance_datagen::array::rand_vec::<Float32Type>(Dimension::from(128)),
4494 )
4495 .col("i", lance_datagen::array::step::<Int32Type>())
4496 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
4497 .await
4498 .unwrap();
4499
4500 create_scalar_index(&mut dataset, "i", false).await;
4502
4503 let options = CompactionOptions {
4504 target_rows_per_fragment: 2_000,
4505 defer_index_remap: true,
4506 ..Default::default()
4507 };
4508
4509 let plan = plan_compaction(&dataset, &options).await.unwrap();
4510 let tasks = plan.tasks();
4511
4512 let mut dataset_clone = dataset.clone();
4513
4514 let rewrite_result = rewrite_files(Cow::Borrowed(&dataset), tasks[0].clone(), &options)
4516 .await
4517 .unwrap();
4518
4519 commit_compaction(
4520 &mut dataset,
4521 Vec::from([rewrite_result]),
4522 Arc::new(DatasetIndexRemapperOptions::default()),
4523 &options,
4524 )
4525 .await
4526 .unwrap();
4527
4528 let Some(frag_reuse_index_meta) = dataset
4530 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
4531 .await
4532 .unwrap()
4533 else {
4534 panic!("Fragment reuse index must be available");
4535 };
4536 let frag_reuse_details = load_frag_reuse_index_details(&dataset, &frag_reuse_index_meta)
4537 .await
4538 .unwrap();
4539 assert_eq!(frag_reuse_details.versions.len(), 1);
4540
4541 let rewrite_result2 =
4543 rewrite_files(Cow::Borrowed(&dataset_clone), tasks[1].clone(), &options)
4544 .await
4545 .unwrap();
4546 let result = commit_compaction(
4547 &mut dataset_clone,
4548 Vec::from([rewrite_result2]),
4549 Arc::new(DatasetIndexRemapperOptions::default()),
4550 &options,
4551 )
4552 .await;
4553 assert!(matches!(result, Err(Error::RetryableCommitConflict { .. })));
4554 }
4555
4556 #[tokio::test]
4557 async fn test_read_bitmap_index_with_defer_index_remap() {
4558 let mut dataset = lance_datagen::gen_batch()
4560 .col(
4561 "vec",
4562 lance_datagen::array::rand_vec::<Float32Type>(Dimension::from(128)),
4563 )
4564 .col(
4565 "category",
4566 lance_datagen::array::cycle::<Int32Type>(vec![1, 2, 3]),
4567 )
4568 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
4569 .await
4570 .unwrap();
4571
4572 let count1 = dataset
4574 .count_rows(Some("category = 1".to_owned()))
4575 .await
4576 .unwrap();
4577 let count2 = dataset
4578 .count_rows(Some("category = 2".to_owned()))
4579 .await
4580 .unwrap();
4581 let count3 = dataset
4582 .count_rows(Some("category = 3".to_owned()))
4583 .await
4584 .unwrap();
4585
4586 let index_name = Some("category_idx".into());
4588 dataset
4589 .create_index(
4590 &["category"],
4591 IndexType::Bitmap,
4592 index_name.clone(),
4593 &ScalarIndexParams::default(),
4594 false,
4595 )
4596 .await
4597 .unwrap();
4598 let indices = dataset.load_indices().await.unwrap();
4599 let original_index = indices
4600 .iter()
4601 .find(|idx| idx.name == "category_idx")
4602 .unwrap();
4603
4604 let options = CompactionOptions {
4606 target_rows_per_fragment: 2_000,
4607 defer_index_remap: true,
4608 ..Default::default()
4609 };
4610
4611 let metrics = compact_files(&mut dataset, options, None).await.unwrap();
4612 assert!(metrics.fragments_removed > 0);
4613 assert!(metrics.fragments_added > 0);
4614
4615 let Some(current_index) = dataset.load_index_by_name("category_idx").await.unwrap() else {
4617 panic!("category index must be available");
4618 };
4619 assert_eq!(current_index.uuid, original_index.uuid);
4620
4621 assert_eq!(
4623 dataset
4624 .count_rows(Some("category = 1".to_owned()))
4625 .await
4626 .unwrap(),
4627 count1
4628 );
4629 assert_eq!(
4630 dataset
4631 .count_rows(Some("category = 2".to_owned()))
4632 .await
4633 .unwrap(),
4634 count2
4635 );
4636 assert_eq!(
4637 dataset
4638 .count_rows(Some("category = 3".to_owned()))
4639 .await
4640 .unwrap(),
4641 count3
4642 );
4643
4644 let mut scanner = dataset.scan();
4646 scanner.filter("category = 1").unwrap();
4647 scanner.project::<String>(&[]).unwrap().with_row_id();
4648 let plan = scanner.explain_plan(false).await.unwrap();
4649 assert!(
4650 plan.contains("ScalarIndexQuery: query=[category = 1]@category_idx(Bitmap)"),
4651 "Expected index query in plan: {}",
4652 plan
4653 );
4654 }
4655
4656 #[tokio::test]
4657 async fn test_read_btree_index_with_defer_index_remap() {
4658 let mut dataset = lance_datagen::gen_batch()
4660 .col(
4661 "vec",
4662 lance_datagen::array::rand_vec::<Float32Type>(Dimension::from(128)),
4663 )
4664 .col("id", lance_datagen::array::step::<Int32Type>())
4665 .into_ram_dataset(FragmentCount::from(110), FragmentRowCount::from(1000))
4666 .await
4667 .unwrap();
4668
4669 let count_low = dataset
4671 .count_rows(Some("id < 1000".to_owned()))
4672 .await
4673 .unwrap();
4674 let count_mid = dataset
4675 .count_rows(Some("id >= 2000 and id < 3000".to_owned()))
4676 .await
4677 .unwrap();
4678 let count_high = dataset
4679 .count_rows(Some("id >= 5000".to_owned()))
4680 .await
4681 .unwrap();
4682
4683 let index_name = Some("id_idx".into());
4685 dataset
4686 .create_index(
4687 &["id"],
4688 IndexType::BTree,
4689 index_name.clone(),
4690 &ScalarIndexParams::default(),
4691 false,
4692 )
4693 .await
4694 .unwrap();
4695 let indices = dataset.load_indices().await.unwrap();
4696 let original_index = indices.iter().find(|idx| idx.name == "id_idx").unwrap();
4697
4698 let options = CompactionOptions {
4700 target_rows_per_fragment: 50_000,
4701 defer_index_remap: true,
4702 ..Default::default()
4703 };
4704
4705 let metrics = compact_files(&mut dataset, options, None).await.unwrap();
4706 assert!(metrics.fragments_removed > 0);
4707 assert!(metrics.fragments_added > 0);
4708
4709 let Some(current_index) = dataset.load_index_by_name("id_idx").await.unwrap() else {
4711 panic!("id index must be available");
4712 };
4713 assert_eq!(current_index.uuid, original_index.uuid);
4714
4715 assert_eq!(
4717 dataset
4718 .count_rows(Some("id < 1000".to_owned()))
4719 .await
4720 .unwrap(),
4721 count_low
4722 );
4723 assert_eq!(
4724 dataset
4725 .count_rows(Some("id >= 2000 and id < 3000".to_owned()))
4726 .await
4727 .unwrap(),
4728 count_mid
4729 );
4730 assert_eq!(
4731 dataset
4732 .count_rows(Some("id >= 5000".to_owned()))
4733 .await
4734 .unwrap(),
4735 count_high
4736 );
4737
4738 let mut scanner = dataset.scan();
4740 scanner.filter("id >= 2000 and id < 3000").unwrap();
4741 scanner.project::<String>(&[]).unwrap().with_row_id();
4742 let plan = scanner.explain_plan(false).await.unwrap();
4743 assert!(
4744 plan.contains("ScalarIndexQuery: query=[id >= 2000 && id < 3000]@id_idx(BTree)"),
4745 "Expected scalar index query in plan: {}",
4746 plan
4747 );
4748 }
4749
4750 #[rstest]
4751 #[case(IndexRemapMode::Compact)]
4752 #[case(IndexRemapMode::Direct)]
4753 #[tokio::test]
4754 async fn test_btree_index_remap_after_compaction(#[case] index_remap_mode: IndexRemapMode) {
4755 let mut dataset = lance_datagen::gen_batch()
4756 .col(
4757 "vec",
4758 lance_datagen::array::rand_vec::<Float32Type>(Dimension::from(32)),
4759 )
4760 .col("id", lance_datagen::array::step::<Int32Type>())
4761 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
4762 .await
4763 .unwrap();
4764
4765 dataset.delete("id % 10 == 0").await.unwrap();
4768
4769 dataset
4770 .create_index(
4771 &["id"],
4772 IndexType::BTree,
4773 Some("id_idx".into()),
4774 &ScalarIndexParams::default(),
4775 false,
4776 )
4777 .await
4778 .unwrap();
4779
4780 let count_low = dataset
4781 .count_rows(Some("id < 1000".to_owned()))
4782 .await
4783 .unwrap();
4784 let count_mid = dataset
4785 .count_rows(Some("id >= 2000 and id < 3000".to_owned()))
4786 .await
4787 .unwrap();
4788 let count_high = dataset
4789 .count_rows(Some("id >= 5000".to_owned()))
4790 .await
4791 .unwrap();
4792
4793 let options = CompactionOptions {
4794 target_rows_per_fragment: 50_000,
4795 index_remap_mode,
4796 ..Default::default()
4797 };
4798 let metrics = compact_files(&mut dataset, options, None).await.unwrap();
4799 assert!(metrics.fragments_removed > 0);
4800 assert!(metrics.fragments_added > 0);
4801
4802 let mut scanner = dataset.scan();
4804 scanner.filter("id >= 2000 and id < 3000").unwrap();
4805 scanner.project::<String>(&[]).unwrap().with_row_id();
4806 let plan = scanner.explain_plan(false).await.unwrap();
4807 assert!(
4808 plan.contains("ScalarIndexQuery: query=[id >= 2000 && id < 3000]@id_idx(BTree)"),
4809 "Expected scalar index query in plan: {}",
4810 plan
4811 );
4812
4813 assert_eq!(
4816 dataset
4817 .count_rows(Some("id < 1000".to_owned()))
4818 .await
4819 .unwrap(),
4820 count_low
4821 );
4822 assert_eq!(
4823 dataset
4824 .count_rows(Some("id >= 2000 and id < 3000".to_owned()))
4825 .await
4826 .unwrap(),
4827 count_mid
4828 );
4829 assert_eq!(
4830 dataset
4831 .count_rows(Some("id >= 5000".to_owned()))
4832 .await
4833 .unwrap(),
4834 count_high
4835 );
4836 }
4837
4838 #[rstest]
4839 #[case(IndexRemapMode::Compact)]
4840 #[case(IndexRemapMode::Direct)]
4841 #[tokio::test]
4842 async fn test_ivf_pq_index_remap_after_compaction(#[case] index_remap_mode: IndexRemapMode) {
4843 use arrow_array::cast::AsArray;
4844 use lance_index::vector::pq::PQBuildParams;
4845
4846 const DIM: u32 = 32;
4847 let mut dataset = lance_datagen::gen_batch()
4848 .col("id", lance_datagen::array::step::<Int32Type>())
4849 .col(
4850 "vec",
4851 lance_datagen::array::rand_vec::<Float32Type>(Dimension::from(DIM)),
4852 )
4853 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
4854 .await
4855 .unwrap();
4856
4857 let params = VectorIndexParams::with_ivf_pq_params(
4858 DistanceType::L2,
4859 small_ivf(),
4860 PQBuildParams {
4861 max_iters: 2,
4862 num_sub_vectors: 2,
4863 ..Default::default()
4864 },
4865 );
4866 dataset
4867 .create_index(
4868 &["vec"],
4869 IndexType::Vector,
4870 Some("vec_idx".into()),
4871 ¶ms,
4872 false,
4873 )
4874 .await
4875 .unwrap();
4876 let original_uuid = dataset
4877 .load_index_by_name("vec_idx")
4878 .await
4879 .unwrap()
4880 .unwrap()
4881 .uuid;
4882
4883 dataset.delete("id % 10 == 0").await.unwrap();
4886
4887 let mut survivors: Vec<(i32, Vec<f32>)> = Vec::new();
4890 {
4891 let mut scanner = dataset.scan();
4892 scanner.project(&["id", "vec"]).unwrap();
4893 let batches = scanner
4894 .try_into_stream()
4895 .await
4896 .unwrap()
4897 .try_collect::<Vec<_>>()
4898 .await
4899 .unwrap();
4900 for batch in &batches {
4901 let ids = batch["id"].as_primitive::<Int32Type>();
4902 let vecs = batch["vec"].as_fixed_size_list();
4903 for i in 0..batch.num_rows() {
4904 let v = vecs.value(i);
4905 survivors.push((
4906 ids.value(i),
4907 v.as_primitive::<Float32Type>().values().to_vec(),
4908 ));
4909 }
4910 }
4911 }
4912 let surviving_ids: std::collections::HashSet<i32> =
4913 survivors.iter().map(|(id, _)| *id).collect();
4914 let step = (survivors.len() / 16).max(1);
4915 let queries: Vec<Vec<f32>> = survivors
4916 .iter()
4917 .step_by(step)
4918 .map(|(_, v)| v.clone())
4919 .collect();
4920 let k = 10;
4921 let mut baseline: Vec<Vec<i32>> = Vec::new();
4922 for q in &queries {
4923 baseline.push(vector_knn_ids(&dataset, q, k).await);
4924 }
4925
4926 let metrics = compact_files(
4929 &mut dataset,
4930 CompactionOptions {
4931 target_rows_per_fragment: 50_000,
4932 index_remap_mode,
4933 ..Default::default()
4934 },
4935 None,
4936 )
4937 .await
4938 .unwrap();
4939 assert!(metrics.fragments_removed > 0);
4940 assert!(metrics.fragments_added > 0);
4941
4942 assert_ne!(
4944 dataset
4945 .load_index_by_name("vec_idx")
4946 .await
4947 .unwrap()
4948 .unwrap()
4949 .uuid,
4950 original_uuid,
4951 "vector index must be physically remapped inline"
4952 );
4953
4954 for (i, q) in queries.iter().enumerate() {
4958 let after = vector_knn_ids(&dataset, q, k).await;
4959 for id in &after {
4960 assert!(
4961 surviving_ids.contains(id),
4962 "KNN returned id {id} that is not a surviving row (query #{i}, mode {index_remap_mode:?})"
4963 );
4964 }
4965 let overlap = after.iter().filter(|id| baseline[i].contains(id)).count();
4966 assert!(
4967 overlap >= 8,
4968 "KNN top-{k} diverged after compaction: overlap {overlap} < 8 (query #{i}, mode {index_remap_mode:?})"
4969 );
4970 }
4971 }
4972
4973 #[rstest]
4974 #[case(IndexRemapMode::Compact)]
4975 #[case(IndexRemapMode::Direct)]
4976 #[tokio::test]
4977 async fn test_inverted_index_remap_after_compaction(#[case] index_remap_mode: IndexRemapMode) {
4978 use arrow_array::cast::AsArray;
4979
4980 let mut dataset = lance_datagen::gen_batch()
4981 .col("id", lance_datagen::array::step::<Int32Type>())
4982 .col("doc", lance_datagen::array::random_sentence(1, 100, false))
4983 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
4984 .await
4985 .unwrap();
4986
4987 dataset
4988 .create_index(
4989 &["doc"],
4990 IndexType::Inverted,
4991 Some("doc_idx".into()),
4992 &InvertedIndexParams::default(),
4993 false,
4994 )
4995 .await
4996 .unwrap();
4997 let original_uuid = dataset
4998 .load_index_by_name("doc_idx")
4999 .await
5000 .unwrap()
5001 .unwrap()
5002 .uuid;
5003
5004 let words: Vec<String> = {
5006 let mut scanner = dataset.scan();
5007 scanner
5008 .project(&["doc"])
5009 .unwrap()
5010 .limit(Some(1), None)
5011 .unwrap();
5012 let batches = scanner
5013 .try_into_stream()
5014 .await
5015 .unwrap()
5016 .try_collect::<Vec<_>>()
5017 .await
5018 .unwrap();
5019 let mut words: Vec<String> = batches[0]["doc"]
5020 .as_string::<i32>()
5021 .value(0)
5022 .split_whitespace()
5023 .map(|s| s.to_string())
5024 .collect();
5025 words.sort();
5026 words.dedup();
5027 words.truncate(3);
5028 words
5029 };
5030 assert!(!words.is_empty(), "sampled document must contain words");
5031
5032 dataset.delete("id % 10 == 0").await.unwrap();
5035
5036 let mut before = Vec::new();
5039 for word in &words {
5040 let mut scanner = dataset.scan();
5041 scanner
5042 .full_text_search(FullTextSearchQuery::new(word.clone()))
5043 .unwrap();
5044 scanner.project::<String>(&[]).unwrap().with_row_id();
5045 before.push(scanner.count_rows().await.unwrap());
5046 }
5047
5048 let options = CompactionOptions {
5051 target_rows_per_fragment: 50_000,
5052 index_remap_mode,
5053 ..Default::default()
5054 };
5055 let metrics = compact_files(&mut dataset, options, None).await.unwrap();
5056 assert!(metrics.fragments_removed > 0);
5057 assert!(metrics.fragments_added > 0);
5058
5059 assert_ne!(
5061 dataset
5062 .load_index_by_name("doc_idx")
5063 .await
5064 .unwrap()
5065 .unwrap()
5066 .uuid,
5067 original_uuid,
5068 "inverted index must be physically remapped inline (mode {index_remap_mode:?})"
5069 );
5070
5071 let mut scanner = dataset.scan();
5073 scanner
5074 .full_text_search(FullTextSearchQuery::new(words[0].clone()))
5075 .unwrap();
5076 scanner.project::<String>(&[]).unwrap().with_row_id();
5077 let plan = scanner.explain_plan(true).await.unwrap();
5078 assert!(
5079 plan.contains("MatchQuery"),
5080 "Expected inverted index scan in plan: {}",
5081 plan
5082 );
5083
5084 for (word, expected) in words.iter().zip(before) {
5087 let mut scanner = dataset.scan();
5088 scanner
5089 .full_text_search(FullTextSearchQuery::new(word.clone()))
5090 .unwrap();
5091 scanner.project::<String>(&[]).unwrap().with_row_id();
5092 assert_eq!(
5093 scanner.count_rows().await.unwrap(),
5094 expected,
5095 "full-text count for {word:?} changed after compaction (mode {index_remap_mode:?})"
5096 );
5097 }
5098 }
5099
5100 #[tokio::test]
5101 async fn test_read_inverted_index_with_defer_index_remap() {
5102 let mut words_gen = lance_datagen::array::random_sentence(1, 100, true);
5104 let doc_col = words_gen
5105 .generate_default(lance_datagen::RowCount::from(6000))
5106 .unwrap();
5107
5108 let batch = RecordBatch::try_new(
5109 Schema::new(vec![Field::new("doc", DataType::LargeUtf8, false)]).into(),
5110 vec![doc_col.clone()],
5111 )
5112 .unwrap();
5113 let schema_ref = batch.schema();
5114 let stream = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema_ref);
5115 let mut dataset = Dataset::write(
5116 stream,
5117 "memory://test/table",
5118 Some(WriteParams {
5119 max_rows_per_file: 1_000, ..Default::default()
5121 }),
5122 )
5123 .await
5124 .unwrap();
5125
5126 let large_string_array = doc_col.as_any().downcast_ref::<LargeStringArray>().unwrap();
5129 let sample_words: Vec<String> = large_string_array
5130 .value(0)
5131 .split_whitespace()
5132 .take(10)
5133 .map(|s| s.to_string())
5134 .collect();
5135 let test_word1 = &sample_words[0];
5136 let test_word2 = &sample_words[1];
5137 let test_word3 = &sample_words[2];
5138
5139 let index_name = Some("doc_idx".into());
5141 dataset
5142 .create_index(
5143 &["doc"],
5144 IndexType::Inverted,
5145 index_name.clone(),
5146 &InvertedIndexParams::default(),
5147 false,
5148 )
5149 .await
5150 .unwrap();
5151 let indices = dataset.load_indices().await.unwrap();
5152 let original_index = indices.iter().find(|idx| idx.name == "doc_idx").unwrap();
5153
5154 let options = CompactionOptions {
5156 target_rows_per_fragment: 2_000,
5157 defer_index_remap: true,
5158 ..Default::default()
5159 };
5160
5161 let metrics = compact_files(&mut dataset, options, None).await.unwrap();
5162 assert!(metrics.fragments_removed > 0);
5163 assert!(metrics.fragments_added > 0);
5164
5165 let Some(current_index) = dataset.load_index_by_name("doc_idx").await.unwrap() else {
5167 panic!("doc index must be available");
5168 };
5169 assert_eq!(current_index.uuid, original_index.uuid);
5170
5171 let mut scanner = dataset.scan();
5173 scanner
5174 .full_text_search(FullTextSearchQuery::new(test_word1.clone()))
5175 .unwrap();
5176 scanner.project::<String>(&[]).unwrap().with_row_id();
5177 let count1 = scanner.count_rows().await.unwrap();
5178 scanner = dataset.scan();
5179 scanner
5180 .full_text_search(FullTextSearchQuery::new(test_word2.clone()))
5181 .unwrap();
5182 scanner.project::<String>(&[]).unwrap().with_row_id();
5183 let count2 = scanner.count_rows().await.unwrap();
5184 scanner = dataset.scan();
5185 scanner
5186 .full_text_search(FullTextSearchQuery::new(test_word3.clone()))
5187 .unwrap();
5188 scanner.project::<String>(&[]).unwrap().with_row_id();
5189 let count3 = scanner.count_rows().await.unwrap();
5190
5191 let mut scanner = dataset.scan();
5193 scanner
5194 .full_text_search(FullTextSearchQuery::new(test_word1.clone()))
5195 .unwrap();
5196 scanner.project::<String>(&[]).unwrap().with_row_id();
5197 let plan = scanner.explain_plan(true).await.unwrap();
5198 assert!(
5199 plan.contains("MatchQuery"),
5200 "Expected inverted index scan in plan: {}",
5201 plan
5202 );
5203 assert!(
5204 !plan.contains("LanceScan"),
5205 "Expected no fragment scan in plan: {}",
5206 plan
5207 );
5208
5209 dataset
5211 .create_index(
5212 &["doc"],
5213 IndexType::Inverted,
5214 index_name.clone(),
5215 &InvertedIndexParams::default(),
5216 true,
5217 )
5218 .await
5219 .unwrap();
5220
5221 let mut scanner = dataset.scan();
5223 scanner
5224 .full_text_search(FullTextSearchQuery::new(test_word1.clone()))
5225 .unwrap();
5226 scanner.project::<String>(&[]).unwrap().with_row_id();
5227 assert_eq!(scanner.count_rows().await.unwrap(), count1);
5228 scanner = dataset.scan();
5229 scanner
5230 .full_text_search(FullTextSearchQuery::new(test_word2.clone()))
5231 .unwrap();
5232 scanner.project::<String>(&[]).unwrap().with_row_id();
5233 assert_eq!(scanner.count_rows().await.unwrap(), count2);
5234 scanner = dataset.scan();
5235 scanner
5236 .full_text_search(FullTextSearchQuery::new(test_word3.clone()))
5237 .unwrap();
5238 scanner.project::<String>(&[]).unwrap().with_row_id();
5239 assert_eq!(scanner.count_rows().await.unwrap(), count3);
5240 }
5241
5242 #[tokio::test]
5250 async fn test_read_inverted_index_with_defer_index_remap_and_deletions() {
5251 const ROWS: i32 = 1200;
5255 const DELETED: i32 = 400;
5256
5257 let ids = Int32Array::from_iter_values(0..ROWS);
5260 let docs = LargeStringArray::from_iter_values((0..ROWS).map(|_| "lance apple orange"));
5261 let batch = RecordBatch::try_new(
5262 Schema::new(vec![
5263 Field::new("id", DataType::Int32, false),
5264 Field::new("doc", DataType::LargeUtf8, false),
5265 ])
5266 .into(),
5267 vec![Arc::new(ids) as ArrayRef, Arc::new(docs) as ArrayRef],
5268 )
5269 .unwrap();
5270 let schema_ref = batch.schema();
5271 let stream = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema_ref);
5272 let mut dataset = Dataset::write(
5273 stream,
5274 "memory://test/table",
5275 Some(WriteParams {
5276 max_rows_per_file: 200, ..Default::default()
5278 }),
5279 )
5280 .await
5281 .unwrap();
5282
5283 dataset
5284 .create_index(
5285 &["doc"],
5286 IndexType::Inverted,
5287 Some("doc_idx".into()),
5288 &InvertedIndexParams::default(),
5289 false,
5290 )
5291 .await
5292 .unwrap();
5293
5294 dataset.delete(&format!("id < {DELETED}")).await.unwrap();
5297 compact_files(
5298 &mut dataset,
5299 CompactionOptions {
5300 target_rows_per_fragment: 2_000,
5301 defer_index_remap: true,
5302 ..Default::default()
5303 },
5304 None,
5305 )
5306 .await
5307 .unwrap();
5308 assert!(
5309 dataset
5310 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
5311 .await
5312 .unwrap()
5313 .is_some(),
5314 "deferred compaction must leave a fragment-reuse index"
5315 );
5316
5317 async fn search_ids(dataset: &Dataset) -> Vec<i32> {
5320 let mut scanner = dataset.scan();
5321 scanner
5322 .full_text_search(FullTextSearchQuery::new("lance".to_owned()))
5323 .unwrap();
5324 scanner.project::<&str>(&["id"]).unwrap();
5325 let batches = scanner
5326 .try_into_stream()
5327 .await
5328 .unwrap()
5329 .try_collect::<Vec<_>>()
5330 .await
5331 .unwrap();
5332 let mut ids: Vec<i32> = batches
5333 .iter()
5334 .flat_map(|b| {
5335 b.column_by_name("id")
5336 .unwrap()
5337 .as_any()
5338 .downcast_ref::<Int32Array>()
5339 .unwrap()
5340 .values()
5341 .to_vec()
5342 })
5343 .collect();
5344 ids.sort_unstable();
5345 ids
5346 }
5347
5348 let expected = (DELETED..ROWS).collect::<Vec<_>>();
5349
5350 let during = search_ids(&dataset).await;
5352 assert_eq!(
5353 during, expected,
5354 "FRI-window FTS must return exactly the surviving rows (no resurrection, no loss, no stale rows)"
5355 );
5356
5357 remapping::remap_column_index(&mut dataset, &["doc"], Some("doc_idx".into()))
5359 .await
5360 .unwrap();
5361 cleanup_frag_reuse_index(&mut dataset).await.unwrap();
5362 let after = search_ids(&dataset).await;
5363 assert_eq!(
5364 after, expected,
5365 "FTS must stay correct after physical remap + fragment-reuse trim"
5366 );
5367 }
5368
5369 #[tokio::test]
5370 async fn test_read_ngram_index_with_defer_index_remap() {
5371 let mut words_gen = lance_datagen::array::random_sentence(1, 100, true);
5373 let doc_col = words_gen
5374 .generate_default(lance_datagen::RowCount::from(6000))
5375 .unwrap();
5376
5377 let batch = RecordBatch::try_new(
5378 Schema::new(vec![Field::new("doc", DataType::LargeUtf8, false)]).into(),
5379 vec![doc_col.clone()],
5380 )
5381 .unwrap();
5382 let schema_ref = batch.schema();
5383 let stream = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema_ref);
5384 let mut dataset = Dataset::write(
5385 stream,
5386 "memory://test/table",
5387 Some(WriteParams {
5388 max_rows_per_file: 1_000, ..Default::default()
5390 }),
5391 )
5392 .await
5393 .unwrap();
5394
5395 let large_string_array = doc_col.as_any().downcast_ref::<LargeStringArray>().unwrap();
5398 let sample_words: Vec<String> = large_string_array
5399 .value(0)
5400 .split_whitespace()
5401 .take(10)
5402 .map(|s| s.to_string())
5403 .collect();
5404 let test_word1 = &sample_words[0];
5405 let test_word2 = &sample_words[1];
5406 let test_word3 = &sample_words[2];
5407
5408 let index_name = Some("doc_idx".into());
5410 dataset
5411 .create_index(
5412 &["doc"],
5413 IndexType::NGram,
5414 index_name.clone(),
5415 &ScalarIndexParams::default(),
5416 false,
5417 )
5418 .await
5419 .unwrap();
5420 let indices = dataset.load_indices().await.unwrap();
5421 let original_index = indices.iter().find(|idx| idx.name == "doc_idx").unwrap();
5422
5423 let count1 = dataset
5425 .count_rows(Some(format!("contains(doc, '{}')", test_word1)))
5426 .await
5427 .unwrap();
5428 let count2 = dataset
5429 .count_rows(Some(format!("contains(doc, '{}')", test_word2)))
5430 .await
5431 .unwrap();
5432 let count3 = dataset
5433 .count_rows(Some(format!("contains(doc, '{}')", test_word3)))
5434 .await
5435 .unwrap();
5436
5437 let options = CompactionOptions {
5439 target_rows_per_fragment: 2_000,
5440 defer_index_remap: true,
5441 ..Default::default()
5442 };
5443
5444 let metrics = compact_files(&mut dataset, options, None).await.unwrap();
5445 assert!(metrics.fragments_removed > 0);
5446 assert!(metrics.fragments_added > 0);
5447
5448 let Some(current_index) = dataset.load_index_by_name("doc_idx").await.unwrap() else {
5450 panic!("doc index must be available");
5451 };
5452 assert_eq!(current_index.uuid, original_index.uuid);
5453
5454 assert_eq!(
5456 dataset
5457 .count_rows(Some(format!("contains(doc, '{}')", test_word1)))
5458 .await
5459 .unwrap(),
5460 count1
5461 );
5462 assert_eq!(
5463 dataset
5464 .count_rows(Some(format!("contains(doc, '{}')", test_word2)))
5465 .await
5466 .unwrap(),
5467 count2
5468 );
5469 assert_eq!(
5470 dataset
5471 .count_rows(Some(format!("contains(doc, '{}')", test_word3)))
5472 .await
5473 .unwrap(),
5474 count3
5475 );
5476
5477 let mut scanner = dataset.scan();
5479 scanner
5480 .filter(&format!("contains(doc, '{}')", test_word1))
5481 .unwrap();
5482 scanner.project::<String>(&[]).unwrap().with_row_id();
5483 let plan = scanner.explain_plan(false).await.unwrap();
5484 assert!(
5485 plan.contains("ScalarIndexQuery: query=[contains(doc, Utf8"),
5486 "Expected scalar index query in plan: {}",
5487 plan
5488 );
5489 }
5490
5491 #[tokio::test]
5492 async fn test_read_label_list_index_with_defer_index_remap() {
5493 let mut dataset = lance_datagen::gen_batch()
5495 .col(
5496 "vec",
5497 lance_datagen::array::rand_vec::<Float32Type>(Dimension::from(128)),
5498 )
5499 .col(
5500 "labels",
5501 lance_datagen::array::rand_list_any(
5502 lance_datagen::array::cycle::<Int64Type>(vec![1, 2, 3, 4, 5]),
5503 false,
5504 ),
5505 )
5506 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
5507 .await
5508 .unwrap();
5509
5510 let count1 = dataset
5512 .count_rows(Some("array_has_any(labels, [1])".to_owned()))
5513 .await
5514 .unwrap();
5515 let count2 = dataset
5516 .count_rows(Some("array_has_any(labels, [5])".to_owned()))
5517 .await
5518 .unwrap();
5519 let count3 = dataset
5520 .count_rows(Some("array_has_any(labels, [10])".to_owned()))
5521 .await
5522 .unwrap();
5523
5524 let index_name = Some("labels_idx".into());
5526 dataset
5527 .create_index(
5528 &["labels"],
5529 IndexType::LabelList,
5530 index_name.clone(),
5531 &ScalarIndexParams::default(),
5532 false,
5533 )
5534 .await
5535 .unwrap();
5536 let indices = dataset.load_indices().await.unwrap();
5537 let original_index = indices.iter().find(|idx| idx.name == "labels_idx").unwrap();
5538
5539 let options = CompactionOptions {
5541 target_rows_per_fragment: 2000,
5542 defer_index_remap: true,
5543 ..Default::default()
5544 };
5545 let metrics = compact_files(&mut dataset, options, None).await.unwrap();
5546 assert!(metrics.fragments_removed > 0);
5547 assert!(metrics.fragments_added > 0);
5548
5549 let indices = dataset.load_indices().await.unwrap();
5551 let current_index = indices.iter().find(|idx| idx.name == "labels_idx").unwrap();
5552 assert_eq!(current_index.uuid, original_index.uuid);
5553
5554 assert_eq!(
5556 dataset
5557 .count_rows(Some("array_has_any(labels, [1])".to_owned()))
5558 .await
5559 .unwrap(),
5560 count1
5561 );
5562 assert_eq!(
5563 dataset
5564 .count_rows(Some("array_has_any(labels, [5])".to_owned()))
5565 .await
5566 .unwrap(),
5567 count2
5568 );
5569 assert_eq!(
5570 dataset
5571 .count_rows(Some("array_has_any(labels, [10])".to_owned()))
5572 .await
5573 .unwrap(),
5574 count3
5575 );
5576
5577 let mut scanner = dataset.scan();
5579 scanner.filter("array_has_any(labels, [1])").unwrap();
5580 scanner.project::<String>(&[]).unwrap().with_row_id();
5581 let plan = scanner.explain_plan(false).await.unwrap();
5582 assert!(
5583 plan.contains(
5584 "ScalarIndexQuery: query=[array_has_any(labels, List([1]))]@labels_idx(LabelList)",
5585 ),
5586 "Expected scalar index query in plan: {}",
5587 plan
5588 );
5589 }
5590
5591 #[tokio::test]
5592 async fn test_read_ivf_pq_index_v3_with_defer_index_remap() {
5593 let mut dataset = lance_datagen::gen_batch()
5595 .col(
5596 "vec",
5597 lance_datagen::array::rand_vec::<Float32Type>(Dimension::from(128)),
5598 )
5599 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
5600 .await
5601 .unwrap();
5602
5603 let query_vec1: PrimitiveArray<Float32Type> =
5605 PrimitiveArray::from_iter_values(std::iter::repeat_n(0.0, 128));
5606 let query_vec2: PrimitiveArray<Float32Type> =
5607 PrimitiveArray::from_iter_values(std::iter::repeat_n(1.1, 128));
5608 let query_vec3: PrimitiveArray<Float32Type> =
5609 PrimitiveArray::from_iter_values(std::iter::repeat_n(2.2, 128));
5610
5611 let mut scanner = dataset.scan();
5613 scanner.nearest("vec", &query_vec1, 10).unwrap();
5614 scanner.project::<String>(&[]).unwrap().with_row_id();
5615 let results1 = scanner
5616 .try_into_stream()
5617 .await
5618 .unwrap()
5619 .try_collect::<Vec<_>>()
5620 .await
5621 .unwrap();
5622 let count1 = results1.len();
5623
5624 scanner = dataset.scan();
5625 scanner.nearest("vec", &query_vec2, 10).unwrap();
5626 scanner.project::<String>(&[]).unwrap().with_row_id();
5627 let results2 = scanner
5628 .try_into_stream()
5629 .await
5630 .unwrap()
5631 .try_collect::<Vec<_>>()
5632 .await
5633 .unwrap();
5634 let count2 = results2.len();
5635
5636 scanner = dataset.scan();
5637 scanner.nearest("vec", &query_vec3, 10).unwrap();
5638 scanner.project::<String>(&[]).unwrap().with_row_id();
5639 let results3 = scanner
5640 .try_into_stream()
5641 .await
5642 .unwrap()
5643 .try_collect::<Vec<_>>()
5644 .await
5645 .unwrap();
5646 let count3 = results3.len();
5647
5648 let index_name = Some("vec_idx".into());
5650 dataset
5651 .create_index(
5652 &["vec"],
5653 IndexType::Vector,
5654 index_name.clone(),
5655 &VectorIndexParams {
5656 metric_type: DistanceType::L2,
5657 stages: vec![
5658 StageParams::Ivf(IvfBuildParams {
5659 max_iters: 2,
5660 num_partitions: Some(2),
5661 sample_rate: 2,
5662 ..Default::default()
5663 }),
5664 StageParams::PQ(PQBuildParams {
5665 max_iters: 2,
5666 num_sub_vectors: 2,
5667 ..Default::default()
5668 }),
5669 ],
5670 version: crate::index::vector::IndexFileVersion::V3,
5671 skip_transpose: false,
5672 runtime_hints: Default::default(),
5673 },
5674 false,
5675 )
5676 .await
5677 .unwrap();
5678 let indices = dataset.load_indices().await.unwrap();
5679 let original_index = indices.iter().find(|idx| idx.name == "vec_idx").unwrap();
5680
5681 let options = CompactionOptions {
5683 target_rows_per_fragment: 2_000,
5684 defer_index_remap: true,
5685 ..Default::default()
5686 };
5687
5688 let metrics = compact_files(&mut dataset, options, None).await.unwrap();
5689 assert!(metrics.fragments_removed > 0);
5690 assert!(metrics.fragments_added > 0);
5691
5692 let Some(current_index) = dataset.load_index_by_name("vec_idx").await.unwrap() else {
5694 panic!("vec index must be available");
5695 };
5696 assert_eq!(current_index.uuid, original_index.uuid);
5697
5698 let mut scanner = dataset.scan();
5700 scanner.nearest("vec", &query_vec1, 10).unwrap();
5701 scanner.project::<String>(&[]).unwrap().with_row_id();
5702 let new_results1 = scanner
5703 .try_into_stream()
5704 .await
5705 .unwrap()
5706 .try_collect::<Vec<_>>()
5707 .await
5708 .unwrap();
5709 assert_eq!(new_results1.len(), count1);
5710
5711 scanner = dataset.scan();
5712 scanner.nearest("vec", &query_vec2, 10).unwrap();
5713 scanner.project::<String>(&[]).unwrap().with_row_id();
5714 let new_results2 = scanner
5715 .try_into_stream()
5716 .await
5717 .unwrap()
5718 .try_collect::<Vec<_>>()
5719 .await
5720 .unwrap();
5721 assert_eq!(new_results2.len(), count2);
5722
5723 scanner = dataset.scan();
5724 scanner.nearest("vec", &query_vec3, 10).unwrap();
5725 scanner.project::<String>(&[]).unwrap().with_row_id();
5726 let new_results3 = scanner
5727 .try_into_stream()
5728 .await
5729 .unwrap()
5730 .try_collect::<Vec<_>>()
5731 .await
5732 .unwrap();
5733 assert_eq!(new_results3.len(), count3);
5734
5735 let mut scanner = dataset.scan();
5737 scanner.nearest("vec", &query_vec1, 10).unwrap();
5738 scanner.project::<String>(&[]).unwrap().with_row_id();
5739 let plan = scanner.explain_plan(false).await.unwrap();
5740 assert!(
5741 plan.contains("ANNSubIndex"),
5742 "Expected vector index scan in plan: {}",
5743 plan
5744 );
5745 assert!(
5746 !plan.contains("LanceScan"),
5747 "Expected no fragment scan in plan: {}",
5748 plan
5749 );
5750 }
5751
5752 #[tokio::test]
5753 async fn test_read_ivf_rq_index_v3_with_defer_index_remap() {
5754 use arrow_array::cast::AsArray;
5755 use lance_index::vector::bq::RQBuildParams;
5756
5757 let mut dataset = lance_datagen::gen_batch()
5758 .col(
5759 "vec",
5760 lance_datagen::array::rand_vec::<Float32Type>(Dimension::from(128)),
5761 )
5762 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
5763 .await
5764 .unwrap();
5765
5766 let stored: Vec<Vec<f32>> = {
5767 let mut scanner = dataset.scan();
5768 scanner.project(&["vec"]).unwrap();
5769 let batches = scanner
5770 .try_into_stream()
5771 .await
5772 .unwrap()
5773 .try_collect::<Vec<_>>()
5774 .await
5775 .unwrap();
5776 let mut out = Vec::new();
5777 for batch in &batches {
5778 let vecs = batch["vec"].as_fixed_size_list();
5779 for i in 0..batch.num_rows() {
5780 let values = vecs.value(i);
5781 let values = values.as_primitive::<Float32Type>();
5782 out.push(values.values().to_vec());
5783 }
5784 }
5785 out
5786 };
5787
5788 let index_name = Some("vec_idx".into());
5789 dataset
5790 .create_index(
5791 &["vec"],
5792 IndexType::Vector,
5793 index_name.clone(),
5794 &VectorIndexParams {
5795 metric_type: DistanceType::L2,
5796 stages: vec![
5797 StageParams::Ivf(IvfBuildParams {
5798 max_iters: 2,
5799 num_partitions: Some(2),
5800 sample_rate: 2,
5801 ..Default::default()
5802 }),
5803 StageParams::RQ(RQBuildParams::new(1)),
5804 ],
5805 version: crate::index::vector::IndexFileVersion::V3,
5806 skip_transpose: false,
5807 runtime_hints: Default::default(),
5808 },
5809 false,
5810 )
5811 .await
5812 .unwrap();
5813 let indices = dataset.load_indices().await.unwrap();
5814 let original_index = indices.iter().find(|idx| idx.name == "vec_idx").unwrap();
5815
5816 let options = CompactionOptions {
5817 target_rows_per_fragment: 2_000,
5818 defer_index_remap: true,
5819 ..Default::default()
5820 };
5821 let metrics = compact_files(&mut dataset, options, None).await.unwrap();
5822 assert!(metrics.fragments_removed > 0);
5823 assert!(metrics.fragments_added > 0);
5824
5825 let Some(current_index) = dataset.load_index_by_name("vec_idx").await.unwrap() else {
5826 panic!("vec index must be available");
5827 };
5828 assert_eq!(current_index.uuid, original_index.uuid);
5829
5830 let frag_reuse_present = dataset
5831 .load_indices()
5832 .await
5833 .unwrap()
5834 .iter()
5835 .any(|idx| idx.name == FRAG_REUSE_INDEX_NAME);
5836 assert!(
5837 frag_reuse_present,
5838 "defer_index_remap must record a {} index",
5839 FRAG_REUSE_INDEX_NAME
5840 );
5841
5842 let sample_step = (stored.len() / 8).max(1);
5843 let mut checked = 0;
5844 for query in stored.iter().step_by(sample_step) {
5845 let query_vec = PrimitiveArray::<Float32Type>::from_iter_values(query.iter().copied());
5846 let mut scanner = dataset.scan();
5847 scanner.nearest("vec", &query_vec, 5).unwrap();
5848 scanner.project(&["vec"]).unwrap().with_row_id();
5849 let batches = scanner
5850 .try_into_stream()
5851 .await
5852 .unwrap()
5853 .try_collect::<Vec<_>>()
5854 .await
5855 .unwrap();
5856 assert!(!batches.is_empty(), "query returned no batches");
5857 let top = &batches[0];
5858 assert!(top.num_rows() > 0, "query returned empty top batch");
5859 let top_vec = top["vec"].as_fixed_size_list().value(0);
5860 let top_vec = top_vec.as_primitive::<Float32Type>();
5861 assert_eq!(
5862 top_vec.values(),
5863 query.as_slice(),
5864 "top-1 self-recall returned a different vector than the query"
5865 );
5866 checked += 1;
5867 }
5868 assert!(checked > 0, "expected to check at least one stored vector");
5869 }
5870
5871 async fn vector_knn_ids(dataset: &Dataset, query: &[f32], k: usize) -> Vec<i32> {
5885 use arrow_array::cast::AsArray;
5886 use arrow_array::types::{Float32Type, Int32Type};
5887 let qa = PrimitiveArray::<Float32Type>::from_iter_values(query.iter().copied());
5888 let mut scanner = dataset.scan();
5889 scanner.nearest("vec", &qa, k).unwrap();
5890 scanner.project(&["id"]).unwrap();
5891 let batches = scanner
5892 .try_into_stream()
5893 .await
5894 .unwrap()
5895 .try_collect::<Vec<_>>()
5896 .await
5897 .unwrap();
5898 let mut ids = Vec::new();
5899 for b in &batches {
5900 ids.extend(b["id"].as_primitive::<Int32Type>().values().iter().copied());
5901 }
5902 ids
5903 }
5904
5905 async fn check_vector_defer_compaction(
5906 params: VectorIndexParams,
5907 delete_predicate: Option<&str>,
5908 k: usize,
5909 min_overlap: usize,
5910 ) {
5911 use arrow_array::cast::AsArray;
5912 use arrow_array::types::{Float32Type, Int32Type};
5913 use lance_datagen::Dimension;
5914
5915 const DIM: u32 = 32;
5916 let mut dataset = lance_datagen::gen_batch()
5917 .col("id", lance_datagen::array::step::<Int32Type>())
5918 .col(
5919 "vec",
5920 lance_datagen::array::rand_vec::<Float32Type>(Dimension::from(DIM)),
5921 )
5922 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
5923 .await
5924 .unwrap();
5925
5926 dataset
5927 .create_index(
5928 &["vec"],
5929 IndexType::Vector,
5930 Some("vec_idx".into()),
5931 ¶ms,
5932 false,
5933 )
5934 .await
5935 .unwrap();
5936 let original_uuid = dataset
5937 .load_index_by_name("vec_idx")
5938 .await
5939 .unwrap()
5940 .unwrap()
5941 .uuid;
5942
5943 if let Some(pred) = delete_predicate {
5944 dataset.delete(pred).await.unwrap();
5945 }
5946
5947 let mut survivors: Vec<(i32, Vec<f32>)> = Vec::new();
5949 {
5950 let mut scanner = dataset.scan();
5951 scanner.project(&["id", "vec"]).unwrap();
5952 let batches = scanner
5953 .try_into_stream()
5954 .await
5955 .unwrap()
5956 .try_collect::<Vec<_>>()
5957 .await
5958 .unwrap();
5959 for batch in &batches {
5960 let ids = batch["id"].as_primitive::<Int32Type>();
5961 let vecs = batch["vec"].as_fixed_size_list();
5962 for i in 0..batch.num_rows() {
5963 let v = vecs.value(i);
5964 let v = v.as_primitive::<Float32Type>().values().to_vec();
5965 survivors.push((ids.value(i), v));
5966 }
5967 }
5968 }
5969 assert!(!survivors.is_empty());
5970 let surviving_ids: std::collections::HashSet<i32> =
5971 survivors.iter().map(|(id, _)| *id).collect();
5972
5973 let step = (survivors.len() / 16).max(1);
5975 let queries: Vec<(i32, Vec<f32>)> = survivors.iter().step_by(step).cloned().collect();
5976 let mut baseline: Vec<Vec<i32>> = Vec::new();
5977 for (_, q) in &queries {
5978 baseline.push(vector_knn_ids(&dataset, q, k).await);
5979 }
5980
5981 let metrics = compact_files(
5983 &mut dataset,
5984 CompactionOptions {
5985 target_rows_per_fragment: 2_000,
5986 defer_index_remap: true,
5987 ..Default::default()
5988 },
5989 None,
5990 )
5991 .await
5992 .unwrap();
5993 assert!(metrics.fragments_removed > 0);
5994 assert!(
5995 dataset
5996 .load_indices()
5997 .await
5998 .unwrap()
5999 .iter()
6000 .any(|idx| idx.name == FRAG_REUSE_INDEX_NAME),
6001 "deferred compaction must record a frag-reuse index"
6002 );
6003 assert_eq!(
6004 dataset
6005 .load_index_by_name("vec_idx")
6006 .await
6007 .unwrap()
6008 .unwrap()
6009 .uuid,
6010 original_uuid,
6011 "index must not be physically remapped yet (FRI window)"
6012 );
6013
6014 for (i, (_, q)) in queries.iter().enumerate() {
6016 let after = vector_knn_ids(&dataset, q, k).await;
6017 for id in &after {
6018 assert!(
6019 surviving_ids.contains(id),
6020 "KNN returned id {id} that is not a surviving row (query #{i})"
6021 );
6022 }
6023 let overlap = after.iter().filter(|id| baseline[i].contains(id)).count();
6024 assert!(
6025 overlap >= min_overlap,
6026 "KNN top-{k} diverged after deferred compaction: overlap {overlap} < {min_overlap} (query #{i})"
6027 );
6028 }
6029 }
6030
6031 fn small_ivf() -> lance_index::vector::ivf::IvfBuildParams {
6032 lance_index::vector::ivf::IvfBuildParams {
6033 max_iters: 2,
6034 num_partitions: Some(2),
6035 sample_rate: 2,
6036 ..Default::default()
6037 }
6038 }
6039
6040 #[tokio::test]
6041 async fn test_ivf_flat_defer_compaction_with_deletions() {
6042 let params = VectorIndexParams::with_ivf_flat_params(DistanceType::L2, small_ivf());
6043 check_vector_defer_compaction(params, Some("id < 1500"), 10, 10).await;
6045 }
6046
6047 #[tokio::test]
6048 async fn test_ivf_hnsw_sq_defer_compaction_merge_only() {
6049 use lance_index::vector::{hnsw::builder::HnswBuildParams, sq::builder::SQBuildParams};
6050 let params = VectorIndexParams::with_ivf_hnsw_sq_params(
6051 DistanceType::L2,
6052 small_ivf(),
6053 HnswBuildParams::default(),
6054 SQBuildParams::default(),
6055 );
6056 check_vector_defer_compaction(params, None, 10, 9).await;
6058 }
6059
6060 #[tokio::test]
6067 async fn test_ivf_pq_defer_compaction_with_deletions() {
6068 use lance_index::vector::pq::PQBuildParams;
6069 let params = VectorIndexParams::with_ivf_pq_params(
6070 DistanceType::L2,
6071 small_ivf(),
6072 PQBuildParams {
6073 max_iters: 2,
6074 num_sub_vectors: 2,
6075 ..Default::default()
6076 },
6077 );
6078 check_vector_defer_compaction(params, Some("id < 1500"), 10, 8).await;
6079 }
6080
6081 #[tokio::test]
6082 async fn test_ivf_sq_defer_compaction_with_deletions() {
6083 use lance_index::vector::sq::builder::SQBuildParams;
6084 let params = VectorIndexParams::with_ivf_sq_params(
6085 DistanceType::L2,
6086 small_ivf(),
6087 SQBuildParams::default(),
6088 );
6089 check_vector_defer_compaction(params, Some("id < 1500"), 10, 8).await;
6090 }
6091
6092 #[tokio::test]
6093 async fn test_ivf_rq_defer_compaction_with_deletions() {
6094 use lance_index::vector::bq::RQBuildParams;
6095 let params = VectorIndexParams::with_ivf_rq_params(
6096 DistanceType::L2,
6097 small_ivf(),
6098 RQBuildParams::new(1),
6099 );
6100 check_vector_defer_compaction(params, Some("id < 1500"), 10, 8).await;
6101 }
6102
6103 async fn check_vector_remap_and_trim(
6109 params: VectorIndexParams,
6110 k: usize,
6111 window_overlap: usize,
6112 post_remap_overlap: Option<usize>,
6113 ) {
6114 use arrow_array::cast::AsArray;
6115 use arrow_array::types::{Float32Type, Int32Type};
6116 use lance_datagen::Dimension;
6117
6118 const DIM: u32 = 32;
6119 let mut dataset = lance_datagen::gen_batch()
6120 .col("id", lance_datagen::array::step::<Int32Type>())
6121 .col(
6122 "vec",
6123 lance_datagen::array::rand_vec::<Float32Type>(Dimension::from(DIM)),
6124 )
6125 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
6126 .await
6127 .unwrap();
6128 dataset
6129 .create_index(
6130 &["vec"],
6131 IndexType::Vector,
6132 Some("vec_idx".into()),
6133 ¶ms,
6134 false,
6135 )
6136 .await
6137 .unwrap();
6138 let original_uuid = dataset
6139 .load_index_by_name("vec_idx")
6140 .await
6141 .unwrap()
6142 .unwrap()
6143 .uuid;
6144
6145 let mut rows: Vec<Vec<f32>> = Vec::new();
6147 {
6148 let mut scanner = dataset.scan();
6149 scanner.project(&["vec"]).unwrap();
6150 let batches = scanner
6151 .try_into_stream()
6152 .await
6153 .unwrap()
6154 .try_collect::<Vec<_>>()
6155 .await
6156 .unwrap();
6157 for batch in &batches {
6158 let vecs = batch["vec"].as_fixed_size_list();
6159 for i in 0..batch.num_rows() {
6160 let v = vecs.value(i);
6161 rows.push(v.as_primitive::<Float32Type>().values().to_vec());
6162 }
6163 }
6164 }
6165 let step = (rows.len() / 16).max(1);
6166 let queries: Vec<Vec<f32>> = rows.iter().step_by(step).cloned().collect();
6167 let mut baseline: Vec<Vec<i32>> = Vec::new();
6168 for q in &queries {
6169 baseline.push(vector_knn_ids(&dataset, q, k).await);
6170 }
6171
6172 let metrics = compact_files(
6174 &mut dataset,
6175 CompactionOptions {
6176 target_rows_per_fragment: 2_000,
6177 defer_index_remap: true,
6178 ..Default::default()
6179 },
6180 None,
6181 )
6182 .await
6183 .unwrap();
6184 assert!(metrics.fragments_removed > 0);
6185 assert_eq!(
6186 dataset
6187 .load_index_by_name("vec_idx")
6188 .await
6189 .unwrap()
6190 .unwrap()
6191 .uuid,
6192 original_uuid,
6193 "index must not be physically remapped yet (FRI window)"
6194 );
6195 for (i, q) in queries.iter().enumerate() {
6196 let window = vector_knn_ids(&dataset, q, k).await;
6197 let overlap = window.iter().filter(|id| baseline[i].contains(id)).count();
6198 assert!(
6199 overlap >= window_overlap,
6200 "FRI-window KNN diverged: overlap {overlap} < {window_overlap} (query #{i})"
6201 );
6202 }
6203
6204 remapping::remap_column_index(&mut dataset, &["vec"], Some("vec_idx".into()))
6206 .await
6207 .unwrap();
6208 cleanup_frag_reuse_index(&mut dataset).await.unwrap();
6209
6210 let remapped_uuid = dataset
6211 .load_index_by_name("vec_idx")
6212 .await
6213 .unwrap()
6214 .unwrap()
6215 .uuid;
6216 assert_ne!(
6217 remapped_uuid, original_uuid,
6218 "index should have been physically remapped"
6219 );
6220 if let Some(meta) = dataset
6221 .load_index_by_name(FRAG_REUSE_INDEX_NAME)
6222 .await
6223 .unwrap()
6224 {
6225 let versions = load_frag_reuse_index_details(&dataset, &meta)
6226 .await
6227 .unwrap()
6228 .versions
6229 .len();
6230 assert_eq!(versions, 0, "frag-reuse index must trim to zero versions");
6231 }
6232
6233 for (i, q) in queries.iter().enumerate() {
6234 let after = vector_knn_ids(&dataset, q, k).await;
6235 assert!(
6237 !after.is_empty(),
6238 "post-remap KNN returned no rows (query #{i})"
6239 );
6240 if let Some(min_overlap) = post_remap_overlap {
6243 let overlap = after.iter().filter(|id| baseline[i].contains(id)).count();
6244 assert!(
6245 overlap >= min_overlap,
6246 "post-remap KNN diverged: overlap {overlap} < {min_overlap} (query #{i})"
6247 );
6248 }
6249 }
6250 }
6251
6252 #[tokio::test]
6253 async fn test_ivf_flat_remap_and_trim() {
6254 let params = VectorIndexParams::with_ivf_flat_params(DistanceType::L2, small_ivf());
6255 check_vector_remap_and_trim(params, 10, 8, Some(8)).await;
6256 }
6257
6258 #[tokio::test]
6265 async fn test_ivf_pq_remap_and_trim() {
6266 use lance_index::vector::pq::PQBuildParams;
6267 let params = VectorIndexParams::with_ivf_pq_params(
6268 DistanceType::L2,
6269 small_ivf(),
6270 PQBuildParams {
6271 max_iters: 2,
6272 num_sub_vectors: 2,
6273 ..Default::default()
6274 },
6275 );
6276 check_vector_remap_and_trim(params, 10, 8, Some(8)).await;
6277 }
6278
6279 #[tokio::test]
6280 async fn test_ivf_sq_remap_and_trim() {
6281 use lance_index::vector::sq::builder::SQBuildParams;
6282 let params = VectorIndexParams::with_ivf_sq_params(
6283 DistanceType::L2,
6284 small_ivf(),
6285 SQBuildParams::default(),
6286 );
6287 check_vector_remap_and_trim(params, 10, 8, Some(8)).await;
6288 }
6289
6290 #[tokio::test]
6291 async fn test_ivf_rq_remap_and_trim() {
6292 use lance_index::vector::bq::RQBuildParams;
6293 let params = VectorIndexParams::with_ivf_rq_params(
6294 DistanceType::L2,
6295 small_ivf(),
6296 RQBuildParams::new(1),
6297 );
6298 check_vector_remap_and_trim(params, 10, 8, Some(8)).await;
6299 }
6300
6301 #[tokio::test]
6302 async fn test_ivf_hnsw_sq_remap_and_trim() {
6303 use lance_index::vector::{hnsw::builder::HnswBuildParams, sq::builder::SQBuildParams};
6304 let params = VectorIndexParams::with_ivf_hnsw_sq_params(
6305 DistanceType::L2,
6306 small_ivf(),
6307 HnswBuildParams::default(),
6308 SQBuildParams::default(),
6309 );
6310 check_vector_remap_and_trim(params, 10, 7, None).await;
6312 }
6313
6314 #[tokio::test]
6315 async fn test_ivf_hnsw_pq_remap_and_trim() {
6316 use lance_index::vector::{hnsw::builder::HnswBuildParams, pq::PQBuildParams};
6317 let params = VectorIndexParams::with_ivf_hnsw_pq_params(
6318 DistanceType::L2,
6319 small_ivf(),
6320 HnswBuildParams::default(),
6321 PQBuildParams {
6322 max_iters: 2,
6323 num_sub_vectors: 2,
6324 ..Default::default()
6325 },
6326 );
6327 check_vector_remap_and_trim(params, 10, 7, None).await;
6328 }
6329
6330 #[tokio::test]
6340 async fn test_bitmap_index_defer_compaction_with_deletions() {
6341 use arrow_array::cast::AsArray;
6342 use arrow_array::types::Int32Type;
6343 let mut dataset = lance_datagen::gen_batch()
6344 .col("id", lance_datagen::array::step::<Int32Type>())
6345 .col(
6346 "category",
6347 lance_datagen::array::cycle::<Int32Type>(vec![1, 2, 3]),
6348 )
6349 .into_ram_dataset(FragmentCount::from(6), FragmentRowCount::from(1000))
6350 .await
6351 .unwrap();
6352 dataset
6353 .create_index(
6354 &["category"],
6355 IndexType::Bitmap,
6356 Some("category_idx".into()),
6357 &ScalarIndexParams::default(),
6358 false,
6359 )
6360 .await
6361 .unwrap();
6362 dataset.delete("id < 1500").await.unwrap();
6363 let metrics = compact_files(
6364 &mut dataset,
6365 CompactionOptions {
6366 target_rows_per_fragment: 2_000,
6367 defer_index_remap: true,
6368 ..Default::default()
6369 },
6370 None,
6371 )
6372 .await
6373 .unwrap();
6374 assert!(metrics.fragments_removed > 0);
6375 assert!(
6376 dataset
6377 .load_indices()
6378 .await
6379 .unwrap()
6380 .iter()
6381 .any(|idx| idx.name == FRAG_REUSE_INDEX_NAME),
6382 "deferred compaction must record a frag-reuse index"
6383 );
6384
6385 let mut scanner = dataset.scan();
6386 scanner.filter("category = 3").unwrap();
6387 scanner.project(&["id"]).unwrap();
6388 let batches = scanner
6389 .try_into_stream()
6390 .await
6391 .unwrap()
6392 .try_collect::<Vec<_>>()
6393 .await
6394 .unwrap();
6395 let mut returned = 0;
6396 for b in &batches {
6397 for id in b["id"].as_primitive::<Int32Type>().values() {
6398 assert!(
6399 *id >= 1500,
6400 "bitmap returned deleted id {id} in the FRI window"
6401 );
6402 returned += 1;
6403 }
6404 }
6405 assert!(returned > 0, "expected surviving category=3 rows");
6406 }
6407
6408 #[tokio::test]
6415 async fn test_default_compaction_planner() {
6416 let test_dir = TempStrDir::default();
6417 let test_uri = &test_dir;
6418
6419 let data = sample_data();
6420 let schema = data.schema();
6421
6422 let reader = RecordBatchIterator::new(vec![Ok(data.clone())], schema.clone());
6424 let write_params = WriteParams {
6425 max_rows_per_file: 2000,
6426 ..Default::default()
6427 };
6428 let dataset = Dataset::write(reader, test_uri, Some(write_params))
6429 .await
6430 .unwrap();
6431
6432 assert_eq!(dataset.get_fragments().len(), 5);
6433
6434 let options = CompactionOptions {
6436 target_rows_per_fragment: 5000,
6437 materialize_deletions_threshold: 2.0,
6438 ..Default::default()
6439 };
6440
6441 let planner = DefaultCompactionPlanner::new(options);
6442 let plan = planner.plan(&dataset).await.unwrap();
6443
6444 assert!(!plan.tasks.is_empty());
6446 assert_eq!(plan.read_version, dataset.manifest.version);
6447 assert!(!plan.options.materialize_deletions);
6449 }
6450
6451 #[test]
6452 fn test_from_dataset_config() {
6453 let config = HashMap::from([
6454 (
6455 "lance.compaction.target_rows_per_fragment".to_string(),
6456 "500000".to_string(),
6457 ),
6458 (
6459 "lance.compaction.max_rows_per_group".to_string(),
6460 "2048".to_string(),
6461 ),
6462 (
6463 "lance.compaction.max_bytes_per_file".to_string(),
6464 "1000000".to_string(),
6465 ),
6466 (
6467 "lance.compaction.materialize_deletions".to_string(),
6468 "false".to_string(),
6469 ),
6470 (
6471 "lance.compaction.materialize_deletions_threshold".to_string(),
6472 "0.25".to_string(),
6473 ),
6474 (
6475 "lance.compaction.defer_index_remap".to_string(),
6476 "true".to_string(),
6477 ),
6478 (
6479 "lance.compaction.batch_size".to_string(),
6480 "4096".to_string(),
6481 ),
6482 (
6483 "lance.compaction.io_buffer_size".to_string(),
6484 "1073741824".to_string(),
6485 ),
6486 (
6487 "lance.compaction.compaction_mode".to_string(),
6488 "try_binary_copy".to_string(),
6489 ),
6490 (
6491 "lance.compaction.binary_copy_read_batch_bytes".to_string(),
6492 "8388608".to_string(),
6493 ),
6494 (
6495 "lance.compaction.index_remap_mode".to_string(),
6496 "compact".to_string(),
6497 ),
6498 ]);
6499
6500 let opts = CompactionOptions::from_dataset_config(&config).unwrap();
6501 assert_eq!(opts.target_rows_per_fragment, 500_000);
6502 assert_eq!(opts.max_rows_per_group, 2048);
6503 assert_eq!(opts.max_bytes_per_file, Some(1_000_000));
6504 assert!(!opts.materialize_deletions);
6505 assert!((opts.materialize_deletions_threshold - 0.25).abs() < f32::EPSILON);
6506 assert!(opts.defer_index_remap);
6507 assert_eq!(opts.batch_size, Some(4096));
6508 assert_eq!(opts.io_buffer_size, Some(1_073_741_824));
6509 assert_eq!(opts.compaction_mode, Some(CompactionMode::TryBinaryCopy));
6510 assert_eq!(opts.binary_copy_read_batch_bytes, Some(8_388_608));
6511 assert_eq!(opts.index_remap_mode, IndexRemapMode::Compact);
6513 }
6514
6515 #[test]
6516 fn test_from_dataset_config_empty() {
6517 let config = HashMap::new();
6518 let opts = CompactionOptions::from_dataset_config(&config).unwrap();
6519 let defaults = CompactionOptions::default();
6520 assert_eq!(
6521 opts.target_rows_per_fragment,
6522 defaults.target_rows_per_fragment
6523 );
6524 assert_eq!(opts.max_rows_per_group, defaults.max_rows_per_group);
6525 assert_eq!(opts.max_bytes_per_file, defaults.max_bytes_per_file);
6526 assert_eq!(opts.materialize_deletions, defaults.materialize_deletions);
6527 assert_eq!(
6528 opts.materialize_deletions_threshold,
6529 defaults.materialize_deletions_threshold
6530 );
6531 assert_eq!(opts.defer_index_remap, defaults.defer_index_remap);
6532 assert_eq!(opts.index_remap_mode, defaults.index_remap_mode);
6533 assert_eq!(opts.index_remap_mode, IndexRemapMode::Direct);
6534 assert_eq!(opts.batch_size, defaults.batch_size);
6535 assert_eq!(opts.compaction_mode, defaults.compaction_mode);
6536 assert_eq!(
6537 opts.binary_copy_read_batch_bytes,
6538 defaults.binary_copy_read_batch_bytes
6539 );
6540 }
6541
6542 #[test]
6543 fn test_from_dataset_config_partial() {
6544 let config = HashMap::from([(
6545 "lance.compaction.target_rows_per_fragment".to_string(),
6546 "500000".to_string(),
6547 )]);
6548
6549 let opts = CompactionOptions::from_dataset_config(&config).unwrap();
6550 assert_eq!(opts.target_rows_per_fragment, 500_000);
6551 let defaults = CompactionOptions::default();
6553 assert_eq!(opts.max_rows_per_group, defaults.max_rows_per_group);
6554 assert_eq!(opts.max_bytes_per_file, defaults.max_bytes_per_file);
6555 assert_eq!(opts.materialize_deletions, defaults.materialize_deletions);
6556 assert_eq!(opts.defer_index_remap, defaults.defer_index_remap);
6557 assert_eq!(opts.batch_size, defaults.batch_size);
6558 assert_eq!(opts.compaction_mode, defaults.compaction_mode);
6559 assert_eq!(
6560 opts.binary_copy_read_batch_bytes,
6561 defaults.binary_copy_read_batch_bytes
6562 );
6563 }
6564
6565 #[test]
6566 fn test_from_dataset_config_ignores_other_keys() {
6567 let config = HashMap::from([
6568 (
6569 "lance.compaction.target_rows_per_fragment".to_string(),
6570 "500000".to_string(),
6571 ),
6572 (
6573 "lance.auto_cleanup.interval".to_string(),
6574 "3600".to_string(),
6575 ),
6576 ("some.other.key".to_string(), "value".to_string()),
6577 ]);
6578
6579 let opts = CompactionOptions::from_dataset_config(&config).unwrap();
6580 assert_eq!(opts.target_rows_per_fragment, 500_000);
6581 }
6582
6583 #[test]
6584 fn test_from_dataset_config_invalid_value() {
6585 let config = HashMap::from([(
6586 "lance.compaction.target_rows_per_fragment".to_string(),
6587 "not_a_number".to_string(),
6588 )]);
6589
6590 let result = CompactionOptions::from_dataset_config(&config);
6591 let err_msg = result.unwrap_err().to_string();
6592 assert!(err_msg.contains("target_rows_per_fragment"));
6593 assert!(err_msg.contains("not_a_number"));
6594 }
6595
6596 #[test]
6597 fn test_from_dataset_config_invalid_bool() {
6598 let config = HashMap::from([(
6599 "lance.compaction.materialize_deletions".to_string(),
6600 "yes".to_string(),
6601 )]);
6602
6603 let result = CompactionOptions::from_dataset_config(&config);
6604 let err_msg = result.unwrap_err().to_string();
6605 assert!(err_msg.contains("materialize_deletions"));
6606 assert!(err_msg.contains("yes"));
6607 }
6608
6609 #[test]
6610 fn test_from_dataset_config_unknown_compaction_key() {
6611 let config = HashMap::from([(
6613 "lance.compaction.unknown_key".to_string(),
6614 "value".to_string(),
6615 )]);
6616
6617 let opts = CompactionOptions::from_dataset_config(&config).unwrap();
6618 let defaults = CompactionOptions::default();
6620 assert_eq!(
6621 opts.target_rows_per_fragment,
6622 defaults.target_rows_per_fragment
6623 );
6624 }
6625
6626 #[test]
6627 fn test_from_dataset_config_invalid_compaction_mode() {
6628 let config = HashMap::from([(
6629 "lance.compaction.compaction_mode".to_string(),
6630 "invalid_mode".to_string(),
6631 )]);
6632
6633 let result = CompactionOptions::from_dataset_config(&config);
6634 let err_msg = result.unwrap_err().to_string();
6635 assert!(err_msg.contains("invalid_mode"));
6636 }
6637
6638 #[test]
6639 fn test_from_dataset_config_max_overlays_per_fragment() {
6640 let key = "lance.compaction.max_overlays_per_fragment".to_string();
6641
6642 let config = HashMap::from([(key.clone(), "3".to_string())]);
6644 let opts = CompactionOptions::from_dataset_config(&config).unwrap();
6645 assert_eq!(opts.max_overlays_per_fragment, Some(3));
6646
6647 let config = HashMap::from([(key.clone(), "None".to_string())]);
6649 let opts = CompactionOptions::from_dataset_config(&config).unwrap();
6650 assert_eq!(opts.max_overlays_per_fragment, None);
6651
6652 let config = HashMap::from([(key, "not_a_number".to_string())]);
6654 let err_msg = CompactionOptions::from_dataset_config(&config)
6655 .unwrap_err()
6656 .to_string();
6657 assert!(err_msg.contains("max_overlays_per_fragment"));
6658 assert!(err_msg.contains("not_a_number"));
6659 }
6660
6661 #[test]
6662 fn test_apply_dataset_config_overrides() {
6663 let config = HashMap::from([(
6664 "lance.compaction.target_rows_per_fragment".to_string(),
6665 "500000".to_string(),
6666 )]);
6667
6668 let mut opts = CompactionOptions {
6669 max_rows_per_group: 4096,
6670 ..Default::default()
6671 };
6672 opts.apply_dataset_config(&config).unwrap();
6673
6674 assert_eq!(opts.target_rows_per_fragment, 500_000);
6676 assert_eq!(opts.max_rows_per_group, 4096);
6678 }
6679
6680 #[test]
6681 fn test_apply_dataset_config_overwrites_matching_field() {
6682 let config = HashMap::from([(
6683 "lance.compaction.max_rows_per_group".to_string(),
6684 "2048".to_string(),
6685 )]);
6686
6687 let mut opts = CompactionOptions {
6688 max_rows_per_group: 4096,
6689 ..Default::default()
6690 };
6691 opts.apply_dataset_config(&config).unwrap();
6692
6693 assert_eq!(opts.max_rows_per_group, 2048);
6695 }
6696
6697 #[tokio::test]
6698 async fn test_max_source_fragments() {
6699 let test_dir = TempStrDir::default();
6700 let test_uri = &test_dir;
6701
6702 let data = sample_data();
6703 let schema = data.schema();
6704
6705 let write_params = WriteParams {
6707 max_rows_per_file: 100,
6708 ..Default::default()
6709 };
6710 Dataset::write(
6711 RecordBatchIterator::new(vec![Ok(data.slice(0, 100))], schema.clone()),
6712 test_uri,
6713 Some(write_params.clone()),
6714 )
6715 .await
6716 .unwrap();
6717 for i in 1..10 {
6718 let mut append_params = write_params.clone();
6719 append_params.mode = WriteMode::Append;
6720 Dataset::write(
6721 RecordBatchIterator::new(vec![Ok(data.slice(i * 100, 100))], schema.clone()),
6722 test_uri,
6723 Some(append_params),
6724 )
6725 .await
6726 .unwrap();
6727 }
6728
6729 let dataset = Dataset::open(test_uri).await.unwrap();
6730 assert_eq!(dataset.get_fragments().len(), 10);
6731
6732 let opts_no_limit = CompactionOptions {
6735 target_rows_per_fragment: 250,
6736 ..Default::default()
6737 };
6738 let plan_all = plan_compaction(&dataset, &opts_no_limit).await.unwrap();
6739 let total_source_frags: usize = plan_all.tasks().iter().map(|t| t.fragments.len()).sum();
6740 assert_eq!(total_source_frags, 10);
6741 assert!(
6742 plan_all.num_tasks() > 2,
6743 "need multiple tasks to test bounding, got {}",
6744 plan_all.num_tasks()
6745 );
6746
6747 let opts_bounded = CompactionOptions {
6750 target_rows_per_fragment: 250,
6751 max_source_fragments: Some(4),
6752 ..Default::default()
6753 };
6754 let plan_bounded = plan_compaction(&dataset, &opts_bounded).await.unwrap();
6755 let bounded_source_frags: usize =
6756 plan_bounded.tasks().iter().map(|t| t.fragments.len()).sum();
6757 assert!(
6758 bounded_source_frags <= 4,
6759 "expected at most 4 source fragments, got {bounded_source_frags}"
6760 );
6761 assert!(
6762 bounded_source_frags > 0,
6763 "expected at least 1 source fragment in bounded plan"
6764 );
6765 assert!(
6766 plan_bounded.num_tasks() < plan_all.num_tasks(),
6767 "bounded plan ({}) should have fewer tasks than unbounded ({})",
6768 plan_bounded.num_tasks(),
6769 plan_all.num_tasks()
6770 );
6771
6772 let mut dataset = dataset;
6774 compact_files(&mut dataset, opts_bounded, None)
6775 .await
6776 .unwrap();
6777 let after_first = dataset.get_fragments().len();
6778 assert!(
6779 after_first < 10,
6780 "expected fewer than 10 fragments after first compaction, got {after_first}"
6781 );
6782 assert!(
6783 after_first > 1,
6784 "expected partial compaction (not fully compacted), got {after_first}"
6785 );
6786
6787 let opts_bounded = CompactionOptions {
6789 target_rows_per_fragment: 250,
6790 max_source_fragments: Some(4),
6791 ..Default::default()
6792 };
6793 compact_files(&mut dataset, opts_bounded, None)
6794 .await
6795 .unwrap();
6796 let after_second = dataset.get_fragments().len();
6797 assert!(
6798 after_second <= after_first,
6799 "expected progress: {after_second} should be <= {after_first}"
6800 );
6801 }
6802
6803 #[tokio::test]
6804 async fn test_compaction_uses_manifest_config() {
6805 let test_dir = TempStrDir::default();
6806 let test_uri = &test_dir;
6807
6808 let data = sample_data();
6809 let schema = data.schema();
6810
6811 let reader = RecordBatchIterator::new(vec![Ok(data.clone())], schema.clone());
6813 let write_params = WriteParams {
6814 max_rows_per_file: 2000,
6815 ..Default::default()
6816 };
6817 let mut dataset = Dataset::write(reader, test_uri, Some(write_params))
6818 .await
6819 .unwrap();
6820
6821 assert_eq!(dataset.get_fragments().len(), 5);
6822
6823 dataset
6825 .update_config([
6826 ("lance.compaction.target_rows_per_fragment", "5000"),
6827 ("lance.compaction.materialize_deletions_threshold", "2.0"),
6828 ])
6829 .await
6830 .unwrap();
6831
6832 let opts = CompactionOptions::from_dataset_config(&dataset.manifest.config).unwrap();
6834 assert_eq!(opts.target_rows_per_fragment, 5000);
6835 assert!((opts.materialize_deletions_threshold - 2.0).abs() < f32::EPSILON);
6836
6837 let plan = plan_compaction(&dataset, &opts).await.unwrap();
6839 assert!(!plan.tasks.is_empty());
6840 assert_eq!(plan.options.target_rows_per_fragment, 5000);
6841 assert!(!plan.options.materialize_deletions);
6843 }
6844
6845 #[tokio::test]
6853 async fn test_rewrite_fri_vs_create_index_conflict() {
6854 use crate::index::DatasetIndexExt;
6855 use crate::index::vector::VectorIndexParams;
6856 use futures::TryStreamExt;
6857 use lance_datagen::{BatchCount, Dimension, RowCount, array, gen_batch};
6858 use lance_index::IndexType;
6859 use lance_linalg::distance::MetricType;
6860
6861 async fn append_fragment(uri: &str, rows: u64) -> Dataset {
6862 let reader = gen_batch()
6863 .col("vec", array::rand_vec::<Float32Type>(Dimension::from(16)))
6864 .into_reader_rows(RowCount::from(rows), BatchCount::from(1));
6865 let params = WriteParams {
6866 max_rows_per_file: rows as usize,
6867 mode: WriteMode::Append,
6868 ..Default::default()
6869 };
6870 Dataset::write(reader, uri, Some(params)).await.unwrap()
6871 }
6872
6873 let tmpdir = TempStrDir::default();
6874 let uri = format!("file://{}", tmpdir.as_str());
6875
6876 let reader = gen_batch()
6878 .col("vec", array::rand_vec::<Float32Type>(Dimension::from(16)))
6879 .into_reader_rows(RowCount::from(256), BatchCount::from(1));
6880 let mut dataset = Dataset::write(
6881 reader,
6882 &uri,
6883 Some(WriteParams {
6884 max_rows_per_file: 256,
6885 mode: WriteMode::Overwrite,
6886 ..Default::default()
6887 }),
6888 )
6889 .await
6890 .unwrap();
6891 let index_params = VectorIndexParams::ivf_pq(2, 8, 2, MetricType::L2, 50);
6892 dataset
6893 .create_index(&["vec"], IndexType::Vector, None, &index_params, true)
6894 .await
6895 .unwrap();
6896
6897 dataset = append_fragment(&uri, 64).await;
6900 let mut stale = dataset.clone();
6901 dataset = append_fragment(&uri, 64).await;
6902
6903 let options = CompactionOptions {
6905 defer_index_remap: true,
6906 ..Default::default()
6907 };
6908 let plan = plan_compaction(&dataset, &options).await.unwrap();
6909 assert!(!plan.tasks.is_empty());
6910 let snapshot = dataset.clone();
6911 let completed: Vec<RewriteResult> = futures::stream::iter(plan.tasks.into_iter())
6912 .map(|task| rewrite_files(Cow::Borrowed(&snapshot), task, &options))
6913 .buffer_unordered(1)
6914 .try_collect()
6915 .await
6916 .unwrap();
6917
6918 stale
6923 .optimize_indices(&lance_index::optimize::OptimizeOptions::append())
6924 .await
6925 .unwrap();
6926
6927 let err = commit_compaction(
6932 &mut dataset,
6933 completed,
6934 Arc::new(DatasetIndexRemapperOptions::default()),
6935 &options,
6936 )
6937 .await
6938 .expect_err("commit should fail with retryable conflict");
6939 assert!(
6940 matches!(err, Error::RetryableCommitConflict { .. }),
6941 "unexpected error: {err}"
6942 );
6943 }
6944
6945 #[tokio::test]
6963 async fn test_distributed_compact_concurrent_delete_no_resurrection() {
6964 let test_dir = TempStrDir::default();
6965 let test_uri = &test_dir;
6966
6967 let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int64, false)]));
6969 let data = RecordBatch::try_new(
6970 schema.clone(),
6971 vec![Arc::new(Int64Array::from_iter_values(0..4_000))],
6972 )
6973 .unwrap();
6974 let mut dataset_plan = Dataset::write(
6975 RecordBatchIterator::new(vec![Ok(data)], schema.clone()),
6976 test_uri,
6977 Some(WriteParams {
6978 max_rows_per_file: 1_000,
6979 ..Default::default()
6980 }),
6981 )
6982 .await
6983 .unwrap();
6984
6985 assert_eq!(dataset_plan.manifest.version, 1);
6986 assert_eq!(dataset_plan.get_fragments().len(), 4);
6987
6988 let options = CompactionOptions {
6990 target_rows_per_fragment: 10_000,
6991 ..Default::default()
6992 };
6993 let plan = plan_compaction(&dataset_plan, &options).await.unwrap();
6994 assert_eq!(plan.tasks().len(), 1, "expected one compaction task");
6995
6996 let dataset_for_tasks = dataset_plan.clone();
7000 let results: Vec<RewriteResult> = futures::stream::iter(plan.compaction_tasks())
7001 .then(|task| {
7002 let ds = dataset_for_tasks.clone();
7003 async move {
7004 task.execute(&ds).await.unwrap()
7006 }
7007 })
7008 .collect()
7009 .await;
7010 assert_eq!(results.len(), 1);
7011 assert_eq!(
7012 results[0].read_version, 1,
7013 "tasks must carry read_version=1"
7014 );
7015
7016 dataset_plan.delete("a < 1000").await.unwrap();
7019 assert_eq!(dataset_plan.manifest.version, 2);
7020
7021 let mut dataset_commit = Dataset::open(test_uri).await.unwrap();
7024 assert_eq!(
7025 dataset_commit.manifest.version, 2,
7026 "fresh dataset must be at the post-delete version"
7027 );
7028
7029 let commit_result = commit_compaction(
7031 &mut dataset_commit,
7032 results,
7033 Arc::new(IgnoreRemap::default()),
7034 &options,
7035 )
7036 .await;
7037
7038 assert!(
7042 commit_result.is_err(),
7043 "commit_compaction must fail with a conflict error when a concurrent \
7044 DELETE touched the same fragments; got Ok instead — deleted rows were \
7045 silently resurrected"
7046 );
7047 let err_msg = commit_result.unwrap_err().to_string();
7048 assert!(
7049 err_msg.contains("retryable")
7050 || err_msg.contains("conflict")
7051 || err_msg.contains("preempted"),
7052 "expected a retryable conflict error, got: {err_msg}"
7053 );
7054
7055 let latest = Dataset::open(test_uri).await.unwrap();
7057 let row_count = latest
7058 .count_rows(Some("a < 1000".to_string()))
7059 .await
7060 .unwrap();
7061 assert_eq!(
7062 row_count, 0,
7063 "rows deleted before compaction must not be resurrected; found {row_count}"
7064 );
7065 }
7066
7067 fn count_all_files_in(dir: &std::path::Path) -> std::io::Result<usize> {
7068 if !dir.exists() {
7069 return Ok(0);
7070 }
7071 let mut count = 0;
7072 for entry in std::fs::read_dir(dir)? {
7073 let path = entry?.path();
7074 if path.is_dir() {
7075 count += count_all_files_in(&path)?;
7076 } else if path.is_file() {
7077 if path
7079 .file_name()
7080 .and_then(|name| name.to_str())
7081 .is_some_and(|file_name| !file_name.starts_with('.'))
7082 {
7083 count += 1;
7084 }
7085 }
7086 }
7087 Ok(count)
7088 }
7089
7090 fn count_data_files_in(base_dir: &str) -> usize {
7091 let data_dir = std::path::Path::new(base_dir).join("data");
7092 count_all_files_in(&data_dir).unwrap_or(0)
7093 }
7094
7095 #[tokio::test]
7100 async fn test_commit_compaction_cleans_up_data_on_commit_failure() {
7101 use crate::dataset::builder::DatasetBuilder;
7102 use crate::utils::test::FailingProxyStore;
7103 use lance_io::object_store::ObjectStoreParams;
7104
7105 let test_dir = TempStrDir::default();
7106 let test_uri = test_dir.as_str();
7107 let path_prefix = if test_uri.starts_with('/') { "" } else { "/" };
7110 let routed_uri = format!("file-object-store://{path_prefix}{test_uri}");
7111
7112 let data = sample_data();
7113 let reader = RecordBatchIterator::new(vec![Ok(data.slice(0, 200))], data.schema());
7114 Dataset::write(
7115 reader,
7116 &routed_uri,
7117 Some(WriteParams {
7118 max_rows_per_file: 100,
7119 enable_stable_row_ids: true,
7124 ..Default::default()
7125 }),
7126 )
7127 .await
7128 .unwrap();
7129
7130 let baseline_files = count_data_files_in(test_uri);
7131
7132 let failing = Arc::new(FailingProxyStore::new());
7133 failing.fail_after_n("put", "_transactions", 1, "injected commit failure");
7139 failing.fail_after_n(
7140 "put_multipart",
7141 "_transactions",
7142 1,
7143 "injected commit failure",
7144 );
7145
7146 let mut dataset = DatasetBuilder::from_uri(&routed_uri)
7147 .with_read_params(crate::dataset::ReadParams {
7148 store_options: Some(ObjectStoreParams {
7149 object_store_wrapper: Some(failing.clone()),
7150 ..Default::default()
7151 }),
7152 ..Default::default()
7153 })
7154 .load()
7155 .await
7156 .unwrap();
7157
7158 let options = CompactionOptions {
7159 target_rows_per_fragment: 1000,
7160 ..Default::default()
7161 };
7162 let result = compact_files(&mut dataset, options, None).await;
7163 assert!(
7164 result.is_err(),
7165 "Compaction should fail when transaction commit fails"
7166 );
7167
7168 assert_eq!(
7169 count_data_files_in(test_uri),
7170 baseline_files,
7171 "Compaction data files should be cleaned up when commit fails"
7172 );
7173 }
7174
7175 #[tokio::test]
7176 async fn test_commit_compaction_cleans_up_blob_v2_sidecars_on_commit_failure() {
7177 use crate::BlobArrayBuilder;
7178 use crate::dataset::builder::DatasetBuilder;
7179 use crate::utils::test::FailingProxyStore;
7180 use lance_io::object_store::ObjectStoreParams;
7181
7182 let test_dir = TempStrDir::default();
7183 let test_uri = test_dir.as_str();
7184 let path_prefix = if test_uri.starts_with('/') { "" } else { "/" };
7185 let routed_uri = format!("file-object-store://{path_prefix}{test_uri}");
7186
7187 let id_array = Arc::new(Int32Array::from(vec![1, 2])) as ArrayRef;
7188 let packed_data = vec![0u8; 100 * 1024];
7190 let dedicated_data = vec![1u8; 5 * 1024 * 1024];
7191 let mut blob_builder = BlobArrayBuilder::new(2);
7192 blob_builder.push_bytes(&packed_data).unwrap();
7193 blob_builder.push_bytes(&dedicated_data).unwrap();
7194 let blob_array: ArrayRef = blob_builder.finish().unwrap();
7195
7196 let schema = Arc::new(Schema::new(vec![
7197 Field::new("id", DataType::Int32, false),
7198 crate::blob_field("blob", true),
7199 ]));
7200 let batch = RecordBatch::try_new(schema.clone(), vec![id_array, blob_array]).unwrap();
7201 let reader = RecordBatchIterator::new(vec![Ok(batch)], schema.clone());
7202
7203 Dataset::write(
7204 reader,
7205 &routed_uri,
7206 Some(WriteParams {
7207 max_rows_per_file: 1, enable_stable_row_ids: true,
7209 data_storage_version: Some(lance_file::version::LanceFileVersion::V2_2),
7210 ..Default::default()
7211 }),
7212 )
7213 .await
7214 .unwrap();
7215
7216 let baseline_files = count_data_files_in(test_uri);
7217
7218 let failing = Arc::new(FailingProxyStore::new());
7219 failing.fail_after_n("put", "_transactions", 1, "injected commit failure");
7220 failing.fail_after_n(
7221 "put_multipart",
7222 "_transactions",
7223 1,
7224 "injected commit failure",
7225 );
7226
7227 let mut dataset = DatasetBuilder::from_uri(&routed_uri)
7228 .with_read_params(crate::dataset::ReadParams {
7229 store_options: Some(ObjectStoreParams {
7230 object_store_wrapper: Some(failing.clone()),
7231 ..Default::default()
7232 }),
7233 ..Default::default()
7234 })
7235 .load()
7236 .await
7237 .unwrap();
7238
7239 let options = CompactionOptions {
7240 target_rows_per_fragment: 1000,
7241 ..Default::default()
7242 };
7243 let result = compact_files(&mut dataset, options, None).await;
7244 assert!(
7245 result.is_err(),
7246 "Compaction should fail when transaction commit fails"
7247 );
7248
7249 assert_eq!(
7250 count_data_files_in(test_uri),
7251 baseline_files,
7252 "Blob v2 sidecars should be cleaned up when commit fails"
7253 );
7254 }
7255
7256 async fn read_blob_bytes_by_index(
7257 dataset: &Arc<Dataset>,
7258 column: &str,
7259 ) -> Vec<(i32, Option<Vec<u8>>)> {
7260 let mut scanner = dataset.scan();
7261 scanner.with_row_id();
7262 let batch = scanner
7263 .project(&["id", column])
7264 .unwrap()
7265 .try_into_batch()
7266 .await
7267 .unwrap();
7268 let ids = batch
7269 .column_by_name("id")
7270 .unwrap()
7271 .as_primitive::<Int32Type>();
7272 let row_ids = batch
7273 .column_by_name(ROW_ID)
7274 .unwrap()
7275 .as_primitive::<UInt64Type>();
7276
7277 let mut result = Vec::with_capacity(batch.num_rows());
7278 for i in 0..batch.num_rows() {
7279 let row_id = row_ids.value(i);
7280 let id = ids.value(i);
7281 let blobs = dataset.take_blobs(&[row_id], column).await.unwrap();
7282 match blobs.into_iter().next().flatten() {
7283 Some(blob) => {
7284 let data = blob.read().await.unwrap();
7285 result.push((id, Some(data.to_vec())));
7286 }
7287 None => result.push((id, None)),
7288 }
7289 }
7290 result
7291 }
7292
7293 fn mixed_blob_values() -> Vec<(i32, Option<Vec<u8>>)> {
7294 vec![
7295 (0, Some(vec![b'0'; 80])),
7296 (1, None),
7297 (2, Some(Vec::new())),
7298 (3, Some(vec![b'3'; 80])),
7299 (4, Some(vec![b'4'; 80])),
7300 (5, Some(vec![b'5'; 80])),
7301 ]
7302 }
7303
7304 async fn assert_compaction_preserves_blob_values(
7305 mut dataset: Dataset,
7306 expected: &[(i32, Option<Vec<u8>>)],
7307 ) {
7308 assert_eq!(dataset.get_fragments().len(), 3);
7309
7310 let mut before = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "blob").await;
7311 before.sort_by_key(|(id, _)| *id);
7312 assert_eq!(before, expected);
7313
7314 compact_files(
7315 &mut dataset,
7316 CompactionOptions {
7317 target_rows_per_fragment: 1024 * 1024,
7318 ..Default::default()
7319 },
7320 None,
7321 )
7322 .await
7323 .unwrap();
7324
7325 assert_eq!(dataset.get_fragments().len(), 1);
7326
7327 let mut after = read_blob_bytes_by_index(&Arc::new(dataset), "blob").await;
7328 after.sort_by_key(|(id, _)| *id);
7329 assert_eq!(after, expected);
7330 }
7331
7332 #[tokio::test]
7333 async fn test_compact_blob_v1_preserves_null_empty_and_payload_order() {
7334 let test_dir = TempStrDir::default();
7335 let expected = mixed_blob_values();
7336 let schema = Arc::new(Schema::new(vec![
7337 Field::new("id", DataType::Int32, false),
7338 Field::new("blob", DataType::LargeBinary, true)
7339 .with_metadata([(BLOB_META_KEY.to_string(), "true".to_string())].into()),
7340 ]));
7341 let batch = RecordBatch::try_new(
7342 schema.clone(),
7343 vec![
7344 Arc::new(Int32Array::from_iter_values(0..expected.len() as i32)),
7345 Arc::new(LargeBinaryArray::from_iter(
7346 expected.iter().map(|(_, value)| value.as_deref()),
7347 )),
7348 ],
7349 )
7350 .unwrap();
7351 let dataset = Dataset::write(
7352 RecordBatchIterator::new(vec![Ok(batch)], schema),
7353 &test_dir,
7354 Some(WriteParams {
7355 data_storage_version: Some(LanceFileVersion::V2_0),
7356 max_rows_per_file: 2,
7357 ..Default::default()
7358 }),
7359 )
7360 .await
7361 .unwrap();
7362
7363 assert_compaction_preserves_blob_values(dataset, &expected).await;
7364 }
7365
7366 #[tokio::test]
7367 async fn test_compact_blob_v2_preserves_null_empty_and_payload_order() {
7368 use crate::BlobArrayBuilder;
7369
7370 let test_dir = TempStrDir::default();
7371 let expected = mixed_blob_values();
7372 let mut blob_builder = BlobArrayBuilder::new(expected.len());
7373 for (_, value) in &expected {
7374 match value {
7375 Some(value) => blob_builder.push_bytes(value).unwrap(),
7376 None => blob_builder.push_null().unwrap(),
7377 }
7378 }
7379 let schema = Arc::new(Schema::new(vec![
7380 Field::new("id", DataType::Int32, false),
7381 crate::blob_field("blob", true),
7382 ]));
7383 let batch = RecordBatch::try_new(
7384 schema.clone(),
7385 vec![
7386 Arc::new(Int32Array::from_iter_values(0..expected.len() as i32)),
7387 blob_builder.finish().unwrap(),
7388 ],
7389 )
7390 .unwrap();
7391 let dataset = Dataset::write(
7392 RecordBatchIterator::new(vec![Ok(batch)], schema),
7393 &test_dir,
7394 Some(WriteParams {
7395 data_storage_version: Some(LanceFileVersion::V2_2),
7396 max_rows_per_file: 2,
7397 ..Default::default()
7398 }),
7399 )
7400 .await
7401 .unwrap();
7402
7403 assert_compaction_preserves_blob_values(dataset, &expected).await;
7404 }
7405
7406 #[tokio::test]
7407 async fn test_compact_blob_v2_preserves_external_references() {
7408 use crate::BlobArrayBuilder;
7409 use lance_core::utils::tempfile::TempDir;
7410 use lance_table::format::BasePath;
7411
7412 let test_dir = TempDir::default();
7413 let external_dir = TempDir::default();
7414 let external_path = external_dir.std_path().join("external.bin");
7415 std::fs::write(&external_path, b"external-data").unwrap();
7416 let external_uri = format!("file://{}", external_path.display());
7417 let base_uri = format!("file://{}", external_dir.std_path().display());
7418
7419 let mut blob_builder = BlobArrayBuilder::new(2);
7420 blob_builder.push_uri(external_uri.clone()).unwrap();
7421 blob_builder.push_bytes(b"inline-data").unwrap();
7422 let blob_array: ArrayRef = blob_builder.finish().unwrap();
7423
7424 let id_array: ArrayRef = Arc::new(Int32Array::from(vec![0, 1]));
7425 let schema = Arc::new(Schema::new(vec![
7426 Field::new("id", DataType::Int32, false),
7427 crate::blob_field("blob", true),
7428 ]));
7429
7430 let batch = RecordBatch::try_new(schema.clone(), vec![id_array, blob_array]).unwrap();
7431 let reader = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema.clone());
7432
7433 let mut dataset = Dataset::write(
7434 reader,
7435 &test_dir.path_str(),
7436 Some(WriteParams {
7437 data_storage_version: Some(LanceFileVersion::V2_2),
7438 max_rows_per_file: 1,
7439 initial_bases: Some(vec![BasePath {
7440 id: 1,
7441 name: Some("external".to_string()),
7442 path: base_uri,
7443 is_dataset_root: false,
7444 }]),
7445 ..Default::default()
7446 }),
7447 )
7448 .await
7449 .unwrap();
7450
7451 assert_eq!(dataset.get_fragments().len(), 2);
7452
7453 for frag in dataset.get_fragments() {
7454 let rows = frag.physical_rows().await.unwrap();
7455 assert!(rows > 0, "fragment {} should have rows", frag.id());
7456 }
7457
7458 let options = CompactionOptions {
7459 target_rows_per_fragment: 1024 * 1024,
7460 ..Default::default()
7461 };
7462 let plan = plan_compaction(&dataset, &options).await.unwrap();
7463 assert!(
7464 !plan.tasks().is_empty(),
7465 "compaction plan should have tasks, got {} tasks",
7466 plan.tasks().len()
7467 );
7468
7469 compact_files(&mut dataset, options, None).await.unwrap();
7470
7471 assert_eq!(dataset.get_fragments().len(), 1);
7472
7473 let scan_result = dataset
7474 .scan()
7475 .project(&["id", "blob"])
7476 .unwrap()
7477 .try_into_batch()
7478 .await
7479 .unwrap();
7480 assert_eq!(scan_result.num_rows(), 2);
7481
7482 let ids = scan_result
7483 .column_by_name("id")
7484 .unwrap()
7485 .as_primitive::<Int32Type>();
7486 let mut id_values: Vec<i32> = ids.iter().map(|v| v.unwrap()).collect();
7487 id_values.sort();
7488 assert_eq!(id_values, vec![0, 1]);
7489
7490 let mut blob_values = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "blob").await;
7491 blob_values.sort_by_key(|(id, _)| *id);
7492 assert_eq!(
7493 blob_values,
7494 vec![
7495 (0, Some(b"external-data".to_vec())),
7496 (1, Some(b"inline-data".to_vec()))
7497 ]
7498 );
7499 }
7500
7501 #[tokio::test]
7502 async fn test_compact_blob_v2_packed_and_dedicated() {
7503 use crate::BlobArrayBuilder;
7504 use lance_arrow::BLOB_DEDICATED_SIZE_THRESHOLD_META_KEY;
7505 use lance_core::utils::tempfile::TempDir;
7506
7507 let test_dir = TempDir::default();
7508
7509 let inline_data = b"small-inline-blob".as_slice();
7510 let packed_data: Vec<u8> = (0..64 * 1024 + 1024).map(|i| (i % 256) as u8).collect();
7511 let dedicated_data: Vec<u8> = (0..4 * 1024 * 1024 + 512)
7512 .map(|i| ((i + 97) % 256) as u8)
7513 .collect();
7514
7515 let mut blob_builder = BlobArrayBuilder::new(3);
7516 blob_builder.push_bytes(inline_data).unwrap();
7517 blob_builder.push_bytes(&packed_data).unwrap();
7518 blob_builder.push_bytes(&dedicated_data).unwrap();
7519 let blob_array: ArrayRef = blob_builder.finish().unwrap();
7520
7521 let id_array: ArrayRef = Arc::new(Int32Array::from(vec![0, 1, 2]));
7522 let mut blob_field = crate::blob_field("blob", true);
7523 {
7524 let metadata = blob_field.metadata().clone();
7525 let mut new_metadata = metadata;
7526 new_metadata.insert(
7527 BLOB_DEDICATED_SIZE_THRESHOLD_META_KEY.to_string(),
7528 (4 * 1024 * 1024).to_string(),
7529 );
7530 blob_field = blob_field.with_metadata(new_metadata);
7531 }
7532 let schema = Arc::new(Schema::new(vec![
7533 Field::new("id", DataType::Int32, false),
7534 blob_field,
7535 ]));
7536
7537 let batch = RecordBatch::try_new(schema.clone(), vec![id_array, blob_array]).unwrap();
7538 let reader = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema.clone());
7539
7540 let mut dataset = Dataset::write(
7541 reader,
7542 &test_dir.path_str(),
7543 Some(WriteParams {
7544 data_storage_version: Some(LanceFileVersion::V2_2),
7545 max_rows_per_file: 1,
7546 ..Default::default()
7547 }),
7548 )
7549 .await
7550 .unwrap();
7551
7552 assert_eq!(dataset.get_fragments().len(), 3);
7553
7554 compact_files(
7555 &mut dataset,
7556 CompactionOptions {
7557 target_rows_per_fragment: 1024 * 1024,
7558 ..Default::default()
7559 },
7560 None,
7561 )
7562 .await
7563 .unwrap();
7564
7565 assert_eq!(dataset.get_fragments().len(), 1);
7566
7567 let scan_result = dataset
7568 .scan()
7569 .project(&["id", "blob"])
7570 .unwrap()
7571 .try_into_batch()
7572 .await
7573 .unwrap();
7574 assert_eq!(scan_result.num_rows(), 3);
7575
7576 let ids = scan_result
7577 .column_by_name("id")
7578 .unwrap()
7579 .as_primitive::<Int32Type>();
7580 let id_values: Vec<i32> = ids.iter().map(|v| v.unwrap()).collect();
7581 assert_eq!(id_values, vec![0, 1, 2]);
7582
7583 let mut blob_values = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "blob").await;
7584 blob_values.sort_by_key(|(id, _)| *id);
7585 assert_eq!(
7586 blob_values,
7587 vec![
7588 (0, Some(inline_data.to_vec())),
7589 (1, Some(packed_data)),
7590 (2, Some(dedicated_data))
7591 ]
7592 );
7593 }
7594
7595 #[tokio::test]
7596 async fn test_compact_blob_v2_with_null_rows() {
7597 use crate::BlobArrayBuilder;
7598 use lance_core::utils::tempfile::TempDir;
7599
7600 let test_dir = TempDir::default();
7601
7602 let mut blob_builder = BlobArrayBuilder::new(4);
7603 blob_builder.push_bytes(b"inline-0").unwrap();
7604 blob_builder.push_null().unwrap();
7605 blob_builder.push_bytes(b"inline-2").unwrap();
7606 blob_builder.push_null().unwrap();
7607 let blob_array: ArrayRef = blob_builder.finish().unwrap();
7608
7609 let id_array: ArrayRef =
7610 Arc::new(Int32Array::from(vec![Some(0), Some(1), Some(2), Some(3)]));
7611 let schema = Arc::new(Schema::new(vec![
7612 Field::new("id", DataType::Int32, false),
7613 crate::blob_field("blob", true),
7614 ]));
7615
7616 let batch = RecordBatch::try_new(schema.clone(), vec![id_array, blob_array]).unwrap();
7617 let reader = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema.clone());
7618
7619 let mut dataset = Dataset::write(
7620 reader,
7621 &test_dir.path_str(),
7622 Some(WriteParams {
7623 data_storage_version: Some(LanceFileVersion::V2_2),
7624 max_rows_per_file: 2,
7625 ..Default::default()
7626 }),
7627 )
7628 .await
7629 .unwrap();
7630
7631 assert_eq!(dataset.get_fragments().len(), 2);
7632
7633 compact_files(
7634 &mut dataset,
7635 CompactionOptions {
7636 target_rows_per_fragment: 1024 * 1024,
7637 ..Default::default()
7638 },
7639 None,
7640 )
7641 .await
7642 .unwrap();
7643
7644 assert_eq!(dataset.get_fragments().len(), 1);
7645
7646 let scan_result = dataset
7647 .scan()
7648 .project(&["id", "blob"])
7649 .unwrap()
7650 .try_into_batch()
7651 .await
7652 .unwrap();
7653 assert_eq!(scan_result.num_rows(), 4);
7654
7655 let ids = scan_result
7656 .column_by_name("id")
7657 .unwrap()
7658 .as_primitive::<Int32Type>();
7659 let id_values: Vec<i32> = ids.iter().map(|v| v.unwrap()).collect();
7660 assert_eq!(id_values, vec![0, 1, 2, 3]);
7661
7662 let blob_col = scan_result.column_by_name("blob").unwrap();
7663 assert!(
7664 matches!(blob_col.data_type(), DataType::Struct(_)),
7665 "blob column should be a struct after compaction"
7666 );
7667
7668 let mut blob_values = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "blob").await;
7669 blob_values.sort_by_key(|(id, _)| *id);
7670 assert_eq!(
7671 blob_values,
7672 vec![
7673 (0, Some(b"inline-0".to_vec())),
7674 (1, None),
7675 (2, Some(b"inline-2".to_vec())),
7676 (3, None)
7677 ]
7678 );
7679 }
7680
7681 #[tokio::test]
7682 async fn test_compact_blob_v2_deleted_rows_not_resurrected() {
7683 use crate::BlobArrayBuilder;
7684 use lance_core::utils::tempfile::TempDir;
7685
7686 let test_dir = TempDir::default();
7687
7688 let mut blob_builder = BlobArrayBuilder::new(4);
7689 blob_builder.push_bytes(b"blob-0").unwrap();
7690 blob_builder.push_bytes(b"blob-1").unwrap();
7691 blob_builder.push_bytes(b"blob-2").unwrap();
7692 blob_builder.push_bytes(b"blob-3").unwrap();
7693 let blob_array: ArrayRef = blob_builder.finish().unwrap();
7694
7695 let id_array: ArrayRef = Arc::new(Int32Array::from(vec![0, 1, 2, 3]));
7696 let schema = Arc::new(Schema::new(vec![
7697 Field::new("id", DataType::Int32, false),
7698 crate::blob_field("blob", true),
7699 ]));
7700
7701 let batch = RecordBatch::try_new(schema.clone(), vec![id_array, blob_array]).unwrap();
7702 let reader = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema.clone());
7703
7704 let mut dataset = Dataset::write(
7705 reader,
7706 &test_dir.path_str(),
7707 Some(WriteParams {
7708 data_storage_version: Some(LanceFileVersion::V2_2),
7709 max_rows_per_file: 2,
7710 ..Default::default()
7711 }),
7712 )
7713 .await
7714 .unwrap();
7715
7716 assert_eq!(dataset.get_fragments().len(), 2);
7717
7718 dataset.delete("id = 1").await.unwrap();
7719 dataset.delete("id = 2").await.unwrap();
7720
7721 compact_files(
7722 &mut dataset,
7723 CompactionOptions {
7724 target_rows_per_fragment: 1024 * 1024,
7725 materialize_deletions_threshold: 0.0,
7726 ..Default::default()
7727 },
7728 None,
7729 )
7730 .await
7731 .unwrap();
7732
7733 let scan_result = dataset
7734 .scan()
7735 .project(&["id", "blob"])
7736 .unwrap()
7737 .try_into_batch()
7738 .await
7739 .unwrap();
7740 assert_eq!(scan_result.num_rows(), 2);
7741
7742 let ids = scan_result
7743 .column_by_name("id")
7744 .unwrap()
7745 .as_primitive::<Int32Type>();
7746 let mut id_values: Vec<i32> = ids.iter().map(|v| v.unwrap()).collect();
7747 id_values.sort();
7748 assert_eq!(id_values, vec![0, 3]);
7749
7750 let blob_col = scan_result.column_by_name("blob").unwrap();
7751 let struct_arr = blob_col.as_any().downcast_ref::<StructArray>().unwrap();
7752 let kind_col = struct_arr
7753 .column_by_name("kind")
7754 .unwrap()
7755 .as_primitive::<UInt8Type>();
7756
7757 for i in 0..kind_col.len() {
7758 assert!(
7759 !kind_col.is_null(i),
7760 "row {} should have a non-null kind after compaction of deleted rows",
7761 i
7762 );
7763 }
7764
7765 let mut blob_values = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "blob").await;
7766 blob_values.sort_by_key(|(id, _)| *id);
7767 assert_eq!(
7768 blob_values,
7769 vec![(0, Some(b"blob-0".to_vec())), (3, Some(b"blob-3".to_vec()))]
7770 );
7771 }
7772
7773 #[tokio::test]
7774 async fn test_compact_blob_v2_external_and_data_blob_mixed() {
7775 use crate::BlobArrayBuilder;
7776 use lance_arrow::BLOB_DEDICATED_SIZE_THRESHOLD_META_KEY;
7777 use lance_core::utils::tempfile::TempDir;
7778 use lance_table::format::BasePath;
7779
7780 let test_dir = TempDir::default();
7781 let external_dir = TempDir::default();
7782 let external_path = external_dir.std_path().join("external.bin");
7783 std::fs::write(&external_path, b"external-payload").unwrap();
7784 let external_uri = format!("file://{}", external_path.display());
7785 let base_uri = format!("file://{}", external_dir.std_path().display());
7786
7787 let packed_data: Vec<u8> = (0..64 * 1024 + 512).map(|i| (i % 256) as u8).collect();
7788
7789 let mut blob_builder = BlobArrayBuilder::new(4);
7790 blob_builder.push_uri(external_uri.clone()).unwrap();
7791 blob_builder.push_bytes(&packed_data).unwrap();
7792 blob_builder.push_bytes(b"inline-small").unwrap();
7793 blob_builder.push_uri(external_uri.clone()).unwrap();
7794 let blob_array: ArrayRef = blob_builder.finish().unwrap();
7795
7796 let id_array: ArrayRef = Arc::new(Int32Array::from(vec![0, 1, 2, 3]));
7797 let mut blob_field = crate::blob_field("blob", true);
7798 {
7799 let mut new_metadata = blob_field.metadata().clone();
7800 new_metadata.insert(
7801 BLOB_DEDICATED_SIZE_THRESHOLD_META_KEY.to_string(),
7802 (4 * 1024 * 1024).to_string(),
7803 );
7804 blob_field = blob_field.with_metadata(new_metadata);
7805 }
7806 let schema = Arc::new(Schema::new(vec![
7807 Field::new("id", DataType::Int32, false),
7808 blob_field,
7809 ]));
7810
7811 let batch = RecordBatch::try_new(schema.clone(), vec![id_array, blob_array]).unwrap();
7812 let reader = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema.clone());
7813
7814 let mut dataset = Dataset::write(
7815 reader,
7816 &test_dir.path_str(),
7817 Some(WriteParams {
7818 data_storage_version: Some(LanceFileVersion::V2_2),
7819 max_rows_per_file: 2,
7820 initial_bases: Some(vec![BasePath {
7821 id: 1,
7822 name: Some("external".to_string()),
7823 path: base_uri,
7824 is_dataset_root: false,
7825 }]),
7826 ..Default::default()
7827 }),
7828 )
7829 .await
7830 .unwrap();
7831
7832 assert_eq!(dataset.get_fragments().len(), 2);
7833
7834 compact_files(
7835 &mut dataset,
7836 CompactionOptions {
7837 target_rows_per_fragment: 1024 * 1024,
7838 ..Default::default()
7839 },
7840 None,
7841 )
7842 .await
7843 .unwrap();
7844
7845 assert_eq!(dataset.get_fragments().len(), 1);
7846
7847 let mut blob_values = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "blob").await;
7848 blob_values.sort_by_key(|(id, _)| *id);
7849 assert_eq!(
7850 blob_values,
7851 vec![
7852 (0, Some(b"external-payload".to_vec())),
7853 (1, Some(packed_data)),
7854 (2, Some(b"inline-small".to_vec())),
7855 (3, Some(b"external-payload".to_vec()))
7856 ]
7857 );
7858 }
7859
7860 #[tokio::test]
7861 async fn test_compact_blob_v2_multiple_blob_columns() {
7862 use crate::BlobArrayBuilder;
7863 use lance_core::utils::tempfile::TempDir;
7864
7865 let test_dir = TempDir::default();
7866
7867 let mut image_builder = BlobArrayBuilder::new(3);
7868 image_builder.push_bytes(b"image-0").unwrap();
7869 image_builder.push_bytes(b"image-1").unwrap();
7870 image_builder.push_bytes(b"image-2").unwrap();
7871 let image_array: ArrayRef = image_builder.finish().unwrap();
7872
7873 let mut thumb_builder = BlobArrayBuilder::new(3);
7874 thumb_builder.push_bytes(b"thumb-0").unwrap();
7875 thumb_builder.push_null().unwrap();
7876 thumb_builder.push_bytes(b"thumb-2").unwrap();
7877 let thumb_array: ArrayRef = thumb_builder.finish().unwrap();
7878
7879 let id_array: ArrayRef = Arc::new(Int32Array::from(vec![0, 1, 2]));
7880 let schema = Arc::new(Schema::new(vec![
7881 Field::new("id", DataType::Int32, false),
7882 crate::blob_field("image", true),
7883 crate::blob_field("thumbnail", true),
7884 ]));
7885
7886 let batch =
7887 RecordBatch::try_new(schema.clone(), vec![id_array, image_array, thumb_array]).unwrap();
7888 let reader = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema.clone());
7889
7890 let mut dataset = Dataset::write(
7891 reader,
7892 &test_dir.path_str(),
7893 Some(WriteParams {
7894 data_storage_version: Some(LanceFileVersion::V2_2),
7895 max_rows_per_file: 1,
7896 ..Default::default()
7897 }),
7898 )
7899 .await
7900 .unwrap();
7901
7902 assert_eq!(dataset.get_fragments().len(), 3);
7903
7904 compact_files(
7905 &mut dataset,
7906 CompactionOptions {
7907 target_rows_per_fragment: 1024 * 1024,
7908 ..Default::default()
7909 },
7910 None,
7911 )
7912 .await
7913 .unwrap();
7914
7915 assert_eq!(dataset.get_fragments().len(), 1);
7916
7917 let mut image_values = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "image").await;
7918 image_values.sort_by_key(|(id, _)| *id);
7919 assert_eq!(
7920 image_values,
7921 vec![
7922 (0, Some(b"image-0".to_vec())),
7923 (1, Some(b"image-1".to_vec())),
7924 (2, Some(b"image-2".to_vec()))
7925 ]
7926 );
7927
7928 let mut thumb_values =
7929 read_blob_bytes_by_index(&Arc::new(dataset.clone()), "thumbnail").await;
7930 thumb_values.sort_by_key(|(id, _)| *id);
7931 assert_eq!(
7932 thumb_values,
7933 vec![
7934 (0, Some(b"thumb-0".to_vec())),
7935 (1, None),
7936 (2, Some(b"thumb-2".to_vec()))
7937 ]
7938 );
7939 }
7940
7941 #[tokio::test]
7942 async fn test_compact_blob_v2_external_and_null_mixed() {
7943 use crate::BlobArrayBuilder;
7944 use lance_core::utils::tempfile::TempDir;
7945 use lance_table::format::BasePath;
7946
7947 let test_dir = TempDir::default();
7948 let external_dir = TempDir::default();
7949 let external_path = external_dir.std_path().join("mixed-external.bin");
7950 std::fs::write(&external_path, b"external-mixed-data").unwrap();
7951 let external_uri = format!("file://{}", external_path.display());
7952 let base_uri = format!("file://{}", external_dir.std_path().display());
7953
7954 let mut blob_builder = BlobArrayBuilder::new(4);
7955 blob_builder.push_uri(external_uri.clone()).unwrap();
7956 blob_builder.push_null().unwrap();
7957 blob_builder.push_uri(external_uri.clone()).unwrap();
7958 blob_builder.push_null().unwrap();
7959 let blob_array: ArrayRef = blob_builder.finish().unwrap();
7960
7961 let id_array: ArrayRef = Arc::new(Int32Array::from(vec![0, 1, 2, 3]));
7962 let schema = Arc::new(Schema::new(vec![
7963 Field::new("id", DataType::Int32, false),
7964 crate::blob_field("blob", true),
7965 ]));
7966
7967 let batch = RecordBatch::try_new(schema.clone(), vec![id_array, blob_array]).unwrap();
7968 let reader = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema.clone());
7969
7970 let mut dataset = Dataset::write(
7971 reader,
7972 &test_dir.path_str(),
7973 Some(WriteParams {
7974 data_storage_version: Some(LanceFileVersion::V2_2),
7975 max_rows_per_file: 2,
7976 initial_bases: Some(vec![BasePath {
7977 id: 1,
7978 name: Some("external".to_string()),
7979 path: base_uri,
7980 is_dataset_root: false,
7981 }]),
7982 ..Default::default()
7983 }),
7984 )
7985 .await
7986 .unwrap();
7987
7988 assert_eq!(dataset.get_fragments().len(), 2);
7989
7990 compact_files(
7991 &mut dataset,
7992 CompactionOptions {
7993 target_rows_per_fragment: 1024 * 1024,
7994 ..Default::default()
7995 },
7996 None,
7997 )
7998 .await
7999 .unwrap();
8000
8001 assert_eq!(dataset.get_fragments().len(), 1);
8002
8003 let mut blob_values = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "blob").await;
8004 blob_values.sort_by_key(|(id, _)| *id);
8005 assert_eq!(
8006 blob_values,
8007 vec![
8008 (0, Some(b"external-mixed-data".to_vec())),
8009 (1, None),
8010 (2, Some(b"external-mixed-data".to_vec())),
8011 (3, None)
8012 ]
8013 );
8014 }
8015
8016 #[tokio::test]
8017 async fn test_compact_blob_v2_all_null_and_all_external_fragments() {
8018 use crate::BlobArrayBuilder;
8019 use lance_core::utils::tempfile::TempDir;
8020 use lance_table::format::BasePath;
8021
8022 let test_dir = TempDir::default();
8023 let external_dir = TempDir::default();
8024 let external_path = external_dir.std_path().join("all-ext.bin");
8025 std::fs::write(&external_path, b"all-external-data").unwrap();
8026 let external_uri = format!("file://{}", external_path.display());
8027 let base_uri = format!("file://{}", external_dir.std_path().display());
8028
8029 let mut null_builder = BlobArrayBuilder::new(2);
8030 null_builder.push_null().unwrap();
8031 null_builder.push_null().unwrap();
8032 let null_array: ArrayRef = null_builder.finish().unwrap();
8033
8034 let mut ext_builder = BlobArrayBuilder::new(2);
8035 ext_builder.push_uri(external_uri.clone()).unwrap();
8036 ext_builder.push_uri(external_uri.clone()).unwrap();
8037 let ext_array: ArrayRef = ext_builder.finish().unwrap();
8038
8039 let id_null_array: ArrayRef = Arc::new(Int32Array::from(vec![0, 1]));
8040 let null_schema = Arc::new(Schema::new(vec![
8041 Field::new("id", DataType::Int32, false),
8042 crate::blob_field("blob", true),
8043 ]));
8044 let null_batch =
8045 RecordBatch::try_new(null_schema.clone(), vec![id_null_array, null_array]).unwrap();
8046
8047 let id_ext_array: ArrayRef = Arc::new(Int32Array::from(vec![2, 3]));
8048 let ext_schema = Arc::new(Schema::new(vec![
8049 Field::new("id", DataType::Int32, false),
8050 crate::blob_field("blob", true),
8051 ]));
8052 let ext_batch =
8053 RecordBatch::try_new(ext_schema.clone(), vec![id_ext_array, ext_array]).unwrap();
8054
8055 let mut dataset = Dataset::write(
8056 RecordBatchIterator::new(
8057 vec![null_batch, ext_batch].into_iter().map(Ok),
8058 null_schema.clone(),
8059 ),
8060 &test_dir.path_str(),
8061 Some(WriteParams {
8062 data_storage_version: Some(LanceFileVersion::V2_2),
8063 max_rows_per_file: 2,
8064 initial_bases: Some(vec![BasePath {
8065 id: 1,
8066 name: Some("external".to_string()),
8067 path: base_uri,
8068 is_dataset_root: false,
8069 }]),
8070 ..Default::default()
8071 }),
8072 )
8073 .await
8074 .unwrap();
8075
8076 assert_eq!(dataset.get_fragments().len(), 2);
8077
8078 compact_files(
8079 &mut dataset,
8080 CompactionOptions {
8081 target_rows_per_fragment: 1024 * 1024,
8082 ..Default::default()
8083 },
8084 None,
8085 )
8086 .await
8087 .unwrap();
8088
8089 assert_eq!(dataset.get_fragments().len(), 1);
8090
8091 let mut blob_values = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "blob").await;
8092 blob_values.sort_by_key(|(id, _)| *id);
8093 assert_eq!(
8094 blob_values,
8095 vec![
8096 (0, None),
8097 (1, None),
8098 (2, Some(b"all-external-data".to_vec())),
8099 (3, Some(b"all-external-data".to_vec()))
8100 ]
8101 );
8102 }
8103
8104 #[tokio::test]
8105 async fn test_compact_blob_v2_external_with_multiple_base_ids() {
8106 use crate::BlobArrayBuilder;
8107 use lance_core::utils::tempfile::TempDir;
8108 use lance_table::format::BasePath;
8109
8110 let test_dir = TempDir::default();
8111 let base_a_dir = TempDir::default();
8112 let base_b_dir = TempDir::default();
8113
8114 let path_a = base_a_dir.std_path().join("data-a.bin");
8115 std::fs::write(&path_a, b"from-base-a").unwrap();
8116 let uri_a = format!("file://{}", path_a.display());
8117 let base_uri_a = format!("file://{}", base_a_dir.std_path().display());
8118
8119 let path_b = base_b_dir.std_path().join("data-b.bin");
8120 std::fs::write(&path_b, b"from-base-b").unwrap();
8121 let uri_b = format!("file://{}", path_b.display());
8122 let base_uri_b = format!("file://{}", base_b_dir.std_path().display());
8123
8124 let mut blob_builder = BlobArrayBuilder::new(4);
8125 blob_builder.push_uri(uri_a.clone()).unwrap();
8126 blob_builder.push_uri(uri_b).unwrap();
8127 blob_builder.push_bytes(b"inline-data").unwrap();
8128 blob_builder.push_uri(uri_a).unwrap();
8129 let blob_array: ArrayRef = blob_builder.finish().unwrap();
8130
8131 let id_array: ArrayRef = Arc::new(Int32Array::from(vec![0, 1, 2, 3]));
8132 let schema = Arc::new(Schema::new(vec![
8133 Field::new("id", DataType::Int32, false),
8134 crate::blob_field("blob", true),
8135 ]));
8136
8137 let batch = RecordBatch::try_new(schema.clone(), vec![id_array, blob_array]).unwrap();
8138 let reader = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema.clone());
8139
8140 let mut dataset = Dataset::write(
8141 reader,
8142 &test_dir.path_str(),
8143 Some(WriteParams {
8144 data_storage_version: Some(LanceFileVersion::V2_2),
8145 max_rows_per_file: 2,
8146 initial_bases: Some(vec![
8147 BasePath {
8148 id: 1,
8149 name: Some("base_a".to_string()),
8150 path: base_uri_a,
8151 is_dataset_root: false,
8152 },
8153 BasePath {
8154 id: 2,
8155 name: Some("base_b".to_string()),
8156 path: base_uri_b,
8157 is_dataset_root: false,
8158 },
8159 ]),
8160 ..Default::default()
8161 }),
8162 )
8163 .await
8164 .unwrap();
8165
8166 assert_eq!(dataset.get_fragments().len(), 2);
8167
8168 compact_files(
8169 &mut dataset,
8170 CompactionOptions {
8171 target_rows_per_fragment: 1024 * 1024,
8172 ..Default::default()
8173 },
8174 None,
8175 )
8176 .await
8177 .unwrap();
8178
8179 assert_eq!(dataset.get_fragments().len(), 1);
8180
8181 let mut blob_values = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "blob").await;
8182 blob_values.sort_by_key(|(id, _)| *id);
8183 assert_eq!(
8184 blob_values,
8185 vec![
8186 (0, Some(b"from-base-a".to_vec())),
8187 (1, Some(b"from-base-b".to_vec())),
8188 (2, Some(b"inline-data".to_vec())),
8189 (3, Some(b"from-base-a".to_vec()))
8190 ]
8191 );
8192 }
8193
8194 #[tokio::test]
8195 async fn test_compact_blob_v2_large_blobs() {
8196 use crate::BlobArrayBuilder;
8197 use lance_core::utils::tempfile::TempDir;
8198
8199 let test_dir = TempDir::default();
8200
8201 let large_blob_a: Vec<u8> = (0..512 * 1024).map(|i| (i % 256) as u8).collect();
8202 let large_blob_b: Vec<u8> = (0..256 * 1024).map(|i| ((i + 42) % 256) as u8).collect();
8203
8204 let mut blob_builder = BlobArrayBuilder::new(3);
8205 blob_builder.push_bytes(&large_blob_a).unwrap();
8206 blob_builder.push_bytes(&large_blob_b).unwrap();
8207 blob_builder.push_bytes(b"small-blob").unwrap();
8208 let blob_array: ArrayRef = blob_builder.finish().unwrap();
8209
8210 let id_array: ArrayRef = Arc::new(Int32Array::from(vec![0, 1, 2]));
8211 let schema = Arc::new(Schema::new(vec![
8212 Field::new("id", DataType::Int32, false),
8213 crate::blob_field("blob", true),
8214 ]));
8215
8216 let batch = RecordBatch::try_new(schema.clone(), vec![id_array, blob_array]).unwrap();
8217 let reader = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema.clone());
8218
8219 let mut dataset = Dataset::write(
8220 reader,
8221 &test_dir.path_str(),
8222 Some(WriteParams {
8223 data_storage_version: Some(LanceFileVersion::V2_2),
8224 max_rows_per_file: 1,
8225 ..Default::default()
8226 }),
8227 )
8228 .await
8229 .unwrap();
8230
8231 assert_eq!(dataset.get_fragments().len(), 3);
8232
8233 compact_files(
8234 &mut dataset,
8235 CompactionOptions {
8236 target_rows_per_fragment: 1024 * 1024,
8237 ..Default::default()
8238 },
8239 None,
8240 )
8241 .await
8242 .unwrap();
8243
8244 assert_eq!(dataset.get_fragments().len(), 1);
8245
8246 let mut blob_values = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "blob").await;
8247 blob_values.sort_by_key(|(id, _)| *id);
8248 assert_eq!(
8249 blob_values,
8250 vec![
8251 (0, Some(large_blob_a)),
8252 (1, Some(large_blob_b)),
8253 (2, Some(b"small-blob".to_vec()))
8254 ]
8255 );
8256 }
8257
8258 #[tokio::test]
8259 async fn test_compact_blob_v2_blob_kind_reclassification() {
8260 use crate::BlobArrayBuilder;
8261 use lance_arrow::BLOB_DEDICATED_SIZE_THRESHOLD_META_KEY;
8262 use lance_core::utils::tempfile::TempDir;
8263
8264 let test_dir = TempDir::default();
8265
8266 let medium_data: Vec<u8> = (0..32 * 1024).map(|i| (i % 256) as u8).collect();
8267
8268 let mut blob_builder = BlobArrayBuilder::new(2);
8269 blob_builder.push_bytes(&medium_data).unwrap();
8270 blob_builder.push_bytes(&medium_data).unwrap();
8271 let blob_array: ArrayRef = blob_builder.finish().unwrap();
8272
8273 let id_array: ArrayRef = Arc::new(Int32Array::from(vec![0, 1]));
8274 let mut blob_field = crate::blob_field("blob", true);
8275 {
8276 let mut new_metadata = blob_field.metadata().clone();
8277 new_metadata.insert(
8278 BLOB_DEDICATED_SIZE_THRESHOLD_META_KEY.to_string(),
8279 (16 * 1024).to_string(),
8280 );
8281 blob_field = blob_field.with_metadata(new_metadata);
8282 }
8283 let schema = Arc::new(Schema::new(vec![
8284 Field::new("id", DataType::Int32, false),
8285 blob_field,
8286 ]));
8287
8288 let batch = RecordBatch::try_new(schema.clone(), vec![id_array, blob_array]).unwrap();
8289 let reader = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema.clone());
8290
8291 let mut dataset = Dataset::write(
8292 reader,
8293 &test_dir.path_str(),
8294 Some(WriteParams {
8295 data_storage_version: Some(LanceFileVersion::V2_2),
8296 max_rows_per_file: 1,
8297 ..Default::default()
8298 }),
8299 )
8300 .await
8301 .unwrap();
8302
8303 assert_eq!(dataset.get_fragments().len(), 2);
8304
8305 compact_files(
8306 &mut dataset,
8307 CompactionOptions {
8308 target_rows_per_fragment: 1024 * 1024,
8309 ..Default::default()
8310 },
8311 None,
8312 )
8313 .await
8314 .unwrap();
8315
8316 assert_eq!(dataset.get_fragments().len(), 1);
8317
8318 let mut blob_values = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "blob").await;
8319 blob_values.sort_by_key(|(id, _)| *id);
8320 assert_eq!(
8321 blob_values,
8322 vec![
8323 (0, Some(medium_data.clone())),
8324 (1, Some(medium_data.clone()))
8325 ]
8326 );
8327 }
8328
8329 #[tokio::test]
8330 async fn test_compact_blob_v2_multi_batch() {
8331 use crate::BlobArrayBuilder;
8332 use lance_core::utils::tempfile::TempDir;
8333
8334 let test_dir = TempDir::default();
8335
8336 let mut blob_builder = BlobArrayBuilder::new(6);
8337 blob_builder.push_bytes(b"batch-0-row-0").unwrap();
8338 blob_builder.push_bytes(b"batch-0-row-1").unwrap();
8339 blob_builder.push_bytes(b"batch-1-row-0").unwrap();
8340 blob_builder.push_null().unwrap();
8341 blob_builder.push_bytes(b"batch-1-row-2").unwrap();
8342 blob_builder.push_bytes(b"batch-1-row-3").unwrap();
8343 let blob_array: ArrayRef = blob_builder.finish().unwrap();
8344
8345 let id_array: ArrayRef = Arc::new(Int32Array::from(vec![0, 1, 2, 3, 4, 5]));
8346 let schema = Arc::new(Schema::new(vec![
8347 Field::new("id", DataType::Int32, false),
8348 crate::blob_field("blob", true),
8349 ]));
8350
8351 let batch = RecordBatch::try_new(schema.clone(), vec![id_array, blob_array]).unwrap();
8352 let reader = RecordBatchIterator::new(vec![batch].into_iter().map(Ok), schema.clone());
8353
8354 let mut dataset = Dataset::write(
8355 reader,
8356 &test_dir.path_str(),
8357 Some(WriteParams {
8358 data_storage_version: Some(LanceFileVersion::V2_2),
8359 max_rows_per_file: 2,
8360 ..Default::default()
8361 }),
8362 )
8363 .await
8364 .unwrap();
8365
8366 assert_eq!(dataset.get_fragments().len(), 3);
8367
8368 compact_files(
8369 &mut dataset,
8370 CompactionOptions {
8371 target_rows_per_fragment: 1024 * 1024,
8372 batch_size: Some(2),
8373 ..Default::default()
8374 },
8375 None,
8376 )
8377 .await
8378 .unwrap();
8379
8380 assert_eq!(dataset.get_fragments().len(), 1);
8381
8382 let mut blob_values = read_blob_bytes_by_index(&Arc::new(dataset.clone()), "blob").await;
8383 blob_values.sort_by_key(|(id, _)| *id);
8384 assert_eq!(
8385 blob_values,
8386 vec![
8387 (0, Some(b"batch-0-row-0".to_vec())),
8388 (1, Some(b"batch-0-row-1".to_vec())),
8389 (2, Some(b"batch-1-row-0".to_vec())),
8390 (3, None),
8391 (4, Some(b"batch-1-row-2".to_vec())),
8392 (5, Some(b"batch-1-row-3".to_vec()))
8393 ]
8394 );
8395 }
8396 use arrow_array::record_batch;
8402 use lance_file::writer::{FileWriter, FileWriterOptions};
8403 use lance_io::utils::CachedFileSize;
8404 use lance_table::format::DataFile;
8405 use lance_table::format::overlay::{DataOverlayFile, OverlayCoverage};
8406 use std::collections::BTreeMap;
8407
8408 use crate::dataset::DATA_DIR;
8409 use crate::dataset::transaction::DataOverlayGroup;
8410
8411 async fn create_base_dataset(uri: &str) -> Dataset {
8414 let batch = record_batch!(
8415 ("id", Int32, (0..12).collect::<Vec<_>>()),
8416 ("val", Int32, (0..12).map(|v| v * 10).collect::<Vec<_>>())
8417 )
8418 .unwrap();
8419 let schema = batch.schema();
8420 let write_params = WriteParams {
8421 max_rows_per_file: 6,
8422 max_rows_per_group: 6,
8423 data_storage_version: Some(LanceFileVersion::Stable),
8424 ..Default::default()
8425 };
8426 let reader = RecordBatchIterator::new(vec![Ok(batch)], schema.clone());
8427 Dataset::write(reader, uri, Some(write_params))
8428 .await
8429 .unwrap()
8430 }
8431
8432 fn i32_array(values: impl IntoIterator<Item = Option<i32>>) -> ArrayRef {
8433 Arc::new(Int32Array::from_iter(values))
8434 }
8435
8436 fn bitmap(offsets: impl IntoIterator<Item = u32>) -> RoaringBitmap {
8437 RoaringBitmap::from_iter(offsets)
8438 }
8439
8440 async fn commit_overlay(
8443 dataset: Dataset,
8444 fragment_id: u64,
8445 fields: &[i32],
8446 coverage: OverlayCoverage,
8447 columns: Vec<ArrayRef>,
8448 ) -> Dataset {
8449 let read_version = dataset.version().version;
8450 let overlay_schema = dataset.schema().project_by_ids(fields, true);
8451 let filename = format!("{}.lance", Uuid::new_v4());
8452 let path = dataset.base.clone().join(DATA_DIR).join(filename.as_str());
8453 let obj_writer = dataset.object_store.create(&path).await.unwrap();
8454 let mut writer = FileWriter::try_new(
8455 obj_writer,
8456 overlay_schema,
8457 FileWriterOptions {
8458 format_version: Some(LanceFileVersion::Stable),
8459 ..Default::default()
8460 },
8461 )
8462 .unwrap();
8463 let file_version = lance_file::version::ConcreteFileVersion::from(writer.version());
8464 for (column_index, array) in columns.into_iter().enumerate() {
8465 writer.write_column(column_index, array).await.unwrap();
8466 }
8467 let summary = writer.finish().await.unwrap();
8468
8469 let mut data_file = DataFile::new_unstarted(filename, file_version);
8470 data_file.fields = writer
8471 .field_id_to_column_indices()
8472 .iter()
8473 .map(|(f, _)| *f as i32)
8474 .collect::<Vec<_>>()
8475 .into();
8476 data_file.column_indices = writer
8477 .field_id_to_column_indices()
8478 .iter()
8479 .map(|(_, c)| *c as i32)
8480 .collect::<Vec<_>>()
8481 .into();
8482 data_file.file_size_bytes = CachedFileSize::new(summary.size_bytes);
8483
8484 Dataset::commit(
8485 WriteDestination::Dataset(Arc::new(dataset)),
8486 Operation::DataOverlay {
8487 groups: vec![DataOverlayGroup {
8488 fragment_id,
8489 overlays: vec![DataOverlayFile {
8490 data_file,
8491 coverage,
8492 committed_version: 0,
8493 }],
8494 }],
8495 },
8496 Some(read_version),
8497 None,
8498 None,
8499 Arc::new(Default::default()),
8500 false,
8501 )
8502 .await
8503 .unwrap()
8504 }
8505
8506 async fn commit_n_overlays(mut dataset: Dataset, n: u32) -> Dataset {
8510 for i in 0..n {
8511 dataset = commit_overlay(
8512 dataset,
8513 0,
8514 &[1],
8515 OverlayCoverage::dense(bitmap([i])),
8516 vec![i32_array([Some(1000 + i as i32)])],
8517 )
8518 .await;
8519 }
8520 dataset
8521 }
8522
8523 fn overlay_only_options(max_overlays_per_fragment: usize) -> CompactionOptions {
8527 CompactionOptions {
8528 max_overlays_per_fragment: Some(max_overlays_per_fragment),
8529 target_rows_per_fragment: 6,
8530 ..Default::default()
8531 }
8532 }
8533
8534 async fn id_val_map(dataset: &Dataset) -> BTreeMap<i32, Option<i32>> {
8536 let mut scanner = dataset.scan();
8537 scanner.project(&["id", "val"]).unwrap();
8538 let batch = scanner.try_into_batch().await.unwrap();
8539 let mut out = BTreeMap::new();
8540 let ids = batch
8541 .column(0)
8542 .as_any()
8543 .downcast_ref::<Int32Array>()
8544 .unwrap();
8545 let vals = batch
8546 .column(1)
8547 .as_any()
8548 .downcast_ref::<Int32Array>()
8549 .unwrap();
8550 for i in 0..batch.num_rows() {
8551 let v = if vals.is_null(i) {
8552 None
8553 } else {
8554 Some(vals.value(i))
8555 };
8556 out.insert(ids.value(i), v);
8557 }
8558 out
8559 }
8560
8561 #[tokio::test]
8562 async fn test_max_overlays_triggers_full_compaction() {
8563 let dataset = create_base_dataset("memory://").await;
8565 let mut dataset = commit_n_overlays(dataset, 3).await;
8566 assert_eq!(
8567 dataset.get_fragment(0).unwrap().metadata().overlays.len(),
8568 3
8569 );
8570
8571 let metrics = compact_files(&mut dataset, overlay_only_options(2), None)
8573 .await
8574 .unwrap();
8575 assert_eq!(metrics.fragments_removed, 1);
8576 assert_eq!(metrics.fragments_added, 1);
8577
8578 let fragments = dataset.get_fragments();
8579 assert_eq!(fragments.len(), 2);
8580 let compacted = fragments
8583 .iter()
8584 .find(|f| f.id() != 1)
8585 .expect("a new fragment id was assigned");
8586 assert!(compacted.metadata().overlays.is_empty());
8587 assert_eq!(compacted.metadata().files.len(), 1);
8588
8589 let values = id_val_map(&dataset).await;
8591 let expected: BTreeMap<i32, Option<i32>> = (0..12)
8592 .map(|id| {
8593 let v = if id < 3 { 1000 + id } else { id * 10 };
8594 (id, Some(v))
8595 })
8596 .collect();
8597 assert_eq!(values, expected);
8598 }
8599
8600 #[tokio::test]
8601 async fn test_below_threshold_is_a_noop() {
8602 let dataset = create_base_dataset("memory://").await;
8603 let mut dataset = commit_n_overlays(dataset, 2).await;
8604
8605 let metrics = compact_files(&mut dataset, overlay_only_options(2), None)
8607 .await
8608 .unwrap();
8609 assert_eq!(metrics.fragments_removed, 0);
8610 assert_eq!(metrics.fragments_added, 0);
8611 assert_eq!(
8612 dataset.get_fragment(0).unwrap().metadata().overlays.len(),
8613 2
8614 );
8615 }
8616
8617 #[tokio::test]
8618 async fn test_overlay_compaction_materializes_deletions() {
8619 let dataset = create_base_dataset("memory://").await;
8620 let mut dataset = commit_n_overlays(dataset, 3).await;
8621 dataset.delete("id = 2").await.unwrap();
8623 assert!(
8624 dataset
8625 .get_fragment(0)
8626 .unwrap()
8627 .metadata()
8628 .deletion_file
8629 .is_some()
8630 );
8631
8632 compact_files(&mut dataset, overlay_only_options(2), None)
8633 .await
8634 .unwrap();
8635
8636 for fragment in dataset.get_fragments() {
8638 assert!(fragment.metadata().deletion_file.is_none());
8639 assert!(fragment.metadata().overlays.is_empty());
8640 }
8641 let values = id_val_map(&dataset).await;
8642 assert!(!values.contains_key(&2));
8643 assert_eq!(values.get(&0), Some(&Some(1000)));
8645 assert_eq!(values.get(&1), Some(&Some(1001)));
8646 }
8647
8648 #[tokio::test]
8649 async fn test_overlay_compaction_reconciles_stale_index() {
8650 let mut dataset = create_base_dataset("memory://").await;
8651 dataset
8653 .create_index(
8654 &["val"],
8655 IndexType::Scalar,
8656 None,
8657 &ScalarIndexParams::default(),
8658 true,
8659 )
8660 .await
8661 .unwrap();
8662
8663 let mut dataset = commit_n_overlays(dataset, 3).await;
8666
8667 let val_index_before = dataset
8668 .load_indices()
8669 .await
8670 .unwrap()
8671 .iter()
8672 .find(|i| i.fields == vec![1])
8673 .expect("val index present")
8674 .clone();
8675 assert!(
8676 val_index_before
8677 .fragment_bitmap
8678 .as_ref()
8679 .unwrap()
8680 .contains(0)
8681 );
8682
8683 compact_files(&mut dataset, overlay_only_options(2), None)
8684 .await
8685 .unwrap();
8686
8687 let indices = dataset.load_indices().await.unwrap();
8690 let val_index = indices
8691 .iter()
8692 .find(|i| i.fields == vec![1])
8693 .expect("val index present");
8694 let compacted_id = dataset
8695 .get_fragments()
8696 .iter()
8697 .map(|f| f.id() as u32)
8698 .find(|id| *id != 1)
8699 .unwrap();
8700 assert!(
8701 !val_index
8702 .fragment_bitmap
8703 .as_ref()
8704 .unwrap()
8705 .contains(compacted_id),
8706 "stale index must drop the compacted fragment from its coverage"
8707 );
8708
8709 let mut scanner = dataset.scan();
8712 scanner
8713 .filter("val = 1000")
8714 .unwrap()
8715 .project(&["id"])
8716 .unwrap();
8717 let batch = scanner.try_into_batch().await.unwrap();
8718 let ids = batch
8719 .column(0)
8720 .as_any()
8721 .downcast_ref::<Int32Array>()
8722 .unwrap();
8723 assert_eq!(ids.len(), 1);
8724 assert_eq!(ids.value(0), 0);
8725
8726 let mut scanner = dataset.scan();
8727 scanner.filter("val = 0").unwrap().project(&["id"]).unwrap();
8728 let batch = scanner.try_into_batch().await.unwrap();
8729 assert_eq!(batch.num_rows(), 0, "stale value 0 must no longer match");
8730 }
8731}