From 55c66b0e955493c24f4bd4ce751980d104386897 Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Thu, 17 Sep 2026 17:41:16 +0800 Subject: [PATCH] feat(dataset): spill row lineage at compaction Compaction used to rechunk the row ids and the two version sequences in two separate passes and write each back inline. It now computes all three together and places each one through place_row_lineage: inline when its encoding fits the table's budget, otherwise as a column of one lineage file per output fragment. Compaction's output is retry-stable, every value being carried over from the input fragments, so spilling it needs nothing from the commit. Only tables that opt in through lance.row_lineage.spill are affected. Updating rows of a table whose compaction has spilled still returns NotSupported from the commit, which cannot read a data file; the next commit lifts that. Co-authored-by: Will Jones Co-Authored-By: Claude Fable 5.1 --- rust/lance/src/dataset/optimize.rs | 192 ++++++++-------------- rust/lance/src/dataset/rowids/spill.rs | 213 ++++++++++++++++++++++++- 2 files changed, 275 insertions(+), 130 deletions(-) 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); + } }