From 2e62889d72011fc6f972a85723a34d1d3f5fc0e3 Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Thu, 17 Sep 2026 18:18:20 +0800 Subject: [PATCH 1/6] feat(dataset): write compaction's spilled lineage into the fragment's data file Compaction computes the output fragments' row lineage before writing: the row ids and versions carry over from the inputs and the output file sizes are fixed in advance. When the table's policy spills a sequence type, its values ride along as hidden uint64 columns of the batches being written, under the reserved field ids, and each output fragment's metadata marks them as spilled into its own data file. A compacted fragment on a table that opts in has one file instead of two, and a scan that projects _rowid reads it from the file it already has open. The write path lets the columns through: write_fragments sets the three hidden fields aside before the schema is checked against the dataset's and puts them back on the written schema under their reserved ids, and Schema::validate admits exactly those ids under those names. The field id constants move to lance-core for that. Binary-copy compaction cannot add columns to the files it copies, so it keeps writing a separate lineage file, as does the update path. Co-Authored-By: Claude Fable 5.1 --- rust/lance-core/src/datatypes/schema.rs | 5 +- rust/lance-core/src/lib.rs | 27 +++ rust/lance-table/src/format/row_ids.rs | 16 +- rust/lance/src/dataset/optimize.rs | 168 +++++++++++++++-- rust/lance/src/dataset/rowids.rs | 5 +- rust/lance/src/dataset/rowids/spill.rs | 233 ++++++++++++++++++++++-- rust/lance/src/dataset/versions/mod.rs | 31 +++- 7 files changed, 432 insertions(+), 53 deletions(-) diff --git a/rust/lance-core/src/datatypes/schema.rs b/rust/lance-core/src/datatypes/schema.rs index 2e328c88d24..8adff191f51 100644 --- a/rust/lance-core/src/datatypes/schema.rs +++ b/rust/lance-core/src/datatypes/schema.rs @@ -350,7 +350,10 @@ impl Schema { // Check for duplicate field ids let mut seen_ids = HashSet::new(); for field in self.fields_pre_order() { - if field.id < 0 { + // A negative id is reserved for system use; the only ones a schema + // may carry are the hidden row lineage columns a data file stores + // next to the user columns, and only under their own names. + if field.id < 0 && crate::row_lineage_field_id(&field.name) != Some(field.id) { return Err(Error::schema(format!( "Field {} has a negative id {}", field.name, field.id diff --git a/rust/lance-core/src/lib.rs b/rust/lance-core/src/lib.rs index 0872dc97371..2c93f8772a1 100644 --- a/rust/lance-core/src/lib.rs +++ b/rust/lance-core/src/lib.rs @@ -51,6 +51,33 @@ pub static ROW_LAST_UPDATED_AT_VERSION_FIELD: LazyLock = pub static ROW_CREATED_AT_VERSION_FIELD: LazyLock = LazyLock::new(|| ArrowField::new(ROW_CREATED_AT_VERSION, DataType::UInt64, true)); +/// Field id of the hidden `_rowid` column a spilled row id sequence lives in. +/// +/// Field ids are `i32` and every negative value is reserved for system use: +/// `-1` is the unassigned sentinel, `-2` the tombstone for a field superseded by +/// a later data file, and `-3..=-5` the three row lineage columns. +pub const ROW_ID_FIELD_ID: i32 = -3; +/// Field id of the hidden `_row_created_at_version` column a spilled created-at +/// version sequence lives in. +pub const ROW_CREATED_AT_VERSION_FIELD_ID: i32 = -4; +/// Field id of the hidden `_row_last_updated_at_version` column a spilled +/// last-updated-at version sequence lives in. +pub const ROW_LAST_UPDATED_AT_VERSION_FIELD_ID: i32 = -5; + +/// The reserved field id of a row lineage column, by column name. +/// +/// These are the only negative field ids a data file schema may carry: a +/// lineage sequence spilled into the fragment's own data file is written as a +/// column under this id, next to the user columns. +pub fn row_lineage_field_id(column_name: &str) -> Option { + match column_name { + ROW_ID => Some(ROW_ID_FIELD_ID), + ROW_CREATED_AT_VERSION => Some(ROW_CREATED_AT_VERSION_FIELD_ID), + ROW_LAST_UPDATED_AT_VERSION => Some(ROW_LAST_UPDATED_AT_VERSION_FIELD_ID), + _ => None, + } +} + /// Check if a column name is a system column. /// /// System columns are virtual columns that are computed at read time and don't diff --git a/rust/lance-table/src/format/row_ids.rs b/rust/lance-table/src/format/row_ids.rs index a9a8a74e0ce..1b749643f02 100644 --- a/rust/lance-table/src/format/row_ids.rs +++ b/rust/lance-table/src/format/row_ids.rs @@ -12,19 +12,9 @@ use serde::{Deserialize, Deserializer, Serialize, Serializer}; use super::pb; -/// Field id of the hidden `_rowid` column that a spilled row id sequence lives in. -/// -/// Field ids are `i32` and every negative value is reserved for system use: -/// `-1` is the unassigned sentinel, `-2` is -/// [`TOMBSTONE_FIELD_ID`](crate::format::overlay::TOMBSTONE_FIELD_ID), and -/// `-3..=-5` are the three row lineage columns. -pub const ROW_ID_FIELD_ID: i32 = -3; -/// Field id of the hidden `_row_created_at_version` column that a spilled -/// created-at version sequence lives in. -pub const ROW_CREATED_AT_VERSION_FIELD_ID: i32 = -4; -/// Field id of the hidden `_row_last_updated_at_version` column that a spilled -/// last-updated-at version sequence lives in. -pub const ROW_LAST_UPDATED_AT_VERSION_FIELD_ID: i32 = -5; +pub use lance_core::{ + ROW_CREATED_AT_VERSION_FIELD_ID, ROW_ID_FIELD_ID, ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, +}; /// A reference to a part of a file, used by the fragment reuse index details. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)] diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index 53d7ea6a4a1..e519a2b81df 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -91,7 +91,8 @@ use super::fragment::FileFragment; use super::index::{DatasetIndexRemapperOptions, load_indices_for_remapping}; use super::rowids::RowVersionKind; use super::rowids::{ - RowLineage, load_row_id_sequences, load_row_version_sequence, place_row_lineage, + RowLineage, RowLineageSpill, load_row_id_sequences, load_row_version_sequence, + place_row_lineage, plan_row_lineage_spill, }; use super::transaction::{ Operation, RewriteGroup, RewrittenIndex, Transaction, TransactionBuilder, @@ -2569,6 +2570,29 @@ async fn rewrite_files( params.enable_stable_row_ids = true; } + // The output fragments' row lineage is known before anything is written: + // the row ids and versions carry over from the inputs, and the output + // file sizes are fixed above. Computing it here lets the sequences that + // leave the manifest ride along as hidden columns of the files being + // written, rather than need a file of their own. Binary copy cannot add + // columns to the files it copies, so it keeps writing that separate file. + let mut lineage_plan: Option<(Vec, RowLineageSpill)> = None; + if !can_binary_copy && dataset.manifest.uses_stable_row_ids() { + let chunk_sizes = file_row_counts + .iter() + .map(|rows| *rows as u64) + .collect::>(); + let lineages = compute_row_lineage(dataset.as_ref(), &fragments, &chunk_sizes).await?; + let spill = plan_row_lineage_spill(dataset.as_ref(), &lineages)?; + if spill.any() { + let stream = reader + .take() + .expect("reader must be prepared for non-binary-copy path"); + reader = Some(append_row_lineage_columns(stream, spill.columns(&lineages))); + } + lineage_plan = Some((lineages, spill)); + } + if can_binary_copy { new_fragments = versions::rewrite_files_binary_copy( write_version, @@ -2604,12 +2628,18 @@ async fn rewrite_files( row_ids_rx = Some(rx); } } else { + // The hidden lineage columns the stream now carries have to be in the + // schema the writer is given, or it would leave them out of the file. + let mut write_schema = dataset.schema().clone(); + if let Some((_, spill)) = &lineage_plan { + write_schema.fields.extend(spill.schema_fields()?); + } let (frags, _) = write_fragments_internal_with_file_row_counts( write_version, Some(dataset.as_ref()), dataset.object_store.clone(), &dataset.base, - dataset.schema().clone(), + write_schema, reader.expect("reader must be prepared for non-binary-copy path"), params, None, @@ -2639,7 +2669,15 @@ async fn rewrite_files( } else { if dataset.manifest.uses_stable_row_ids() { log::info!("Compaction task {}: rechunking stable row ids", task_id); - rechunk_row_lineage(dataset.as_ref(), &mut new_fragments, &fragments).await?; + match lineage_plan.take() { + Some((lineages, spill)) => { + place_row_lineage_in_files(&mut new_fragments, &lineages, spill)? + } + None => { + rechunk_row_lineage(dataset.as_ref(), &mut new_fragments, &fragments) + .await? + } + } } Ok(None) } @@ -2679,14 +2717,13 @@ async fn rewrite_files( } /// Carry the stable row ids and per-row versions of `old_fragments` over to -/// `new_fragments`, which hold the same live rows in the same order, and place -/// each new fragment's sequences inline or in a spilled column as the table's -/// spill policy and their size call for. -async fn rechunk_row_lineage( +/// output fragments of `chunk_sizes` rows each, in order, which hold the same +/// live rows in the same order. +async fn compute_row_lineage( dataset: &Dataset, - new_fragments: &mut [Fragment], old_fragments: &[Fragment], -) -> Result<()> { + chunk_sizes: &[u64], +) -> Result> { let mut old_sequences = load_row_id_sequences(dataset, old_fragments) .try_collect::>() .await?; @@ -2730,10 +2767,6 @@ async fn rechunk_row_lineage( } } - let chunk_sizes: Vec = new_fragments - .iter() - .map(|frag| frag.physical_rows.unwrap() as u64) - .collect(); debug_assert_eq!( { old_sequences.iter().map(|(_, seq)| seq.len()).sum::() }, { chunk_sizes.iter().sum::() }, @@ -2755,27 +2788,122 @@ async fn rechunk_row_lineage( )?; let new_last_updated_at = lance_table::rowids::version::rechunk_version_sequences( old_last_updated_sequences, - chunk_sizes, + chunk_sizes.iter().copied(), false, )?; - for (((fragment, row_ids), created_at), last_updated_at) in new_fragments - .iter_mut() - .zip(new_row_ids) + Ok(new_row_ids + .into_iter() .zip(new_created_at) .zip(new_last_updated_at) - { - let lineage = RowLineage { + .map(|((row_ids, created_at), last_updated_at)| RowLineage { row_ids, created_at, last_updated_at, - }; + }) + .collect()) +} + +/// Carry the lineage of `old_fragments` over to the already written +/// `new_fragments` and place each one's sequences inline or in a separate +/// lineage file, as the table's spill policy and their size call for. This is +/// the path for output files the lineage could not be written into. +async fn rechunk_row_lineage( + dataset: &Dataset, + new_fragments: &mut [Fragment], + old_fragments: &[Fragment], +) -> Result<()> { + let chunk_sizes: Vec = new_fragments + .iter() + .map(|frag| frag.physical_rows.unwrap() as u64) + .collect(); + let lineages = compute_row_lineage(dataset, old_fragments, &chunk_sizes).await?; + for (fragment, lineage) in new_fragments.iter_mut().zip(lineages) { place_row_lineage(dataset, &lineage).await?.apply(fragment); } + Ok(()) +} +/// Mark each new fragment's spilled sequences as living in its own data file, +/// which was written with them as hidden columns, and place the rest inline. +fn place_row_lineage_in_files( + new_fragments: &mut [Fragment], + lineages: &[RowLineage], + spill: RowLineageSpill, +) -> Result<()> { + if new_fragments.len() != lineages.len() { + return Err(Error::internal(format!( + "compaction wrote {} fragments but planned lineage for {}", + new_fragments.len(), + lineages.len() + ))); + } + for (fragment, lineage) in new_fragments.iter_mut().zip(lineages) { + let [file] = fragment.files.as_slice() else { + return Err(Error::internal(format!( + "compaction wrote {} data files for fragment {}; the row lineage columns \ + are in exactly one", + fragment.files.len(), + fragment.id + ))); + }; + let missing = spill + .field_ids() + .find(|field_id| !file.fields.contains(field_id)); + if let Some(field_id) = missing { + return Err(Error::internal(format!( + "compaction wrote fragment {}'s data file {} without row lineage field {}", + fragment.id, file.path, field_id + ))); + } + spill.place_in_file(lineage).apply(fragment); + } Ok(()) } +/// Append the hidden row lineage columns to the batches of `stream`, in row +/// order, so the file writer stores them next to the user columns. +fn append_row_lineage_columns( + stream: SendableRecordBatchStream, + columns: Vec<(i32, &'static str, Vec)>, +) -> SendableRecordBatchStream { + let mut fields = stream.schema().fields().to_vec(); + for (_, name, _) in &columns { + fields.push(Arc::new(ArrowField::new( + *name, + ArrowDataType::UInt64, + false, + ))); + } + let schema = Arc::new(ArrowSchema::new(fields)); + let values = columns + .into_iter() + .map(|(_, _, values)| arrow_array::UInt64Array::from(values)) + .collect::>(); + let total_rows = values.first().map_or(0, |column| column.len()); + let batch_schema = schema.clone(); + let mut offset = 0usize; + let batches = stream.map(move |batch| { + let batch = batch?; + let rows = batch.num_rows(); + if offset + rows > total_rows { + return Err(datafusion::error::DataFusionError::External(Box::new( + Error::internal(format!( + "compaction read more rows than its row lineage covers ({total_rows})" + )), + ))); + } + let mut arrays = batch.columns().to_vec(); + for column in &values { + arrays.push(Arc::new(column.slice(offset, rows)) as ArrayRef); + } + offset += rows; + RecordBatch::try_new(batch_schema.clone(), arrays) + .map_err(datafusion::error::DataFusionError::from) + }); + Box::pin(RecordBatchStreamAdapter::new(schema, batches)) +} + /// Commit the results of file compaction. /// /// It is not required that all tasks are passed to this method. If some failed, diff --git a/rust/lance/src/dataset/rowids.rs b/rust/lance/src/dataset/rowids.rs index 4a77b45757a..93a12ec14d9 100644 --- a/rust/lance/src/dataset/rowids.rs +++ b/rust/lance/src/dataset/rowids.rs @@ -27,8 +27,9 @@ 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, - place_row_lineage, read_spilled_row_ids, read_spilled_versions, + PlacedRowLineage, RowLineage, RowLineageSpill, SPILL_ROW_LINEAGE_CONFIG_KEY, + inline_row_lineage_max_bytes, place_row_lineage, plan_row_lineage_spill, read_spilled_row_ids, + read_spilled_versions, }; pub(super) use validate::validate_stable_row_ids; diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs index 7c2168e43e2..a5b5a660e1a 100644 --- a/rust/lance/src/dataset/rowids/spill.rs +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -103,6 +103,146 @@ pub fn inline_row_lineage_max_bytes(dataset: &Dataset) -> Result> }) } +/// Which of a compaction task's sequences leave the manifest, decided for all +/// its output fragments together so that every data file the task writes +/// carries the same hidden columns. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct RowLineageSpill { + pub row_ids: bool, + pub created_at: bool, + pub last_updated_at: bool, +} + +impl RowLineageSpill { + pub fn any(&self) -> bool { + self.row_ids || self.created_at || self.last_updated_at + } + + /// The reserved field ids of the columns that spill, in column order. + pub fn field_ids(&self) -> impl Iterator { + [ + (self.row_ids, ROW_ID_FIELD_ID), + (self.created_at, ROW_CREATED_AT_VERSION_FIELD_ID), + (self.last_updated_at, ROW_LAST_UPDATED_AT_VERSION_FIELD_ID), + ] + .into_iter() + .filter(|(spilled, _)| *spilled) + .map(|(_, field_id)| field_id) + } + + /// The hidden columns to write, as `(field id, column name, values)`, + /// concatenated over `lineages` in order. + pub fn columns(&self, lineages: &[RowLineage]) -> Vec<(i32, &'static str, Vec)> { + let mut columns = Vec::with_capacity(3); + if self.row_ids { + let values = lineages.iter().flat_map(|l| l.row_ids.iter()).collect(); + columns.push((ROW_ID_FIELD_ID, ROW_ID, values)); + } + if self.created_at { + let values = lineages + .iter() + .flat_map(|l| l.created_at.versions()) + .collect(); + columns.push(( + ROW_CREATED_AT_VERSION_FIELD_ID, + ROW_CREATED_AT_VERSION, + values, + )); + } + if self.last_updated_at { + let values = lineages + .iter() + .flat_map(|l| l.last_updated_at.versions()) + .collect(); + columns.push(( + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION, + values, + )); + } + columns + } + + /// The hidden columns as write-schema fields, under their reserved ids. + pub fn schema_fields(&self) -> Result> { + [ + (self.row_ids, ROW_ID, ROW_ID_FIELD_ID), + ( + self.created_at, + ROW_CREATED_AT_VERSION, + ROW_CREATED_AT_VERSION_FIELD_ID, + ), + ( + self.last_updated_at, + ROW_LAST_UPDATED_AT_VERSION, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, + ), + ] + .into_iter() + .filter(|(spilled, _, _)| *spilled) + .map(|(_, name, id)| { + let mut field = lance_core::datatypes::Field::try_from(&ArrowField::new( + name, + DataType::UInt64, + false, + ))?; + field.id = id; + Ok(field) + }) + .collect() + } + + /// The placement for one fragment whose own data file was written with + /// the spilled columns: those sequences are marked as spilled, the rest + /// are placed inline. There is no lineage file to add to the fragment. + pub fn place_in_file(&self, lineage: &RowLineage) -> PlacedRowLineage { + PlacedRowLineage { + row_ids: if self.row_ids { + RowIdMeta::Column + } else { + RowIdMeta::Inline(write_row_ids(&lineage.row_ids).into()) + }, + created_at: if self.created_at { + RowDatasetVersionMeta::Column + } else { + RowDatasetVersionMeta::Inline(write_dataset_versions(&lineage.created_at).into()) + }, + last_updated_at: if self.last_updated_at { + RowDatasetVersionMeta::Column + } else { + RowDatasetVersionMeta::Inline( + write_dataset_versions(&lineage.last_updated_at).into(), + ) + }, + file: None, + } + } +} + +/// Decide which sequence types of `lineages` spill: each one whose encoding +/// exceeds the table's inline budget in any of them. Nothing spills on a +/// table that has not opted in. +pub fn plan_row_lineage_spill( + dataset: &Dataset, + lineages: &[RowLineage], +) -> Result { + let Some(limit) = inline_row_lineage_max_bytes(dataset)? else { + return Ok(RowLineageSpill::default()); + }; + let over = |encoded: usize| encoded > limit; + Ok(RowLineageSpill { + row_ids: lineages + .iter() + .any(|l| over(write_row_ids(&l.row_ids).len())), + created_at: lineages + .iter() + .any(|l| over(write_dataset_versions(&l.created_at).len())), + last_updated_at: lineages + .iter() + .any(|l| over(write_dataset_versions(&l.last_updated_at).len())), + }) +} + /// The per-row lineage of one fragment, in row offset order. pub struct RowLineage { pub row_ids: RowIdSequence, @@ -510,7 +650,7 @@ async fn read_spilled_column( mod tests { use super::*; use crate::dataset::cleanup::{CleanupPolicyBuilder, cleanup_old_versions}; - use crate::dataset::optimize::{CompactionOptions, compact_files}; + use crate::dataset::optimize::{CompactionMode, CompactionOptions, compact_files}; use crate::dataset::rowids::{RowVersionKind, load_row_id_sequence, load_row_version_sequence}; use crate::dataset::{ColumnAlteration, UpdateBuilder, WriteMode, WriteParams}; use arrow_array::{Int32Array, RecordBatchIterator}; @@ -925,22 +1065,23 @@ mod tests { ), "compaction must spill every sequence under a zero inline budget, got {metadata:?}" ); - // The three sequences share one lineage file, listed after the user - // data file among the fragment's files and found by field id. - assert_eq!(metadata.files.len(), 2); - let lineage_file = &metadata.files[1]; + // The three columns ride in the fragment's own data file, after the + // user column, so the fragment has no extra file to reference. + assert_eq!(metadata.files.len(), 1); + let data_file = &metadata.files[0]; assert_eq!( - lineage_file.fields.as_ref(), + data_file.fields.as_ref(), [ + 0, ROW_ID_FIELD_ID, ROW_CREATED_AT_VERSION_FIELD_ID, ROW_LAST_UPDATED_AT_VERSION_FIELD_ID ] ); - for field_id in lineage_file.fields.iter() { + for field_id in data_file.fields.iter().filter(|id| **id < 0) { assert_eq!( metadata.row_lineage_file(*field_id).unwrap(), - Some(lineage_file) + Some(data_file) ); } assert_ne!( @@ -978,6 +1119,50 @@ mod tests { assert_eq!(versions_of(&created_at), before.1); } + /// Binary-copy compaction copies the input files page by page and cannot + /// add columns to them, so its spilled lineage goes to a separate file. + #[tokio::test] + async fn binary_copy_compaction_spills_to_a_separate_file() { + let dir = TempStrDir::default(); + let uri = dir.as_str(); + let mut dataset = appended_dataset(uri, 4, 250).await; + spill_everything(&mut dataset).await; + let before = collect_lineage(&dataset).await; + + compact_files( + &mut dataset, + CompactionOptions { + compaction_mode: Some(CompactionMode::ForceBinaryCopy), + ..one_fragment() + }, + None, + ) + .await + .unwrap(); + + let fragments = dataset.get_fragments(); + assert_eq!(fragments.len(), 1); + let metadata = fragments[0].metadata(); + assert!( + matches!(metadata.row_id_meta, Some(RowIdMeta::Column)), + "compaction must spill the row ids under a zero inline budget, got {metadata:?}" + ); + // The copied data file carries only the user column; the lineage + // follows it as a file of its own. + assert_eq!(metadata.files.len(), 2); + assert!(metadata.files[0].fields.iter().all(|field| *field >= 0)); + assert_eq!( + metadata.files[1].fields.as_ref(), + [ + ROW_ID_FIELD_ID, + ROW_CREATED_AT_VERSION_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID + ] + ); + assert_eq!(collect_lineage(&dataset).await, before); + dataset.validate().await.unwrap(); + } + /// Cleanup decides what to delete by walking /// [`Fragment::referenced_lance_files`], so a spilled sequence has to be /// reachable from there. If it were not, an ordinary cleanup would delete a @@ -1443,7 +1628,9 @@ mod tests { enum SpillingWrite { /// An update rewriting the rows with `i >= 500`. Update, - /// A compaction into a single fragment. + /// A binary-copy compaction into a single fragment. A reencoding one + /// would write the lineage next to `i` in the data file, which the + /// schema change keeps for `i` alone. Compaction, } @@ -1488,18 +1675,32 @@ mod tests { dataset = updated.new_dataset.as_ref().clone(); } SpillingWrite::Compaction => { - compact_files(&mut dataset, one_fragment(), None) - .await - .unwrap(); + let options = CompactionOptions { + compaction_mode: Some(CompactionMode::ForceBinaryCopy), + ..one_fragment() + }; + compact_files(&mut dataset, options, None).await.unwrap(); } } + // Each spilled fragment keeps its lineage in a file of its own, which + // nothing but the lineage keeps alive through the schema change. + let spilled = dataset + .manifest + .fragments + .iter() + .filter(|fragment| fragment.has_spilled_row_lineage()) + .collect::>(); assert!( - dataset - .get_fragments() - .iter() - .any(|fragment| fragment.metadata().has_spilled_row_lineage()), + !spilled.is_empty(), "the fixture must spill before the schema change" ); + for fragment in spilled { + assert_eq!(fragment.files.len(), 2, "{fragment:?}"); + assert!( + fragment.files[1].fields.iter().all(|field| *field < 0), + "expected a lineage-only file: {fragment:?}" + ); + } let before = by_key(&collect_rows(&dataset).await); let j = || ColumnAlteration::new("j".to_string()); diff --git a/rust/lance/src/dataset/versions/mod.rs b/rust/lance/src/dataset/versions/mod.rs index 051e933a6bf..e4af2e1c3e2 100644 --- a/rust/lance/src/dataset/versions/mod.rs +++ b/rust/lance/src/dataset/versions/mod.rs @@ -152,12 +152,18 @@ pub async fn write_fragments( _ => normalized_schema, }; let version_name = format!("{version:?}"); - let schema = write::prepare_write_schema( + // A writer that spills row lineage into the fragment's own data file + // carries the hidden columns in its stream. They are not dataset fields: + // set them aside before the schema is checked against the dataset's and + // put them back, under their reserved ids, on the schema that is written. + let (normalized_schema, lineage_fields) = split_row_lineage_fields(normalized_schema); + let mut schema = write::prepare_write_schema( dataset, normalized_schema, ¶ms, schema_compare_options(version), )?; + schema.fields.extend(lineage_fields); match version { ConcreteFileVersion::V1 | ConcreteFileVersion::V2_0 | ConcreteFileVersion::V2_1 => { write::validate_legacy_blob_write_schema(&schema, &version_name)?; @@ -183,6 +189,29 @@ pub async fn write_fragments( Ok((fragments, schema)) } +/// Take the hidden row lineage columns out of a write schema, each keyed to +/// the reserved field id its name maps to. +fn split_row_lineage_fields(schema: Schema) -> (Schema, Vec) { + let (lineage, user): (Vec<_>, Vec<_>) = schema + .fields + .into_iter() + .partition(|field| lance_core::row_lineage_field_id(&field.name).is_some()); + let lineage = lineage + .into_iter() + .map(|mut field| { + field.id = lance_core::row_lineage_field_id(&field.name).unwrap(); + field + }) + .collect(); + ( + Schema { + fields: user, + metadata: schema.metadata, + }, + lineage, + ) +} + #[allow(clippy::too_many_arguments)] pub async fn write_fragments_direct( version: ConcreteFileVersion, From 6dd150be917f7386451f59fdf5780eed1bdacc40 Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Tue, 29 Sep 2026 05:12:54 +0800 Subject: [PATCH 2/6] fix(dataset): read spilled lineage columns by their column index read_spilled_column took the projection field from the file schema's top-level field at the lineage column's physical index. That only holds for lineage-only files. Compaction now writes the lineage columns after the user columns, and a column index counts physical columns: one per leaf from 2.1 on, and one per list or struct as well in 2.0. After a struct or a 2.0 list, the lookup picked another lineage field or none at all, so every reader of a spilled sequence (scan, take, the row id index, commit read-ahead, validate, the next compaction) failed on that fragment. Build the projection field from the reserved id instead, as a non-nullable UInt64 with the reserved name, and keep the column index from the data file. RowLineageSpill::schema_fields builds the written fields through the same helper, so the read and write sides share one definition. A lookup by id in the file schema would not work: a lineage-only file stores its fields under ids 0..2. The reader now opens through open_projected_reader. For a wide compaction output it fetches only the lineage column's metadata instead of decoding every user column's, using the fragment reader's threshold. Lineage-only files and 2.0 files still load the full metadata. Either way, the column index is checked against the file's column count before a reader is built on it, so an entry past the file's columns is still reported as a corrupt file naming the data file, as the old schema lookup did, rather than as invalid input from the reader. compaction_spills_and_reads_back_row_lineage becomes an rstest with a struct column on 2.2 and a list column on 2.0. With all three sequences spilled, both cases failed before this change. spilled_lineage_shares_one_file_and_round_trips also checks the error for an out-of-range column index. Co-Authored-By: Claude Opus 5.5 --- rust/lance/src/dataset/rowids/spill.rs | 330 ++++++++++++++++--------- 1 file changed, 218 insertions(+), 112 deletions(-) diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs index a5b5a660e1a..8b6fe955a74 100644 --- a/rust/lance/src/dataset/rowids/spill.rs +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -26,11 +26,12 @@ use std::sync::Arc; use arrow_array::{Array, ArrayRef, RecordBatch, UInt64Array}; use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema}; -use futures::{FutureExt, TryStreamExt}; +use futures::{FutureExt, StreamExt, TryStreamExt}; use lance_core::datatypes::Schema; use lance_core::{ROW_CREATED_AT_VERSION, ROW_ID, ROW_LAST_UPDATED_AT_VERSION}; use lance_encoding::decoder::{DecoderPlugins, FilterExpression}; -use lance_file::reader::{FileReader, ReaderProjection}; +use lance_file::LanceEncodingsIo; +use lance_file::reader::{FileReader, ProjectedFileReader, ReaderProjection}; use lance_file::version::ConcreteFileVersion; use lance_file::versions; use lance_file::writer::FileWriterOptions; @@ -165,31 +166,7 @@ impl RowLineageSpill { /// The hidden columns as write-schema fields, under their reserved ids. pub fn schema_fields(&self) -> Result> { - [ - (self.row_ids, ROW_ID, ROW_ID_FIELD_ID), - ( - self.created_at, - ROW_CREATED_AT_VERSION, - ROW_CREATED_AT_VERSION_FIELD_ID, - ), - ( - self.last_updated_at, - ROW_LAST_UPDATED_AT_VERSION, - ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, - ), - ] - .into_iter() - .filter(|(spilled, _, _)| *spilled) - .map(|(_, name, id)| { - let mut field = lance_core::datatypes::Field::try_from(&ArrowField::new( - name, - DataType::UInt64, - false, - ))?; - field.id = id; - Ok(field) - }) - .collect() + self.field_ids().map(lineage_field).collect() } /// The placement for one fragment whose own data file was written with @@ -219,6 +196,25 @@ impl RowLineageSpill { } } +/// The hidden column a spilled lineage sequence is stored in: a non-nullable +/// `UInt64` named after the sequence, under its reserved `field_id`. +fn lineage_field(field_id: i32) -> Result { + let name = match field_id { + ROW_ID_FIELD_ID => ROW_ID, + ROW_CREATED_AT_VERSION_FIELD_ID => ROW_CREATED_AT_VERSION, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID => ROW_LAST_UPDATED_AT_VERSION, + _ => { + return Err(Error::internal(format!( + "field id {field_id} is not a row lineage field id" + ))); + } + }; + let mut field = + lance_core::datatypes::Field::try_from(&ArrowField::new(name, DataType::UInt64, false))?; + field.id = field_id; + Ok(field) +} + /// Decide which sequence types of `lineages` spill: each one whose encoding /// exceeds the table's inline budget in any of them. Nothing spills on a /// table that has not opted in. @@ -519,7 +515,7 @@ pub async fn read_spilled_versions( /// Read the hidden `UInt64` column `field_id` of `fragment` in full: one /// value per physical row, from the one data file that carries the id. /// -/// Callers box this future: it drives the full data file reader, and inlined +/// Callers box this future: it drives the data file reader, and inlined /// into the row id index build (reached from `take`, and from there from index /// builds and `optimize_indices`) it makes those futures too deep for the trait /// solver to prove `Send`/`Sync` (E0275). @@ -556,6 +552,21 @@ async fn read_spilled_column( ) })?; + // The projected field is built from the reserved id rather than taken + // from the file schema. A column index counts physical columns -- one per + // leaf, and in 2.0 one per list or struct as well -- so once the lineage + // columns follow nested user columns, as in a compaction output, it is no + // longer the field's position among the file's top-level fields. Nor can + // the file schema be searched by id: a lineage-only file stores its + // fields under ids 0, 1 and 2. + let projection = ReaderProjection { + schema: Arc::new(Schema { + fields: vec![lineage_field(field_id)?], + metadata: Default::default(), + }), + column_indices: vec![column_index], + }; + // Resolved through `data_file_dir` rather than `data_dir` so a shallow // clone, which rewrites `base_id` on every referenced file, still finds it. let path: Path = dataset @@ -569,46 +580,83 @@ async fn read_spilled_column( let file = scheduler .open_file(&path, &data_file.file_size_bytes) .await?; - let reader = FileReader::try_open( - file, - None, - Arc::::default(), - &dataset.metadata_cache.file_metadata_cache(&path), - dataset.file_reader_options.clone().unwrap_or_default(), - ) - .await?; - - // The lineage columns are flat primitives, so the file schema's column - // position is the column index in every file version. - let field = reader - .schema() - .fields - .get(column_index as usize) - .ok_or_else(|| { - Error::corrupt_file_named( + let options = dataset.file_reader_options.clone().unwrap_or_default(); + let cache = dataset.metadata_cache.file_metadata_cache(&path); + let io = + Arc::new(LanceEncodingsIo::new(file.clone()).with_read_chunk_size(options.read_chunk_size)); + // An index past the file's columns means the data file entry and the file + // disagree. The reader would reject the projection as invalid input without + // naming the file, so it is checked here, once the column count is known. + let check_column_index = |num_columns: usize| { + if (column_index as usize) < num_columns { + Ok(()) + } else { + Err(Error::corrupt_file_named( &data_file.path, - format!("spilled row lineage file has no column at index {column_index}"), - ) - })?; - let projection = ReaderProjection { - schema: Arc::new(Schema { - fields: vec![field.clone()], - metadata: Default::default(), - }), - column_indices: vec![column_index], + format!( + "spilled row lineage column {field_id} is at column index {column_index}, \ + but the file has only {num_columns} columns" + ), + )) + } }; + // A compaction output holds the lineage next to every user column, and + // decoding all their metadata to read one column would cost as much as + // opening the file for a scan. Past a few columns, only the one read here + // has its metadata fetched, by the same threshold the fragment reader uses. + let live_columns = data_file + .column_indices + .iter() + .filter(|column_index| **column_index >= 0) + .count(); + let reader = versions::open_projected_reader( + data_file.file_version()?, + &projection, + projection.column_indices.len().saturating_mul(4) < live_columns, + || async { + let metadata_index = FileReader::read_metadata_index(&file).await?; + check_column_index(metadata_index.num_columns() as usize)?; + let reader = ProjectedFileReader::try_open_with_metadata_index( + io.clone(), + path.clone(), + Some(projection.clone()), + Arc::::default(), + Arc::new(metadata_index), + &cache, + options.clone(), + ) + .await?; + Ok(Some(reader)) + }, + || async { + let metadata = FileReader::read_all_metadata(&file).await?; + check_column_index(metadata.column_infos.len())?; + ProjectedFileReader::try_open_with_file_metadata( + io.clone(), + path.clone(), + Some(projection.clone()), + Arc::::default(), + Arc::new(metadata), + &cache, + options.clone(), + ) + .await + }, + ) + .await?; let mut values: Vec = Vec::with_capacity(reader.num_rows() as usize); - let mut stream = reader - .read_stream_projected( + let mut batches = reader + .read_tasks( ReadBatchParams::RangeFull, SPILL_BATCH_ROWS as u32, - 8, - projection, + None, FilterExpression::no_filter(), ) - .await?; - while let Some(batch) = stream.try_next().await? { + .await? + .map(|task| task.task) + .buffered(8); + while let Some(batch) = batches.try_next().await? { let column = batch .column(0) .as_any() @@ -653,7 +701,8 @@ mod tests { use crate::dataset::optimize::{CompactionMode, CompactionOptions, compact_files}; use crate::dataset::rowids::{RowVersionKind, load_row_id_sequence, load_row_version_sequence}; use crate::dataset::{ColumnAlteration, UpdateBuilder, WriteMode, WriteParams}; - use arrow_array::{Int32Array, RecordBatchIterator}; + use arrow_array::builder::{ListBuilder, StringBuilder}; + use arrow_array::{Int32Array, RecordBatchIterator, StringArray, StructArray}; use arrow_schema::Field; use chrono::Utc; use lance_core::utils::tempfile::TempStrDir; @@ -765,6 +814,18 @@ mod tests { versions_of(&last_updated_at), versions_of(&lineage.last_updated_at) ); + + // A data file entry that points past the file's columns is corruption + // of that file, not a bad projection from the caller. + fragment.files[0].column_indices = vec![3, 1, 2].into(); + let error = read_spilled_row_ids(&dataset, &fragment).await.unwrap_err(); + assert!(matches!(error, Error::CorruptFile { .. }), "{error}"); + assert!( + error + .to_string() + .contains("is at column index 3, but the file has only 3 columns"), + "{error}" + ); } #[tokio::test] @@ -926,62 +987,88 @@ mod tests { .await .unwrap(); } + + /// The user columns of a test table: its key column `i`, and possibly one + /// more column next to it. + #[derive(Clone, Copy, Debug)] + enum UserColumns { + KeyOnly, + /// Adds `meta: struct`, which from 2.1 on takes one + /// physical column per child and none for the struct itself. + WithStruct, + /// Adds `tags: list`, which in 2.0 takes one physical column for + /// the list and one for its items. + WithList, + /// Adds `j: int32` holding the key again, so a schema change has a + /// column to rename, drop or cast while `i` still identifies every + /// row. + WithCopy, + } + + /// A batch of `columns` holding `keys` in column `i`. + fn keyed_batch(columns: UserColumns, keys: std::ops::Range) -> RecordBatch { + let mut fields = vec![Field::new("i", DataType::Int32, false)]; + let mut arrays: Vec = vec![Arc::new(Int32Array::from_iter_values(keys.clone()))]; + match columns { + UserColumns::KeyOnly => {} + UserColumns::WithStruct => { + let a: ArrayRef = Arc::new(Int32Array::from_iter_values(keys.clone())); + let b: ArrayRef = Arc::new(StringArray::from_iter_values( + keys.map(|key| format!("b{key}")), + )); + let meta = StructArray::from(vec![ + (Arc::new(Field::new("a", DataType::Int32, true)), a), + (Arc::new(Field::new("b", DataType::Utf8, true)), b), + ]); + fields.push(Field::new("meta", meta.data_type().clone(), true)); + arrays.push(Arc::new(meta)); + } + UserColumns::WithList => { + let mut tags = ListBuilder::new(StringBuilder::new()); + for key in keys { + tags.values().append_value(format!("t{key}")); + tags.append(true); + } + let tags = tags.finish(); + fields.push(Field::new("tags", tags.data_type().clone(), true)); + arrays.push(Arc::new(tags)); + } + UserColumns::WithCopy => { + fields.push(Field::new("j", DataType::Int32, true)); + arrays.push(Arc::new(Int32Array::from_iter_values(keys))); + } + } + RecordBatch::try_new(Arc::new(ArrowSchema::new(fields)), arrays).unwrap() + } + /// A stable-row-id dataset built from `chunks` separate appends, so /// compacting it has several sequences to concatenate. async fn appended_dataset(uri: &str, chunks: i32, rows_per_chunk: i32) -> Dataset { - let schema = test_schema(); - let mut dataset: Option = None; - for chunk in 0..chunks { - let batch = RecordBatch::try_new( - schema.clone(), - vec![Arc::new(Int32Array::from_iter_values( - (chunk * rows_per_chunk)..((chunk + 1) * rows_per_chunk), - ))], - ) - .unwrap(); - let reader = RecordBatchIterator::new(vec![Ok(batch)], schema.clone()); - dataset = Some( - Dataset::write( - reader, - uri, - Some(WriteParams { - enable_stable_row_ids: true, - mode: if chunk == 0 { - WriteMode::Create - } else { - WriteMode::Append - }, - ..Default::default() - }), - ) - .await - .unwrap(), - ); - } - dataset.unwrap() + appended_dataset_with(uri, chunks, rows_per_chunk, UserColumns::KeyOnly, None).await } - /// 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), - ])); + /// [`appended_dataset`] with the given user columns, written in `version`, + /// or the default version when that is `None`. + async fn appended_dataset_with( + uri: &str, + chunks: i32, + rows_per_chunk: i32, + columns: UserColumns, + version: Option, + ) -> Dataset { 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()); + for chunk in 0..chunks { + let keys = (chunk * rows_per_chunk)..((chunk + 1) * rows_per_chunk); + let batch = keyed_batch(columns, keys); + let schema = batch.schema(); + let reader = RecordBatchIterator::new(vec![Ok(batch)], schema); dataset = Some( Dataset::write( reader, uri, Some(WriteParams { enable_stable_row_ids: true, + data_storage_version: version, mode: if chunk == 0 { WriteMode::Create } else { @@ -1028,11 +1115,22 @@ mod tests { ) } + /// Compaction writes the spilled columns after the user columns of its + /// output file. A column index counts physical columns, so after a + /// struct, or a list in 2.0, a lineage column's index no longer matches + /// its position among the file's top-level fields. + #[rstest] + #[case::flat(UserColumns::KeyOnly, None)] + #[case::nested_v2_2(UserColumns::WithStruct, Some(LanceFileVersion::V2_2))] + #[case::list_v2_0(UserColumns::WithList, Some(LanceFileVersion::V2_0))] #[tokio::test] - async fn compaction_spills_and_reads_back_row_lineage() { + async fn compaction_spills_and_reads_back_row_lineage( + #[case] columns: UserColumns, + #[case] version: Option, + ) { let dir = TempStrDir::default(); let uri = dir.as_str(); - let mut dataset = appended_dataset(uri, 4, 250).await; + let mut dataset = appended_dataset_with(uri, 4, 250, columns, version).await; spill_everything(&mut dataset).await; // Four appends at four versions, so the compacted created-at sequence // has four runs rather than one. @@ -1066,18 +1164,27 @@ mod tests { "compaction must spill every sequence under a zero inline budget, got {metadata:?}" ); // The three columns ride in the fragment's own data file, after the - // user column, so the fragment has no extra file to reference. + // user columns, so the fragment has no extra file to reference. Their + // column indices continue from the last physical user column. assert_eq!(metadata.files.len(), 1); let data_file = &metadata.files[0]; + let (user_fields, lineage_fields) = data_file.fields.split_at(data_file.fields.len() - 3); + assert!( + user_fields.iter().all(|field| *field >= 0), + "unexpected data file fields {:?}", + data_file.fields + ); assert_eq!( - data_file.fields.as_ref(), + lineage_fields, [ - 0, ROW_ID_FIELD_ID, ROW_CREATED_AT_VERSION_FIELD_ID, ROW_LAST_UPDATED_AT_VERSION_FIELD_ID ] ); + let (user_columns, lineage_columns) = data_file.column_indices.split_at(user_fields.len()); + let next = user_columns.iter().filter(|column| **column >= 0).count() as i32; + assert_eq!(lineage_columns, [next, next + 1, next + 2]); for field_id in data_file.fields.iter().filter(|id| **id < 0) { assert_eq!( metadata.row_lineage_file(*field_id).unwrap(), @@ -1415,7 +1522,6 @@ mod tests { use crate::dataset::{ MergeInsertBuilder, MergeInsertWriteMode, WhenMatched, WhenNotMatched, }; - use arrow_array::StringArray; let dir = TempStrDir::default(); let uri = dir.as_str(); @@ -1658,7 +1764,7 @@ mod tests { ) { let dir = TempStrDir::default(); let uri = dir.as_str(); - let mut dataset = two_column_dataset(uri).await; + let mut dataset = appended_dataset_with(uri, 4, 250, UserColumns::WithCopy, None).await; spill_everything(&mut dataset).await; match spilling_write { SpillingWrite::Update => { From 0369072fe78edf184007d229c433ac9cc6edb986 Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Tue, 29 Sep 2026 05:15:08 +0800 Subject: [PATCH 3/6] refactor(dataset): key hidden lineage fields on their reserved ids versions::write_fragments runs for every writer: insert, update, both merge_insert paths and compaction. It used to treat any top-level field named _rowid, _row_created_at_version or _row_last_updated_at_version as a hidden lineage column, overwrite its id with the reserved one and write it where readers skip it. A user column with one of those names, from alter_columns or an external caller of write_fragments_internal, would have been dropped from the dataset's view without an error. Matching by name was only needed because the 2.2+ blob promotion's set_field_id renumbers every negative id before the split ran. Split before the promotion instead, and take only fields that already carry the reserved id of their name. RowLineageSpill::schema_fields is the only code that assigns those ids. A field with a reserved id that is not a non-nullable UInt64 is rejected with an internal error, because the readers decode these columns as exactly that type. The unwrap and the second name lookup go away. Schema::validate now exempts reserved ids only for top-level fields, where a data file stores the lineage columns. A nested field that reuses one fails the negative-id or duplicate-id check. Co-Authored-By: Claude Opus 5.5 --- rust/lance-core/src/datatypes/schema.rs | 67 ++++++++++++++++-- rust/lance/src/dataset/versions/mod.rs | 94 ++++++++++++++++++++----- 2 files changed, 139 insertions(+), 22 deletions(-) diff --git a/rust/lance-core/src/datatypes/schema.rs b/rust/lance-core/src/datatypes/schema.rs index 8adff191f51..0a132875ff5 100644 --- a/rust/lance-core/src/datatypes/schema.rs +++ b/rust/lance-core/src/datatypes/schema.rs @@ -347,13 +347,21 @@ impl Schema { } } + // A negative id is reserved for system use; the only ones a schema may + // carry are the hidden row lineage columns a data file stores next to + // the user columns, as top-level fields under their own names. A + // nested field reusing one of these ids fails the duplicate check. + let row_lineage_ids = self + .fields + .iter() + .filter(|field| crate::row_lineage_field_id(&field.name) == Some(field.id)) + .map(|field| field.id) + .collect::>(); + // Check for duplicate field ids let mut seen_ids = HashSet::new(); for field in self.fields_pre_order() { - // A negative id is reserved for system use; the only ones a schema - // may carry are the hidden row lineage columns a data file stores - // next to the user columns, and only under their own names. - if field.id < 0 && crate::row_lineage_field_id(&field.name) != Some(field.id) { + if field.id < 0 && !row_lineage_ids.contains(&field.id) { return Err(Error::schema(format!( "Field {} has a negative id {}", field.name, field.id @@ -1923,6 +1931,57 @@ mod tests { assert!(error.to_string().contains("Duplicate field id 0")); } + #[test] + fn test_validate_admits_row_lineage_ids_only_at_top_level() { + let uint64_field = |name: &str, id: i32| { + let mut field = Field::new_arrow(name, DataType::UInt64, false).unwrap(); + field.id = id; + field + }; + let mut key = Field::new_arrow("i", DataType::Int32, false).unwrap(); + key.id = 0; + + // A hidden lineage column: top-level, under its own name's reserved id. + let schema = Schema { + fields: vec![key.clone(), uint64_field(ROW_ID, crate::ROW_ID_FIELD_ID)], + metadata: HashMap::new(), + }; + schema.validate().unwrap(); + + // Another lineage column's id does not belong to this name. + let schema = Schema { + fields: vec![ + key.clone(), + uint64_field(ROW_ID, crate::ROW_CREATED_AT_VERSION_FIELD_ID), + ], + metadata: HashMap::new(), + }; + let error = schema.validate().unwrap_err(); + assert!(matches!(error, Error::Schema { .. }), "{error}"); + assert!(error.to_string().contains("negative id"), "{error}"); + + // A data file stores lineage columns only at the top level. + let mut parent = Field::new_arrow( + "s", + DataType::Struct(ArrowFields::from(vec![ArrowField::new( + ROW_CREATED_AT_VERSION, + DataType::UInt64, + false, + )])), + true, + ) + .unwrap(); + parent.id = 1; + parent.children[0].id = crate::ROW_CREATED_AT_VERSION_FIELD_ID; + let schema = Schema { + fields: vec![key, parent], + metadata: HashMap::new(), + }; + let error = schema.validate().unwrap_err(); + assert!(matches!(error, Error::Schema { .. }), "{error}"); + assert!(error.to_string().contains("negative id"), "{error}"); + } + #[test] fn test_resolve_quoted_fields() { // Test that top-level fields with dots are rejected during validation diff --git a/rust/lance/src/dataset/versions/mod.rs b/rust/lance/src/dataset/versions/mod.rs index e4af2e1c3e2..79f8f55660d 100644 --- a/rust/lance/src/dataset/versions/mod.rs +++ b/rust/lance/src/dataset/versions/mod.rs @@ -145,6 +145,12 @@ pub async fn write_fragments( target_bases_info: Option>, file_row_counts: Option>, ) -> Result<(Vec, Schema)> { + // A writer that spills row lineage into the fragment's own data file + // carries the hidden columns in its stream. They are not dataset fields: + // set them aside before the schema is checked against the dataset's and + // put them back on the schema that is written. This has to come before + // the blob promotion, which gives every negative field id a new one. + let (normalized_schema, lineage_fields) = split_row_lineage_fields(normalized_schema)?; let normalized_schema = match version { ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 => { write::promote_legacy_blob_schema(&normalized_schema)? @@ -152,11 +158,6 @@ pub async fn write_fragments( _ => normalized_schema, }; let version_name = format!("{version:?}"); - // A writer that spills row lineage into the fragment's own data file - // carries the hidden columns in its stream. They are not dataset fields: - // set them aside before the schema is checked against the dataset's and - // put them back, under their reserved ids, on the schema that is written. - let (normalized_schema, lineage_fields) = split_row_lineage_fields(normalized_schema); let mut schema = write::prepare_write_schema( dataset, normalized_schema, @@ -189,27 +190,39 @@ pub async fn write_fragments( Ok((fragments, schema)) } -/// Take the hidden row lineage columns out of a write schema, each keyed to -/// the reserved field id its name maps to. -fn split_row_lineage_fields(schema: Schema) -> (Schema, Vec) { +/// Take the hidden row lineage columns out of a write schema. +/// +/// A hidden column is a top-level field that already carries the reserved id +/// of its name, as the writers that spill lineage into a data file assign it. +/// A user column that only shares the name has a non-negative id and stays a +/// user column, to be checked against the dataset schema like any other. +fn split_row_lineage_fields(schema: Schema) -> Result<(Schema, Vec)> { let (lineage, user): (Vec<_>, Vec<_>) = schema .fields .into_iter() - .partition(|field| lance_core::row_lineage_field_id(&field.name).is_some()); - let lineage = lineage - .into_iter() - .map(|mut field| { - field.id = lance_core::row_lineage_field_id(&field.name).unwrap(); - field - }) - .collect(); - ( + .partition(|field| lance_core::row_lineage_field_id(&field.name) == Some(field.id)); + // The readers decode these columns as non-nullable `UInt64`, whatever the + // written type, so any other type would only fail once read back. + if let Some(field) = lineage + .iter() + .find(|field| field.nullable || field.data_type() != DataType::UInt64) + { + return Err(Error::internal(format!( + "hidden row lineage column {} (field id {}) must be a non-nullable UInt64, got \ + {} (nullable: {})", + field.name, + field.id, + field.data_type(), + field.nullable + ))); + } + Ok(( Schema { fields: user, metadata: schema.metadata, }, lineage, - ) + )) } #[allow(clippy::too_many_arguments)] @@ -1115,3 +1128,48 @@ pub fn validate_row_stream_read(version: ConcreteFileVersion) -> Result<()> { | ConcreteFileVersion::V2_3 => Ok(()), } } + +#[cfg(test)] +mod tests { + use super::*; + use lance_core::{ROW_CREATED_AT_VERSION, ROW_ID, ROW_ID_FIELD_ID}; + use rstest::rstest; + + fn lance_field(name: &str, data_type: DataType, nullable: bool, id: i32) -> Field { + let mut field = Field::try_from(&ArrowField::new(name, data_type, nullable)).unwrap(); + field.id = id; + field + } + + #[rstest] + #[case::uint64(DataType::UInt64, false, true)] + #[case::not_uint64(DataType::Int64, false, false)] + #[case::nullable(DataType::UInt64, true, false)] + fn split_row_lineage_fields_keys_on_reserved_ids( + #[case] data_type: DataType, + #[case] nullable: bool, + #[case] valid: bool, + ) { + let key = lance_field("i", DataType::Int32, false, 0); + // A user column that only shares a lineage column's name keeps its own + // id, so it stays with the fields checked against the dataset schema. + let named_like_lineage = lance_field(ROW_CREATED_AT_VERSION, DataType::UInt64, true, 1); + let hidden = lance_field(ROW_ID, data_type, nullable, ROW_ID_FIELD_ID); + let schema = Schema { + fields: vec![key.clone(), named_like_lineage.clone(), hidden.clone()], + metadata: HashMap::new(), + }; + + let result = split_row_lineage_fields(schema); + if !valid { + let error = result.unwrap_err(); + assert!(matches!(error, Error::Internal { .. }), "{error}"); + assert!(error.to_string().contains("non-nullable UInt64"), "{error}"); + return; + } + let (user, lineage) = result.unwrap(); + assert_eq!(user.fields, vec![key, named_like_lineage]); + // Split out with its reserved id intact. + assert_eq!(lineage, vec![hidden]); + } +} From 8fcfd06114919093f19df83334f2d9cb87c015eb Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Tue, 29 Sep 2026 05:19:01 +0800 Subject: [PATCH 4/6] fix(dataset): follow written file sizes when placing compaction lineage Compaction computed its output fragments' row lineage from the planned file row counts before writing, then paired it with the fragments the writer produced and checked only how many there were. A byte limit (max_bytes_per_file) makes the writer close a file early and spread the remaining rows over files of other sizes, so a stable-row-id compaction failed with a fragment count mismatch or, when the count happened to match, placed row ids shifted onto the wrong rows. This held for every stable-row-id table, including ones that never opted in to spilling. Only a table that can spill now plans its lineage before the write. Every other table carries its lineage over after the write, sized by the written physical rows, as it did before lineage could go into the data file. Placement follows the written sizes: when they differ from the plan, the inline sequences are cut again at them, the written total must equal the planned one, and every inline sequence must match its fragment's physical rows. The hidden columns need nothing, since they were appended row by row. The spill plan is now decided per sequence type. A type goes into the data files only when it is over the budget in every output fragment. When it is over in only some, the task writes no hidden columns and places each fragment's lineage on its own after the write, so a Range sequence that fits in a few manifest bytes no longer costs a read per row. The spilled sequences are released as soon as their column values are built, and each batch gets its own copy of those values. A slice shared the whole task's buffer, which the writer counts when sizing pages, so it cut a page for every batch. RowLineageSpill and the plan are no longer public API, and one table of column names and field ids drives the field ids, the write schema and the stream schema. Co-Authored-By: Claude Opus 5.5 --- rust/lance/src/dataset/optimize.rs | 629 +++++++++++++++++++++---- rust/lance/src/dataset/rowids.rs | 6 +- rust/lance/src/dataset/rowids/spill.rs | 230 ++++++--- 3 files changed, 708 insertions(+), 157 deletions(-) diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index e519a2b81df..06f203bd925 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -91,8 +91,8 @@ use super::fragment::FileFragment; use super::index::{DatasetIndexRemapperOptions, load_indices_for_remapping}; use super::rowids::RowVersionKind; use super::rowids::{ - RowLineage, RowLineageSpill, load_row_id_sequences, load_row_version_sequence, - place_row_lineage, plan_row_lineage_spill, + RowLineage, RowLineagePlan, RowLineageSpill, inline_row_lineage_max_bytes, + load_row_id_sequences, load_row_version_sequence, place_row_lineage, plan_row_lineage_spill, }; use super::transaction::{ Operation, RewriteGroup, RewrittenIndex, Transaction, TransactionBuilder, @@ -135,6 +135,7 @@ use lance_index::frag_reuse::{FRAG_REUSE_INDEX_NAME, FragReuseGroup}; use lance_index::is_system_index; use lance_index::metrics::NoOpMetricsCollector; use lance_table::format::{Fragment, IndexMetadata, RowDatasetVersionSequence}; +use lance_table::rowids::RowIdSequence; use roaring::{RoaringBitmap, RoaringTreemap}; use serde::{Deserialize, Serialize}; use tracing::{info, warn}; @@ -2570,30 +2571,7 @@ async fn rewrite_files( params.enable_stable_row_ids = true; } - // The output fragments' row lineage is known before anything is written: - // the row ids and versions carry over from the inputs, and the output - // file sizes are fixed above. Computing it here lets the sequences that - // leave the manifest ride along as hidden columns of the files being - // written, rather than need a file of their own. Binary copy cannot add - // columns to the files it copies, so it keeps writing that separate file. - let mut lineage_plan: Option<(Vec, RowLineageSpill)> = None; - if !can_binary_copy && dataset.manifest.uses_stable_row_ids() { - let chunk_sizes = file_row_counts - .iter() - .map(|rows| *rows as u64) - .collect::>(); - let lineages = compute_row_lineage(dataset.as_ref(), &fragments, &chunk_sizes).await?; - let spill = plan_row_lineage_spill(dataset.as_ref(), &lineages)?; - if spill.any() { - let stream = reader - .take() - .expect("reader must be prepared for non-binary-copy path"); - reader = Some(append_row_lineage_columns(stream, spill.columns(&lineages))); - } - lineage_plan = Some((lineages, spill)); - } - - if can_binary_copy { + let lineage_plan = if can_binary_copy { new_fragments = versions::rewrite_files_binary_copy( write_version, dataset.as_ref(), @@ -2627,27 +2605,64 @@ async fn rewrite_files( let _ = tx.send(captured); row_ids_rx = Some(rx); } + None } else { - // The hidden lineage columns the stream now carries have to be in the - // schema the writer is given, or it would leave them out of the file. + let mut stream = reader.expect("reader must be prepared for non-binary-copy path"); let mut write_schema = dataset.schema().clone(); - if let Some((_, spill)) = &lineage_plan { - write_schema.fields.extend(spill.schema_fields()?); - } + // On a table that can spill, the output fragments' lineage is computed + // before anything is written: the row ids and versions carry over from + // the inputs, so a sequence that leaves the manifest can ride along as + // a hidden column of the file being written rather than need a file of + // its own. Every other table carries its lineage over after the write, + // cut at the row counts actually written, and holds none of it while + // writing. Binary copy cannot add columns to the files it copies, so + // it always takes that path. + let spill_budget = if dataset.manifest.uses_stable_row_ids() { + inline_row_lineage_max_bytes(dataset.as_ref())? + } else { + None + }; + let lineage_plan = match spill_budget { + Some(limit) => { + let planned_rows = file_row_counts + .iter() + .map(|rows| *rows as u64) + .collect::>(); + let mut lineages = + compute_row_lineage(dataset.as_ref(), &fragments, &planned_rows).await?; + let plan = plan_row_lineage_spill(limit, &lineages); + if let RowLineagePlan::InFile(spill) = plan + && spill.any() + { + let fields = spill.schema_fields()?; + let columns = spill.take_columns(&mut lineages); + stream = append_row_lineage_columns(stream, &fields, columns); + // The writer stores only the columns its schema lists. + write_schema.fields.extend(fields); + } + Some(PlannedRowLineage { + lineages, + planned_rows, + plan, + }) + } + None => None, + }; let (frags, _) = write_fragments_internal_with_file_row_counts( write_version, Some(dataset.as_ref()), dataset.object_store.clone(), &dataset.base, write_schema, - reader.expect("reader must be prepared for non-binary-copy path"), + stream, params, None, Some(file_row_counts), ) .await?; new_fragments = frags; - } + lineage_plan + }; log::info!("Compaction task {}: file written", task_id); @@ -2669,9 +2684,10 @@ async fn rewrite_files( } else { if dataset.manifest.uses_stable_row_ids() { log::info!("Compaction task {}: rechunking stable row ids", task_id); - match lineage_plan.take() { - Some((lineages, spill)) => { - place_row_lineage_in_files(&mut new_fragments, &lineages, spill)? + match lineage_plan { + Some(planned) => { + place_planned_row_lineage(dataset.as_ref(), &mut new_fragments, planned) + .await? } None => { rechunk_row_lineage(dataset.as_ref(), &mut new_fragments, &fragments) @@ -2807,7 +2823,8 @@ async fn compute_row_lineage( /// Carry the lineage of `old_fragments` over to the already written /// `new_fragments` and place each one's sequences inline or in a separate /// lineage file, as the table's spill policy and their size call for. This is -/// the path for output files the lineage could not be written into. +/// the path for tables that cannot spill, and for binary copy, whose output +/// files cannot take the lineage as extra columns. async fn rechunk_row_lineage( dataset: &Dataset, new_fragments: &mut [Fragment], @@ -2824,63 +2841,225 @@ async fn rechunk_row_lineage( Ok(()) } -/// Mark each new fragment's spilled sequences as living in its own data file, -/// which was written with them as hidden columns, and place the rest inline. +/// The row lineage a compaction task computed for its planned output +/// fragments before writing them, and how the task places it. +struct PlannedRowLineage { + /// Each planned output fragment's lineage, in order. The sequences the + /// plan writes into the data files are empty: their values went into the + /// written columns. + lineages: Vec, + /// Each planned output fragment's row count, which the writer only + /// follows until a byte limit closes a file early. + planned_rows: Vec, + plan: RowLineagePlan, +} + +/// Place the lineage compaction planned before the write on the fragments it +/// wrote. +async fn place_planned_row_lineage( + dataset: &Dataset, + new_fragments: &mut [Fragment], + planned: PlannedRowLineage, +) -> Result<()> { + let PlannedRowLineage { + lineages, + planned_rows, + plan, + } = planned; + match plan { + RowLineagePlan::InFile(spill) => { + place_row_lineage_in_files(new_fragments, lineages, &planned_rows, spill) + } + RowLineagePlan::PerFragment => { + let lineages = fit_row_lineage_to_written( + new_fragments, + lineages, + &planned_rows, + RowLineageSpill::default(), + )?; + for (fragment, lineage) in new_fragments.iter_mut().zip(lineages) { + place_row_lineage(dataset, &lineage).await?.apply(fragment); + } + Ok(()) + } + } +} + +/// Mark each written fragment's spilled sequences as living in its own data +/// file, which was written with them as hidden columns, and place the rest +/// inline. fn place_row_lineage_in_files( new_fragments: &mut [Fragment], - lineages: &[RowLineage], + lineages: Vec, + planned_rows: &[u64], spill: RowLineageSpill, ) -> Result<()> { - if new_fragments.len() != lineages.len() { + let lineages = fit_row_lineage_to_written(new_fragments, lineages, planned_rows, spill)?; + for (fragment, lineage) in new_fragments.iter_mut().zip(&lineages) { + if spill.any() { + let [file] = fragment.files.as_slice() else { + return Err(Error::internal(format!( + "compaction wrote {} data files for fragment {}; the row lineage columns \ + are in exactly one", + fragment.files.len(), + fragment.id + ))); + }; + let missing = spill + .field_ids() + .find(|field_id| !file.fields.contains(field_id)); + if let Some(field_id) = missing { + return Err(Error::internal(format!( + "compaction wrote fragment {}'s data file {} without row lineage field {}", + fragment.id, file.path, field_id + ))); + } + } + spill.place_in_file(lineage).apply(fragment); + } + Ok(()) +} + +/// Fit the lineage planned for output fragments of `planned_rows` rows each to +/// the fragments the writer actually produced, one lineage per fragment. +/// +/// A byte limit can close a file before its planned row count, after which the +/// writer spreads the remaining rows over files of other sizes, so the +/// sequences kept inline are cut again at the written row counts. The types +/// `in_file` names need nothing: their columns were appended to the written +/// rows one by one, so they follow any split, and their sequences are empty. +/// +/// The in-file plan was made for the planned sizes. A written fragment is +/// never larger than the largest planned one, but it can join the tail of one +/// planned fragment to the head of the next, so a sequence that plan keeps +/// inline can end up over the budget, up to about twice it. That is acceptable +/// because the budget only bounds manifest growth and is not a format limit. +/// The per-fragment plan decides again on the cut sequences. +fn fit_row_lineage_to_written( + new_fragments: &[Fragment], + lineages: Vec, + planned_rows: &[u64], + in_file: RowLineageSpill, +) -> Result> { + let written_rows = new_fragments + .iter() + .map(|fragment| match fragment.physical_rows { + Some(rows) => Ok(rows as u64), + None => Err(Error::internal(format!( + "compaction wrote fragment {} without a physical row count", + fragment.id + ))), + }) + .collect::>>()?; + let planned_total = planned_rows.iter().sum::(); + let written_total = written_rows.iter().sum::(); + if written_total != planned_total { return Err(Error::internal(format!( - "compaction wrote {} fragments but planned lineage for {}", - new_fragments.len(), - lineages.len() + "compaction wrote {written_total} rows, in fragments of {written_rows:?} rows, \ + but planned row lineage for {planned_total} rows, in fragments of \ + {planned_rows:?} rows" ))); } - for (fragment, lineage) in new_fragments.iter_mut().zip(lineages) { - let [file] = fragment.files.as_slice() else { - return Err(Error::internal(format!( - "compaction wrote {} data files for fragment {}; the row lineage columns \ - are in exactly one", - fragment.files.len(), - fragment.id - ))); + + let lineages = if written_rows.as_slice() == planned_rows { + lineages + } else { + let cut_error = |error: Error| { + Error::internal(format!( + "compaction could not cut the row lineage planned for fragments of \ + {planned_rows:?} rows at the written fragments of {written_rows:?} rows: \ + {error}" + )) }; - let missing = spill - .field_ids() - .find(|field_id| !file.fields.contains(field_id)); - if let Some(field_id) = missing { + let mut row_ids = Vec::with_capacity(lineages.len()); + let mut created_at = Vec::with_capacity(lineages.len()); + let mut last_updated_at = Vec::with_capacity(lineages.len()); + for lineage in lineages { + row_ids.push(lineage.row_ids); + created_at.push(lineage.created_at); + last_updated_at.push(lineage.last_updated_at); + } + let fragment_count = written_rows.len(); + let written = || written_rows.iter().copied(); + let row_ids = if in_file.row_ids { + vec![RowIdSequence::new(); fragment_count] + } else { + lance_table::rowids::rechunk_sequences(row_ids, written(), false).map_err(cut_error)? + }; + let created_at = if in_file.created_at { + vec![RowDatasetVersionSequence::new(); fragment_count] + } else { + lance_table::rowids::version::rechunk_version_sequences(created_at, written(), false) + .map_err(cut_error)? + }; + let last_updated_at = if in_file.last_updated_at { + vec![RowDatasetVersionSequence::new(); fragment_count] + } else { + lance_table::rowids::version::rechunk_version_sequences( + last_updated_at, + written(), + false, + ) + .map_err(cut_error)? + }; + row_ids + .into_iter() + .zip(created_at) + .zip(last_updated_at) + .map(|((row_ids, created_at), last_updated_at)| RowLineage { + row_ids, + created_at, + last_updated_at, + }) + .collect::>() + }; + + if lineages.len() != new_fragments.len() { + return Err(Error::internal(format!( + "compaction has row lineage for {} fragments but wrote {}", + lineages.len(), + new_fragments.len() + ))); + } + for ((fragment, &rows), lineage) in new_fragments.iter().zip(&written_rows).zip(&lineages) { + let inline_lengths = [ + (in_file.row_ids, "row ids", lineage.row_ids.len()), + ( + in_file.created_at, + "created-at versions", + lineage.created_at.len(), + ), + ( + in_file.last_updated_at, + "last-updated-at versions", + lineage.last_updated_at.len(), + ), + ]; + let mismatch = inline_lengths + .into_iter() + .find(|&(is_in_file, _, len)| !is_in_file && len != rows); + if let Some((_, sequence, len)) = mismatch { return Err(Error::internal(format!( - "compaction wrote fragment {}'s data file {} without row lineage field {}", - fragment.id, file.path, field_id + "compaction has {len} inline {sequence} for fragment {} of {rows} physical rows", + fragment.id ))); } - spill.place_in_file(lineage).apply(fragment); } - Ok(()) + Ok(lineages) } -/// Append the hidden row lineage columns to the batches of `stream`, in row -/// order, so the file writer stores them next to the user columns. +/// Append the hidden row lineage columns, `fields` with the matching entry of +/// `columns` as each one's values, to the batches of `stream`, in row order, +/// so the file writer stores them next to the user columns. fn append_row_lineage_columns( stream: SendableRecordBatchStream, - columns: Vec<(i32, &'static str, Vec)>, + fields: &[LanceField], + columns: Vec>, ) -> SendableRecordBatchStream { - let mut fields = stream.schema().fields().to_vec(); - for (_, name, _) in &columns { - fields.push(Arc::new(ArrowField::new( - *name, - ArrowDataType::UInt64, - false, - ))); - } - let schema = Arc::new(ArrowSchema::new(fields)); - let values = columns - .into_iter() - .map(|(_, _, values)| arrow_array::UInt64Array::from(values)) - .collect::>(); - let total_rows = values.first().map_or(0, |column| column.len()); + let mut arrow_fields = stream.schema().fields().to_vec(); + arrow_fields.extend(fields.iter().map(|field| Arc::new(ArrowField::from(field)))); + let schema = Arc::new(ArrowSchema::new(arrow_fields)); + let total_rows = columns.first().map_or(0, |values| values.len()); let batch_schema = schema.clone(); let mut offset = 0usize; let batches = stream.map(move |batch| { @@ -2894,8 +3073,12 @@ fn append_row_lineage_columns( ))); } let mut arrays = batch.columns().to_vec(); - for column in &values { - arrays.push(Arc::new(column.slice(offset, rows)) as ArrayRef); + // Each batch gets a buffer of its own. A slice would share the whole + // task's buffer, and the writer, which sizes pages by the memory an + // array holds, would then cut a tiny page for every batch. + for values in &columns { + let batch_values = values[offset..offset + rows].to_vec(); + arrays.push(Arc::new(arrow_array::UInt64Array::from(batch_values)) as ArrayRef); } offset += rows; RecordBatch::try_new(batch_schema.clone(), arrays) @@ -3423,6 +3606,10 @@ mod tests { use crate::dataset::WriteDestination; use crate::dataset::index::frag_reuse::cleanup_frag_reuse_index; use crate::dataset::optimize::remapping::{transpose_row_addrs, transpose_row_ids_from_digest}; + use crate::dataset::rowids::read_spilled_row_ids; + use crate::dataset::rowids::{ + INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY, SPILL_ROW_LINEAGE_CONFIG_KEY, + }; use crate::dataset::scanner::ColumnOrdering; use crate::index::DatasetIndexExt; use crate::index::frag_reuse::{load_frag_reuse_index_details, open_frag_reuse_index}; @@ -3439,6 +3626,7 @@ mod tests { use lance_arrow::BLOB_META_KEY; use lance_core::Error; use lance_core::ROW_ID; + use lance_core::ROW_ID_FIELD_ID; use lance_core::utils::address::RowAddress; use lance_core::utils::tempfile::TempStrDir; use lance_datagen::Dimension; @@ -3452,7 +3640,9 @@ mod tests { use lance_index::vector::pq::PQBuildParams; use lance_index::{Index, IndexType}; use lance_linalg::distance::{DistanceType, MetricType}; + use lance_table::format::RowIdMeta; use lance_table::io::manifest::read_manifest_indexes; + use lance_table::rowids::read_row_ids; use lance_testing::datagen::{BatchGenerator, IncrementingInt32, RandomVector}; use rstest::rstest; use std::collections::HashSet; @@ -4645,6 +4835,287 @@ mod tests { assert_eq!(second_metrics, CompactionMetrics::default()); } + /// A byte limit can close a compaction output file before its planned row + /// count, and the writer then spreads the remaining rows over files of + /// other sizes. The lineage placed inline has to follow the written row + /// counts, and a written total that differs from the planned one is an + /// error rather than lineage shifted onto the wrong rows. + #[rstest] + #[case::extra_file( + &[10, 10], + &[9, 6, 5], + RowLineageSpill { row_ids: true, ..Default::default() }, + None + )] + #[case::shifted( + &[101, 100, 100], + &[100, 101, 100], + RowLineageSpill { row_ids: true, ..Default::default() }, + None + )] + #[case::shifted_all_inline( + &[101, 100, 100], + &[100, 101, 100], + RowLineageSpill::default(), + None + )] + #[case::mismatched_totals( + &[10, 10], + &[10, 9], + RowLineageSpill { row_ids: true, ..Default::default() }, + Some("wrote 19 rows") + )] + fn place_row_lineage_in_files_follows_written_sizes( + #[case] planned_rows: &[u64], + #[case] written_rows: &[usize], + #[case] spill: RowLineageSpill, + #[case] expected_error: Option<&str>, + ) { + let total_rows = planned_rows.iter().sum::(); + // Every row has a version of its own, so a sequence cut at the wrong + // row cannot pass for the right one. + let versions = (1..=total_rows).collect::>(); + let mut lineages = Vec::with_capacity(planned_rows.len()); + let mut start = 0_u64; + for rows in planned_rows { + let end = start + rows; + let created_at = + RowDatasetVersionSequence::from_versions(&versions[start as usize..end as usize]); + lineages.push(RowLineage { + row_ids: RowIdSequence::from(start..end), + last_updated_at: created_at.clone(), + created_at, + }); + start = end; + } + let columns = spill.take_columns(&mut lineages); + if spill.row_ids { + assert_eq!(columns, vec![(0..total_rows).collect::>()]); + } else { + assert!(columns.is_empty(), "{columns:?}"); + } + // Each written file carries the user column and the spilled columns. + let fields = std::iter::once(0) + .chain(spill.field_ids()) + .collect::>(); + let column_indices = (0..fields.len() as i32).collect::>(); + let mut fragments = written_rows + .iter() + .enumerate() + .map(|(id, rows)| { + Fragment::new(id as u64) + .with_file( + format!("{id}.lance"), + fields.clone(), + column_indices.clone(), + ConcreteFileVersion::V2_0, + None, + ) + .with_physical_rows(*rows) + }) + .collect::>(); + + let placed = place_row_lineage_in_files(&mut fragments, lineages, planned_rows, spill); + if let Some(expected) = expected_error { + let error = placed.unwrap_err(); + let Error::Internal { message, .. } = &error else { + panic!("expected an internal error, got {error}"); + }; + assert!(message.contains(expected), "{message}"); + return; + } + placed.unwrap(); + let mut placed_row_ids = Vec::with_capacity(versions.len()); + let mut placed_created_at = Vec::with_capacity(versions.len()); + let mut placed_last_updated_at = Vec::with_capacity(versions.len()); + for fragment in &fragments { + let physical_rows = fragment.physical_rows.unwrap() as u64; + match &fragment.row_id_meta { + Some(RowIdMeta::Column) if spill.row_ids => {} + Some(RowIdMeta::Inline(bytes)) if !spill.row_ids => { + let row_ids = read_row_ids(bytes).unwrap(); + assert_eq!(row_ids.len(), physical_rows); + placed_row_ids.extend(row_ids.iter()); + } + other => panic!( + "fragment {} has row id meta {other:?} under {spill:?}", + fragment.id + ), + } + for (meta, collected) in [ + (&fragment.created_at_version_meta, &mut placed_created_at), + ( + &fragment.last_updated_at_version_meta, + &mut placed_last_updated_at, + ), + ] { + let sequence = meta.as_ref().unwrap().load_sequence().unwrap(); + assert_eq!(sequence.len(), physical_rows); + collected.extend(sequence.versions()); + } + } + if !spill.row_ids { + assert_eq!(placed_row_ids, (0..total_rows).collect::>()); + } + assert_eq!(placed_created_at, versions); + assert_eq!(placed_last_updated_at, versions); + } + + /// A sequence type over the budget in only some output fragments sends the + /// task down the per-fragment path: once the lineage is cut at the written + /// sizes, a fragment whose row ids fit keeps them inline, and only the ones + /// over the budget get a lineage file of their own. + #[tokio::test] + async fn per_fragment_plan_spills_only_the_fragments_over_budget() { + let mut dataset = lance_datagen::gen_batch() + .col("i", lance_datagen::array::step::()) + .into_ram_dataset_with_params( + FragmentCount::from(1), + FragmentRowCount::from(10), + Some(WriteParams { + enable_stable_row_ids: true, + max_rows_per_file: 10, + ..Default::default() + }), + ) + .await + .unwrap(); + dataset + .update_config([ + (SPILL_ROW_LINEAGE_CONFIG_KEY, "true"), + (INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY, "100"), + ]) + .await + .unwrap(); + // A range encodes to a few bytes at any length. The scattered ids, a + // permutation of 500..1000 with no runs, take about two bytes a row, + // so any 250 of them are well over the 100-byte budget. + let scattered = (0..500_u64) + .map(|i| 500 + (i * 7919) % 500) + .collect::>(); + let row_ids = [ + RowIdSequence::from(0..500), + RowIdSequence::from(scattered.as_slice()), + ]; + let lineages = row_ids + .into_iter() + .map(|row_ids| { + let created_at = + RowDatasetVersionSequence::from_uniform_row_count(row_ids.len(), 1); + RowLineage { + row_ids, + last_updated_at: created_at.clone(), + created_at, + } + }) + .collect::>(); + let plan = plan_row_lineage_spill(100, &lineages); + assert_eq!(plan, RowLineagePlan::PerFragment); + // Written as [250, 500, 250] rather than the planned [500, 500], so + // the middle fragment joins the tail of the range to the head of the + // scattered ids. + let mut fragments = [250, 500, 250] + .into_iter() + .enumerate() + .map(|(id, rows)| { + Fragment::new(id as u64) + .with_file( + format!("{id}.lance"), + vec![0], + vec![0], + ConcreteFileVersion::V2_0, + None, + ) + .with_physical_rows(rows) + }) + .collect::>(); + + place_planned_row_lineage( + &dataset, + &mut fragments, + PlannedRowLineage { + lineages, + planned_rows: vec![500, 500], + plan, + }, + ) + .await + .unwrap(); + + assert!( + matches!(fragments[0].row_id_meta, Some(RowIdMeta::Inline(_))), + "a fragment whose row ids fit must keep them inline, got {:?}", + fragments[0].row_id_meta + ); + assert_eq!(fragments[0].files.len(), 1); + for fragment in &fragments[1..] { + assert!( + matches!(fragment.row_id_meta, Some(RowIdMeta::Column)), + "fragment {}'s row ids are over the budget, got {:?}", + fragment.id, + fragment.row_id_meta + ); + assert_eq!(fragment.files.len(), 2); + assert_eq!(fragment.files[1].fields.as_ref(), [ROW_ID_FIELD_ID]); + } + let mut placed_row_ids = Vec::with_capacity(1000); + for fragment in &fragments { + let row_ids = match &fragment.row_id_meta { + Some(RowIdMeta::Inline(bytes)) => read_row_ids(bytes).unwrap(), + _ => read_spilled_row_ids(&dataset, fragment).await.unwrap(), + }; + placed_row_ids.extend(row_ids.iter()); + } + assert_eq!( + placed_row_ids, + (0..500) + .chain(scattered.iter().copied()) + .collect::>() + ); + } + + /// The writer sizes a column's pages by the memory its arrays hold, so a + /// lineage array sharing one buffer with the whole task would make it cut + /// a tiny page for every batch. + #[tokio::test] + async fn append_row_lineage_columns_gives_each_batch_its_own_buffer() { + let batch = arrow_array::record_batch!(("i", Int32, [0, 1, 2, 3])).unwrap(); + let schema = batch.schema(); + let batches = vec![batch; 3] + .into_iter() + .map(Ok::<_, datafusion::error::DataFusionError>); + let input: SendableRecordBatchStream = Box::pin(RecordBatchStreamAdapter::new( + schema, + futures::stream::iter(batches), + )); + let spill = RowLineageSpill { + row_ids: true, + ..Default::default() + }; + let fields = spill.schema_fields().unwrap(); + + let output = append_row_lineage_columns(input, &fields, vec![(0..12).collect()]) + .try_collect::>() + .await + .unwrap(); + + assert_eq!(output.len(), 3); + for (index, batch) in output.iter().enumerate() { + let start = index as u64 * 4; + let row_ids = batch[ROW_ID].as_primitive::(); + assert_eq!( + row_ids.values().to_vec(), + (start..start + 4).collect::>() + ); + assert!( + row_ids.get_buffer_memory_size() <= batch.num_rows() * 8, + "batch {index}'s row ids hold {} bytes for {} rows", + row_ids.get_buffer_memory_size(), + batch.num_rows() + ); + } + } + #[rstest] #[tokio::test] async fn test_compact_data_files( diff --git a/rust/lance/src/dataset/rowids.rs b/rust/lance/src/dataset/rowids.rs index 93a12ec14d9..6474d39df7b 100644 --- a/rust/lance/src/dataset/rowids.rs +++ b/rust/lance/src/dataset/rowids.rs @@ -27,10 +27,10 @@ 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, RowLineageSpill, SPILL_ROW_LINEAGE_CONFIG_KEY, - inline_row_lineage_max_bytes, place_row_lineage, plan_row_lineage_spill, read_spilled_row_ids, - read_spilled_versions, + PlacedRowLineage, RowLineage, SPILL_ROW_LINEAGE_CONFIG_KEY, inline_row_lineage_max_bytes, + place_row_lineage, read_spilled_row_ids, read_spilled_versions, }; +pub(crate) use spill::{RowLineagePlan, RowLineageSpill, plan_row_lineage_spill}; pub(super) use validate::validate_stable_row_ids; /// Load a row id sequence from the given dataset and fragment. diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs index 8b6fe955a74..ab139f43278 100644 --- a/rust/lance/src/dataset/rowids/spill.rs +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -104,9 +104,20 @@ pub fn inline_row_lineage_max_bytes(dataset: &Dataset) -> Result> }) } -/// Which of a compaction task's sequences leave the manifest, decided for all -/// its output fragments together so that every data file the task writes -/// carries the same hidden columns. +/// The hidden row lineage columns by name and reserved field id, in the order +/// a data file written by compaction stores them. +const LINEAGE_COLUMNS: [(&str, i32); 3] = [ + (ROW_ID, ROW_ID_FIELD_ID), + (ROW_CREATED_AT_VERSION, ROW_CREATED_AT_VERSION_FIELD_ID), + ( + ROW_LAST_UPDATED_AT_VERSION, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, + ), +]; + +/// Which of a compaction task's sequences leave the manifest as hidden columns +/// of the data files it writes. Every file the task writes carries the same +/// columns, so this holds for all of its output fragments. #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] pub struct RowLineageSpill { pub row_ids: bool, @@ -119,52 +130,52 @@ impl RowLineageSpill { self.row_ids || self.created_at || self.last_updated_at } - /// The reserved field ids of the columns that spill, in column order. + /// The reserved field ids of the spilled columns, in column order. pub fn field_ids(&self) -> impl Iterator { - [ - (self.row_ids, ROW_ID_FIELD_ID), - (self.created_at, ROW_CREATED_AT_VERSION_FIELD_ID), - (self.last_updated_at, ROW_LAST_UPDATED_AT_VERSION_FIELD_ID), - ] - .into_iter() - .filter(|(spilled, _)| *spilled) - .map(|(_, field_id)| field_id) + let flags = [self.row_ids, self.created_at, self.last_updated_at]; + LINEAGE_COLUMNS + .into_iter() + .zip(flags) + .filter_map(|((_, field_id), spilled)| spilled.then_some(field_id)) } - /// The hidden columns to write, as `(field id, column name, values)`, - /// concatenated over `lineages` in order. - pub fn columns(&self, lineages: &[RowLineage]) -> Vec<(i32, &'static str, Vec)> { - let mut columns = Vec::with_capacity(3); - if self.row_ids { - let values = lineages.iter().flat_map(|l| l.row_ids.iter()).collect(); - columns.push((ROW_ID_FIELD_ID, ROW_ID, values)); - } - if self.created_at { - let values = lineages - .iter() - .flat_map(|l| l.created_at.versions()) - .collect(); - columns.push(( - ROW_CREATED_AT_VERSION_FIELD_ID, - ROW_CREATED_AT_VERSION, - values, - )); - } - if self.last_updated_at { - let values = lineages - .iter() - .flat_map(|l| l.last_updated_at.versions()) - .collect(); - columns.push(( - ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, - ROW_LAST_UPDATED_AT_VERSION, - values, - )); - } - columns + /// Move the spilled sequences out of `lineages` as the values of the + /// hidden columns, one vector per column in column order, concatenated + /// over `lineages`. The spilled sequences are left empty: nothing places + /// them inline, and their run-length form would otherwise stay alive + /// through the write next to the values built from it. + pub fn take_columns(&self, lineages: &mut [RowLineage]) -> Vec> { + let total_rows = lineages + .iter() + .map(|lineage| lineage.row_ids.len() as usize) + .sum::(); + self.field_ids() + .map(|field_id| { + let mut values = Vec::with_capacity(total_rows); + for lineage in lineages.iter_mut() { + match field_id { + ROW_ID_FIELD_ID => { + let row_ids = std::mem::take(&mut lineage.row_ids); + values.extend(row_ids.iter()); + } + ROW_CREATED_AT_VERSION_FIELD_ID => { + let created_at = std::mem::take(&mut lineage.created_at); + values.extend(created_at.versions()); + } + // The only other id `field_ids` yields. + _ => { + let last_updated_at = std::mem::take(&mut lineage.last_updated_at); + values.extend(last_updated_at.versions()); + } + } + } + values + }) + .collect() } - /// The hidden columns as write-schema fields, under their reserved ids. + /// The spilled columns as write-schema fields, in column order and under + /// their reserved ids. pub fn schema_fields(&self) -> Result> { self.field_ids().map(lineage_field).collect() } @@ -197,46 +208,77 @@ impl RowLineageSpill { } /// The hidden column a spilled lineage sequence is stored in: a non-nullable -/// `UInt64` named after the sequence, under its reserved `field_id`. +/// `UInt64` named after the sequence in [`LINEAGE_COLUMNS`], under its +/// reserved `field_id`. fn lineage_field(field_id: i32) -> Result { - let name = match field_id { - ROW_ID_FIELD_ID => ROW_ID, - ROW_CREATED_AT_VERSION_FIELD_ID => ROW_CREATED_AT_VERSION, - ROW_LAST_UPDATED_AT_VERSION_FIELD_ID => ROW_LAST_UPDATED_AT_VERSION, - _ => { - return Err(Error::internal(format!( - "field id {field_id} is not a row lineage field id" - ))); - } - }; + let (name, _) = LINEAGE_COLUMNS + .into_iter() + .find(|(_, id)| *id == field_id) + .ok_or_else(|| { + Error::internal(format!( + "field id {field_id} is not a reserved row lineage field id" + )) + })?; let mut field = lance_core::datatypes::Field::try_from(&ArrowField::new(name, DataType::UInt64, false))?; field.id = field_id; Ok(field) } -/// Decide which sequence types of `lineages` spill: each one whose encoding -/// exceeds the table's inline budget in any of them. Nothing spills on a -/// table that has not opted in. -pub fn plan_row_lineage_spill( - dataset: &Dataset, - lineages: &[RowLineage], -) -> Result { - let Some(limit) = inline_row_lineage_max_bytes(dataset)? else { - return Ok(RowLineageSpill::default()); - }; - let over = |encoded: usize| encoded > limit; - Ok(RowLineageSpill { - row_ids: lineages - .iter() - .any(|l| over(write_row_ids(&l.row_ids).len())), - created_at: lineages +/// How a compaction task places the lineage of the fragments it writes. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum RowLineagePlan { + /// Each sequence type the spill names is over the inline budget in every + /// output fragment and is written into the data files as a hidden column; + /// every other type is under the budget in all of them and stays inline. + InFile(RowLineageSpill), + /// Some sequence type is over the budget in some output fragments but not + /// in others. A column in every file would also spill the sequences that + /// fit inline, and their readers would then pay a read per row for what a + /// few bytes of manifest hold. The task writes no hidden columns instead, + /// and after the write places each fragment's lineage on its own with + /// [`place_row_lineage`], which spills only what is over the budget, to a + /// separate lineage file. + PerFragment, +} + +/// Plan how the lineage of a compaction task's output fragments, `lineages`, +/// is placed when every encoded sequence over `limit` bytes has to leave the +/// manifest (see [`inline_row_lineage_max_bytes`]). +pub fn plan_row_lineage_spill(limit: usize, lineages: &[RowLineage]) -> RowLineagePlan { + let over = |encoded: Vec| encoded.len() > limit; + let row_ids = over_in_all_or_none(lineages.iter().map(|l| over(write_row_ids(&l.row_ids)))); + let created_at = over_in_all_or_none( + lineages .iter() - .any(|l| over(write_dataset_versions(&l.created_at).len())), - last_updated_at: lineages + .map(|l| over(write_dataset_versions(&l.created_at))), + ); + let last_updated_at = over_in_all_or_none( + lineages .iter() - .any(|l| over(write_dataset_versions(&l.last_updated_at).len())), - }) + .map(|l| over(write_dataset_versions(&l.last_updated_at))), + ); + match (row_ids, created_at, last_updated_at) { + (Some(row_ids), Some(created_at), Some(last_updated_at)) => { + RowLineagePlan::InFile(RowLineageSpill { + row_ids, + created_at, + last_updated_at, + }) + } + _ => RowLineagePlan::PerFragment, + } +} + +/// Whether one sequence type is over the budget in every output fragment +/// (`Some(true)`) or in none (`Some(false)`), given each output's verdict in +/// `over`; `None` when the outputs disagree. It stops encoding at the first +/// disagreement. +fn over_in_all_or_none(mut over: impl Iterator) -> Option { + let Some(first) = over.next() else { + return Some(false); + }; + over.all(|next| next == first).then_some(first) } /// The per-row lineage of one fragment, in row offset order. @@ -975,6 +1017,44 @@ mod tests { assert!(placed.file.is_none()); } + /// Compaction writes a sequence type into its data files only when it is + /// over the budget in every output fragment. A type over it in only some + /// of them sends the task down the per-fragment path, so the sequences + /// that fit stay inline. + #[rstest] + #[case::all_over( + vec![scattered_row_ids(500), scattered_row_ids(500)], + RowLineagePlan::InFile(RowLineageSpill { row_ids: true, ..Default::default() }) + )] + #[case::none_over( + vec![RowIdSequence::from(0..500), RowIdSequence::from(500..1000)], + RowLineagePlan::InFile(RowLineageSpill::default()) + )] + #[case::mixed( + vec![RowIdSequence::from(0..500), scattered_row_ids(500)], + RowLineagePlan::PerFragment + )] + fn plan_row_lineage_spill_decides_per_kind( + #[case] row_ids: Vec, + #[case] expected: RowLineagePlan, + ) { + // Single-run version sequences encode to a few bytes, far under the + // budget, so the row ids alone decide. + let lineages = row_ids + .into_iter() + .map(|row_ids| { + let created_at = + RowDatasetVersionSequence::from_uniform_row_count(row_ids.len(), 1); + RowLineage { + row_ids, + last_updated_at: created_at.clone(), + created_at, + } + }) + .collect::>(); + assert_eq!(plan_row_lineage_spill(100, &lineages), expected); + } + /// Opt the table into spilling, at a zero inline budget so every sequence /// spills regardless of size: reaching the natural 200 KiB threshold needs /// ~25k scattered rows, more than these tests need to prove. From c0541e31eb7a6c7e1040d7123c986bdc12d23d9e Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Tue, 29 Sep 2026 06:16:27 +0800 Subject: [PATCH 5/6] fix(dataset): reclaim lineage carriers whose user columns are gone Compaction writes spilled row lineage into the data file of each fragment it writes, and schema-only commits keep that file because it holds the lineage's only copy. Once every user column in it is dropped, cast or replaced, the file carries nothing but lineage and dead user bytes, and two things went wrong. FileFragment::validate treated any file with a non-negative field id as user data and opened it. After drop_columns the carrier still lists the dropped ids, shares no field with the schema, and validate() reported the table as corrupt. A file with no schema field that carries one of the fragment's spilled sequences is now left unopened and checked through the sequences it carries. Every other file is opened as before, so a stale file of user fields the schema no longer has is still reported. Nothing rewrote the file either: the planner picks fragments by size, deletions and overlays only, so the dead bytes stayed referenced for good. A fragment with a file that holds no schema field, holds a dropped or tombstoned user field, and carries lineage now compacts on its own, which moves the lineage next to the live columns and lets the file go. Neither layout compaction writes matches -- its data files hold live fields, and the lineage-only files of binary copy and update hold no user field -- so the rewrite cannot repeat. Binary copy cannot rewrite such a file, so under ForceBinaryCopy the planner leaves the fragment alone rather than plan a task that fails the whole run. Co-Authored-By: Claude Opus 5.5 --- rust/lance/src/dataset/fragment.rs | 35 +++- rust/lance/src/dataset/optimize.rs | 95 ++++++++++- rust/lance/src/dataset/rowids/spill.rs | 219 ++++++++++++++++++++++++- 3 files changed, 337 insertions(+), 12 deletions(-) diff --git a/rust/lance/src/dataset/fragment.rs b/rust/lance/src/dataset/fragment.rs index 1a59e33eae1..875ed8a483f 100644 --- a/rust/lance/src/dataset/fragment.rs +++ b/rust/lance/src/dataset/fragment.rs @@ -1754,7 +1754,10 @@ impl FileFragment { /// Verifies: /// * All field ids in the fragment are distinct /// * Within each data file, field ids are in increasing order - /// * All data files exist and have the same length + /// * All data files holding user data exist and have the same length. A + /// file kept only for the spilled row lineage it carries is not opened; + /// [`Dataset::validate`] reads that lineage back, which checks its + /// length. /// * Field ids are distinct between data files. /// * Deletion file exists and has rowids in the correct range /// * `Fragment.physical_rows` matches length of file @@ -1821,14 +1824,28 @@ impl FileFragment { data_file.validate(&self.dataset.data_file_dir(data_file)?)?; } - // A file holding only row lineage columns has no dataset field to open - // it by; its length is checked against `physical_rows` when the - // sequences it carries are validated. - let user_data_files = self - .metadata - .files - .iter() - .filter(|data_file| data_file.fields.iter().any(|field| *field >= 0)); + // A file that holds no field of the dataset schema is not opened when + // it holds no user field at all, or when the fragment keeps it for a + // spilled row lineage sequence it carries -- a data file whose user + // columns were all dropped or replaced after compaction wrote the + // lineage next to them. The sequences it carries are checked against + // `physical_rows` when they are validated. Any other file is opened, + // so a file listing only user fields the schema does not have is + // still reported. + let schema_field_ids = self + .dataset + .schema() + .fields_pre_order() + .map(|field| field.id) + .collect::>(); + let spilled_field_ids = self.metadata.spilled_row_lineage_field_ids(); + let user_data_files = self.metadata.files.iter().filter(|data_file| { + let fields = &data_file.fields; + let holds_schema_field = fields.iter().any(|id| schema_field_ids.contains(id)); + let holds_user_field = fields.iter().any(|id| *id >= 0); + let holds_spilled_lineage = fields.iter().any(|id| spilled_field_ids.contains(id)); + holds_schema_field || (holds_user_field && !holds_spilled_lineage) + }); let get_lengths = user_data_files.clone().map(|data_file| async move { let data_file_dir = self.dataset.data_file_dir(data_file)?; let reader = self diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index 06f203bd925..55e930bc160 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -134,7 +134,11 @@ use lance_file::version::{ConcreteFileVersion, LanceFileVersion}; use lance_index::frag_reuse::{FRAG_REUSE_INDEX_NAME, FragReuseGroup}; use lance_index::is_system_index; use lance_index::metrics::NoOpMetricsCollector; -use lance_table::format::{Fragment, IndexMetadata, RowDatasetVersionSequence}; +use lance_table::format::overlay::TOMBSTONE_FIELD_ID; +use lance_table::format::{ + Fragment, IndexMetadata, ROW_CREATED_AT_VERSION_FIELD_ID, ROW_ID_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, RowDatasetVersionSequence, +}; use lance_table::rowids::RowIdSequence; use roaring::{RoaringBitmap, RoaringTreemap}; use serde::{Deserialize, Serialize}; @@ -885,6 +889,19 @@ impl CompactionPlanner for DefaultCompactionPlanner { .collect::>() }; + let schema_field_ids = dataset + .schema() + .fields_pre_order() + .map(|field| field.id) + .collect::>(); + // Only reencoding reclaims a lineage carrier whose user columns are + // gone: its fields are not the schema's, so binary copy never applies + // to it, and under `ForceBinaryCopy` its task would fail the whole run. + let is_binary_copy_only = matches!( + self.options.compaction_mode(), + CompactionMode::ForceBinaryCopy + ); + let mut candidate_bins: Vec = Vec::new(); let mut current_bin: Option = None; let mut i = 0; @@ -910,6 +927,12 @@ impl CompactionPlanner for DefaultCompactionPlanner { // Too many overlays: fully compact this fragment on its own, // regardless of its size or deletion count. Some(CompactionCandidacy::CompactItself) + } else if !is_binary_copy_only + && pins_dead_user_columns_for_lineage(&fragment, &schema_field_ids) + { + // Only a rewrite moves the row lineage out of a file whose + // user columns are all gone and lets the file go. + Some(CompactionCandidacy::CompactItself) } else if self.options.materialize_deletions && metrics.deletion_percentage() > self.options.materialize_deletions_threshold { @@ -1111,6 +1134,43 @@ async fn collect_metrics(fragment: &FileFragment) -> Result { }) } +/// Whether one of `fragment`'s data files is referenced only for the row +/// lineage columns it carries, while the user columns next to them have all +/// been dropped or replaced. `schema_field_ids` holds the ids of the schema's +/// fields. +/// +/// Compaction writes spilled lineage into the data file of each fragment it +/// writes. Once every user column in that file is dropped, cast or replaced, +/// the lineage alone keeps the file, and its dead bytes, referenced: cleanup +/// cannot remove it and nothing else rewrites it. A lineage-only file, which +/// binary copy and update write, holds no dead user column and never matches; +/// nor does any file compaction writes, so compacting a fragment that matches +/// cannot leave one that matches again. +fn pins_dead_user_columns_for_lineage( + fragment: &Fragment, + schema_field_ids: &HashSet, +) -> bool { + fragment.files.iter().any(|file| { + let mut holds_dead_user_column = false; + let mut holds_lineage = false; + for field_id in file.fields.iter() { + if schema_field_ids.contains(field_id) { + return false; + } + match *field_id { + ROW_ID_FIELD_ID + | ROW_CREATED_AT_VERSION_FIELD_ID + | ROW_LAST_UPDATED_AT_VERSION_FIELD_ID => holds_lineage = true, + TOMBSTONE_FIELD_ID => holds_dead_user_column = true, + // A user field the schema no longer has: dropping a column + // leaves the ids of the files it keeps as they are. + other => holds_dead_user_column |= other >= 0, + } + } + holds_dead_user_column && holds_lineage + }) +} + /// Truncates a planned task list to the configured per-run source budgets /// (`max_source_fragments`, `max_source_rows`, `max_source_bytes`). /// @@ -5074,6 +5134,39 @@ mod tests { ); } + /// A fragment compacts on its own when one of its files is kept only for + /// the row lineage it carries next to user columns that are gone. Neither + /// layout compaction writes may qualify, or every compaction would plan + /// another. Field -2 is the tombstone and -3..=-5 are the lineage columns; + /// the schema holds fields 0 and 1, and 5 and 6 were dropped. + #[rstest] + #[case::in_file_carrier(vec![vec![0, 1, -3, -4, -5]], false)] + #[case::lineage_only_file(vec![vec![0, 1], vec![-3, -4, -5]], false)] + #[case::tombstoned_carrier(vec![vec![-2, -2, -3, -4, -5], vec![0, 1]], true)] + #[case::carrier_of_dropped_columns(vec![vec![5, 6, -3, -4, -5], vec![0, 1]], true)] + #[case::dead_file_without_lineage(vec![vec![-2, 6], vec![0, 1]], false)] + fn compaction_reclaims_lineage_carriers_of_dead_user_columns( + #[case] files: Vec>, + #[case] expected: bool, + ) { + let mut fragment = Fragment::new(0); + for (index, fields) in files.into_iter().enumerate() { + let column_indices = (0..fields.len() as i32).collect(); + fragment.add_file( + format!("{index}.lance"), + fields, + column_indices, + ConcreteFileVersion::V2_2, + None, + ); + } + let schema_field_ids = HashSet::from([0, 1]); + assert_eq!( + pins_dead_user_columns_for_lineage(&fragment, &schema_field_ids), + expected + ); + } + /// The writer sizes a column's pages by the memory its arrays hold, so a /// lineage array sharing one buffer with the whole task would make it cut /// a tiny page for every batch. diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs index ab139f43278..0d23a9bd7f4 100644 --- a/rust/lance/src/dataset/rowids/spill.rs +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -740,9 +740,15 @@ async fn read_spilled_column( mod tests { use super::*; use crate::dataset::cleanup::{CleanupPolicyBuilder, cleanup_old_versions}; - use crate::dataset::optimize::{CompactionMode, CompactionOptions, compact_files}; + use crate::dataset::fragment::FileFragment; + use crate::dataset::optimize::{ + CompactionMode, CompactionOptions, compact_files, plan_compaction, + }; use crate::dataset::rowids::{RowVersionKind, load_row_id_sequence, load_row_version_sequence}; - use crate::dataset::{ColumnAlteration, UpdateBuilder, WriteMode, WriteParams}; + use crate::dataset::transaction::Operation; + use crate::dataset::{ + ColumnAlteration, NewColumnTransform, UpdateBuilder, WriteMode, WriteParams, + }; use arrow_array::builder::{ListBuilder, StringBuilder}; use arrow_array::{Int32Array, RecordBatchIterator, StringArray, StructArray}; use arrow_schema::Field; @@ -751,6 +757,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 lance_table::format::overlay::TOMBSTONE_FIELD_ID; use rstest::rstest; /// A sequence with no runs to exploit, which is what a globally shuffled @@ -1348,6 +1355,11 @@ mod tests { ); assert_eq!(collect_lineage(&dataset).await, before); dataset.validate().await.unwrap(); + + // The lineage-only file holds no dead user column, so it is no reason + // to compact the fragment again. + let plan = plan_compaction(&dataset, &one_fragment()).await.unwrap(); + assert_eq!(plan.num_tasks(), 0); } /// Cleanup decides what to delete by walking @@ -1913,4 +1925,207 @@ mod tests { let reopened = Dataset::open(uri).await.unwrap(); assert_eq!(by_key(&collect_rows(&reopened).await), expected); } + + /// How [`dropping_every_column_of_a_lineage_carrier_keeps_lineage_until_compaction`] + /// takes every user column away from the data file that compaction wrote + /// the lineage into. + #[derive(Debug, Clone, Copy)] + enum CarrierChange { + /// `k` is added in a file of its own and `i` and `j` are dropped. A + /// drop leaves the files it keeps as they are, so the data file still + /// lists the ids of `i` and `j`, which the schema no longer has. + Drop, + /// `i` and `j` are cast. A cast rewrites the columns under new field + /// ids into a new file and leaves their old ids, now dead, in the data + /// file. + Cast, + /// `i` and `j` are rewritten under their own field ids by a + /// `DataReplacement`, which tombstones them in the data file. + Replace, + } + + /// Compaction writes the lineage into the fragment's data file, which + /// then outlives its user columns. It stays, as the lineage's only copy, + /// and neither reads nor `validate` open it for user data. The next + /// compaction rewrites the fragment, which moves the lineage next to the + /// live columns and lets the old file go. Binary copy cannot rewrite it, + /// so a compaction restricted to binary copy leaves the fragment alone + /// rather than planning a task that fails. + #[rstest] + #[case::drop(CarrierChange::Drop)] + #[case::cast(CarrierChange::Cast)] + #[case::replace(CarrierChange::Replace)] + #[tokio::test] + async fn dropping_every_column_of_a_lineage_carrier_keeps_lineage_until_compaction( + #[case] change: CarrierChange, + ) { + let dir = TempStrDir::default(); + let uri = dir.as_str(); + let mut dataset = appended_dataset_with(uri, 4, 250, UserColumns::WithCopy, None).await; + spill_everything(&mut dataset).await; + compact_files(&mut dataset, one_fragment(), None) + .await + .unwrap(); + let (row_ids, created_at, _) = collect_lineage(&dataset).await; + + // The ids the data file lists for `i` and `j` once the change is made. + let carrier_user_fields = match change { + CarrierChange::Drop => { + // `k` lives in a file of its own, so dropping `i` and `j` + // leaves the compacted file with no field in the schema. + dataset + .add_columns( + NewColumnTransform::SqlExpressions(vec![("k".into(), "i + 1".into())]), + None, + None, + ) + .await + .unwrap(); + dataset.drop_columns(&["i", "j"]).await.unwrap(); + [0, 1] + } + CarrierChange::Cast => { + let alterations = [ + ColumnAlteration::new("i".into()).cast_to(DataType::Int64), + ColumnAlteration::new("j".into()).cast_to(DataType::Int64), + ]; + dataset.alter_columns(&alterations).await.unwrap(); + [0, 1] + } + CarrierChange::Replace => { + let batch = keyed_batch(UserColumns::WithCopy, 0..1_000); + let fragments = dataset.get_fragments(); + let replacement = fragments[0] + .write_columns(futures::stream::iter([Ok(batch)]), dataset.schema()) + .await + .unwrap(); + let read_version = dataset.manifest.version; + let operation = Operation::DataReplacement { + replacements: vec![replacement], + }; + dataset = Dataset::commit( + Arc::new(dataset), + operation, + Some(read_version), + None, + None, + Arc::new(Default::default()), + false, + ) + .await + .unwrap(); + [TOMBSTONE_FIELD_ID, TOMBSTONE_FIELD_ID] + } + }; + + let fragments = dataset.get_fragments(); + assert_eq!(fragments.len(), 1); + let metadata = fragments[0].metadata(); + let carrier = metadata + .row_lineage_file(ROW_ID_FIELD_ID) + .unwrap() + .expect("the schema change must keep the lineage carrier"); + let (user_fields, lineage_fields) = carrier.fields.split_at(2); + assert_eq!(user_fields, carrier_user_fields, "{metadata:?}"); + assert_eq!( + lineage_fields, + [ + ROW_ID_FIELD_ID, + ROW_CREATED_AT_VERSION_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID + ] + ); + assert!( + carrier + .fields + .iter() + .all(|field_id| dataset.schema().field_by_id(*field_id).is_none()), + "the carrier must have lost every user column: {metadata:?}" + ); + dataset.validate().await.unwrap(); + // A change may stamp every row as updated, so only the row ids and + // created-at versions are compared. + let changed = collect_lineage(&dataset).await; + assert_eq!((&changed.0, &changed.1), (&row_ids, &created_at)); + + let binary_copy_only = CompactionOptions { + compaction_mode: Some(CompactionMode::ForceBinaryCopy), + ..one_fragment() + }; + let plan = plan_compaction(&dataset, &binary_copy_only).await.unwrap(); + assert_eq!(plan.num_tasks(), 0); + + compact_files(&mut dataset, one_fragment(), None) + .await + .unwrap(); + + let fragments = dataset.get_fragments(); + assert_eq!(fragments.len(), 1); + let metadata = fragments[0].metadata(); + let mut expected_fields = dataset.schema().field_ids(); + expected_fields.extend([ + ROW_ID_FIELD_ID, + ROW_CREATED_AT_VERSION_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, + ]); + assert_eq!(metadata.files.len(), 1, "{metadata:?}"); + assert_eq!( + metadata.files[0].fields.as_ref(), + expected_fields.as_slice() + ); + assert_eq!(collect_lineage(&dataset).await, changed); + dataset.validate().await.unwrap(); + let reopened = Dataset::open(uri).await.unwrap(); + assert_eq!(collect_lineage(&reopened).await, changed); + + // What compaction wrote holds no dead user column, so it plans no + // further rewrite. + let plan = plan_compaction(&dataset, &one_fragment()).await.unwrap(); + assert_eq!(plan.num_tasks(), 0); + } + + /// `FileFragment::validate` leaves a file without a schema field unopened + /// only when the fragment keeps it for a spilled sequence it carries, as + /// it keeps a lineage carrier whose user columns are gone. Any other file + /// whose user fields are all outside the schema is opened and reported, as + /// it always was, even when it lists a lineage id the fragment does not + /// spill. + #[rstest] + #[case::stale_user_file(vec![7], false)] + #[case::dead_lineage_carrier(vec![7, ROW_ID_FIELD_ID], true)] + #[case::unspilled_lineage_id(vec![7, ROW_ID_FIELD_ID], false)] + #[tokio::test] + async fn validate_skips_only_files_kept_for_spilled_lineage( + #[case] fields: Vec, + #[case] spills_row_ids: bool, + ) { + let dir = TempStrDir::default(); + let dataset = tiny_dataset(dir.as_str()).await; + let mut fragment = dataset.manifest.fragments[0].clone(); + // A second entry for the fragment's data file, under field ids the + // schema does not have. + let mut extra = fragment.files[0].clone(); + extra.column_indices = (0..fields.len() as i32).collect(); + extra.fields = fields.into(); + fragment.files.push(extra); + if spills_row_ids { + fragment.row_id_meta = Some(RowIdMeta::Column); + } + + let result = FileFragment::new(Arc::new(dataset), fragment) + .validate() + .await; + if spills_row_ids { + result.unwrap(); + } else { + let error = result.unwrap_err(); + assert!(matches!(error, Error::CorruptFile { .. }), "{error}"); + assert!( + error + .to_string() + .contains("did not have any fields in common with the dataset schema"), + "{error}" + ); + } + } } From ea175d58e9270a59c683d38a30bfe64d9d3a1348 Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Tue, 29 Sep 2026 06:21:27 +0800 Subject: [PATCH 6/6] test(dataset): cover the in-file lineage layout across layouts and follow-ups Compaction now writes spilled lineage into the data file of each fragment it writes, but the tests only covered one output fragment with no deleted rows, never compacted such a fragment again, and the cleanup test had quietly stopped covering a lineage-only file. - compaction_spills_and_reads_back_row_lineage adds a task that writes two fragments, so the hidden columns have to break where the writer ends a file, and a compaction of deleted rows. Every case checks that each fragment's sequences hold its physical rows, and takes sampled row ids through the row id index after a cold reopen. Four appends of 250 under a 300-row target plan two single-output tasks, so the two-output case compacts three appends. - compact_twice_reads_back_in_file_lineage compacts a fragment whose lineage is in its data file: the scan must not return the hidden columns the task appends itself. - cleanup_keeps_a_live_spilled_file runs for the in-file layout and for binary copy's lineage-only file, and pins each layout before cleaning up. - schema_change_keeps_spilled_lineage also casts a column after its binary-copy compaction, and points to the test covering a data file that carries the lineage next to user columns. Co-Authored-By: Claude Opus 5.5 --- rust/lance/src/dataset/rowids/spill.rs | 313 ++++++++++++++++++------- 1 file changed, 230 insertions(+), 83 deletions(-) diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs index 0d23a9bd7f4..e9e5fff33a2 100644 --- a/rust/lance/src/dataset/rowids/spill.rs +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -750,6 +750,8 @@ mod tests { ColumnAlteration, NewColumnTransform, UpdateBuilder, WriteMode, WriteParams, }; use arrow_array::builder::{ListBuilder, StringBuilder}; + use arrow_array::cast::AsArray; + use arrow_array::types::Int32Type; use arrow_array::{Int32Array, RecordBatchIterator, StringArray, StructArray}; use arrow_schema::Field; use chrono::Utc; @@ -1202,81 +1204,132 @@ mod tests { ) } - /// Compaction writes the spilled columns after the user columns of its - /// output file. A column index counts physical columns, so after a + /// The table [`compaction_spills_and_reads_back_row_lineage`] compacts, + /// and the fragments the compaction writes from it. + #[derive(Clone, Copy, Debug)] + enum CompactionShape { + /// Four appends of 250 rows, compacted into one fragment. + OneOutput, + /// Three appends of 250 rows under a 300-row target. They make one + /// task, which writes them as two fragments of 375 rows, so the + /// hidden columns have to break at the row the writer ends a file. + TwoOutputs, + /// Four appends of 250 rows with every seventh row deleted, compacted + /// into one fragment of the 857 left. A deleted row leaves the lineage + /// sequences at the offset it leaves the data. + WithDeletions, + } + + /// Compaction writes the spilled columns after the user columns of each + /// file it writes. A column index counts physical columns, so after a /// struct, or a list in 2.0, a lineage column's index no longer matches /// its position among the file's top-level fields. #[rstest] - #[case::flat(UserColumns::KeyOnly, None)] - #[case::nested_v2_2(UserColumns::WithStruct, Some(LanceFileVersion::V2_2))] - #[case::list_v2_0(UserColumns::WithList, Some(LanceFileVersion::V2_0))] + #[case::flat(UserColumns::KeyOnly, None, CompactionShape::OneOutput)] + #[case::nested_v2_2( + UserColumns::WithStruct, + Some(LanceFileVersion::V2_2), + CompactionShape::OneOutput + )] + #[case::list_v2_0( + UserColumns::WithList, + Some(LanceFileVersion::V2_0), + CompactionShape::OneOutput + )] + #[case::multi_output(UserColumns::KeyOnly, None, CompactionShape::TwoOutputs)] + #[case::with_deletions(UserColumns::KeyOnly, None, CompactionShape::WithDeletions)] #[tokio::test] async fn compaction_spills_and_reads_back_row_lineage( #[case] columns: UserColumns, #[case] version: Option, + #[case] shape: CompactionShape, ) { + let (appends, target_rows_per_fragment, output_rows) = match shape { + CompactionShape::OneOutput => (4, 1_000, vec![1_000]), + CompactionShape::TwoOutputs => (3, 300, vec![375, 375]), + CompactionShape::WithDeletions => (4, 1_000, vec![857]), + }; let dir = TempStrDir::default(); let uri = dir.as_str(); - let mut dataset = appended_dataset_with(uri, 4, 250, columns, version).await; + let mut dataset = appended_dataset_with(uri, appends, 250, columns, version).await; spill_everything(&mut dataset).await; - // Four appends at four versions, so the compacted created-at sequence - // has four runs rather than one. - let before = collect_lineage(&dataset).await; - assert_eq!( - before - .1 - .iter() - .collect::>() - .len(), - 4 - ); + if matches!(shape, CompactionShape::WithDeletions) { + dataset.delete("i % 7 = 0").await.unwrap(); + } + // One version per append, so the compacted created-at sequence has a + // run per append rather than one. + let before = collect_rows(&dataset).await; + let created_at_versions = before + .iter() + .map(|(_, _, created_at, _)| *created_at) + .collect::>(); + assert_eq!(created_at_versions.len(), appends as usize); - compact_files(&mut dataset, one_fragment(), None) - .await - .unwrap(); + compact_files( + &mut dataset, + CompactionOptions { + target_rows_per_fragment, + ..Default::default() + }, + None, + ) + .await + .unwrap(); let fragments = dataset.get_fragments(); - assert_eq!(fragments.len(), 1); - let metadata = fragments[0].metadata(); - assert!( - matches!(metadata.row_id_meta, Some(RowIdMeta::Column)) - && matches!( - metadata.created_at_version_meta, - Some(RowDatasetVersionMeta::Column) - ) - && matches!( - metadata.last_updated_at_version_meta, - Some(RowDatasetVersionMeta::Column) - ), - "compaction must spill every sequence under a zero inline budget, got {metadata:?}" - ); - // The three columns ride in the fragment's own data file, after the - // user columns, so the fragment has no extra file to reference. Their - // column indices continue from the last physical user column. - assert_eq!(metadata.files.len(), 1); - let data_file = &metadata.files[0]; - let (user_fields, lineage_fields) = data_file.fields.split_at(data_file.fields.len() - 3); - assert!( - user_fields.iter().all(|field| *field >= 0), - "unexpected data file fields {:?}", - data_file.fields - ); - assert_eq!( - lineage_fields, - [ - ROW_ID_FIELD_ID, - ROW_CREATED_AT_VERSION_FIELD_ID, - ROW_LAST_UPDATED_AT_VERSION_FIELD_ID - ] - ); - let (user_columns, lineage_columns) = data_file.column_indices.split_at(user_fields.len()); - let next = user_columns.iter().filter(|column| **column >= 0).count() as i32; - assert_eq!(lineage_columns, [next, next + 1, next + 2]); - for field_id in data_file.fields.iter().filter(|id| **id < 0) { + let written_rows = fragments + .iter() + .map(|fragment| fragment.metadata().physical_rows.unwrap()) + .collect::>(); + assert_eq!(written_rows, output_rows); + for fragment in &fragments { + let metadata = fragment.metadata(); + assert!( + matches!(metadata.row_id_meta, Some(RowIdMeta::Column)) + && matches!( + metadata.created_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ) + && matches!( + metadata.last_updated_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ), + "compaction must spill every sequence under a zero budget, got {metadata:?}" + ); + // The three columns ride in the fragment's own data file, after + // the user columns, so the fragment has no extra file to + // reference. Their column indices continue from the last physical + // user column. + assert_eq!(metadata.files.len(), 1); + let data_file = &metadata.files[0]; + let (user_fields, lineage_fields) = + data_file.fields.split_at(data_file.fields.len() - 3); + assert!( + user_fields.iter().all(|field| *field >= 0), + "unexpected data file fields {:?}", + data_file.fields + ); assert_eq!( - metadata.row_lineage_file(*field_id).unwrap(), - Some(data_file) + lineage_fields, + [ + ROW_ID_FIELD_ID, + ROW_CREATED_AT_VERSION_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID + ] ); + let (user_columns, lineage_columns) = + data_file.column_indices.split_at(user_fields.len()); + let next = user_columns.iter().filter(|column| **column >= 0).count() as i32; + assert_eq!(lineage_columns, [next, next + 1, next + 2]); + for field_id in data_file.fields.iter().filter(|id| **id < 0) { + assert_eq!( + metadata.row_lineage_file(*field_id).unwrap(), + Some(data_file) + ); + } + // Each fragment's columns hold its own rows and no others. + let row_ids = load_row_id_sequence(&dataset, metadata).await.unwrap(); + assert_eq!(Some(row_ids.len() as usize), metadata.physical_rows); } assert_ne!( dataset.manifest.reader_feature_flags & FLAG_UNSTABLE_SPILLED_ROW_LINEAGE, @@ -1291,7 +1344,7 @@ mod tests { // The lineage survives the rewrite and is still served through the // ordinary scan path, now from the data file columns. - assert_eq!(collect_lineage(&dataset).await, before); + assert_eq!(collect_rows(&dataset).await, before); // `validate_stable_row_ids` reads every fragment's sequences back and // checks them against the fragment length, so this covers the loaders // independently of the scan. @@ -1299,18 +1352,39 @@ mod tests { // Re-opened cold, so nothing is served from this process's caches. let reopened = Dataset::open(uri).await.unwrap(); - assert_eq!(collect_lineage(&reopened).await, before); - let fragment = &reopened.get_fragments()[0]; - let row_ids = load_row_id_sequence(&reopened, fragment.metadata()) - .await - .unwrap(); - assert_eq!(row_ids.iter().collect::>(), before.0); - let created_at = - load_row_version_sequence(&reopened, fragment.metadata(), RowVersionKind::CreatedAt) + assert_eq!(collect_rows(&reopened).await, before); + let mut row_ids = Vec::with_capacity(before.len()); + let mut created_at = Vec::with_capacity(before.len()); + for fragment in reopened.get_fragments() { + let metadata = fragment.metadata(); + let row_id_sequence = load_row_id_sequence(&reopened, metadata).await.unwrap(); + row_ids.extend(row_id_sequence.iter()); + let kind = RowVersionKind::CreatedAt; + let created_at_sequence = load_row_version_sequence(&reopened, metadata, kind) .await .unwrap() .expect("a compacted fragment carries created-at versions"); - assert_eq!(versions_of(&created_at), before.1); + created_at.extend(created_at_sequence.versions()); + } + let expected_row_ids = before.iter().map(|(_, row_id, _, _)| *row_id); + assert_eq!(row_ids, expected_row_ids.collect::>()); + let expected_created_at = before.iter().map(|(_, _, created, _)| *created); + assert_eq!(created_at, expected_created_at.collect::>()); + + // A take by row id resolves each id through the index built from the + // loaded sequences, so a sequence cut at the wrong row returns the + // wrong key. + let (sample_keys, sample_ids): (Vec, Vec) = before + .iter() + .step_by(61) + .map(|(key, row_id, _, _)| (*key, *row_id)) + .unzip(); + let taken = reopened + .take_rows(&sample_ids, reopened.schema().project(&["i"]).unwrap()) + .await + .unwrap(); + let taken_keys = taken["i"].as_primitive::().values().to_vec(); + assert_eq!(taken_keys, sample_keys); } /// Binary-copy compaction copies the input files page by page and cannot @@ -1362,30 +1436,100 @@ mod tests { assert_eq!(plan.num_tasks(), 0); } - /// Cleanup decides what to delete by walking - /// [`Fragment::referenced_lance_files`], so a spilled sequence has to be - /// reachable from there. If it were not, an ordinary cleanup would delete a - /// live file and leave the fragment claiming row ids it can no longer read. + /// A second compaction reads the lineage back from the columns the first + /// one wrote into the data file. Its scan must not return those columns + /// with the user data: the task appends the lineage columns itself, and a + /// second copy would clash by name. #[tokio::test] - async fn cleanup_keeps_a_live_spilled_file() { + async fn compact_twice_reads_back_in_file_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 = collect_lineage(&dataset).await; - compact_files(&mut dataset, one_fragment(), None) .await .unwrap(); + // Deleting every seventh row makes the compacted fragment compact + // again on its own, which masks the sequences read back from its file. + dataset.delete("i % 7 = 0").await.unwrap(); + let before = collect_rows(&dataset).await; + + let metrics = compact_files(&mut dataset, one_fragment(), None) + .await + .unwrap(); + assert_eq!(metrics.fragments_removed, 1); + + let fragments = dataset.get_fragments(); + assert_eq!(fragments.len(), 1); + let metadata = fragments[0].metadata(); + assert_eq!(metadata.physical_rows, Some(before.len())); + assert_eq!(metadata.files.len(), 1, "{metadata:?}"); + assert_eq!( + metadata.files[0].fields.as_ref(), + [ + 0, + ROW_ID_FIELD_ID, + ROW_CREATED_AT_VERSION_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID + ] + ); + assert_eq!(collect_rows(&dataset).await, before); + dataset.validate().await.unwrap(); + + let reopened = Dataset::open(uri).await.unwrap(); + assert_eq!(collect_rows(&reopened).await, before); + } + + /// Cleanup decides what to delete by walking + /// [`Fragment::referenced_lance_files`], so the file carrying a spilled + /// sequence has to be reachable from there. If it were not, an ordinary + /// cleanup would delete a live file and leave the fragment claiming row + /// ids it can no longer read. A reencoding compaction writes the lineage + /// into the fragment's data file, which cleanup keeps for the user columns + /// anyway; binary copy writes a file holding nothing but lineage, which + /// only the lineage keeps. + #[rstest] + #[case::in_file(CompactionMode::Reencode)] + #[case::lineage_file(CompactionMode::ForceBinaryCopy)] + #[tokio::test] + async fn cleanup_keeps_a_live_spilled_file(#[case] mode: CompactionMode) { + let dir = TempStrDir::default(); + let uri = dir.as_str(); + let mut dataset = appended_dataset(uri, 4, 250).await; + spill_everything(&mut dataset).await; + let before = collect_lineage(&dataset).await; + + compact_files( + &mut dataset, + CompactionOptions { + compaction_mode: Some(mode), + ..one_fragment() + }, + None, + ) + .await + .unwrap(); - let spilled = dataset.get_fragments()[0] - .metadata() + // Pin the layout, so each case keeps covering the file it is named + // after. + let fragments = dataset.get_fragments(); + let metadata = fragments[0].metadata(); + let carrier = metadata .row_lineage_file(ROW_ID_FIELD_ID) .unwrap() - .expect("compaction must spill under a zero inline budget") - .path - .clone(); - let on_disk = std::path::Path::new(uri).join("data").join(&spilled); + .expect("compaction must spill under a zero inline budget"); + if matches!(mode, CompactionMode::ForceBinaryCopy) { + assert_eq!(metadata.files.len(), 2, "{metadata:?}"); + assert_eq!(carrier, &metadata.files[1]); + assert!( + carrier.fields.iter().all(|field| *field < 0), + "binary copy must write a lineage-only file: {metadata:?}" + ); + } else { + assert_eq!(metadata.files.len(), 1, "{metadata:?}"); + assert_eq!(carrier, &metadata.files[0]); + } + let on_disk = std::path::Path::new(uri).join("data").join(&carrier.path); assert!(on_disk.exists(), "no spilled file written at {on_disk:?}"); // Everything written so far is older than this instant, so the @@ -1843,12 +1987,15 @@ mod tests { /// 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. + /// keep the file carrying the lineage, which is its only copy. A data file + /// holding the lineage next to user columns the change removes is covered + /// by [`dropping_every_column_of_a_lineage_carrier_keeps_lineage_until_compaction`]. #[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)] + #[case::compact_then_cast(SpillingWrite::Compaction, SchemaChange::Cast)] #[tokio::test] async fn schema_change_keeps_spilled_lineage( #[case] spilling_write: SpillingWrite,