From efa125ee0d64783659ff065856fd73f1af77100c Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Thu, 17 Sep 2026 17:49:28 +0800 Subject: [PATCH] feat(table): load spilled row lineage ahead of a commit Building a manifest is synchronous and cannot read a data file, yet two commit-time paths need existing lineage: resolving which rows an update rewrote so each row's created-at version carries over, and overlaying a partial column rewrite's patched offsets onto the fragment's last-updated-at sequence. The commit path in lance now reads every spilled sequence of the current manifest ahead of each build attempt and hands them over in ManifestBuildConfig::spilled_row_lineage; the two paths consult that map for a spilled fragment and still refuse if a sequence they need is missing. UpdateBuilder, merge_insert in both write modes and externally assembled Operation::Updates therefore work on a spilled table unchanged, producing inline lineage for the rewritten rows; the next compaction spills it again. Co-Authored-By: Claude Fable 5.1 --- rust/lance-table/src/format/manifest.rs | 7 + rust/lance-table/src/rowids/version.rs | 54 +++- .../src/transaction/manifest_build.rs | 3 + .../src/transaction/row_version.rs | 155 +++++++---- .../src/transaction/test_support.rs | 1 + rust/lance/src/dataset.rs | 1 + rust/lance/src/dataset/rowids.rs | 60 ++++- rust/lance/src/dataset/rowids/spill.rs | 255 +++++++++++++++++- rust/lance/src/dataset/write/merge_insert.rs | 7 +- rust/lance/src/io/commit.rs | 32 ++- 10 files changed, 499 insertions(+), 76 deletions(-) diff --git a/rust/lance-table/src/format/manifest.rs b/rust/lance-table/src/format/manifest.rs index 6ea623c60da..86d98838457 100644 --- a/rust/lance-table/src/format/manifest.rs +++ b/rust/lance-table/src/format/manifest.rs @@ -750,6 +750,13 @@ pub struct ManifestBuildConfig { /// It bypasses the "cannot enable stable row ids on existing dataset" guard and /// sets `manifest.next_row_id` to the provided value before activating the flag. pub migration_next_row_id: Option, + /// Row lineage sequences of the current manifest's fragments that live + /// outside the manifest, read ahead of the build. An update that rewrites + /// rows needs the existing row ids and created-at versions to carry each + /// row's lineage over, and a partial column rewrite needs the existing + /// last-updated-at versions; the build cannot read a data file itself. Only + /// consulted for fragments whose sequences are spilled. + pub spilled_row_lineage: std::sync::Arc, } #[derive(Debug, Clone, Copy, PartialEq, Eq)] diff --git a/rust/lance-table/src/rowids/version.rs b/rust/lance-table/src/rowids/version.rs index 1b2ed5509d1..a269b1cc7f3 100644 --- a/rust/lance-table/src/rowids/version.rs +++ b/rust/lance-table/src/rowids/version.rs @@ -7,6 +7,7 @@ //! update version for each row in a Lance dataset, enabling efficient //! cross-version diff operations. +use std::collections::HashMap; use std::{ops::Range, sync::Arc}; use lance_core::Error; @@ -21,6 +22,25 @@ use crate::format::{Fragment, pb}; use crate::rowids::segment::U64Segment; use crate::rowids::{RowIdSequence, read_row_ids}; +/// One fragment's row lineage sequences that live outside the manifest, read +/// ahead of a commit by the caller. +/// +/// Building a manifest is synchronous and has no object store, so it cannot +/// read a sequence spilled to a data file column (`RowIdMeta::Column`, +/// `RowDatasetVersionMeta::Column`). A caller that can do IO loads them first +/// and hands them over in [`ManifestBuildConfig`](crate::format::ManifestBuildConfig). +/// Only the spilled sequences need to be present; inline ones are decoded on +/// the spot. +#[derive(Debug, Clone, Default)] +pub struct LoadedRowLineage { + pub row_ids: Option>, + pub created_at: Option>, + pub last_updated_at: Option>, +} + +/// [`LoadedRowLineage`] per fragment id. +pub type SpilledRowLineage = HashMap; + /// A run of identical versions over a contiguous span of row positions. /// /// Span is expressed as a U64Segment over row offsets (0..N within a fragment), @@ -695,11 +715,15 @@ pub fn refresh_row_latest_update_meta_for_full_frag_rewrite_cols( /// `updated_offsets` are local row offsets (within the fragment) that have been updated. /// Existing version metadata is preserved and only the updated positions are set to `current_version`. /// If no existing metadata is present, positions default to `prev_version`. +/// +/// A fragment whose existing versions are spilled has them read from +/// `spilled`; the refreshed sequence is placed inline. pub fn refresh_row_latest_update_meta_for_partial_frag_rewrite_cols( fragment: &mut Fragment, updated_offsets: &[usize], current_version: u64, prev_version: u64, + spilled: &SpilledRowLineage, ) -> Result<()> { // Determine row count for fragment let row_count_u64: u64 = if let Some(pr) = fragment.physical_rows { @@ -727,17 +751,26 @@ pub fn refresh_row_latest_update_meta_for_partial_frag_rewrite_cols( // Build base version vector from existing meta or previous dataset version let mut base_versions: Vec = Vec::with_capacity(row_count_u64 as usize); if let Some(meta) = fragment.last_updated_at_version_meta.as_ref() { - if matches!(meta, RowDatasetVersionMeta::Column) { + let base_seq = if matches!(meta, RowDatasetVersionMeta::Column) { // The existing versions of the rows this update leaves alone - // live in a data file, which this commit-time path cannot read. - // Defaulting them would silently rewrite their lineage. - return Err(Error::not_supported(format!( - "fragment {} stores its last-updated-at versions outside the manifest; \ - partially rewriting its columns is not supported yet", - fragment.id - ))); - } - if let Ok(base_seq) = meta.load_sequence() { + // live in a data file, which this path cannot read. Defaulting + // them would silently rewrite their lineage, so the caller has + // to have read them ahead of time. + let loaded = spilled + .get(&fragment.id) + .and_then(|lineage| lineage.last_updated_at.clone()) + .ok_or_else(|| { + Error::not_supported(format!( + "fragment {} stores its last-updated-at versions outside the \ + manifest and they were not loaded ahead of the commit", + fragment.id + )) + })?; + Some(loaded.as_ref().clone()) + } else { + meta.load_sequence().ok() + }; + if let Some(base_seq) = base_seq { base_versions.extend(base_seq.versions().take(row_count_u64 as usize)); base_versions.resize(row_count_u64 as usize, prev_version); } else { @@ -901,6 +934,7 @@ mod tests { &[ROWS - 1], 3, 1, + &Default::default(), ) .unwrap(); diff --git a/rust/lance-table/src/transaction/manifest_build.rs b/rust/lance-table/src/transaction/manifest_build.rs index f18975d49c8..378237a0d8f 100644 --- a/rust/lance-table/src/transaction/manifest_build.rs +++ b/rust/lance-table/src/transaction/manifest_build.rs @@ -795,6 +795,7 @@ impl Transaction { &offsets, new_version, prev_version, + &config.spilled_row_lineage, )?; } } @@ -821,6 +822,7 @@ impl Transaction { existing_fragments, new_fragments.as_mut_slice(), new_version, + &config.spilled_row_lineage, )?; } @@ -1339,6 +1341,7 @@ impl Transaction { &covered_offsets, new_version, 1, + &config.spilled_row_lineage, )?; } } diff --git a/rust/lance-table/src/transaction/row_version.rs b/rust/lance-table/src/transaction/row_version.rs index 2ae085b9dc3..05c74c16487 100644 --- a/rust/lance-table/src/transaction/row_version.rs +++ b/rust/lance-table/src/transaction/row_version.rs @@ -10,10 +10,11 @@ //! fragment and offset it came from, which is what most of this module does. use crate::format::{Fragment, RowDatasetVersionMeta, RowDatasetVersionSequence, RowIdMeta}; -use crate::rowids::version::build_version_meta; +use crate::rowids::version::{SpilledRowLineage, build_version_meta}; use crate::rowids::{RowIdSequence, read_row_ids, write_row_ids}; use crate::transaction::Transaction; use lance_core::{Error, Result}; +use std::borrow::Cow; use std::cmp::Ordering; use std::collections::{HashMap, HashSet}; @@ -52,10 +53,15 @@ fn resolve_created_at_version( /// For each new fragment produced by an update, set `created_at_version_meta` /// (preserved from the original rows) and `last_updated_at_version_meta`. +/// +/// `spilled` supplies the row ids and created-at versions of existing +/// fragments that keep them outside the manifest; see +/// [`SpilledRowLineage`]. pub(super) fn resolve_update_version_metadata( existing_fragments: &[Fragment], new_fragments: &mut [Fragment], new_version: u64, + spilled: &SpilledRowLineage, ) -> Result<()> { // Collect only the row IDs we actually need to resolve, those appearing in new_fragments // with inline metadata. This bounds the lookup map to O(updated rows) instead of O(all dataset rows) @@ -91,14 +97,20 @@ pub(super) fn resolve_update_version_metadata( // cannot be read here. Skipping such a fragment would make its rows // look freshly inserted and stamp them with a new created-at version. let seq = match &frag.row_id_meta { - Some(RowIdMeta::Inline(data)) => read_row_ids(data).ok(), - Some(RowIdMeta::Column) => { - return Err(Error::not_supported(format!( - "fragment {} stores its row ids outside the manifest; updating rows \ - of a table with spilled row ids is not supported yet", - frag.id - ))); - } + Some(RowIdMeta::Inline(data)) => read_row_ids(data).ok().map(Cow::Owned), + Some(RowIdMeta::Column) => Some(Cow::Borrowed( + spilled + .get(&frag.id) + .and_then(|lineage| lineage.row_ids.as_deref()) + .ok_or_else(|| { + Error::not_supported(format!( + "fragment {} stores its row ids outside the manifest and they \ + were not loaded ahead of the commit; cannot resolve which rows \ + this update rewrote", + frag.id + )) + })?, + )), None => None, }; if let Some(seq) = seq { @@ -138,15 +150,20 @@ pub(super) fn resolve_update_version_metadata( }; if matches!(meta, RowDatasetVersionMeta::Column) { // The rewritten rows' original created-at versions live in a data - // file, which this commit-time path cannot read. Defaulting them - // would silently rewrite their lineage. - return Err(Error::not_supported(format!( - "fragment {} stores its created-at versions outside the manifest; \ - updating rows it holds is not supported yet", - frag.id - ))); - } - if let Ok(seq) = meta.load_sequence() { + // file, which this path cannot read. Defaulting them would silently + // rewrite their lineage, so the caller has to have read them. + let loaded = spilled + .get(&frag.id) + .and_then(|lineage| lineage.created_at.clone()) + .ok_or_else(|| { + Error::not_supported(format!( + "fragment {} stores its created-at versions outside the manifest \ + and they were not loaded ahead of the commit", + frag.id + )) + })?; + version_cache.insert(frag.id, loaded.as_ref().clone()); + } else if let Ok(seq) = meta.load_sequence() { version_cache.insert(frag.id, seq); } } @@ -388,19 +405,27 @@ mod tests { } #[test] - fn test_resolve_update_versions_refuses_spilled_lineage() { - // Resolving lineage here means reading row ids and versions back, which - // this commit-time path cannot do for a spilled sequence. Skipping such - // a fragment would make its rows look freshly inserted, so it refuses. - let inline_ids = |range: std::ops::Range| { - Some(RowIdMeta::Inline( - write_row_ids(&RowIdSequence::from(range)).into(), - )) - }; - let fragment = |id: u64, rows: usize| Fragment { - id, - physical_rows: Some(rows), - row_id_meta: None, + fn test_resolve_update_versions_reads_spilled_source_lineage_from_config() { + // The rewritten rows' sources are found by scanning existing row ids, + // which a spilled fragment does not expose here: the commit's caller + // reads them ahead of time. Without that, skipping the fragment would + // make its rows look freshly inserted, so it is an error instead. + let existing = vec![Fragment { + id: 1, + physical_rows: Some(50), + row_id_meta: Some(RowIdMeta::Column), + files: vec![], + overlays: vec![], + deletion_file: None, + last_updated_at_version_meta: Some(inline_versions(50, 2)), + created_at_version_meta: Some(RowDatasetVersionMeta::Column), + }]; + let new_fragment = |ids: std::ops::Range| Fragment { + id: 2, + physical_rows: Some(2), + row_id_meta: Some(RowIdMeta::Inline( + write_row_ids(&RowIdSequence::from(ids)).into(), + )), files: vec![], overlays: vec![], deletion_file: None, @@ -408,39 +433,51 @@ mod tests { created_at_version_meta: None, }; - // An existing fragment with spilled row ids, when the update rewrote rows. - let existing = vec![Fragment { - row_id_meta: Some(RowIdMeta::Column), - created_at_version_meta: Some(inline_versions(50, 2)), - ..fragment(1, 50) - }]; - let mut new_fragments = vec![Fragment { - row_id_meta: inline_ids(10..12), - ..fragment(2, 2) - }]; - let error = resolve_update_version_metadata(&existing, &mut new_fragments, 9).unwrap_err(); - assert!(matches!(error, Error::NotSupported { .. }), "{error:?}"); - - // An existing fragment with spilled created-at versions. - let existing = vec![Fragment { - row_id_meta: inline_ids(0..50), - created_at_version_meta: Some(RowDatasetVersionMeta::Column), - ..fragment(1, 50) - }]; - let mut new_fragments = vec![Fragment { - row_id_meta: inline_ids(10..12), - ..fragment(2, 2) - }]; - let error = resolve_update_version_metadata(&existing, &mut new_fragments, 9).unwrap_err(); + let mut new_fragments = vec![new_fragment(10..12)]; + let error = + resolve_update_version_metadata(&existing, &mut new_fragments, 9, &Default::default()) + .unwrap_err(); assert!(matches!(error, Error::NotSupported { .. }), "{error:?}"); - // A new fragment whose own row ids are spilled. - let mut new_fragments = vec![Fragment { + // A new fragment whose own row ids are spilled cannot be resolved. + let mut spilled_new = vec![Fragment { row_id_meta: Some(RowIdMeta::Column), - ..fragment(2, 50) + ..new_fragment(0..2) }]; - let error = resolve_update_version_metadata(&[], &mut new_fragments, 9).unwrap_err(); + let error = resolve_update_version_metadata(&[], &mut spilled_new, 9, &Default::default()) + .unwrap_err(); assert!(matches!(error, Error::NotSupported { .. }), "{error:?}"); + + // Fragment 1 holds ids 100..150 and rows 10 and 11 were created at + // versions 3 and 4 respectively; that is what the rewrite must keep. + let mut created_at_versions = vec![2; 50]; + created_at_versions[10] = 3; + created_at_versions[11] = 4; + let spilled = SpilledRowLineage::from([( + 1, + crate::rowids::version::LoadedRowLineage { + row_ids: Some(Arc::new(RowIdSequence::from(100..150))), + created_at: Some(Arc::new(RowDatasetVersionSequence::from_versions( + &created_at_versions, + ))), + last_updated_at: None, + }, + )]); + let mut new_fragments = vec![new_fragment(110..112)]; + resolve_update_version_metadata(&existing, &mut new_fragments, 9, &spilled).unwrap(); + let versions = |meta: &Option| { + meta.as_ref() + .unwrap() + .load_sequence() + .unwrap() + .versions() + .collect::>() + }; + assert_eq!(versions(&new_fragments[0].created_at_version_meta), [3, 4]); + assert_eq!( + versions(&new_fragments[0].last_updated_at_version_meta), + [9, 9] + ); } #[test] diff --git a/rust/lance-table/src/transaction/test_support.rs b/rust/lance-table/src/transaction/test_support.rs index 895abfbb746..a7c5799429d 100644 --- a/rust/lance-table/src/transaction/test_support.rs +++ b/rust/lance-table/src/transaction/test_support.rs @@ -30,6 +30,7 @@ pub fn default_build_config() -> ManifestBuildConfig { storage_format: None, disable_transaction_file: false, migration_next_row_id: None, + spilled_row_lineage: Default::default(), } } diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index ec9f113cdee..426c97f1a4a 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -4155,6 +4155,7 @@ impl ManifestWriteConfig { storage_format: self.storage_format.clone(), disable_transaction_file: self.disable_transaction_file, migration_next_row_id: self.migration_next_row_id, + spilled_row_lineage: Default::default(), } } } diff --git a/rust/lance/src/dataset/rowids.rs b/rust/lance/src/dataset/rowids.rs index 71718232909..78644267fdd 100644 --- a/rust/lance/src/dataset/rowids.rs +++ b/rust/lance/src/dataset/rowids.rs @@ -17,7 +17,10 @@ use lance_table::{ ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, RowDatasetVersionMeta, RowDatasetVersionSequence, RowIdMeta, }, - rowids::{FragmentRowIdIndex, RowIdIndex, RowIdSequence, read_row_ids}, + rowids::{ + FragmentRowIdIndex, RowIdIndex, RowIdSequence, read_row_ids, + version::{LoadedRowLineage, SpilledRowLineage}, + }, }; use std::sync::Arc; @@ -127,6 +130,61 @@ pub async fn load_row_version_sequence( } } +/// Read ahead every lineage sequence of `fragments` that lives outside the +/// manifest, for a commit that will need to consult them. +/// +/// Building a manifest is synchronous and cannot read a data file, so the +/// commit path calls this first, over the whole manifest, and hands the result +/// over in `ManifestBuildConfig::spilled_row_lineage`. Only spilled sequences +/// are loaded; when none is, this returns an empty map without IO. +pub async fn load_spilled_row_lineage<'a>( + dataset: &Dataset, + fragments: impl IntoIterator, +) -> Result> { + // A `for` loop rather than `map`: a closure returning a future that borrows + // its argument trips the higher-ranked lifetime check on the outer future. + let mut loads = Vec::new(); + for fragment in fragments { + if fragment.has_spilled_row_lineage() { + loads.push(load_fragment_spilled_lineage(dataset, fragment)); + } + } + let loaded: SpilledRowLineage = futures::stream::iter(loads) + .buffer_unordered(dataset.object_store.io_parallelism()) + .try_collect() + .await?; + Ok(Arc::new(loaded)) +} + +/// The spilled sequences of one fragment, for [`load_spilled_row_lineage`]. +async fn load_fragment_spilled_lineage( + dataset: &Dataset, + fragment: &Fragment, +) -> Result<(u64, LoadedRowLineage)> { + let row_ids = match &fragment.row_id_meta { + Some(RowIdMeta::Column) => Some(load_row_id_sequence(dataset, fragment).await?), + _ => None, + }; + let mut versions = [None, None]; + for (slot, kind) in versions + .iter_mut() + .zip([RowVersionKind::CreatedAt, RowVersionKind::LastUpdatedAt]) + { + if let Some(RowDatasetVersionMeta::Column) = kind.meta(fragment) { + *slot = load_row_version_sequence(dataset, fragment, kind).await?; + } + } + let [created_at, last_updated_at] = versions; + Ok(( + fragment.id, + LoadedRowLineage { + row_ids, + created_at, + last_updated_at, + }, + )) +} + /// Load row id sequences from the given dataset and fragments. /// /// Returned as a vector of (fragment_id, sequence) pairs. These are not diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs index 4fcdc6b62ef..91ec1476035 100644 --- a/rust/lance/src/dataset/rowids/spill.rs +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -416,7 +416,7 @@ mod tests { 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 crate::dataset::{UpdateBuilder, WriteMode, WriteParams}; use arrow_array::{Int32Array, RecordBatchIterator}; use arrow_schema::Field; use chrono::Utc; @@ -820,4 +820,257 @@ mod tests { let reopened = Dataset::open(uri).await.unwrap(); assert_eq!(collect_lineage(&reopened).await, before); } + /// Every row's key, row id, created-at and last-updated-at version. + async fn collect_rows(dataset: &Dataset) -> Vec<(i32, u64, u64, u64)> { + let mut scanner = dataset.scan(); + scanner + .project(&[ + "i", + ROW_ID, + ROW_CREATED_AT_VERSION, + ROW_LAST_UPDATED_AT_VERSION, + ]) + .unwrap(); + let batch = scanner.try_into_batch().await.unwrap(); + let u64s = |name: &str| { + batch + .column_by_name(name) + .unwrap() + .as_any() + .downcast_ref::() + .unwrap() + .values() + .to_vec() + }; + let keys = batch + .column_by_name("i") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap() + .values() + .to_vec(); + let (ids, created, updated) = ( + u64s(ROW_ID), + u64s(ROW_CREATED_AT_VERSION), + u64s(ROW_LAST_UPDATED_AT_VERSION), + ); + keys.into_iter() + .zip(ids) + .zip(created) + .zip(updated) + .map(|(((key, id), created), updated)| (key, id, created, updated)) + .collect() + } + + fn by_key(rows: &[(i32, u64, u64, u64)]) -> std::collections::BTreeMap { + rows.iter() + .map(|(key, id, created, updated)| (*key, (*id, *created, *updated))) + .collect() + } + + /// Resolving the rewritten rows' original created-at versions happens at + /// commit time, inside `lance-table`, which cannot read a data file. The + /// commit path reads the spilled sequences ahead of the build, so an + /// update on a spilled table keeps every row's lineage the way it does on + /// an inline one. + #[tokio::test] + async fn updating_rows_with_spilled_lineage_keeps_their_created_at() { + let dir = TempStrDir::default(); + let uri = dir.as_str(); + let mut dataset = appended_dataset(uri, 4, 250).await; + spill_everything(&mut dataset).await; + compact_files(&mut dataset, one_fragment(), None) + .await + .unwrap(); + let before = by_key(&collect_rows(&dataset).await); + // Row 700 came in the third append, so its created-at is not the + // default a reader would fall back to. + let (id_700, created_700, _) = before[&700]; + assert_eq!(created_700, 3); + + let updated = UpdateBuilder::new(Arc::new(dataset)) + .update_where("i = 700") + .unwrap() + .set("i", "7000") + .unwrap() + .build() + .unwrap() + .execute() + .await + .unwrap(); + let updated = updated.new_dataset.as_ref(); + let update_version = updated.version().version; + let after = by_key(&collect_rows(updated).await); + + // The rewritten row keeps its id and its created-at, and is stamped + // with the update's version; every other row is untouched. + assert_eq!(after[&7000], (id_700, created_700, update_version)); + assert!(!after.contains_key(&700)); + for (key, lineage) in before.iter().filter(|(key, _)| **key != 700) { + assert_eq!(after[key], *lineage, "row {key} must be untouched"); + } + updated.validate().await.unwrap(); + } + + /// merge_insert rewrites the matched rows and appends the inserted ones in + /// one fragment; the former keep their lineage, the latter start at the + /// commit version. Both resolve through the read-ahead spilled sequences. + #[tokio::test] + async fn merge_insert_on_spilled_table_keeps_matched_lineage() { + use crate::dataset::{MergeInsertBuilder, WhenMatched, WhenNotMatched}; + + let dir = TempStrDir::default(); + let uri = dir.as_str(); + let mut dataset = appended_dataset(uri, 4, 250).await; + spill_everything(&mut dataset).await; + compact_files(&mut dataset, one_fragment(), None) + .await + .unwrap(); + let before = by_key(&collect_rows(&dataset).await); + + // Keys 300 and 700 exist and were appended at different versions; + // 5000 does not exist and is inserted. + let schema = test_schema(); + let source = RecordBatch::try_new( + schema.clone(), + vec![Arc::new(Int32Array::from(vec![300, 700, 5000]))], + ) + .unwrap(); + let (merged, stats) = MergeInsertBuilder::try_new(Arc::new(dataset), vec!["i".into()]) + .unwrap() + .when_matched(WhenMatched::UpdateAll) + .when_not_matched(WhenNotMatched::InsertAll) + .try_build() + .unwrap() + .execute_reader(Box::new(RecordBatchIterator::new([Ok(source)], schema))) + .await + .unwrap(); + assert_eq!((stats.num_updated_rows, stats.num_inserted_rows), (2, 1)); + let merge_version = merged.version().version; + let after = by_key(&collect_rows(&merged).await); + + for key in [300, 700] { + let (id, created, _) = before[&key]; + assert_eq!(after[&key], (id, created, merge_version), "row {key}"); + } + let (_, created, updated) = after[&5000]; + assert_eq!((created, updated), (merge_version, merge_version)); + for (key, lineage) in before.iter().filter(|(key, _)| ![300, 700].contains(key)) { + assert_eq!(after[key], *lineage, "row {key} must be untouched"); + } + merged.validate().await.unwrap(); + } + + /// A partial column rewrite patches an existing fragment in place and + /// stamps only the patched rows' last-updated-at, which means overlaying + /// the fragment's existing sequence; on a spilled fragment that sequence + /// is read ahead of the commit and the refreshed one goes back inline. + #[tokio::test] + async fn partial_column_rewrite_on_spilled_fragment_stamps_only_matched_rows() { + use crate::dataset::{ + MergeInsertBuilder, MergeInsertWriteMode, WhenMatched, WhenNotMatched, + }; + use arrow_array::StringArray; + + let dir = TempStrDir::default(); + let uri = dir.as_str(); + // A third column keeps the patch source a strict subset of the schema, + // which is what makes this an in-place column rewrite. + let schema = Arc::new(ArrowSchema::new(vec![ + Field::new("i", DataType::Int32, false), + Field::new("tag", DataType::Utf8, true), + Field::new("other", DataType::Utf8, true), + ])); + let mut dataset: Option = None; + for chunk in 0..4 { + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int32Array::from_iter_values( + (chunk * 250)..((chunk + 1) * 250), + )), + Arc::new(StringArray::from(vec!["t"; 250])), + Arc::new(StringArray::from(vec!["o"; 250])), + ], + ) + .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(), + ); + } + let mut dataset = dataset.unwrap(); + spill_everything(&mut dataset).await; + compact_files(&mut dataset, one_fragment(), None) + .await + .unwrap(); + assert!( + matches!( + dataset.get_fragments()[0] + .metadata() + .last_updated_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ), + "the fixture must start with spilled last-updated-at versions" + ); + let before = by_key(&collect_rows(&dataset).await); + + let source_schema = Arc::new(ArrowSchema::new(vec![ + Field::new("i", DataType::Int32, false), + Field::new("tag", DataType::Utf8, true), + ])); + let source = RecordBatch::try_new( + source_schema.clone(), + vec![ + Arc::new(Int32Array::from(vec![300, 700])), + Arc::new(StringArray::from(vec!["patched"; 2])), + ], + ) + .unwrap(); + let (patched, stats) = MergeInsertBuilder::try_new(Arc::new(dataset), vec!["i".into()]) + .unwrap() + .when_matched(WhenMatched::UpdateAll) + .when_not_matched(WhenNotMatched::DoNothing) + .write_mode(MergeInsertWriteMode::RewriteColumns) + .try_build() + .unwrap() + .execute_reader(Box::new(RecordBatchIterator::new( + [Ok(source)], + source_schema, + ))) + .await + .unwrap(); + assert_eq!(stats.num_updated_rows, 2); + assert_eq!( + patched.get_fragments().len(), + 1, + "in-place patches add no fragment" + ); + let patch_version = patched.version().version; + let after = by_key(&collect_rows(&patched).await); + + for key in [300, 700] { + let (id, created, _) = before[&key]; + assert_eq!(after[&key], (id, created, patch_version), "row {key}"); + } + for (key, lineage) in before.iter().filter(|(key, _)| ![300, 700].contains(key)) { + assert_eq!(after[key], *lineage, "row {key} must be untouched"); + } + patched.validate().await.unwrap(); + } } diff --git a/rust/lance/src/dataset/write/merge_insert.rs b/rust/lance/src/dataset/write/merge_insert.rs index 8bed2cb2879..6f57a21ec3f 100644 --- a/rust/lance/src/dataset/write/merge_insert.rs +++ b/rust/lance/src/dataset/write/merge_insert.rs @@ -47,7 +47,7 @@ use super::{ CommitBuilder, TargetBaseInfo, WriteMode, WriteParams, validate_and_resolve_target_bases_with_primary, write_fragments_internal, }; -use crate::dataset::rowids::get_row_id_index; +use crate::dataset::rowids::{get_row_id_index, load_spilled_row_lineage}; use crate::dataset::transaction::UpdateMode::{RewriteColumns, RewriteRows}; use crate::dataset::utils::CapturedRowIds; use crate::index::DatasetIndexExt; @@ -1875,11 +1875,16 @@ impl MergeInsertJob { updated_offsets.sort_unstable(); updated_offsets.dedup(); + // The fragment's existing versions may be spilled to a + // data file, which the refresh cannot read itself. + let spilled_lineage = + load_spilled_row_lineage(&dataset, [&updated_fragment]).await?; lance_table::rowids::version::refresh_row_latest_update_meta_for_partial_frag_rewrite_cols( &mut updated_fragment, &updated_offsets, current_version, dataset.manifest.version, + &spilled_lineage, )?; } diff --git a/rust/lance/src/io/commit.rs b/rust/lance/src/io/commit.rs index e5120faeffc..4bc284288a4 100644 --- a/rust/lance/src/io/commit.rs +++ b/rust/lance/src/io/commit.rs @@ -36,8 +36,9 @@ use lance_io::utils::CachedFileSize; use lance_select::RowAddrTreeMap; use lance_table::feature_flags::ensure_can_write_manifest; use lance_table::format::{ - DETACHED_VERSION_MASK, DeletionFile, Fragment, IndexMetadata, Manifest, WriterVersion, - is_detached_version, list_index_files_with_sizes, operation_may_change_schema, pb, + DETACHED_VERSION_MASK, DeletionFile, Fragment, IndexMetadata, Manifest, ManifestBuildConfig, + WriterVersion, is_detached_version, list_index_files_with_sizes, operation_may_change_schema, + pb, }; use lance_table::io::commit::{ CommitConfig, CommitError, CommitHandler, ManifestLocation, ManifestNamingScheme, @@ -50,6 +51,7 @@ use super::ObjectStore; use crate::Dataset; use crate::dataset::cleanup::auto_cleanup_hook; use crate::dataset::fragment::FileFragment; +use crate::dataset::rowids::load_spilled_row_lineage; use crate::dataset::transaction::{Operation, Transaction}; use crate::dataset::{ ManifestWriteConfig, NewTransactionResult, TRANSACTIONS_DIR, load_new_transactions, @@ -1217,7 +1219,7 @@ pub(crate) async fn do_commit_detached_transaction( Some(dataset.manifest.as_ref()), load_all_indices(dataset).await?.as_ref().clone(), &transaction_file, - &write_config.to_build_config(), + &build_config_for_attempt(dataset, transaction, write_config).await?, )?, }; @@ -1375,6 +1377,28 @@ pub(crate) async fn commit_detached_transaction( } /// Load new transactions and sort them by version in ascending order (oldest to newest) +/// The build config for one commit attempt against `dataset`'s current manifest. +/// +/// An operation that carries rows' lineage over from existing fragments needs +/// their sequences at build time, and the build cannot read the ones spilled to +/// data files; they are read here, per attempt, so a rebase onto a newer +/// manifest sees that manifest's fragments. +async fn build_config_for_attempt( + dataset: &Dataset, + transaction: &Transaction, + write_config: &ManifestWriteConfig, +) -> Result { + let mut config = write_config.to_build_config(); + if matches!( + transaction.operation, + Operation::Update { .. } | Operation::DataOverlay { .. } + ) { + config.spilled_row_lineage = + load_spilled_row_lineage(dataset, dataset.manifest.fragments.iter()).await?; + } + Ok(config) +} + async fn load_and_sort_new_transactions( dataset: &Dataset, ) -> Result<(Dataset, Vec<(u64, Arc)>)> { @@ -1569,7 +1593,7 @@ pub(crate) async fn commit_transaction( Some(dataset.manifest.as_ref()), load_all_indices(&dataset).await?.as_ref().clone(), transaction_file, - &write_config.to_build_config(), + &build_config_for_attempt(&dataset, &transaction, write_config).await?, read_version_state, )?, };