From be440e280d0cdb74160bf42326051dd2056a9030 Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Thu, 17 Sep 2026 18:03:13 +0800 Subject: [PATCH 1/5] feat(dataset): spill row lineage when updating rows UpdateBuilder reads the rewritten rows' created-at versions alongside their row ids while it scans them, and places each new fragment's lineage through place_row_lineage: inline on a table that has not opted in, as before, and spilled on one that has, so an update-heavy table no longer regrows its manifest between compactions. The values it places are the ones a commit conflict cannot change; the last-updated-at version is the commit's, and the commit now stamps it on any new fragment whose created-at versions the writer placed, spilled or inline, instead of resolving them again from the existing fragments. Legacy V1 datasets, whose reader does not serve the version columns, keep leaving the created-at lookup to the commit. Co-Authored-By: Claude Fable 5.1 --- .../src/transaction/row_version.rs | 44 ++++++- rust/lance/src/dataset/rowids/spill.rs | 112 +++++++++++++++++ rust/lance/src/dataset/utils.rs | 89 ++++++++++---- rust/lance/src/dataset/write/update.rs | 116 +++++++++++++++--- 4 files changed, 314 insertions(+), 47 deletions(-) diff --git a/rust/lance-table/src/transaction/row_version.rs b/rust/lance-table/src/transaction/row_version.rs index 05c74c16487..0f350b7f105 100644 --- a/rust/lance-table/src/transaction/row_version.rs +++ b/rust/lance-table/src/transaction/row_version.rs @@ -169,14 +169,25 @@ pub(super) fn resolve_update_version_metadata( } for fragment in new_fragments.iter_mut() { + // A writer that read the rewritten rows already knows their created-at + // versions and may place them itself, spilled or inline. The + // last-updated-at version is this commit's, which only the commit + // knows (a conflict retry moves it), so it is stamped here regardless + // of what the writer left. + if fragment.created_at_version_meta.is_some() { + fragment.last_updated_at_version_meta = build_version_meta(fragment, new_version); + continue; + } let row_ids = match &fragment.row_id_meta { Some(RowIdMeta::Inline(data)) => read_row_ids(data).ok(), Some(RowIdMeta::Column) => { // Resolving the versions needs the row ids, which this - // commit-time path cannot read back from a data file. + // commit-time path cannot read back from a data file; a writer + // that spills them has to place the created-at versions too. return Err(Error::not_supported(format!( - "fragment {} stores its row ids outside the manifest; committing it \ - through an update is not supported yet", + "fragment {} stores its row ids outside the manifest but carries no \ + created-at version metadata; a writer that spills row ids must place \ + the created-at versions with them", fragment.id ))); } @@ -404,6 +415,33 @@ mod tests { assert_eq!(fragments[0].row_id_meta, Some(spilled)); } + #[test] + fn test_resolve_update_versions_keeps_writer_placed_created_at() { + // A writer that placed a new fragment's created-at versions, spilled + // or inline, keeps them; the commit only stamps last-updated-at with + // its own version, which the writer could not know. + let mut new_fragments = 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, 7)), + created_at_version_meta: Some(RowDatasetVersionMeta::Column), + }]; + + resolve_update_version_metadata(&[], &mut new_fragments, 9, &Default::default()).unwrap(); + assert_eq!( + new_fragments[0].created_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ); + assert_eq!( + new_fragments[0].last_updated_at_version_meta, + Some(inline_versions(50, 9)) + ); + } + #[test] fn test_resolve_update_versions_reads_spilled_source_lineage_from_config() { // The rewritten rows' sources are found by scanning existing row ids, diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs index 91ec1476035..e1e0fa22814 100644 --- a/rust/lance/src/dataset/rowids/spill.rs +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -1073,4 +1073,116 @@ mod tests { } patched.validate().await.unwrap(); } + /// With the table opted in, an update spills the lineage of the rows it + /// rewrote at write time: their ids and created-at versions are known + /// before the commit and cannot change on a retry. The commit stamps the + /// last-updated-at version; the writer's placeholder is replaced. + #[tokio::test] + async fn update_on_an_opted_in_table_spills_the_rewritten_rows_lineage() { + 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 = by_key(&collect_rows(&dataset).await); + + let updated = UpdateBuilder::new(Arc::new(dataset)) + .update_where("i >= 500") + .unwrap() + .set("i", "i + 10000") + .unwrap() + .build() + .unwrap() + .execute() + .await + .unwrap(); + let updated = updated.new_dataset.as_ref(); + let update_version = updated.version().version; + + let rewritten = updated + .get_fragments() + .into_iter() + .map(|fragment| fragment.metadata().clone()) + .find(|metadata| matches!(metadata.row_id_meta, Some(RowIdMeta::Column))) + .expect("the rewritten rows' fragment must spill its row ids"); + assert!( + matches!( + rewritten.created_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ), + "the created-at versions must spill with the row ids" + ); + // The lineage file is one of the fragment's files, after its data. + assert_eq!(rewritten.files.len(), 2); + assert!( + rewritten.files[1].fields.iter().all(|field| *field < 0), + "the lineage file holds only lineage columns: {:?}", + rewritten.files[1].fields + ); + assert!( + matches!( + rewritten.last_updated_at_version_meta, + Some(RowDatasetVersionMeta::Inline(_)) + ), + "the commit stamps last-updated-at inline, got {:?}", + rewritten.last_updated_at_version_meta + ); + + let after = by_key(&collect_rows(updated).await); + for (key, (id, created, updated_at)) in before.iter() { + if *key >= 500 { + assert_eq!( + after[&(key + 10000)], + (*id, *created, update_version), + "row {key}" + ); + } else { + assert_eq!(after[key], (*id, *created, *updated_at), "row {key}"); + } + } + updated.validate().await.unwrap(); + } + + /// Without the opt-in the same update places everything inline, as every + /// release has, and the lineage it carries is the same. + #[tokio::test] + async fn update_on_a_table_that_did_not_opt_in_keeps_lineage_inline() { + let dir = TempStrDir::default(); + let uri = dir.as_str(); + let dataset = appended_dataset(uri, 4, 250).await; + let before = by_key(&collect_rows(&dataset).await); + + let updated = UpdateBuilder::new(Arc::new(dataset)) + .update_where("i >= 500") + .unwrap() + .set("i", "i + 10000") + .unwrap() + .build() + .unwrap() + .execute() + .await + .unwrap(); + let updated = updated.new_dataset.as_ref(); + let update_version = updated.version().version; + + for fragment in updated.get_fragments() { + assert!( + !fragment.metadata().has_spilled_row_lineage(), + "fragment {} spilled without the table opting in", + fragment.id() + ); + } + let after = by_key(&collect_rows(updated).await); + for (key, (id, created, updated_at)) in before.iter() { + if *key >= 500 { + assert_eq!( + after[&(key + 10000)], + (*id, *created, update_version), + "row {key}" + ); + } else { + assert_eq!(after[key], (*id, *created, *updated_at), "row {key}"); + } + } + updated.validate().await.unwrap(); + } } diff --git a/rust/lance/src/dataset/utils.rs b/rust/lance/src/dataset/utils.rs index 254718464e9..5b23ab6b4a1 100644 --- a/rust/lance/src/dataset/utils.rs +++ b/rust/lance/src/dataset/utils.rs @@ -14,37 +14,42 @@ use lance_arrow::json::{ arrow_json_to_lance_json, convert_json_columns, convert_lance_json_to_arrow, has_arrow_json_fields, has_json_fields, lance_json_to_arrow_json, }; -use lance_core::ROW_ID; +use lance_core::{ROW_CREATED_AT_VERSION, ROW_ID}; use lance_table::rowids::{RowIdIndex, RowIdSequence}; use roaring::RoaringTreemap; use std::borrow::Cow; use std::sync::Arc; use std::sync::mpsc::Receiver; +fn u64_values<'a>(batch: &'a RecordBatch, column: usize, what: &str) -> &'a [u64] { + let array = batch.column(column); + array + .as_any() + .downcast_ref::() + .unwrap_or_else(|| panic!("{what} had an unexpected type: {}", array.data_type())) + .values() +} + fn extract_row_ids( row_ids: &mut CapturedRowIds, batch: RecordBatch, row_id_idx: usize, - non_row_id_projection: &[usize], + created_at_idx: Option, + data_projection: &[usize], ) -> DFResult { - let row_ids_arr = batch.column(row_id_idx); - let row_ids_itr = row_ids_arr - .as_any() - .downcast_ref::() - .unwrap_or_else(|| { - panic!( - "Row ids had an unexpected type: {}", - row_ids_arr.data_type() - ) - }) - .values(); - row_ids.capture(row_ids_itr)?; - Ok(batch.project(non_row_id_projection)?) + row_ids.capture(u64_values(&batch, row_id_idx, "Row ids"))?; + if let Some(created_at_idx) = created_at_idx { + row_ids.capture_created_at(u64_values(&batch, created_at_idx, "Created-at versions")); + } + Ok(batch.project(data_projection)?) } /// Given a stream that includes a row id column, return a stream that will /// capture the row id. At completion of the stream, the captured row ids can /// be received from the returned receiver. +/// +/// A `_row_created_at_version` column, if the stream carries one, is captured +/// alongside the row ids and removed from the output the same way. pub fn make_rowid_capture_stream( mut target: SendableRecordBatchStream, stable_row_ids: bool, @@ -57,14 +62,17 @@ pub fn make_rowid_capture_stream( let (row_id_idx, _) = schema .column_with_name(ROW_ID) .expect("Received a batch without row ids"); - let non_row_ids_cols = (0..schema.fields.len()) - .filter(|col| *col != row_id_idx) + let created_at_idx = schema + .column_with_name(ROW_CREATED_AT_VERSION) + .map(|(idx, _)| idx); + let data_cols = (0..schema.fields.len()) + .filter(|col| *col != row_id_idx && Some(*col) != created_at_idx) .collect::>(); - let output_schema = Arc::new(schema.project(&non_row_ids_cols)?); + let output_schema = Arc::new(schema.project(&data_cols)?); let stream = futures::stream::poll_fn(move |cx| match target.poll_next_unpin(cx) { std::task::Poll::Ready(Some(Ok(batch))) => { - let res = extract_row_ids(&mut row_ids, batch, row_id_idx, &non_row_ids_cols); + let res = extract_row_ids(&mut row_ids, batch, row_id_idx, created_at_idx, &data_cols); std::task::Poll::Ready(Some(res)) } std::task::Poll::Ready(Some(Err(err))) => std::task::Poll::Ready(Some(Err(err))), @@ -84,13 +92,21 @@ pub fn make_rowid_capture_stream( #[derive(Debug)] pub enum CapturedRowIds { AddressStyle(RoaringTreemap), - SequenceStyle(RowIdSequence), + SequenceStyle { + row_ids: RowIdSequence, + /// The created-at version of each captured row, in capture order, when + /// the stream carried them; empty otherwise. + created_at: Vec, + }, } impl CapturedRowIds { pub fn new(stable_row_ids: bool) -> Self { if stable_row_ids { - Self::SequenceStyle(RowIdSequence::new()) + Self::SequenceStyle { + row_ids: RowIdSequence::new(), + created_at: Vec::new(), + } } else { Self::AddressStyle(RoaringTreemap::new()) } @@ -103,16 +119,39 @@ impl CapturedRowIds { ids.append(row_ids.iter().cloned()) .map_err(|e| datafusion::error::DataFusionError::Execution(e.to_string()))?; } - Self::SequenceStyle(sequence) => { + Self::SequenceStyle { + row_ids: sequence, .. + } => { sequence.extend(row_ids.into()); } } Ok(()) } + /// Record the created-at versions of the rows just passed to [`Self::capture`]. + pub fn capture_created_at(&mut self, versions: &[u64]) { + if let Self::SequenceStyle { created_at, .. } = self { + created_at.extend_from_slice(versions); + } + } + pub fn row_id_sequence(&self) -> Option<&RowIdSequence> { match self { - Self::SequenceStyle(sequence) => Some(sequence), + Self::SequenceStyle { row_ids, .. } => Some(row_ids), + _ => None, + } + } + + /// The captured rows' created-at versions, one per captured row id, when + /// the stream carried them. + pub fn created_at_versions(&self) -> Option<&[u64]> { + match self { + Self::SequenceStyle { + row_ids, + created_at, + } if created_at.len() as u64 == row_ids.len() && !created_at.is_empty() => { + Some(created_at) + } _ => None, } } @@ -120,7 +159,9 @@ impl CapturedRowIds { pub fn row_addrs(&self, index: Option<&RowIdIndex>) -> Result> { match self { Self::AddressStyle(addrs) => Ok(Cow::Borrowed(addrs)), - Self::SequenceStyle(sequence) => { + Self::SequenceStyle { + row_ids: sequence, .. + } => { let mut treemap = RoaringTreemap::new(); let Some(index) = index else { panic!("RowIdIndex required for sequence style row ids") diff --git a/rust/lance/src/dataset/write/update.rs b/rust/lance/src/dataset/write/update.rs index ee08ef22b74..0c4ab9ac1f9 100644 --- a/rust/lance/src/dataset/write/update.rs +++ b/rust/lance/src/dataset/write/update.rs @@ -8,7 +8,7 @@ use std::time::Duration; use super::cleanup_data_fragments; use super::retry::{RetryConfig, RetryExecutor, execute_with_retry}; use super::{CommitBuilder, WriteParams, write_fragments_internal}; -use crate::dataset::rowids::get_row_id_index; +use crate::dataset::rowids::{RowLineage, get_row_id_index, place_row_lineage}; use crate::dataset::transaction::UpdateMode::RewriteRows; use crate::dataset::transaction::{Operation, Transaction}; use crate::dataset::utils::make_rowid_capture_stream; @@ -28,10 +28,10 @@ use lance_arrow::json::{JsonArray, is_json_field}; use lance_core::datatypes::BlobHandling; use lance_core::error::{InvalidInputSnafu, box_error}; use lance_core::utils::tokio::get_num_compute_intensive_cpus; -use lance_core::{ROW_ADDR_FIELD, ROW_ID_FIELD, ROW_OFFSET_FIELD}; +use lance_core::{ROW_ADDR_FIELD, ROW_CREATED_AT_VERSION, ROW_ID_FIELD, ROW_OFFSET_FIELD}; use lance_file::version::{ConcreteFileVersion, LanceFileVersion}; use lance_select::RowAddrTreeMap; -use lance_table::format::{Fragment, RowIdMeta}; +use lance_table::format::{Fragment, RowDatasetVersionSequence, RowIdMeta}; use roaring::RoaringTreemap; use snafu::ResultExt; @@ -336,6 +336,28 @@ impl UpdateJob { scanner.with_row_address(); } scanner.with_row_id(); + // The rewritten rows keep their created-at versions, so read them here + // where the source rows are in hand rather than have the commit look + // them up again. The legacy file reader does not serve the version + // columns; those datasets keep leaving the lookup to the commit. + if self.dataset.manifest.uses_stable_row_ids() + && self + .dataset + .manifest + .data_storage_format + .lance_file_format() + != ConcreteFileVersion::V1 + { + let columns = self + .dataset + .schema() + .fields + .iter() + .map(|field| field.name.as_str()) + .chain([ROW_CREATED_AT_VERSION]) + .collect::>(); + scanner.project(&columns)?; + } if let Some(expr) = &self.condition { scanner.filter_expr(expr.clone()); @@ -471,23 +493,18 @@ impl UpdateJob { .map_err(|err| Error::internal(format!("Failed to receive row ids: {}", err)))?; if let Some(row_id_sequence) = removed_row_ids.row_id_sequence() { - let fragment_sizes = new_fragments - .iter() - .map(|f| f.physical_rows.unwrap() as u64); - let sequences = lance_table::rowids::rechunk_sequences( - [row_id_sequence.clone()], - fragment_sizes, - false, - ) - .map_err(|e| { - Error::internal(format!( - "Captured row ids not equal to number of rows written: {}", - e - )) - })?; - for (fragment, sequence) in new_fragments.iter_mut().zip(sequences) { - let serialized = lance_table::rowids::write_row_ids(&sequence); - fragment.row_id_meta = Some(RowIdMeta::Inline(serialized.into())); + let placed = self + .place_rewritten_lineage(&mut new_fragments, &removed_row_ids, row_id_sequence) + .await; + if let Err(e) = placed { + cleanup_data_fragments( + &self.dataset.object_store, + &self.dataset.base, + None, + &new_fragments, + ) + .await; + return Err(e); } } @@ -524,6 +541,65 @@ impl UpdateJob { }) } + /// Give each new fragment the row ids and created-at versions its rows + /// carried before the rewrite, placed inline or spilled as the table's + /// policy calls for. The last-updated-at version is the commit's to stamp; + /// a single-run placeholder keeps the metadata complete until then. + async fn place_rewritten_lineage( + &self, + new_fragments: &mut [Fragment], + captured: &crate::dataset::utils::CapturedRowIds, + row_id_sequence: &lance_table::rowids::RowIdSequence, + ) -> Result<()> { + let fragment_sizes = new_fragments + .iter() + .map(|f| f.physical_rows.unwrap() as u64) + .collect::>(); + let sequences = lance_table::rowids::rechunk_sequences( + [row_id_sequence.clone()], + fragment_sizes.iter().copied(), + false, + ) + .map_err(|e| { + Error::internal(format!( + "Captured row ids not equal to number of rows written: {}", + e + )) + })?; + // Without the created-at versions in hand the commit resolves them from + // the existing fragments, as it always has; the row ids stay inline so + // it can read them. + let Some(created_at) = captured.created_at_versions() else { + for (fragment, row_ids) in new_fragments.iter_mut().zip(sequences) { + let serialized = lance_table::rowids::write_row_ids(&row_ids); + fragment.row_id_meta = Some(RowIdMeta::Inline(serialized.into())); + } + return Ok(()); + }; + let placeholder_version = self.dataset.manifest.version + 1; + let mut offset = 0; + for ((fragment, row_ids), size) in + new_fragments.iter_mut().zip(sequences).zip(fragment_sizes) + { + let size = size as usize; + let lineage = RowLineage { + row_ids, + created_at: RowDatasetVersionSequence::from_versions( + &created_at[offset..offset + size], + ), + last_updated_at: RowDatasetVersionSequence::from_uniform_row_count( + size as u64, + placeholder_version, + ), + }; + offset += size; + place_row_lineage(&self.dataset, &lineage) + .await? + .apply(fragment); + } + Ok(()) + } + async fn commit_impl( &self, dataset: Arc, From 1b0e9d4c74b45acbdae536f550801d322b9414e3 Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Tue, 29 Sep 2026 05:14:41 +0800 Subject: [PATCH 2/5] fix(table): keep row lineage carriers through schema-only commits A spilled row lineage sequence lives only in the data file whose fields carry its reserved id (-3, -4 or -5), and those ids are never in the dataset schema. Three commit paths decide whether a data file is still live by asking whether any of its fields is in the schema: Project (drop_columns, and alter_columns renames and nullability changes), DataReplacement, and the alter_columns cast. Each of them dropped the carrier while the fragment kept its Column marker. After an update on an opted-in table, a single rename or drop left a manifest whose next row-id read failed, and cleanup would then delete the only copy of the rows' ids and created-at versions. Add Fragment::spilled_row_lineage_field_ids and have the three retain sites also keep a file that carries one of the fragment's spilled ids. Add Fragment::validate_row_lineage_carriers and call it on every fragment in build_manifest before the manifest is assembled. It requires each Column arm to have exactly one v2 carrier, so any future path that loses a carrier fails the commit instead of publishing lineage no reader can load. The check reads metadata only and skips fragments that spill nothing. It covers every fragment, not only the ones the operation touched, so it also refuses to build on a manifest that already lost a carrier. Co-Authored-By: Claude Opus 5.5 --- rust/lance-table/src/format/fragment.rs | 60 ++++- .../src/transaction/manifest_build.rs | 224 +++++++++++++++++- rust/lance/src/dataset/rowids/spill.rs | 129 +++++++++- rust/lance/src/dataset/schema_evolution.rs | 5 +- 4 files changed, 407 insertions(+), 11 deletions(-) diff --git a/rust/lance-table/src/format/fragment.rs b/rust/lance-table/src/format/fragment.rs index 9d5763e88dd..4f4b98cd5a0 100644 --- a/rust/lance-table/src/format/fragment.rs +++ b/rust/lance-table/src/format/fragment.rs @@ -13,7 +13,10 @@ use object_store::path::Path; use serde::{Deserialize, Deserializer, Serialize, Serializer}; use super::overlay::{DataOverlayFile, TOMBSTONE_FIELD_ID, sort_overlays_newest_last}; -use super::row_ids::RowIdMeta; +use super::row_ids::{ + ROW_CREATED_AT_VERSION_FIELD_ID, ROW_ID_FIELD_ID, ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, + RowIdMeta, +}; use crate::format::pb; use crate::rowids::version::{ @@ -600,6 +603,61 @@ impl Fragment { Ok(file) } + /// The reserved field ids of the row lineage sequences this fragment + /// marks as spilled: [`ROW_ID_FIELD_ID`] when its row ids are, and + /// likewise [`ROW_CREATED_AT_VERSION_FIELD_ID`] and + /// [`ROW_LAST_UPDATED_AT_VERSION_FIELD_ID`] for its versions. + /// + /// These ids are never in the dataset schema, so code that decides whether + /// a data file is still needed by its schema fields has to keep a file + /// carrying one of them as well: that file is the sequence's only copy. + pub fn spilled_row_lineage_field_ids(&self) -> Vec { + let mut field_ids = Vec::new(); + if matches!(self.row_id_meta, Some(RowIdMeta::Column)) { + field_ids.push(ROW_ID_FIELD_ID); + } + if matches!( + self.created_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ) { + field_ids.push(ROW_CREATED_AT_VERSION_FIELD_ID); + } + if matches!( + self.last_updated_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ) { + field_ids.push(ROW_LAST_UPDATED_AT_VERSION_FIELD_ID); + } + field_ids + } + + /// Check that every sequence this fragment marks as spilled has exactly + /// one carrier among [`Self::files`], and that the carrier is a v2 file, + /// the only version that can hold the columns. + /// + /// This reads metadata only. Committing a fragment that fails it would + /// publish lineage no reader can load, and once cleanup removed the + /// unreferenced carrier that lineage would be lost for good. + pub(crate) fn validate_row_lineage_carriers(&self) -> Result<()> { + for field_id in self.spilled_row_lineage_field_ids() { + let Some(file) = self.row_lineage_file(field_id)? else { + return Err(Error::internal(format!( + "cannot commit fragment {}: it marks row lineage field {} as spilled but \ + none of its data files carries it", + self.id, field_id + ))); + }; + if file.file_version()? == ConcreteFileVersion::V1 { + return Err(Error::internal(format!( + "cannot commit fragment {}: its spilled row lineage field {} is carried by \ + legacy v1 data file {}, which cannot hold row lineage columns", + self.id, field_id, file.path + ))); + } + } + Ok(()) + } + pub fn from_json(json: &str) -> Result { let fragment: Self = serde_json::from_str(json)?; Ok(fragment) diff --git a/rust/lance-table/src/transaction/manifest_build.rs b/rust/lance-table/src/transaction/manifest_build.rs index 42feeb2e4da..5320e8f1543 100644 --- a/rust/lance-table/src/transaction/manifest_build.rs +++ b/rust/lance-table/src/transaction/manifest_build.rs @@ -1202,15 +1202,18 @@ impl Transaction { // We might have removed all fields for certain data files, so // we should remove the data files that are no longer relevant. + // A file carrying a spilled row lineage sequence stays: its + // reserved ids are never in the schema, and it is the only copy. let remaining_field_ids = schema .fields_pre_order() .map(|f| f.id) .collect::>(); for fragment in final_fragments.iter_mut() { + let spilled = fragment.spilled_row_lineage_field_ids(); fragment.files.retain(|file| { - file.fields - .iter() - .any(|field_id| remaining_field_ids.contains(field_id)) + file.fields.iter().any(|field_id| { + remaining_field_ids.contains(field_id) || spilled.contains(field_id) + }) }); } @@ -1367,14 +1370,22 @@ impl Transaction { // the dataset schema: a file kept alive only by // tombstones or by ids the schema no longer defines is // unreachable to readers, uncollectable by cleanup, - // and reported corrupt by validate(). + // and reported corrupt by validate(). The exception is + // a file carrying one of the fragment's spilled row + // lineage sequences: their reserved ids are never in + // the schema, the file holds their only copy, and the + // replaced user columns in it are dead space until + // compaction. let live_ids = schema .fields_pre_order() .map(|field| field.id) .collect::>(); - new_frag - .files - .retain(|file| file.fields.iter().any(|f| live_ids.contains(f))); + let spilled = new_frag.spilled_row_lineage_field_ids(); + new_frag.files.retain(|file| { + file.fields + .iter() + .any(|f| live_ids.contains(f) || spilled.contains(f)) + }); new_frag.files.push(new_file.clone()); } @@ -1550,6 +1561,17 @@ impl Transaction { } } + // A spilled row lineage sequence lives only in its carrier file. An + // operation that dropped or duplicated that file would otherwise go + // unnoticed until the next read, by which time cleanup may have + // deleted the only copy. Every fragment is checked, not only the ones + // this operation touched: the check reads metadata only, skips + // fragments that spill nothing, and also refuses to build on a + // manifest that already lost a carrier. + for fragment in &final_fragments { + fragment.validate_row_lineage_carriers()?; + } + let user_requested_version = match (&config.storage_format, config.use_legacy_format) { (Some(storage_format), _) => Some(storage_format.lance_file_format()), (None, Some(true)) => Some(ConcreteFileVersion::V1), @@ -1902,7 +1924,8 @@ mod tests { use crate::format::overlay::OverlayCoverage; use crate::format::pb; use crate::format::{ - DeletionFile, DeletionFileType, RowDatasetVersionMeta, RowDatasetVersionSequence, RowIdMeta, + DeletionFile, DeletionFileType, ROW_CREATED_AT_VERSION_FIELD_ID, ROW_ID_FIELD_ID, + RowDatasetVersionMeta, RowDatasetVersionSequence, RowIdMeta, }; use crate::rowids::{RowIdSequence, write_row_ids}; use crate::transaction::test_support::{ @@ -4400,6 +4423,191 @@ mod tests { assert_eq!(fragment.overlays[0].data_file.fields.as_ref(), &[5]); } + /// A projection keeps only the data files that still share a field with + /// the schema. The reserved ids of spilled row lineage are never in it, so + /// without an exception the file carrying them, their only copy, would go. + #[test] + fn project_keeps_file_carrying_spilled_row_lineage() { + let mut fragment = Fragment::new(0); + fragment.physical_rows = Some(4); + fragment.row_id_meta = Some(RowIdMeta::Column); + fragment.created_at_version_meta = Some(RowDatasetVersionMeta::Column); + fragment.files = vec![ + DataFile::new( + "d.lance", + vec![0, 1], + vec![0, 1], + ConcreteFileVersion::V2_0, + None, + None, + ), + DataFile::new( + "l.lance", + vec![ROW_ID_FIELD_ID, ROW_CREATED_AT_VERSION_FIELD_ID], + vec![0, 1], + ConcreteFileVersion::V2_0, + None, + None, + ), + ]; + let schema = ArrowSchema::new(vec![ + ArrowField::new("a", DataType::Int32, false), + ArrowField::new("b", DataType::Int32, false), + ]); + let mut manifest = Manifest::new( + LanceSchema::try_from(&schema).unwrap(), + Arc::new(vec![fragment]), + DataStorageFormat::new(ConcreteFileVersion::V2_0), + HashMap::new(), + ); + manifest.reader_feature_flags = FLAG_STABLE_ROW_IDS; + manifest.writer_feature_flags = FLAG_STABLE_ROW_IDS; + + // Drop field 1, as `drop_columns(["b"])` would. + let transaction = Transaction::new( + manifest.version, + Operation::Project { + schema: manifest.schema.project_by_ids(&[0], true), + preserves_nullability: true, + }, + None, + ); + let (result, _) = transaction + .build_manifest(Some(&manifest), vec![], "txn", &default_build_config()) + .unwrap(); + + let fragment = &result.fragments[0]; + let paths = fragment + .files + .iter() + .map(|file| file.path.as_str()) + .collect::>(); + assert_eq!(paths, ["d.lance", "l.lance"]); + for field_id in [ROW_ID_FIELD_ID, ROW_CREATED_AT_VERSION_FIELD_ID] { + assert_eq!( + fragment + .row_lineage_file(field_id) + .unwrap() + .map(|file| file.path.as_str()), + Some("l.lance"), + "field {field_id}" + ); + } + } + + /// Replacing part of a wider file tombstones the replaced fields where + /// they live and drops any file left without a schema field. The file + /// carrying the fragment's spilled row ids may have none, yet it is their + /// only copy and must stay next to the new file that answers for the user + /// fields. The carrier is either a file of its own, as an update writes + /// it, or a file that also holds user columns, all of which are replaced. + #[rstest::rstest] + #[case::lineage_only_file( + vec![("v.lance", vec![3, 4, 5, 6]), ("l.lance", vec![ROW_ID_FIELD_ID])], + vec![3, 4], + vec![ + ("v.lance", vec![TOMBSTONE_FIELD_ID, TOMBSTONE_FIELD_ID, 5, 6]), + ("l.lance", vec![ROW_ID_FIELD_ID]), + ("v-new.lance", vec![3, 4]), + ], + "l.lance" + )] + #[case::in_file_carrier( + vec![("v.lance", vec![3, 4, 5, 6, ROW_ID_FIELD_ID])], + vec![3, 4, 5, 6], + vec![ + ( + "v.lance", + vec![ + TOMBSTONE_FIELD_ID, + TOMBSTONE_FIELD_ID, + TOMBSTONE_FIELD_ID, + TOMBSTONE_FIELD_ID, + ROW_ID_FIELD_ID, + ], + ), + ("v-new.lance", vec![3, 4, 5, 6]), + ], + "v.lance" + )] + fn data_replacement_keeps_lineage_carrier( + #[case] files: Vec<(&str, Vec)>, + #[case] replaced_fields: Vec, + #[case] expected_files: Vec<(&str, Vec)>, + #[case] carrier: &str, + ) { + let mut fragment = Fragment::new(0); + fragment.row_id_meta = Some(RowIdMeta::Column); + fragment.files = files + .into_iter() + .map(|(path, fields)| { + let column_indices = (0..fields.len() as i32).collect(); + DataFile::new( + path, + fields, + column_indices, + ConcreteFileVersion::V2_0, + None, + None, + ) + }) + .collect(); + + let fragment = replace_fields(fragment, replaced_fields, 1, 1).unwrap(); + let files = fragment + .files + .iter() + .map(|file| (file.path.as_str(), file.fields.to_vec())) + .collect::>(); + assert_eq!(files, expected_files); + assert_eq!( + fragment + .row_lineage_file(ROW_ID_FIELD_ID) + .unwrap() + .map(|file| file.path.as_str()), + Some(carrier) + ); + } + + /// Whatever hands the commit a fragment whose spilled row ids have no v2 + /// carrier -- an operation that dropped the file, or a writer that never + /// attached it -- the commit refuses it rather than publish lineage that + /// no reader can load. + #[rstest::rstest] + #[case::no_carrier(vec![0], ConcreteFileVersion::V2_0)] + #[case::v1_carrier(vec![0, ROW_ID_FIELD_ID], ConcreteFileVersion::V1)] + fn build_manifest_rejects_spilled_arm_without_carrier( + #[case] fields: Vec, + #[case] version: ConcreteFileVersion, + ) { + let manifest = sample_manifest(); + let mut fragment = Fragment::new(0); + fragment.physical_rows = Some(4); + fragment.row_id_meta = Some(RowIdMeta::Column); + // Column indices do not matter to the commit, only ids and version. + let data_file = DataFile::new("d.lance", fields, vec![], version, None, None); + fragment.files = vec![data_file]; + let transaction = Transaction::new( + manifest.version, + Operation::Merge { + fragments: vec![fragment], + schema: manifest.schema.clone(), + preserves_nullability: true, + }, + None, + ); + + let err = transaction + .build_manifest(Some(&manifest), vec![], "txn", &default_build_config()) + .unwrap_err(); + assert!(matches!(err, Error::Internal { .. }), "{err:?}"); + let message = err.to_string(); + assert!( + message.contains("fragment 0") && message.contains(&ROW_ID_FIELD_ID.to_string()), + "{message}" + ); + } + #[test] fn test_data_overlay_build_manifest_merges_duplicate_groups() { // Two groups targeting the same fragment must both survive (a HashMap diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs index e1e0fa22814..b656fc73247 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::{UpdateBuilder, WriteMode, WriteParams}; + use crate::dataset::{ColumnAlteration, UpdateBuilder, WriteMode, WriteParams}; use arrow_array::{Int32Array, RecordBatchIterator}; use arrow_schema::Field; use chrono::Utc; @@ -424,6 +424,7 @@ mod tests { 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; + use rstest::rstest; /// A sequence with no runs to exploit, which is what a globally shuffled /// table produces and what forces the spill path. @@ -650,6 +651,42 @@ mod tests { dataset.unwrap() } + /// Like [`appended_dataset`] with four appends of 250 rows, plus a column + /// `j` next to the key `i`, so a schema change has a column to rename, + /// drop or cast while `i` still identifies every row. + async fn two_column_dataset(uri: &str) -> Dataset { + let schema = Arc::new(ArrowSchema::new(vec![ + Field::new("i", DataType::Int32, false), + Field::new("j", DataType::Int32, true), + ])); + let mut dataset: Option = None; + for chunk in 0..4 { + let keys = Int32Array::from_iter_values((chunk * 250)..((chunk + 1) * 250)); + let batch = + RecordBatch::try_new(schema.clone(), vec![Arc::new(keys.clone()), Arc::new(keys)]) + .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, @@ -1185,4 +1222,94 @@ mod tests { } updated.validate().await.unwrap(); } + + /// The write that spills the fixture's lineage ahead of a schema change. + /// Both leave it in a lineage-only file next to the fragment's data. + #[derive(Debug, Clone, Copy)] + enum SpillingWrite { + /// An update rewriting the rows with `i >= 500`. + Update, + /// A compaction into a single fragment. + Compaction, + } + + /// A change to column `j` that adds or removes no rows. + #[derive(Debug, Clone, Copy)] + enum SchemaChange { + Rename, + Drop, + Cast, + } + + /// A schema change keeps only the data files that still hold a schema + /// field, and the reserved ids of spilled lineage never are one. Renames + /// and drops commit a projection and a cast rewrites the column; each must + /// keep the file carrying the lineage, which is its only copy. + #[rstest] + #[case::update_then_rename(SpillingWrite::Update, SchemaChange::Rename)] + #[case::update_then_drop(SpillingWrite::Update, SchemaChange::Drop)] + #[case::update_then_cast(SpillingWrite::Update, SchemaChange::Cast)] + #[case::compact_then_drop(SpillingWrite::Compaction, SchemaChange::Drop)] + #[tokio::test] + async fn schema_change_keeps_spilled_lineage( + #[case] spilling_write: SpillingWrite, + #[case] change: SchemaChange, + ) { + let dir = TempStrDir::default(); + let uri = dir.as_str(); + let mut dataset = two_column_dataset(uri).await; + spill_everything(&mut dataset).await; + match spilling_write { + SpillingWrite::Update => { + let updated = UpdateBuilder::new(Arc::new(dataset)) + .update_where("i >= 500") + .unwrap() + .set("i", "i + 10000") + .unwrap() + .build() + .unwrap() + .execute() + .await + .unwrap(); + dataset = updated.new_dataset.as_ref().clone(); + } + SpillingWrite::Compaction => { + compact_files(&mut dataset, one_fragment(), None) + .await + .unwrap(); + } + } + assert!( + dataset + .get_fragments() + .iter() + .any(|fragment| fragment.metadata().has_spilled_row_lineage()), + "the fixture must spill before the schema change" + ); + let before = by_key(&collect_rows(&dataset).await); + + let j = || ColumnAlteration::new("j".to_string()); + let changed = match change { + SchemaChange::Rename => dataset.alter_columns(&[j().rename("k".to_string())]).await, + SchemaChange::Drop => dataset.drop_columns(&["j"]).await, + SchemaChange::Cast => dataset.alter_columns(&[j().cast_to(DataType::Int64)]).await, + }; + changed.unwrap(); + + let mut expected = before; + if matches!(change, SchemaChange::Cast) { + // A cast rewrites `j` in every fragment, which the commit records + // as an update of every row; ids and created-at stay. + let cast_version = dataset.version().version; + for (_, _, updated_at) in expected.values_mut() { + *updated_at = cast_version; + } + } + assert_eq!(by_key(&collect_rows(&dataset).await), expected); + dataset.validate().await.unwrap(); + + // Re-opened cold, so the lineage is read back from the committed files. + let reopened = Dataset::open(uri).await.unwrap(); + assert_eq!(by_key(&collect_rows(&reopened).await), expected); + } } diff --git a/rust/lance/src/dataset/schema_evolution.rs b/rust/lance/src/dataset/schema_evolution.rs index 6b2afd1986a..b57f12c80b4 100644 --- a/rust/lance/src/dataset/schema_evolution.rs +++ b/rust/lance/src/dataset/schema_evolution.rs @@ -1119,10 +1119,13 @@ pub(super) async fn alter_columns( .collect::>() .into(); } + // A file carrying a spilled row lineage sequence stays: its + // reserved ids are never in the schema, and it is the only copy. + let spilled = frag.spilled_row_lineage_field_ids(); frag.files.retain(|f| { f.fields .iter() - .any(|field| schema_field_ids.contains(field)) + .any(|field| schema_field_ids.contains(field) || spilled.contains(field)) }); Ok(frag) }) From 5ae7de0137a4bb9bbcfc97cfb9bd50df4662f7d1 Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Tue, 29 Sep 2026 05:12:11 +0800 Subject: [PATCH 3/5] perf(table): skip created-at lookup for writer-placed update lineage The commit resolved every Update new fragment's created-at versions from the existing fragments before noticing that a writer had already placed them, and then discarded the result: it collected the rewritten row ids, walked the row ids of every overlapping existing fragment and decoded the source fragments' created-at sequences. It also read every spilled lineage sequence of the table ahead of each attempt, whether or not the build would look at it. Writer-placed lineage is now defined narrowly: spilled row ids or spilled created-at versions. Those are the only values the commit can neither resolve nor rewrite, and only a writer carrying rows' lineage over produces them. Such fragments skip the lookup and only have their last-updated-at versions stamped; inline created-at versions next to spilled row ids must cover every physical row. Inline created-at versions next to inline row ids are resolved again, as before this stack, so a hand-built transaction cannot commit stale ones and the documented Python contract still holds. The commit skips the spilled-lineage read-ahead for an update whose new fragments all carry writer-placed lineage and that patched no offsets in place. Co-Authored-By: Claude Opus 5.5 --- rust/lance-table/src/transaction.rs | 1 + .../src/transaction/row_version.rs | 216 ++++++++++++++++-- rust/lance/src/io/commit.rs | 25 +- 3 files changed, 221 insertions(+), 21 deletions(-) diff --git a/rust/lance-table/src/transaction.rs b/rust/lance-table/src/transaction.rs index 8d5c6439773..9a000afb8d1 100644 --- a/rust/lance-table/src/transaction.rs +++ b/rust/lance-table/src/transaction.rs @@ -48,6 +48,7 @@ pub use operation::{ TaggedRewriteAssembly, UpdateMode, UpdatedFragmentOffsets, reordered_sources, }; pub use prepare::{FragReuseUpdate, PreparedIndices}; +pub use row_version::has_writer_placed_lineage; pub use update_map::{ UpdateMap, UpdateMapEntry, translate_config_updates, translate_schema_metadata_updates, }; diff --git a/rust/lance-table/src/transaction/row_version.rs b/rust/lance-table/src/transaction/row_version.rs index 0f350b7f105..6524e4a0709 100644 --- a/rust/lance-table/src/transaction/row_version.rs +++ b/rust/lance-table/src/transaction/row_version.rs @@ -51,9 +51,40 @@ fn resolve_created_at_version( .unwrap_or(UNKNOWN_CREATED_AT_VERSION) } +/// Whether an update's new `fragment` carries row lineage its writer placed +/// outside the manifest, which the commit keeps instead of resolving. +/// +/// That is spilled row ids or spilled created-at versions. The commit cannot +/// read either back, and only a writer carrying the lineage of rows it read +/// over into the fragment produces them. Inline created-at versions next to +/// inline row ids do not count: the commit resolves those again from the row +/// ids, whatever the caller supplied, so a hand-built transaction cannot +/// commit stale ones. +/// +/// ``` +/// use lance_table::format::{Fragment, RowIdMeta}; +/// use lance_table::transaction::has_writer_placed_lineage; +/// +/// let mut fragment = Fragment::new(0); +/// assert!(!has_writer_placed_lineage(&fragment)); +/// fragment.row_id_meta = Some(RowIdMeta::Column); +/// assert!(has_writer_placed_lineage(&fragment)); +/// ``` +pub fn has_writer_placed_lineage(fragment: &Fragment) -> bool { + matches!(fragment.row_id_meta, Some(RowIdMeta::Column)) + || matches!( + fragment.created_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ) +} + /// For each new fragment produced by an update, set `created_at_version_meta` /// (preserved from the original rows) and `last_updated_at_version_meta`. /// +/// A fragment with writer-placed lineage (see [`has_writer_placed_lineage`]) +/// keeps its created-at versions and only has its last-updated-at versions +/// stamped. +/// /// `spilled` supplies the row ids and created-at versions of existing /// fragments that keep them outside the manifest; see /// [`SpilledRowLineage`]. @@ -65,8 +96,12 @@ pub(super) fn resolve_update_version_metadata( ) -> 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) + // A fragment with writer-placed lineage needs no lookup; when every new + // fragment has it, the set stays empty and the existing fragments are not + // scanned at all. let needed_row_ids: HashSet = new_fragments .iter() + .filter(|f| !has_writer_placed_lineage(f)) .filter_map(|f| match &f.row_id_meta { Some(RowIdMeta::Inline(data)) => read_row_ids(data).ok(), _ => None, @@ -169,18 +204,41 @@ pub(super) fn resolve_update_version_metadata( } for fragment in new_fragments.iter_mut() { - // A writer that read the rewritten rows already knows their created-at - // versions and may place them itself, spilled or inline. The - // last-updated-at version is this commit's, which only the commit - // knows (a conflict retry moves it), so it is stamped here regardless - // of what the writer left. - if fragment.created_at_version_meta.is_some() { - fragment.last_updated_at_version_meta = build_version_meta(fragment, new_version); - continue; - } - let row_ids = match &fragment.row_id_meta { - Some(RowIdMeta::Inline(data)) => read_row_ids(data).ok(), - Some(RowIdMeta::Column) => { + let row_ids = match (&fragment.row_id_meta, &fragment.created_at_version_meta) { + // Writer-placed lineage (see `has_writer_placed_lineage`): the + // writer read the rewritten rows and placed their created-at + // versions itself. The last-updated-at version is this commit's, + // which only the commit knows (a conflict retry moves it), so it + // is stamped here whatever the writer left. + (_, Some(RowDatasetVersionMeta::Column)) => { + fragment.last_updated_at_version_meta = build_version_meta(fragment, new_version); + continue; + } + (Some(RowIdMeta::Column), Some(created_at)) => { + // Unlike a spilled sequence, an inline one can be checked + // against the fragment here without IO. + let placed_rows = created_at + .load_sequence() + .map_err(|error| { + Error::invalid_input(format!( + "fragment {} carries created-at version metadata that does not \ + decode: {error}", + fragment.id + )) + })? + .len(); + let physical_rows = fragment.physical_rows.unwrap_or(0) as u64; + if placed_rows != physical_rows { + return Err(Error::invalid_input(format!( + "fragment {} carries {placed_rows} created-at versions for \ + {physical_rows} physical rows", + fragment.id + ))); + } + fragment.last_updated_at_version_meta = build_version_meta(fragment, new_version); + continue; + } + (Some(RowIdMeta::Column), None) => { // Resolving the versions needs the row ids, which this // commit-time path cannot read back from a data file; a writer // that spills them has to place the created-at versions too. @@ -191,7 +249,10 @@ pub(super) fn resolve_update_version_metadata( fragment.id ))); } - None => None, + // Inline created-at versions next to inline row ids are resolved + // again below, overwriting whatever the caller supplied. + (Some(RowIdMeta::Inline(data)), _) => read_row_ids(data).ok(), + (None, _) => None, }; if let Some(row_ids) = row_ids { @@ -353,6 +414,7 @@ mod tests { created_at_versions, default_build_config, last_updated_at_versions, make_stable_row_id_manifest, update_txn, }; + use rstest::rstest; use std::sync::Arc; #[test] @@ -417,9 +479,10 @@ mod tests { #[test] fn test_resolve_update_versions_keeps_writer_placed_created_at() { - // A writer that placed a new fragment's created-at versions, spilled - // or inline, keeps them; the commit only stamps last-updated-at with - // its own version, which the writer could not know. + // A writer that placed a new fragment's lineage outside the manifest + // keeps its created-at versions; the commit only stamps + // last-updated-at with its own version, which the writer could not + // know. let mut new_fragments = vec![Fragment { id: 1, physical_rows: Some(50), @@ -442,6 +505,127 @@ mod tests { ); } + #[test] + fn test_resolve_update_versions_skips_source_scan_for_spilled_new_fragments() { + // A new fragment whose writer spilled its created-at versions needs no + // lookup, so the commit does not scan the existing fragments' row ids + // at all: a spilled source that was not loaded ahead is no obstacle. + 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 mut new_fragments = vec![Fragment { + id: 2, + physical_rows: Some(2), + row_id_meta: Some(RowIdMeta::Inline( + write_row_ids(&RowIdSequence::from(10..12)).into(), + )), + files: vec![], + overlays: vec![], + deletion_file: None, + last_updated_at_version_meta: None, + created_at_version_meta: Some(RowDatasetVersionMeta::Column), + }]; + + resolve_update_version_metadata(&existing, &mut new_fragments, 9, &Default::default()) + .unwrap(); + assert_eq!( + new_fragments[0].created_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ); + assert_eq!( + new_fragments[0].last_updated_at_version_meta, + Some(inline_versions(2, 9)) + ); + } + + #[test] + fn test_resolve_update_versions_recomputes_inline_created_at() { + // Inline created-at versions next to inline row ids are not + // writer-placed: the commit resolves them from the row ids again, so a + // hand-built transaction cannot commit stale ones. + let row_ids = Some(RowIdMeta::Inline( + write_row_ids(&RowIdSequence::from(0..10)).into(), + )); + let existing = vec![Fragment { + id: 1, + physical_rows: Some(10), + row_id_meta: row_ids.clone(), + files: vec![], + overlays: vec![], + deletion_file: None, + last_updated_at_version_meta: Some(inline_versions(10, 3)), + created_at_version_meta: Some(inline_versions(10, 3)), + }]; + let mut new_fragments = vec![Fragment { + id: 2, + physical_rows: Some(10), + row_id_meta: row_ids, + files: vec![], + overlays: vec![], + deletion_file: None, + last_updated_at_version_meta: None, + created_at_version_meta: Some(inline_versions(10, 7)), + }]; + + resolve_update_version_metadata(&existing, &mut new_fragments, 9, &Default::default()) + .unwrap(); + assert_eq!( + new_fragments[0].created_at_version_meta, + Some(inline_versions(10, 3)) + ); + assert_eq!( + new_fragments[0].last_updated_at_version_meta, + Some(inline_versions(10, 9)) + ); + } + + #[rstest] + #[case::covers_every_row(50)] + #[case::too_short(40)] + fn test_resolve_update_versions_checks_placed_inline_created_at(#[case] placed_rows: u64) { + // Inline created-at versions next to spilled row ids are writer-placed + // and kept as they are, so they have to cover every row. + let mut fragments = 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: None, + created_at_version_meta: Some(inline_versions(placed_rows, 3)), + }]; + + let result = resolve_update_version_metadata(&[], &mut fragments, 9, &Default::default()); + if placed_rows == 50 { + result.unwrap(); + assert_eq!( + fragments[0].created_at_version_meta, + Some(inline_versions(50, 3)) + ); + assert_eq!( + fragments[0].last_updated_at_version_meta, + Some(inline_versions(50, 9)) + ); + } else { + let error = result.unwrap_err(); + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error + .to_string() + .contains("40 created-at versions for 50 physical rows"), + "{error}" + ); + } + } + #[test] fn test_resolve_update_versions_reads_spilled_source_lineage_from_config() { // The rewritten rows' sources are found by scanning existing row ids, diff --git a/rust/lance/src/io/commit.rs b/rust/lance/src/io/commit.rs index 61a15dcf287..107180ba1d9 100644 --- a/rust/lance/src/io/commit.rs +++ b/rust/lance/src/io/commit.rs @@ -44,7 +44,7 @@ use lance_table::io::commit::{ CommitConfig, CommitError, CommitHandler, ManifestLocation, ManifestNamingScheme, }; use lance_table::io::manifest::read_manifest; -use lance_table::transaction::{FragReuseUpdate, PreparedIndices}; +use lance_table::transaction::{FragReuseUpdate, PreparedIndices, has_writer_placed_lineage}; use rand::{Rng, rng}; use roaring::RoaringBitmap; @@ -1437,10 +1437,25 @@ async fn build_config_for_attempt( write_config: &ManifestWriteConfig, ) -> Result { let mut config = write_config.to_build_config(); - if matches!( - transaction.operation, - Operation::Update { .. } | Operation::DataOverlay { .. } - ) { + let reads_existing_lineage = match &transaction.operation { + // An update resolves its new fragments' created-at versions from the + // existing fragments unless their writer placed them, and refreshes + // the existing last-updated-at versions of the offsets it rewrote in + // place (`updated_fragment_offsets`). + Operation::Update { + new_fragments, + updated_fragment_offsets, + .. + } => { + !new_fragments.iter().all(has_writer_placed_lineage) + || updated_fragment_offsets + .as_ref() + .is_some_and(|offsets| !offsets.0.is_empty()) + } + Operation::DataOverlay { .. } => true, + _ => false, + }; + if reads_existing_lineage { config.spilled_row_lineage = load_spilled_row_lineage(dataset, dataset.manifest.fragments.iter()).await?; } From bf8679695c3049f1b0577f6e05ba7d2d9b1e626d Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Tue, 29 Sep 2026 05:19:11 +0800 Subject: [PATCH 4/5] perf(dataset): place update lineage only when it spills UpdateBuilder read every rewritten row's created-at version on any v2 table with stable row ids, held them as a Vec for the whole write, and placed them on every new fragment, inline or spilled, next to a placeholder last-updated-at sequence that could itself spill into the lineage file. Tables that never opted in paid for a created-at column on every scanned row and a larger transaction, and a malformed lance.row_lineage.inline_max_bytes was only reported after the whole rewrite had been written. The update now resolves the inline budget once, before any IO, and reads created-at versions only when the table may spill. A crate-private place_carried_row_lineage places what a rewritten fragment carries over: when neither the row ids nor the created-at versions exceed the budget, only the row ids are placed, inline, and the commit resolves created-at as it always has; once either spills, both are placed and only the over-budget ones go to the lineage file. Last-updated-at is never placed, since the commit stamps it. The public place_row_lineage is unchanged. Created-at versions are captured run-length encoded, extending a run across batch boundaries, and split per output fragment with rechunk_version_sequences. A capture that does not match the rows written is now an internal error instead of a silent fallback to the commit. Co-Authored-By: Claude Opus 5.5 --- rust/lance/src/dataset/rowids.rs | 1 + rust/lance/src/dataset/rowids/spill.rs | 178 ++++++++++++++++- rust/lance/src/dataset/utils.rs | 91 +++++++-- rust/lance/src/dataset/write/update.rs | 264 ++++++++++++++++++++----- 4 files changed, 463 insertions(+), 71 deletions(-) diff --git a/rust/lance/src/dataset/rowids.rs b/rust/lance/src/dataset/rowids.rs index 78644267fdd..4a77b45757a 100644 --- a/rust/lance/src/dataset/rowids.rs +++ b/rust/lance/src/dataset/rowids.rs @@ -24,6 +24,7 @@ use lance_table::{ }; use std::sync::Arc; +pub(crate) use spill::place_carried_row_lineage; pub use spill::{ DEFAULT_INLINE_ROW_LINEAGE_MAX_BYTES, INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY, PlacedRowLineage, RowLineage, SPILL_ROW_LINEAGE_CONFIG_KEY, inline_row_lineage_max_bytes, diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs index b656fc73247..6c04707ebb7 100644 --- a/rust/lance/src/dataset/rowids/spill.rs +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -61,10 +61,13 @@ const SPILL_BATCH_ROWS: usize = 64 * 1024; /// single-run version sequence. pub const DEFAULT_INLINE_ROW_LINEAGE_MAX_BYTES: usize = 200 * 1024; -/// Table config key that turns spilling on: `"true"` lets compaction move -/// oversized lineage sequences out of the manifest. Absent or anything else, -/// every sequence stays inline however large it grows, which is what every -/// released build does; a table that never sets it stays readable by them. +/// Table config key that turns spilling on: `"true"` lets compaction and +/// update move the oversized lineage sequences they carry over from existing +/// rows out of the manifest -- row ids and created-at versions, and at +/// compaction last-updated-at versions too; an update's last-updated-at +/// version is the commit's and stays inline. Absent or anything else, every +/// sequence stays inline however large it grows, which is what every released +/// build does; a table that never sets it stays readable by them. pub const SPILL_ROW_LINEAGE_CONFIG_KEY: &str = "lance.row_lineage.spill"; /// Table config key overriding [`DEFAULT_INLINE_ROW_LINEAGE_MAX_BYTES`], as a @@ -210,6 +213,99 @@ pub async fn place_row_lineage( }) } +/// The lineage an update carries over into one of its new fragments, as +/// [`place_carried_row_lineage`] placed it. +/// +/// `pub` rather than `pub(crate)` only because `clippy::redundant_pub_crate` +/// fires inside the private `spill` module; the crate-private re-export of +/// [`place_carried_row_lineage`] keeps it off the public API. +pub struct CarriedRowLineage { + row_ids: RowIdMeta, + /// `None` when neither sequence spilled: the commit then resolves the + /// created-at versions from the inline row ids, as on a table that has not + /// opted in. `Some`, inline or spilled, once either sequence spilled. + created_at: Option, + /// The lineage file holding the spilled sequences, `None` when nothing + /// spilled. + file: Option, +} + +impl CarriedRowLineage { + /// Put the placement on `fragment`: its row ids, its created-at versions + /// when they were placed, and the lineage file, if any, as one more of its + /// data files. The last-updated-at versions are left for the commit. + pub fn apply(self, fragment: &mut Fragment) { + fragment.row_id_meta = Some(self.row_ids); + fragment.created_at_version_meta = self.created_at; + fragment.last_updated_at_version_meta = None; + if let Some(file) = self.file { + fragment.files.push(file); + } + } +} + +/// Place the row ids and created-at versions that rewritten rows carry over +/// into a new fragment, spilling each whose encoding exceeds `limit` into a +/// hidden column of one new data file. +/// +/// When neither exceeds `limit`, only the row ids are placed, inline, and the +/// created-at versions are left for the commit to resolve from them, as it +/// does on a table that has not opted in. Once either spills the commit can no +/// longer resolve them, since it cannot read spilled row ids, so the created-at +/// versions are placed as well, spilled or inline. +/// +/// The last-updated-at version is never placed: it is the commit's, and a +/// conflict retry moves it. +pub async fn place_carried_row_lineage( + dataset: &Dataset, + limit: usize, + row_ids: &RowIdSequence, + created_at: &RowDatasetVersionSequence, +) -> Result { + let inline_row_ids = write_row_ids(row_ids); + let inline_created_at = write_dataset_versions(created_at); + let spill_row_ids = inline_row_ids.len() > limit; + let spill_created_at = inline_created_at.len() > limit; + if !spill_row_ids && !spill_created_at { + return Ok(CarriedRowLineage { + row_ids: RowIdMeta::Inline(inline_row_ids.into()), + created_at: None, + file: None, + }); + } + + // Materialized before the write, as in `place_row_lineage`, so the future + // stays `Send`. + let mut columns: Vec<(i32, &str, ArrayRef)> = Vec::with_capacity(2); + if spill_row_ids { + let ids = UInt64Array::from(row_ids.iter().collect::>()); + columns.push((ROW_ID_FIELD_ID, ROW_ID, Arc::new(ids))); + } + if spill_created_at { + let versions = UInt64Array::from(created_at.versions().collect::>()); + columns.push(( + ROW_CREATED_AT_VERSION_FIELD_ID, + ROW_CREATED_AT_VERSION, + Arc::new(versions), + )); + } + let file = write_lineage_file(dataset, &columns).await?; + + Ok(CarriedRowLineage { + row_ids: if spill_row_ids { + RowIdMeta::Column + } else { + RowIdMeta::Inline(inline_row_ids.into()) + }, + created_at: Some(if spill_created_at { + RowDatasetVersionMeta::Column + } else { + RowDatasetVersionMeta::Inline(inline_created_at.into()) + }), + file: Some(file), + }) +} + /// Write `columns` as the hidden columns of one new data file and return the /// [`DataFile`] that locates them, listing the columns' field ids in order. async fn write_lineage_file( @@ -577,6 +673,80 @@ mod tests { ); } + /// What an update carries over is placed in full only once something + /// spills: the commit resolves created-at versions from inline row ids but + /// cannot from spilled ones, and it always stamps last-updated-at itself. + #[rstest] + #[case::nothing_spills(false, false)] + #[case::row_ids_spill(true, false)] + #[case::created_at_spills(false, true)] + #[case::both_spill(true, true)] + #[tokio::test] + async fn carried_lineage_places_created_at_once_anything_spills( + #[case] spill_row_ids: bool, + #[case] spill_created_at: bool, + ) { + let dir = TempStrDir::default(); + let dataset = tiny_dataset(dir.as_str()).await; + // Under a 100-byte budget a scattered or alternating sequence of 1,000 + // rows spills, and a range or a single run stays inline. + let row_ids = if spill_row_ids { + scattered_row_ids(1_000) + } else { + RowIdSequence::from(0..1_000) + }; + let created_at = if spill_created_at { + alternating_versions(1_000, 1) + } else { + RowDatasetVersionSequence::from_uniform_row_count(1_000, 1) + }; + + let mut fragment = Fragment::new(42); + fragment.physical_rows = Some(1_000); + place_carried_row_lineage(&dataset, 100, &row_ids, &created_at) + .await + .unwrap() + .apply(&mut fragment); + + assert_eq!( + matches!(fragment.row_id_meta, Some(RowIdMeta::Column)), + spill_row_ids + ); + assert_eq!(fragment.last_updated_at_version_meta, None); + let placed_row_ids = load_row_id_sequence(&dataset, &fragment).await.unwrap(); + assert_eq!( + placed_row_ids.iter().collect::>(), + row_ids.iter().collect::>() + ); + if !spill_row_ids && !spill_created_at { + assert_eq!(fragment.created_at_version_meta, None); + assert!(fragment.files.is_empty()); + } else { + let spilled = [ + (spill_row_ids, ROW_ID_FIELD_ID), + (spill_created_at, ROW_CREATED_AT_VERSION_FIELD_ID), + ] + .into_iter() + .filter_map(|(spills, field_id)| spills.then_some(field_id)) + .collect::>(); + assert_eq!(fragment.files.len(), 1); + assert_eq!(fragment.files[0].fields.as_ref(), spilled.as_slice()); + assert_eq!( + matches!( + fragment.created_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ), + spill_created_at + ); + let placed_created_at = + load_row_version_sequence(&dataset, &fragment, RowVersionKind::CreatedAt) + .await + .unwrap() + .expect("created-at versions are placed once anything spills"); + assert_eq!(versions_of(&placed_created_at), versions_of(&created_at)); + } + } + /// The format allows the columns only in v2 files, so a legacy v1 table /// keeps everything inline however it is configured. #[tokio::test] diff --git a/rust/lance/src/dataset/utils.rs b/rust/lance/src/dataset/utils.rs index 5b23ab6b4a1..376eed126a0 100644 --- a/rust/lance/src/dataset/utils.rs +++ b/rust/lance/src/dataset/utils.rs @@ -15,6 +15,8 @@ use lance_arrow::json::{ has_arrow_json_fields, has_json_fields, lance_json_to_arrow_json, }; use lance_core::{ROW_CREATED_AT_VERSION, ROW_ID}; +use lance_table::format::{RowDatasetVersionRun, RowDatasetVersionSequence}; +use lance_table::rowids::segment::U64Segment; use lance_table::rowids::{RowIdIndex, RowIdSequence}; use roaring::RoaringTreemap; use std::borrow::Cow; @@ -65,6 +67,13 @@ pub fn make_rowid_capture_stream( let created_at_idx = schema .column_with_name(ROW_CREATED_AT_VERSION) .map(|(idx, _)| idx); + // Started here, rather than on the first batch, so a stream that carries + // the column but no rows still reports an empty capture. + if created_at_idx.is_some() + && let CapturedRowIds::SequenceStyle { created_at, .. } = &mut row_ids + { + *created_at = Some(RowDatasetVersionSequence::new()); + } let data_cols = (0..schema.fields.len()) .filter(|col| *col != row_id_idx && Some(*col) != created_at_idx) .collect::>(); @@ -94,9 +103,9 @@ pub enum CapturedRowIds { AddressStyle(RoaringTreemap), SequenceStyle { row_ids: RowIdSequence, - /// The created-at version of each captured row, in capture order, when - /// the stream carried them; empty otherwise. - created_at: Vec, + /// The created-at versions of the captured rows, in capture order, + /// when the stream carried them; `None` when it did not. + created_at: Option, }, } @@ -105,7 +114,7 @@ impl CapturedRowIds { if stable_row_ids { Self::SequenceStyle { row_ids: RowIdSequence::new(), - created_at: Vec::new(), + created_at: None, } } else { Self::AddressStyle(RoaringTreemap::new()) @@ -130,8 +139,36 @@ impl CapturedRowIds { /// Record the created-at versions of the rows just passed to [`Self::capture`]. pub fn capture_created_at(&mut self, versions: &[u64]) { - if let Self::SequenceStyle { created_at, .. } = self { - created_at.extend_from_slice(versions); + let Self::SequenceStyle { + created_at: Some(sequence), + .. + } = self + else { + return; + }; + // Run-length encoded as they arrive, and a run that carries on from + // the previous batch is extended rather than restarted: a run per + // batch would bloat the inline metadata of a large rewrite. + for run in versions.chunk_by(|a, b| a == b) { + let (version, rows) = (run[0], run.len() as u64); + let start = match sequence.runs.last_mut() { + Some(RowDatasetVersionRun { + span: U64Segment::Range(span), + version: last_version, + }) => { + if *last_version == version { + span.end += rows; + continue; + } + span.end + } + // Every run pushed below spans a range, so this is the first. + _ => 0, + }; + sequence.runs.push(RowDatasetVersionRun { + span: U64Segment::Range(start..start + rows), + version, + }); } } @@ -142,17 +179,12 @@ impl CapturedRowIds { } } - /// The captured rows' created-at versions, one per captured row id, when - /// the stream carried them. - pub fn created_at_versions(&self) -> Option<&[u64]> { + /// The captured rows' created-at versions, in capture order, when the + /// stream carried them. + pub fn created_at_sequence(&self) -> Option<&RowDatasetVersionSequence> { match self { - Self::SequenceStyle { - row_ids, - created_at, - } if created_at.len() as u64 == row_ids.len() && !created_at.is_empty() => { - Some(created_at) - } - _ => None, + Self::SequenceStyle { created_at, .. } => created_at.as_ref(), + Self::AddressStyle(_) => None, } } @@ -360,3 +392,30 @@ impl SchemaAdapter { )) } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn capture_created_at_extends_a_run_across_batches() { + let mut captured = CapturedRowIds::SequenceStyle { + row_ids: RowIdSequence::new(), + created_at: Some(RowDatasetVersionSequence::new()), + }; + captured.capture_created_at(&[1, 1, 2]); + captured.capture_created_at(&[]); + captured.capture_created_at(&[2, 3]); + // One run per version, the one split by the batch boundary included, + // with spans absolute over the whole capture. + assert_eq!( + captured.created_at_sequence(), + Some(&RowDatasetVersionSequence::from_versions(&[1, 1, 2, 2, 3])) + ); + + // A stream that did not carry the column captures nothing. + let mut not_carried = CapturedRowIds::new(true); + not_carried.capture_created_at(&[1, 2]); + assert_eq!(not_carried.created_at_sequence(), None); + } +} diff --git a/rust/lance/src/dataset/write/update.rs b/rust/lance/src/dataset/write/update.rs index 0c4ab9ac1f9..b9b4bb4fd1a 100644 --- a/rust/lance/src/dataset/write/update.rs +++ b/rust/lance/src/dataset/write/update.rs @@ -8,7 +8,9 @@ use std::time::Duration; use super::cleanup_data_fragments; use super::retry::{RetryConfig, RetryExecutor, execute_with_retry}; use super::{CommitBuilder, WriteParams, write_fragments_internal}; -use crate::dataset::rowids::{RowLineage, get_row_id_index, place_row_lineage}; +use crate::dataset::rowids::{ + get_row_id_index, inline_row_lineage_max_bytes, place_carried_row_lineage, +}; use crate::dataset::transaction::UpdateMode::RewriteRows; use crate::dataset::transaction::{Operation, Transaction}; use crate::dataset::utils::make_rowid_capture_stream; @@ -32,6 +34,8 @@ use lance_core::{ROW_ADDR_FIELD, ROW_CREATED_AT_VERSION, ROW_ID_FIELD, ROW_OFFSE use lance_file::version::{ConcreteFileVersion, LanceFileVersion}; use lance_select::RowAddrTreeMap; use lance_table::format::{Fragment, RowDatasetVersionSequence, RowIdMeta}; +use lance_table::rowids::version::rechunk_version_sequences; +use lance_table::rowids::{RowIdSequence, rechunk_sequences, write_row_ids}; use roaring::RoaringTreemap; use snafu::ResultExt; @@ -316,6 +320,15 @@ impl UpdateJob { } async fn execute_impl(self) -> Result { + // Resolved before any IO, so a malformed budget fails the update + // before it rewrites the rows rather than after. Only a table with + // stable row ids carries lineage over, so no other table depends on + // the setting. + let spill_budget = if self.dataset.manifest.uses_stable_row_ids() { + inline_row_lineage_max_bytes(&self.dataset)? + } else { + None + }; let mut scanner = self.dataset.scan(); let legacy_blob_ids = self .dataset @@ -336,18 +349,12 @@ impl UpdateJob { scanner.with_row_address(); } scanner.with_row_id(); - // The rewritten rows keep their created-at versions, so read them here - // where the source rows are in hand rather than have the commit look - // them up again. The legacy file reader does not serve the version - // columns; those datasets keep leaving the lookup to the commit. - if self.dataset.manifest.uses_stable_row_ids() - && self - .dataset - .manifest - .data_storage_format - .lance_file_format() - != ConcreteFileVersion::V1 - { + // The rewritten rows keep their created-at versions. On a table that + // may spill, read them here, where the source rows are in hand: once + // their row ids spill the commit can no longer look them up. Elsewhere + // the commit looks them up from the inline row ids, and reading them + // would only build a column for every scanned row. + if spill_budget.is_some() { let columns = self .dataset .schema() @@ -494,7 +501,12 @@ impl UpdateJob { if let Some(row_id_sequence) = removed_row_ids.row_id_sequence() { let placed = self - .place_rewritten_lineage(&mut new_fragments, &removed_row_ids, row_id_sequence) + .place_rewritten_lineage( + &mut new_fragments, + row_id_sequence, + removed_row_ids.created_at_sequence(), + spill_budget, + ) .await; if let Err(e) = placed { cleanup_data_fragments( @@ -541,59 +553,50 @@ impl UpdateJob { }) } - /// Give each new fragment the row ids and created-at versions its rows - /// carried before the rewrite, placed inline or spilled as the table's - /// policy calls for. The last-updated-at version is the commit's to stamp; - /// a single-run placeholder keeps the metadata complete until then. + /// Give each new fragment the row ids its rows carried before the rewrite + /// and, where the lineage spills under `spill_budget`, their created-at + /// versions too; see [`place_carried_row_lineage`]. The last-updated-at + /// version is the commit's to stamp. + /// + /// `created_at` is `None` when the scan did not read the created-at + /// versions, which it only does when the table may spill. async fn place_rewritten_lineage( &self, new_fragments: &mut [Fragment], - captured: &crate::dataset::utils::CapturedRowIds, - row_id_sequence: &lance_table::rowids::RowIdSequence, + row_ids: &RowIdSequence, + created_at: Option<&RowDatasetVersionSequence>, + spill_budget: Option, ) -> Result<()> { let fragment_sizes = new_fragments .iter() .map(|f| f.physical_rows.unwrap() as u64) .collect::>(); - let sequences = lance_table::rowids::rechunk_sequences( - [row_id_sequence.clone()], - fragment_sizes.iter().copied(), - false, - ) - .map_err(|e| { + let row_ids = rechunk_sequences([row_ids.clone()], fragment_sizes.iter().copied(), false) + .map_err(|e| { Error::internal(format!( "Captured row ids not equal to number of rows written: {}", e )) })?; - // Without the created-at versions in hand the commit resolves them from - // the existing fragments, as it always has; the row ids stay inline so - // it can read them. - let Some(created_at) = captured.created_at_versions() else { - for (fragment, row_ids) in new_fragments.iter_mut().zip(sequences) { - let serialized = lance_table::rowids::write_row_ids(&row_ids); - fragment.row_id_meta = Some(RowIdMeta::Inline(serialized.into())); + let (Some(limit), Some(created_at)) = (spill_budget, created_at) else { + // Nothing can spill: the row ids stay inline and the commit + // resolves the created-at versions from them, as it always has. + for (fragment, row_ids) in new_fragments.iter_mut().zip(row_ids) { + fragment.row_id_meta = Some(RowIdMeta::Inline(write_row_ids(&row_ids).into())); } return Ok(()); }; - let placeholder_version = self.dataset.manifest.version + 1; - let mut offset = 0; - for ((fragment, row_ids), size) in - new_fragments.iter_mut().zip(sequences).zip(fragment_sizes) + let created_at = + rechunk_version_sequences([created_at.clone()], fragment_sizes.iter().copied(), false) + .map_err(|e| { + Error::internal(format!( + "Captured created-at versions not equal to number of rows written: {e}" + )) + })?; + for ((fragment, row_ids), created_at) in + new_fragments.iter_mut().zip(row_ids).zip(created_at) { - let size = size as usize; - let lineage = RowLineage { - row_ids, - created_at: RowDatasetVersionSequence::from_versions( - &created_at[offset..offset + size], - ), - last_updated_at: RowDatasetVersionSequence::from_uniform_row_count( - size as u64, - placeholder_version, - ), - }; - offset += size; - place_row_lineage(&self.dataset, &lineage) + place_carried_row_lineage(&self.dataset, limit, &row_ids, &created_at) .await? .apply(fragment); } @@ -764,14 +767,19 @@ mod tests { use super::*; + use crate::dataset::rowids::{ + INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY, SPILL_ROW_LINEAGE_CONFIG_KEY, + read_spilled_row_ids, read_spilled_versions, + }; use crate::dataset::{WriteDestination, WriteMode}; use crate::index::DatasetIndexExt; use crate::index::vector::VectorIndexParams; - use crate::utils::test::{DatagenExt, FragmentCount, FragmentRowCount}; + use crate::utils::test::{DatagenExt, FailingProxyStore, FragmentCount, FragmentRowCount}; use arrow::{ array::AsArray, datatypes::{Int64Type, UInt32Type}, }; + use arrow_array::record_batch; use arrow_array::types::{Float32Type, Int32Type}; use arrow_array::{ Int64Array, RecordBatchIterator, StringArray, StructArray, UInt32Array, UInt64Array, @@ -790,6 +798,7 @@ mod tests { use lance_io::object_store::ObjectStoreParams; use lance_linalg::distance::MetricType; use lance_table::feature_flags::FLAG_MIXED_DATA_FILE_VERSIONS; + use lance_table::format::ROW_CREATED_AT_VERSION_FIELD_ID; use object_store::throttle::ThrottleConfig; use rstest::rstest; use tokio::sync::Barrier; @@ -2338,6 +2347,159 @@ mod tests { ); } + /// On a table that cannot spill, or when nothing an update carries over + /// exceeds the inline budget, the new fragments carry only inline row ids: + /// the commit resolves their created-at versions, as it always has. + #[rstest] + #[case::not_opted_in(false)] + #[case::under_budget(true)] + #[tokio::test] + async fn update_leaves_inline_created_at_to_the_commit(#[case] opted_in: bool) { + let (dataset, _test_dir) = make_test_dataset(LanceFileVersion::V2_0, true).await; + let mut dataset = dataset.as_ref().clone(); + if opted_in { + dataset + .update_config([(SPILL_ROW_LINEAGE_CONFIG_KEY, "true")]) + .await + .unwrap(); + } + let update_data = UpdateBuilder::new(Arc::new(dataset)) + .update_where("id >= 15") + .unwrap() + .set("name", "'bar'") + .unwrap() + .build() + .unwrap() + .execute_impl() + .await + .unwrap(); + + assert!(!update_data.new_fragments.is_empty()); + for fragment in &update_data.new_fragments { + assert!( + matches!(fragment.row_id_meta, Some(RowIdMeta::Inline(_))), + "{fragment:?}" + ); + assert_eq!(fragment.created_at_version_meta, None); + assert_eq!(fragment.last_updated_at_version_meta, None); + assert_eq!(fragment.files.len(), 1, "no lineage file: {fragment:?}"); + } + } + + /// A malformed inline budget fails the update before it writes anything. + /// Every write to the data directory fails here, so an update that + /// rewrote the rows first would report that failure instead. + #[tokio::test] + async fn update_rejects_malformed_inline_max_bytes_before_writing() { + let test_dir = TempStrDir::default(); + let test_uri = test_dir.as_str(); + // Prefix `/` so Windows drive letters (e.g. `C:`) don't get parsed as + // the URL authority. + let path_prefix = if test_uri.starts_with('/') { "" } else { "/" }; + let routed_uri = format!("file-object-store://{path_prefix}{test_uri}"); + let batch = record_batch!(("id", Int64, [0, 1, 2, 3, 4, 5])).unwrap(); + let schema = batch.schema(); + let write_params = WriteParams { + enable_stable_row_ids: true, + ..Default::default() + }; + let batches = RecordBatchIterator::new([Ok(batch)], schema); + let mut dataset = Dataset::write(batches, &routed_uri, Some(write_params)) + .await + .unwrap(); + dataset + .update_config([ + (SPILL_ROW_LINEAGE_CONFIG_KEY, "true"), + (INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY, "200KB"), + ]) + .await + .unwrap(); + + let failing = Arc::new(FailingProxyStore::new()); + failing.fail_when("put", "/data/", "injected data write failure"); + failing.fail_when("put_multipart", "/data/", "injected data write failure"); + let dataset = DatasetBuilder::from_uri(&routed_uri) + .with_read_params(ReadParams { + store_options: Some(ObjectStoreParams { + object_store_wrapper: Some(failing), + ..Default::default() + }), + ..Default::default() + }) + .load() + .await + .unwrap(); + + let error = UpdateBuilder::new(Arc::new(dataset)) + .update_where("id >= 3") + .unwrap() + .set("id", "id + 100") + .unwrap() + .build() + .unwrap() + .execute() + .await + .unwrap_err(); + assert!(matches!(error, Error::InvalidInput { .. }), "{error:?}"); + assert!( + error + .to_string() + .contains(INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY), + "{error}" + ); + } + + /// The captured lineage is split across the new fragments by their row + /// counts, and a created-at capture that does not cover every written row + /// is an error rather than a silent fallback to the commit. + #[tokio::test] + async fn place_rewritten_lineage_splits_lineage_by_output_fragment() { + let (dataset, _test_dir) = make_test_dataset(LanceFileVersion::V2_0, true).await; + let job = UpdateBuilder::new(dataset) + .set("name", "'bar'") + .unwrap() + .build() + .unwrap(); + let output_fragments = || { + [(1, 3), (2, 2)].map(|(id, rows)| { + let mut fragment = Fragment::new(id); + fragment.physical_rows = Some(rows); + fragment + }) + }; + let row_ids = RowIdSequence::from([10u64, 11, 12, 20, 21].as_slice()); + + // A zero budget spills everything, so each fragment's share is read + // back from its own lineage file. + let mut fragments = output_fragments(); + let created_at = RowDatasetVersionSequence::from_versions(&[1, 1, 2, 3, 3]); + job.place_rewritten_lineage(&mut fragments, &row_ids, Some(&created_at), Some(0)) + .await + .unwrap(); + let mut placed = Vec::new(); + for fragment in &fragments { + let ids = read_spilled_row_ids(&job.dataset, fragment).await.unwrap(); + let versions = + read_spilled_versions(&job.dataset, fragment, ROW_CREATED_AT_VERSION_FIELD_ID) + .await + .unwrap(); + placed.push(( + ids.iter().collect::>(), + versions.versions().collect::>(), + )); + } + assert_eq!(placed[0], (vec![10, 11, 12], vec![1, 1, 2])); + assert_eq!(placed[1], (vec![20, 21], vec![3, 3])); + + let mut fragments = output_fragments(); + let too_few = RowDatasetVersionSequence::from_versions(&[1, 1, 2, 3]); + let error = job + .place_rewritten_lineage(&mut fragments, &row_ids, Some(&too_few), Some(0)) + .await + .unwrap_err(); + assert!(matches!(error, Error::Internal { .. }), "{error:?}"); + } + #[tokio::test] async fn test_update_with_blob() { use arrow_array::LargeBinaryArray; From 14ea9db8a80dbbcd2113d14519eb199e1f41caae Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Tue, 29 Sep 2026 05:21:34 +0800 Subject: [PATCH 5/5] test(dataset): tighten update lineage tests The two update tests were copies that differed only in the opt-in and a few layout checks, and neither could tell carried-over lineage from a spilled last-updated-at placeholder: the opted-in one only checked that the lineage file's fields were negative. They are now one rstest that pins the lineage file to exactly the row id and created-at columns, checks the feature flags this update is the first to raise, and reads the lineage back from a cold open. The spilled-source update test drew its one rewritten row from the middle of a created-at run and still described commit-time resolution, which UpdateBuilder no longer relies on once its lineage spills. It now deletes rows first and rewrites a selection across the compacted fragment's run boundaries, checking every row, once with the writer placing the lineage and once under a budget that leaves the created-at lookup to the commit and its read-ahead. Co-Authored-By: Claude Opus 5.5 --- rust/lance/src/dataset/rowids/spill.rs | 230 +++++++++++++++---------- 1 file changed, 137 insertions(+), 93 deletions(-) diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs index 6c04707ebb7..7c2168e43e2 100644 --- a/rust/lance/src/dataset/rowids/spill.rs +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -1076,13 +1076,22 @@ mod tests { .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. + /// Updating rows whose lineage is spilled keeps every row's created-at + /// version. The update's scan reads each rewritten row's created-at + /// version, spilled column included. When what the update carries over + /// spills again, the writer places those versions and the commit only + /// stamps last-updated-at; when it fits inline, the commit resolves them + /// from the row ids, reading the source fragment's spilled sequences ahead + /// of the build. Deleted rows and a selection across the compacted + /// fragment's created-at runs mean a version read at the wrong offset + /// would show up as a neighbour's. + #[rstest] + #[case::writer_spills(true)] + #[case::commit_resolves(false)] #[tokio::test] - async fn updating_rows_with_spilled_lineage_keeps_their_created_at() { + async fn updating_rows_with_spilled_lineage_keeps_their_created_at( + #[case] update_spills: bool, + ) { let dir = TempStrDir::default(); let uri = dir.as_str(); let mut dataset = appended_dataset(uri, 4, 250).await; @@ -1090,16 +1099,34 @@ mod tests { compact_files(&mut dataset, one_fragment(), None) .await .unwrap(); + // Deleted rows move the compacted fragment's scan positions away from + // its physical offsets, which its lineage is indexed by. + dataset.delete("i % 7 = 0").await.unwrap(); + if !update_spills { + // A budget the update's carried-over lineage fits in, which leaves + // resolving the created-at versions to the commit. + dataset + .update_config([(INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY, "1000000")]) + .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 selected = |key: i32| (240..260).contains(&key) || (745..755).contains(&key); + // The selection spans the run boundaries at 250, between the first + // and second appends, and at 750, between the third and fourth. + assert_eq!( + before + .iter() + .filter(|(key, _)| selected(**key)) + .map(|(_, (_, created, _))| *created) + .collect::>(), + std::collections::BTreeSet::from([1, 2, 3, 4]) + ); let updated = UpdateBuilder::new(Arc::new(dataset)) - .update_where("i = 700") + .update_where("(i >= 240 AND i < 260) OR (i >= 745 AND i < 755)") .unwrap() - .set("i", "7000") + .set("i", "i + 10000") .unwrap() .build() .unwrap() @@ -1108,14 +1135,39 @@ mod tests { .unwrap(); let updated = updated.new_dataset.as_ref(); let update_version = updated.version().version; - let after = by_key(&collect_rows(updated).await); + let rewritten = updated + .manifest + .fragments + .last() + .expect("the update adds a fragment"); + assert_eq!( + matches!(rewritten.row_id_meta, Some(RowIdMeta::Column)), + update_spills, + "{rewritten:?}" + ); + assert_eq!( + matches!( + rewritten.created_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ), + update_spills, + "{rewritten:?}" + ); - // The rewritten row keeps its id and its created-at, and is stamped + // Each 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"); + let after = by_key(&collect_rows(updated).await); + assert_eq!(after.len(), before.len()); + for (key, (id, created, updated_at)) in before.iter() { + if selected(*key) { + assert_eq!( + after[&(key + 10000)], + (*id, *created, update_version), + "row {key}" + ); + } else { + assert_eq!(after[key], (*id, *created, *updated_at), "row {key}"); + } } updated.validate().await.unwrap(); } @@ -1280,82 +1332,24 @@ mod tests { } patched.validate().await.unwrap(); } - /// With the table opted in, an update spills the lineage of the rows it - /// rewrote at write time: their ids and created-at versions are known - /// before the commit and cannot change on a retry. The commit stamps the - /// last-updated-at version; the writer's placeholder is replaced. + + /// An update carries the rewritten rows' ids and created-at versions over + /// and the commit stamps their last-updated-at version. On an opted-in + /// table what the update carries over spills at write time -- it is known + /// before the commit and a retry cannot change it -- while the commit's + /// stamp stays inline. Without the opt-in everything stays inline, as + /// every release has, and the lineage is the same. + #[rstest] + #[case::opted_in(true)] + #[case::not_opted_in(false)] #[tokio::test] - async fn update_on_an_opted_in_table_spills_the_rewritten_rows_lineage() { + async fn update_places_the_rewritten_rows_lineage(#[case] opted_in: bool) { 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 = by_key(&collect_rows(&dataset).await); - - let updated = UpdateBuilder::new(Arc::new(dataset)) - .update_where("i >= 500") - .unwrap() - .set("i", "i + 10000") - .unwrap() - .build() - .unwrap() - .execute() - .await - .unwrap(); - let updated = updated.new_dataset.as_ref(); - let update_version = updated.version().version; - - let rewritten = updated - .get_fragments() - .into_iter() - .map(|fragment| fragment.metadata().clone()) - .find(|metadata| matches!(metadata.row_id_meta, Some(RowIdMeta::Column))) - .expect("the rewritten rows' fragment must spill its row ids"); - assert!( - matches!( - rewritten.created_at_version_meta, - Some(RowDatasetVersionMeta::Column) - ), - "the created-at versions must spill with the row ids" - ); - // The lineage file is one of the fragment's files, after its data. - assert_eq!(rewritten.files.len(), 2); - assert!( - rewritten.files[1].fields.iter().all(|field| *field < 0), - "the lineage file holds only lineage columns: {:?}", - rewritten.files[1].fields - ); - assert!( - matches!( - rewritten.last_updated_at_version_meta, - Some(RowDatasetVersionMeta::Inline(_)) - ), - "the commit stamps last-updated-at inline, got {:?}", - rewritten.last_updated_at_version_meta - ); - - let after = by_key(&collect_rows(updated).await); - for (key, (id, created, updated_at)) in before.iter() { - if *key >= 500 { - assert_eq!( - after[&(key + 10000)], - (*id, *created, update_version), - "row {key}" - ); - } else { - assert_eq!(after[key], (*id, *created, *updated_at), "row {key}"); - } + if opted_in { + spill_everything(&mut dataset).await; } - updated.validate().await.unwrap(); - } - - /// Without the opt-in the same update places everything inline, as every - /// release has, and the lineage it carries is the same. - #[tokio::test] - async fn update_on_a_table_that_did_not_opt_in_keeps_lineage_inline() { - let dir = TempStrDir::default(); - let uri = dir.as_str(); - let dataset = appended_dataset(uri, 4, 250).await; let before = by_key(&collect_rows(&dataset).await); let updated = UpdateBuilder::new(Arc::new(dataset)) @@ -1371,13 +1365,57 @@ mod tests { let updated = updated.new_dataset.as_ref(); let update_version = updated.version().version; - for fragment in updated.get_fragments() { + if opted_in { + let rewritten = updated + .get_fragments() + .into_iter() + .map(|fragment| fragment.metadata().clone()) + .find(|metadata| matches!(metadata.row_id_meta, Some(RowIdMeta::Column))) + .expect("the rewritten rows' fragment must spill its row ids"); assert!( - !fragment.metadata().has_spilled_row_lineage(), - "fragment {} spilled without the table opting in", - fragment.id() + matches!( + rewritten.created_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ), + "the created-at versions must spill with the row ids" + ); + // The lineage file is one of the fragment's files, after its data, + // and holds only what the update carried over. + assert_eq!(rewritten.files.len(), 2); + assert_eq!( + rewritten.files[1].fields.as_ref(), + [ROW_ID_FIELD_ID, ROW_CREATED_AT_VERSION_FIELD_ID] ); + assert!( + matches!( + rewritten.last_updated_at_version_meta, + Some(RowDatasetVersionMeta::Inline(_)) + ), + "the commit stamps last-updated-at inline, got {:?}", + rewritten.last_updated_at_version_meta + ); + // The source fragments were plain appends, so this update is the + // commit that first spills anything and has to raise the flag. + assert_ne!( + updated.manifest.reader_feature_flags & FLAG_UNSTABLE_SPILLED_ROW_LINEAGE, + 0, + "a spilled sequence must set the reader feature flag" + ); + assert_ne!( + updated.manifest.writer_feature_flags & FLAG_UNSTABLE_SPILLED_ROW_LINEAGE, + 0, + "a spilled sequence must set the writer feature flag" + ); + } else { + for fragment in updated.get_fragments() { + assert!( + !fragment.metadata().has_spilled_row_lineage(), + "fragment {} spilled without the table opting in", + fragment.id() + ); + } } + let after = by_key(&collect_rows(updated).await); for (key, (id, created, updated_at)) in before.iter() { if *key >= 500 { @@ -1391,6 +1429,12 @@ mod tests { } } updated.validate().await.unwrap(); + + // Re-opened cold, so the lineage is read through the committed + // manifest rather than from this process's caches. + let reopened = Dataset::open(uri).await.unwrap(); + assert_eq!(by_key(&collect_rows(&reopened).await), after); + reopened.validate().await.unwrap(); } /// The write that spills the fixture's lineage ahead of a schema change.