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.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/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-table/src/transaction/row_version.rs b/rust/lance-table/src/transaction/row_version.rs index 05c74c16487..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,55 @@ pub(super) fn resolve_update_version_metadata( } for fragment in new_fragments.iter_mut() { - 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. + // 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 ))); } - 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 { @@ -342,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] @@ -404,6 +477,155 @@ 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 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), + 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_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/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 91ec1476035..7c2168e43e2 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( @@ -416,7 +512,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 +520,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. @@ -576,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] @@ -650,6 +821,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, @@ -869,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; @@ -883,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() @@ -901,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(); } @@ -1073,4 +1332,198 @@ mod tests { } patched.validate().await.unwrap(); } + + /// 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_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; + if opted_in { + 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; + + 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!( + 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 { + 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(); + + // 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. + /// 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) }) diff --git a/rust/lance/src/dataset/utils.rs b/rust/lance/src/dataset/utils.rs index 254718464e9..376eed126a0 100644 --- a/rust/lance/src/dataset/utils.rs +++ b/rust/lance/src/dataset/utils.rs @@ -14,37 +14,44 @@ 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::format::{RowDatasetVersionRun, RowDatasetVersionSequence}; +use lance_table::rowids::segment::U64Segment; 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 +64,24 @@ 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); + // 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::>(); - 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 +101,21 @@ pub fn make_rowid_capture_stream( #[derive(Debug)] pub enum CapturedRowIds { AddressStyle(RoaringTreemap), - SequenceStyle(RowIdSequence), + SequenceStyle { + row_ids: RowIdSequence, + /// The created-at versions of the captured rows, in capture order, + /// when the stream carried them; `None` when it did not. + created_at: Option, + }, } 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: None, + } } else { Self::AddressStyle(RoaringTreemap::new()) } @@ -103,24 +128,72 @@ 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]) { + 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, + }); + } + } + 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, in capture order, when the + /// stream carried them. + pub fn created_at_sequence(&self) -> Option<&RowDatasetVersionSequence> { + match self { + Self::SequenceStyle { created_at, .. } => created_at.as_ref(), + Self::AddressStyle(_) => None, + } + } + 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") @@ -319,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 ee08ef22b74..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::get_row_id_index; +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; @@ -28,10 +30,12 @@ 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 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,6 +349,22 @@ impl UpdateJob { scanner.with_row_address(); } scanner.with_row_id(); + // 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() + .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 +500,23 @@ 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, + row_id_sequence, + removed_row_ids.created_at_sequence(), + spill_budget, + ) + .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 +553,56 @@ impl UpdateJob { }) } + /// 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], + 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 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 + )) + })?; + 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 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) + { + place_carried_row_lineage(&self.dataset, limit, &row_ids, &created_at) + .await? + .apply(fragment); + } + Ok(()) + } + async fn commit_impl( &self, dataset: Arc, @@ -688,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, @@ -714,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; @@ -2262,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; 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?; }