diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index 08d2ab2913f..7b5711106d5 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -89,7 +89,10 @@ use std::sync::Arc; use super::fragment::FileFragment; use super::index::{DatasetIndexRemapperOptions, load_indices_for_remapping}; -use super::rowids::load_row_id_sequences; +use super::rowids::RowVersionKind; +use super::rowids::{ + RowLineage, load_row_id_sequences, load_row_version_sequence, place_row_lineage, +}; use super::transaction::{ Operation, RewriteGroup, RewrittenIndex, Transaction, TransactionBuilder, }; @@ -130,7 +133,7 @@ use lance_file::version::{ConcreteFileVersion, LanceFileVersion}; use lance_index::frag_reuse::{FRAG_REUSE_INDEX_NAME, FragReuseGroup}; use lance_index::is_system_index; use lance_index::metrics::NoOpMetricsCollector; -use lance_table::format::{Fragment, IndexMetadata, RowIdMeta}; +use lance_table::format::{Fragment, IndexMetadata, RowDatasetVersionSequence}; use roaring::{RoaringBitmap, RoaringTreemap}; use serde::{Deserialize, Serialize}; use tracing::{info, warn}; @@ -2626,13 +2629,7 @@ async fn rewrite_files( } else { if dataset.manifest.uses_stable_row_ids() { log::info!("Compaction task {}: rechunking stable row ids", task_id); - rechunk_stable_row_ids(dataset.as_ref(), &mut new_fragments, &fragments).await?; - recalc_versions_for_rewritten_fragments( - dataset.as_ref(), - &mut new_fragments, - &fragments, - ) - .await?; + rechunk_row_lineage(dataset.as_ref(), &mut new_fragments, &fragments).await?; } Ok(None) } @@ -2671,7 +2668,11 @@ async fn rewrite_files( }) } -async fn rechunk_stable_row_ids( +/// Carry the stable row ids and per-row versions of `old_fragments` over to +/// `new_fragments`, which hold the same live rows in the same order, and place +/// each new fragment's sequences inline or in a spilled column as the table's +/// spill policy and their size call for. +async fn rechunk_row_lineage( dataset: &Dataset, new_fragments: &mut [Fragment], old_fragments: &[Fragment], @@ -2687,146 +2688,79 @@ async fn rechunk_stable_row_ids( .expect("Fragment not found") }); + // Load old per-row version sequences, defaulting the way readers do: + // created-at to version 1, last-updated-at to created-at. + let mut old_created_at_sequences = Vec::with_capacity(old_fragments.len()); + let mut old_last_updated_sequences = Vec::with_capacity(old_fragments.len()); + for (frag, (_, row_ids)) in old_fragments.iter().zip(old_sequences.iter()) { + let created_at = + match load_row_version_sequence(dataset, frag, RowVersionKind::CreatedAt).await? { + Some(sequence) => sequence.as_ref().clone(), + None => RowDatasetVersionSequence::from_uniform_row_count(row_ids.len(), 1), + }; + let last_updated_at = + match load_row_version_sequence(dataset, frag, RowVersionKind::LastUpdatedAt).await? { + Some(sequence) => sequence.as_ref().clone(), + None => created_at.clone(), + }; + old_created_at_sequences.push(created_at); + old_last_updated_sequences.push(last_updated_at); + } + // Need to remove deleted rows - futures::stream::iter(old_sequences.iter_mut().zip(old_fragments.iter())) - .map(Ok) - .try_for_each(|((_, seq), frag)| async move { - if let Some(deletion_file) = &frag.deletion_file { - let deletions = read_dataset_deletion_file(dataset, frag.id, deletion_file).await?; - - let mut new_seq = seq.as_ref().clone(); - new_seq.mask(deletions.to_sorted_iter())?; - *seq = Arc::new(new_seq); - } - Ok::<(), crate::Error>(()) - }) - .await?; + for (index, frag) in old_fragments.iter().enumerate() { + if let Some(deletion_file) = &frag.deletion_file { + let deletions = read_dataset_deletion_file(dataset, frag.id, deletion_file).await?; + + let mut new_seq = old_sequences[index].1.as_ref().clone(); + new_seq.mask(deletions.to_sorted_iter())?; + old_sequences[index].1 = Arc::new(new_seq); + old_created_at_sequences[index].mask(deletions.to_sorted_iter())?; + old_last_updated_sequences[index].mask(deletions.to_sorted_iter())?; + } + } + let chunk_sizes: Vec = new_fragments + .iter() + .map(|frag| frag.physical_rows.unwrap() as u64) + .collect(); debug_assert_eq!( { old_sequences.iter().map(|(_, seq)| seq.len()).sum::() }, - { - new_fragments - .iter() - .map(|frag| frag.physical_rows.unwrap() as u64) - .sum::() - }, + { chunk_sizes.iter().sum::() }, "{:?}", old_sequences ); - let new_sequences = lance_table::rowids::rechunk_sequences( + let new_row_ids = lance_table::rowids::rechunk_sequences( old_sequences .into_iter() .map(|(_, seq)| seq.as_ref().clone()), - new_fragments - .iter() - .map(|frag| frag.physical_rows.unwrap() as u64), + chunk_sizes.iter().copied(), false, )?; - - for (fragment, sequence) in new_fragments.iter_mut().zip(new_sequences) { - // TODO: if large enough, serialize to separate file - let serialized = lance_table::rowids::write_row_ids(&sequence); - fragment.row_id_meta = Some(RowIdMeta::Inline(serialized.into())); - } - - Ok(()) -} - -/// After row id rechunking, preserve per-row latest update versions by masking deletions and rechunking -async fn recalc_versions_for_rewritten_fragments( - dataset: &Dataset, - new_fragments: &mut [Fragment], - old_fragments: &[Fragment], -) -> Result<()> { - // Load old per-row last_updated_at version sequences - let mut old_last_updated_sequences: Vec = - Vec::with_capacity(old_fragments.len()); - // Load old per-row created_at version sequences - let mut old_created_at_sequences: Vec = - Vec::with_capacity(old_fragments.len()); - - for frag in old_fragments.iter() { - let row_count = if let Some(row_id_meta) = &frag.row_id_meta { - match row_id_meta { - RowIdMeta::Inline(data) => lance_table::rowids::read_row_ids(data)?.len(), - RowIdMeta::Column => frag.physical_rows.unwrap_or(0) as u64, - } - } else { - frag.physical_rows.unwrap_or(0) as u64 - }; - - // Load created_at sequence (default to version 1 if missing) - let mut created_at_seq = if let Some(version_meta) = &frag.created_at_version_meta { - version_meta.load_sequence().map_err(|e| { - Error::internal(format!("Failed to load created_at version sequence: {}", e)) - })? - } else { - // Default: treat all rows as created at version 1 - lance_table::format::RowDatasetVersionSequence::from_uniform_row_count(row_count, 1) - }; - - // Load last_updated_at sequence (default to same as created_at sequence) - let mut last_updated_seq = if let Some(version_meta) = &frag.last_updated_at_version_meta { - version_meta.load_sequence().map_err(|e| { - Error::internal(format!( - "Failed to load last_updated_at version sequence: {}", - e - )) - })? - } else { - created_at_seq.clone() - }; - - // Apply deletion mask if present (positions are local offsets) - if let Some(deletion_file) = &frag.deletion_file { - let deletions = read_dataset_deletion_file(dataset, frag.id, deletion_file).await?; - last_updated_seq.mask(deletions.to_sorted_iter())?; - created_at_seq.mask(deletions.to_sorted_iter())?; - } - - old_last_updated_sequences.push(last_updated_seq); - old_created_at_sequences.push(created_at_seq); - } - - // Ensure row counts match new fragments total - let old_total: u64 = old_last_updated_sequences.iter().map(|s| s.len()).sum(); - let new_total: u64 = new_fragments - .iter() - .map(|f| f.physical_rows.unwrap_or(0) as u64) - .sum(); - debug_assert_eq!(old_total, new_total); - - // Rechunk version runs aligned to new fragment sizes - let chunk_sizes: Vec = new_fragments - .iter() - .map(|f| f.physical_rows.unwrap_or(0) as u64) - .collect(); - - let new_last_updated_sequences = lance_table::rowids::version::rechunk_version_sequences( - old_last_updated_sequences, - chunk_sizes.clone(), + let new_created_at = lance_table::rowids::version::rechunk_version_sequences( + old_created_at_sequences, + chunk_sizes.iter().copied(), false, )?; - - let new_created_at_sequences = lance_table::rowids::version::rechunk_version_sequences( - old_created_at_sequences, + let new_last_updated_at = lance_table::rowids::version::rechunk_version_sequences( + old_last_updated_sequences, chunk_sizes, false, )?; - // Set both version metadata on new fragments - for ((fragment, last_updated_seq), created_at_seq) in new_fragments + for (((fragment, row_ids), created_at), last_updated_at) in new_fragments .iter_mut() - .zip(new_last_updated_sequences) - .zip(new_created_at_sequences) + .zip(new_row_ids) + .zip(new_created_at) + .zip(new_last_updated_at) { - fragment.last_updated_at_version_meta = Some( - lance_table::format::RowDatasetVersionMeta::from_sequence(&last_updated_seq).unwrap(), - ); - fragment.created_at_version_meta = Some( - lance_table::format::RowDatasetVersionMeta::from_sequence(&created_at_seq).unwrap(), - ); + let lineage = RowLineage { + row_ids, + created_at, + last_updated_at, + }; + place_row_lineage(dataset, &lineage).await?.apply(fragment); } Ok(()) diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs index 9fae144d33b..4fcdc6b62ef 100644 --- a/rust/lance/src/dataset/rowids/spill.rs +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -413,11 +413,17 @@ async fn read_spilled_column( #[cfg(test)] mod tests { use super::*; - use crate::dataset::WriteParams; + use crate::dataset::cleanup::{CleanupPolicyBuilder, cleanup_old_versions}; + use crate::dataset::optimize::{CompactionOptions, compact_files}; + use crate::dataset::rowids::{RowVersionKind, load_row_id_sequence, load_row_version_sequence}; + use crate::dataset::{WriteMode, WriteParams}; use arrow_array::{Int32Array, RecordBatchIterator}; use arrow_schema::Field; + use chrono::Utc; use lance_core::utils::tempfile::TempStrDir; + use lance_core::{ROW_CREATED_AT_VERSION, ROW_ID, ROW_LAST_UPDATED_AT_VERSION}; use lance_file::version::LanceFileVersion; + use lance_table::feature_flags::FLAG_UNSTABLE_SPILLED_ROW_LINEAGE; /// A sequence with no runs to exploit, which is what a globally shuffled /// table produces and what forces the spill path. @@ -609,4 +615,209 @@ mod tests { .await .unwrap(); } + /// A stable-row-id dataset built from `chunks` separate appends, so + /// compacting it has several sequences to concatenate. + async fn appended_dataset(uri: &str, chunks: i32, rows_per_chunk: i32) -> Dataset { + let schema = test_schema(); + let mut dataset: Option = None; + for chunk in 0..chunks { + let batch = RecordBatch::try_new( + schema.clone(), + vec![Arc::new(Int32Array::from_iter_values( + (chunk * rows_per_chunk)..((chunk + 1) * rows_per_chunk), + ))], + ) + .unwrap(); + let reader = RecordBatchIterator::new(vec![Ok(batch)], schema.clone()); + dataset = Some( + Dataset::write( + reader, + uri, + Some(WriteParams { + enable_stable_row_ids: true, + mode: if chunk == 0 { + WriteMode::Create + } else { + WriteMode::Append + }, + ..Default::default() + }), + ) + .await + .unwrap(), + ); + } + dataset.unwrap() + } + + fn one_fragment() -> CompactionOptions { + CompactionOptions { + target_rows_per_fragment: 1_000, + ..Default::default() + } + } + + /// The lineage columns of every row, in scan order. + async fn collect_lineage(dataset: &Dataset) -> (Vec, Vec, Vec) { + let mut scanner = dataset.scan(); + scanner + .project(&[ROW_ID, ROW_CREATED_AT_VERSION, ROW_LAST_UPDATED_AT_VERSION]) + .unwrap(); + let batch = scanner.try_into_batch().await.unwrap(); + let column = |name: &str| { + batch + .column_by_name(name) + .unwrap() + .as_any() + .downcast_ref::() + .unwrap() + .values() + .to_vec() + }; + ( + column(ROW_ID), + column(ROW_CREATED_AT_VERSION), + column(ROW_LAST_UPDATED_AT_VERSION), + ) + } + + #[tokio::test] + async fn compaction_spills_and_reads_back_row_lineage() { + let dir = TempStrDir::default(); + let uri = dir.as_str(); + let mut dataset = appended_dataset(uri, 4, 250).await; + spill_everything(&mut dataset).await; + // Four appends at four versions, so the compacted created-at sequence + // has four runs rather than one. + let before = collect_lineage(&dataset).await; + assert_eq!( + before + .1 + .iter() + .collect::>() + .len(), + 4 + ); + + compact_files(&mut dataset, one_fragment(), None) + .await + .unwrap(); + + let fragments = dataset.get_fragments(); + assert_eq!(fragments.len(), 1); + let metadata = fragments[0].metadata(); + assert!( + matches!(metadata.row_id_meta, Some(RowIdMeta::Column)) + && matches!( + metadata.created_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ) + && matches!( + metadata.last_updated_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ), + "compaction must spill every sequence under a zero inline budget, got {metadata:?}" + ); + // The three sequences share one lineage file, listed after the user + // data file among the fragment's files and found by field id. + assert_eq!(metadata.files.len(), 2); + let lineage_file = &metadata.files[1]; + assert_eq!( + lineage_file.fields.as_ref(), + [ + ROW_ID_FIELD_ID, + ROW_CREATED_AT_VERSION_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID + ] + ); + for field_id in lineage_file.fields.iter() { + assert_eq!( + metadata.row_lineage_file(*field_id).unwrap(), + Some(lineage_file) + ); + } + assert_ne!( + dataset.manifest.reader_feature_flags & FLAG_UNSTABLE_SPILLED_ROW_LINEAGE, + 0, + "a spilled sequence must set the reader feature flag" + ); + assert_ne!( + dataset.manifest.writer_feature_flags & FLAG_UNSTABLE_SPILLED_ROW_LINEAGE, + 0, + "a spilled sequence must set the writer feature flag" + ); + + // The lineage survives the rewrite and is still served through the + // ordinary scan path, now from the data file columns. + assert_eq!(collect_lineage(&dataset).await, before); + // `validate_stable_row_ids` reads every fragment's sequences back and + // checks them against the fragment length, so this covers the loaders + // independently of the scan. + dataset.validate().await.unwrap(); + + // Re-opened cold, so nothing is served from this process's caches. + let reopened = Dataset::open(uri).await.unwrap(); + assert_eq!(collect_lineage(&reopened).await, before); + let fragment = &reopened.get_fragments()[0]; + let row_ids = load_row_id_sequence(&reopened, fragment.metadata()) + .await + .unwrap(); + assert_eq!(row_ids.iter().collect::>(), before.0); + let created_at = + load_row_version_sequence(&reopened, fragment.metadata(), RowVersionKind::CreatedAt) + .await + .unwrap() + .expect("a compacted fragment carries created-at versions"); + assert_eq!(versions_of(&created_at), before.1); + } + + /// Cleanup decides what to delete by walking + /// [`Fragment::referenced_lance_files`], so a spilled sequence has to be + /// reachable from there. If it were not, an ordinary cleanup would delete a + /// live file and leave the fragment claiming row ids it can no longer read. + #[tokio::test] + async fn cleanup_keeps_a_live_spilled_file() { + let dir = TempStrDir::default(); + let uri = dir.as_str(); + let mut dataset = appended_dataset(uri, 4, 250).await; + spill_everything(&mut dataset).await; + let before = collect_lineage(&dataset).await; + + compact_files(&mut dataset, one_fragment(), None) + .await + .unwrap(); + + let spilled = dataset.get_fragments()[0] + .metadata() + .row_lineage_file(ROW_ID_FIELD_ID) + .unwrap() + .expect("compaction must spill under a zero inline budget") + .path + .clone(); + let on_disk = std::path::Path::new(uri).join("data").join(&spilled); + assert!(on_disk.exists(), "no spilled file written at {on_disk:?}"); + + // Everything written so far is older than this instant, so the + // pre-compaction versions and their data files are all candidates. + let removed = cleanup_old_versions( + &dataset, + CleanupPolicyBuilder::default() + .before_timestamp(Utc::now()) + .delete_unverified(true) + .build(), + ) + .await + .unwrap(); + assert!( + removed.old_versions > 0, + "expected the pre-compaction versions to be cleaned up" + ); + assert!( + on_disk.exists(), + "cleanup deleted the live spilled row lineage file at {on_disk:?}" + ); + + let reopened = Dataset::open(uri).await.unwrap(); + assert_eq!(collect_lineage(&reopened).await, before); + } }