diff --git a/rust/lance-core/src/datatypes/schema.rs b/rust/lance-core/src/datatypes/schema.rs index 2e328c88d24..0a132875ff5 100644 --- a/rust/lance-core/src/datatypes/schema.rs +++ b/rust/lance-core/src/datatypes/schema.rs @@ -347,10 +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() { - if field.id < 0 { + if field.id < 0 && !row_lineage_ids.contains(&field.id) { return Err(Error::schema(format!( "Field {} has a negative id {}", field.name, field.id @@ -1920,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-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/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 53d7ea6a4a1..55e930bc160 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, 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, @@ -133,7 +134,12 @@ 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}; use tracing::{info, warn}; @@ -883,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; @@ -908,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 { @@ -1109,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`). /// @@ -2569,7 +2631,7 @@ async fn rewrite_files( params.enable_stable_row_ids = true; } - if can_binary_copy { + let lineage_plan = if can_binary_copy { new_fragments = versions::rewrite_files_binary_copy( write_version, dataset.as_ref(), @@ -2603,21 +2665,64 @@ async fn rewrite_files( let _ = tx.send(captured); row_ids_rx = Some(rx); } + None } else { + let mut stream = reader.expect("reader must be prepared for non-binary-copy path"); + let mut write_schema = dataset.schema().clone(); + // 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, - dataset.schema().clone(), - reader.expect("reader must be prepared for non-binary-copy path"), + write_schema, + stream, params, None, Some(file_row_counts), ) .await?; new_fragments = frags; - } + lineage_plan + }; log::info!("Compaction task {}: file written", task_id); @@ -2639,7 +2744,16 @@ 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 { + 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) + .await? + } + } } Ok(None) } @@ -2679,14 +2793,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 +2843,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 +2864,289 @@ 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 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], + 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(()) +} + +/// 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: Vec, + planned_rows: &[u64], + spill: RowLineageSpill, +) -> Result<()> { + 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 {written_total} rows, in fragments of {written_rows:?} rows, \ + but planned row lineage for {planned_total} rows, in fragments of \ + {planned_rows:?} rows" + ))); + } + + 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 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 has {len} inline {sequence} for fragment {} of {rows} physical rows", + fragment.id + ))); + } + } + Ok(lineages) +} + +/// 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, + fields: &[LanceField], + columns: Vec>, +) -> SendableRecordBatchStream { + 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| { + 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(); + // 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) + .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, @@ -3295,6 +3666,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}; @@ -3311,6 +3686,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; @@ -3324,7 +3700,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; @@ -4517,6 +4895,320 @@ 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::>() + ); + } + + /// 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. + #[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 4a77b45757a..6474d39df7b 100644 --- a/rust/lance/src/dataset/rowids.rs +++ b/rust/lance/src/dataset/rowids.rs @@ -30,6 +30,7 @@ pub use spill::{ 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 7c2168e43e2..e9e5fff33a2 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; @@ -103,6 +104,183 @@ pub fn inline_row_lineage_max_bytes(dataset: &Dataset) -> Result> }) } +/// 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, + 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 spilled columns, in column order. + pub fn field_ids(&self) -> impl Iterator { + 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)) + } + + /// 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 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() + } + + /// 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, + } + } +} + +/// The hidden column a spilled lineage sequence is stored in: a non-nullable +/// `UInt64` named after the sequence in [`LINEAGE_COLUMNS`], under its +/// reserved `field_id`. +fn lineage_field(field_id: i32) -> Result { + 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) +} + +/// 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() + .map(|l| over(write_dataset_versions(&l.created_at))), + ); + let last_updated_at = over_in_all_or_none( + lineages + .iter() + .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. pub struct RowLineage { pub row_ids: RowIdSequence, @@ -379,7 +557,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). @@ -416,6 +594,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 @@ -429,46 +622,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() @@ -510,16 +740,26 @@ 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::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 arrow_array::{Int32Array, RecordBatchIterator}; + use crate::dataset::transaction::Operation; + use crate::dataset::{ + 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; use lance_core::utils::tempfile::TempStrDir; use lance_core::{ROW_CREATED_AT_VERSION, ROW_ID, ROW_LAST_UPDATED_AT_VERSION}; use lance_file::version::LanceFileVersion; use lance_table::feature_flags::FLAG_UNSTABLE_SPILLED_ROW_LINEAGE; + use lance_table::format::overlay::TOMBSTONE_FIELD_ID; use rstest::rstest; /// A sequence with no runs to exploit, which is what a globally shuffled @@ -625,6 +865,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] @@ -774,6 +1026,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. @@ -786,62 +1076,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 { @@ -888,60 +1204,132 @@ mod tests { ) } + /// 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, 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() { + 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(uri, 4, 250).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 - ); - - compact_files(&mut dataset, one_fragment(), None) - .await - .unwrap(); + 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, + 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 sequences share one lineage file, listed after the user - // data file among the fragment's files and found by field id. - assert_eq!(metadata.files.len(), 2); - let lineage_file = &metadata.files[1]; - assert_eq!( - lineage_file.fields.as_ref(), - [ - ROW_ID_FIELD_ID, - ROW_CREATED_AT_VERSION_FIELD_ID, - ROW_LAST_UPDATED_AT_VERSION_FIELD_ID - ] - ); - for field_id in lineage_file.fields.iter() { + 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(lineage_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, @@ -956,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. @@ -964,44 +1352,184 @@ 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); } - /// 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. + /// 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 cleanup_keeps_a_live_spilled_file() { + 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(); + + // 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); + } + + /// 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 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; 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 spilled = dataset.get_fragments()[0] - .metadata() + 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(); + + // 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 @@ -1230,7 +1758,6 @@ mod tests { use crate::dataset::{ MergeInsertBuilder, MergeInsertWriteMode, WhenMatched, WhenNotMatched, }; - use arrow_array::StringArray; let dir = TempStrDir::default(); let uri = dir.as_str(); @@ -1443,7 +1970,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, } @@ -1458,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, @@ -1471,7 +2003,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 => { @@ -1488,18 +2020,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()); @@ -1526,4 +2072,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}" + ); + } + } } diff --git a/rust/lance/src/dataset/versions/mod.rs b/rust/lance/src/dataset/versions/mod.rs index 051e933a6bf..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,12 +158,13 @@ pub async fn write_fragments( _ => normalized_schema, }; let version_name = format!("{version:?}"); - let schema = write::prepare_write_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 +190,41 @@ pub async fn write_fragments( Ok((fragments, schema)) } +/// 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) == 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)] pub async fn write_fragments_direct( version: ConcreteFileVersion, @@ -1086,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]); + } +}