diff --git a/rust/lance/src/dataset/fragment.rs b/rust/lance/src/dataset/fragment.rs index 017381d1e73..9e3dd0e89d1 100644 --- a/rust/lance/src/dataset/fragment.rs +++ b/rust/lance/src/dataset/fragment.rs @@ -52,7 +52,7 @@ use lance_io::object_store::ObjectStore; use lance_io::scheduler::{FileScheduler, ScanScheduler, SchedulerConfig}; use lance_io::stream::RecordBatchStream; use lance_io::utils::CachedFileSize; -use lance_table::format::{DataFile, DeletionFile, Fragment, RowDatasetVersionMeta}; +use lance_table::format::{DataFile, DeletionFile, Fragment}; use lance_table::io::deletion::{deletion_file_path, write_deletion_file}; use lance_table::rowids::RowIdSequence; use lance_table::utils::stream::{ @@ -65,7 +65,7 @@ use roaring::RoaringBitmap; use self::write::FragmentCreateBuilder; use super::hash_joiner::HashJoiner; -use super::rowids::load_row_id_sequence; +use super::rowids::{RowVersionKind, load_row_id_sequence, load_row_version_sequence}; use super::scanner::Scanner; use super::updater::Updater; @@ -1155,33 +1155,44 @@ impl FileFragment { futures::future::Either::Right(futures::future::ready(Ok(None))) }; - // The reader builders below decode version sequences from the manifest - // and fall back to version 1 when they cannot; a spilled sequence must - // not fall through to that. - for (wanted, meta) in [ - ( - read_config.with_row_created_at_version, - &self.metadata.created_at_version_meta, - ), - ( - read_config.with_row_last_updated_at_version, - &self.metadata.last_updated_at_version_meta, - ), - ] { - if wanted && matches!(meta, Some(RowDatasetVersionMeta::Column)) { - return Err(Error::not_supported(format!( - "row versions of fragment {} are spilled to a data file column, which \ - this build cannot read", - self.id() - ))); + let version_load = |kind: RowVersionKind, wanted: bool| { + if wanted { + futures::future::Either::Left(load_row_version_sequence( + &self.dataset, + &self.metadata, + kind, + )) + } else { + futures::future::Either::Right(futures::future::ready(Ok(None))) } - } + }; + let last_updated_at_load = version_load( + RowVersionKind::LastUpdatedAt, + read_config.with_row_last_updated_at_version, + ); + let created_at_load = version_load( + RowVersionKind::CreatedAt, + read_config.with_row_created_at_version, + ); - let (opened_files, deletion_vec, row_id_sequence) = - join!(open_files, deletion_vec_load, row_id_load); + let ( + opened_files, + deletion_vec, + row_id_sequence, + last_updated_at_sequence, + created_at_sequence, + ) = join!( + open_files, + deletion_vec_load, + row_id_load, + last_updated_at_load, + created_at_load + ); let opened_files = opened_files?; let deletion_vec = deletion_vec?; let row_id_sequence = row_id_sequence?; + let last_updated_at_sequence = last_updated_at_sequence?; + let created_at_sequence = created_at_sequence?; if opened_files.is_empty() && !read_config.has_system_cols() { return Err(Error::not_found(format!( @@ -1224,10 +1235,10 @@ impl FileFragment { reader.with_row_address(); } if read_config.with_row_last_updated_at_version { - reader.with_row_last_updated_at_version(); + reader.with_row_last_updated_at_version(last_updated_at_sequence); } if read_config.with_row_created_at_version { - reader.with_row_created_at_version(); + reader.with_row_created_at_version(created_at_sequence); } Ok(reader) @@ -1793,7 +1804,15 @@ impl FileFragment { data_file.validate(&self.dataset.data_file_dir(data_file)?)?; } - let get_lengths = self.metadata.files.iter().map(|data_file| async move { + // 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)); + 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 .open_reader(data_file, None, &FragReadConfig::default()) @@ -1814,7 +1833,7 @@ impl FileFragment { let get_lengths = get_lengths?; let expected_length = get_lengths.first().unwrap_or(&0); - for (length, data_file) in get_lengths.iter().zip(self.metadata.files.iter()) { + for (length, data_file) in get_lengths.iter().zip(user_data_files) { if length != expected_length { let path = self .dataset @@ -3447,17 +3466,15 @@ impl FragmentReader { self } - pub(crate) fn with_row_last_updated_at_version(&mut self) -> &mut Self { + /// Emit the `_row_last_updated_at_version` column, served from `sequence`; + /// `None` means the fragment has no version metadata and every row reads + /// as version 1. + pub(crate) fn with_row_last_updated_at_version( + &mut self, + sequence: Option>, + ) -> &mut Self { self.with_row_last_updated_at_version = true; - - // Load the version sequence if not already loaded - if self.last_updated_at_sequence.is_none() - && let Some(meta) = &self.fragment.last_updated_at_version_meta - && let Ok(sequence) = meta.load_sequence() - { - self.last_updated_at_sequence = Some(Arc::new(sequence)); - } - // If no metadata or load fails, sequence remains None (will default to version 1) + self.last_updated_at_sequence = sequence; // Add the version column to the output schema self.output_schema = self @@ -3468,17 +3485,15 @@ impl FragmentReader { self } - pub(crate) fn with_row_created_at_version(&mut self) -> &mut Self { + /// Emit the `_row_created_at_version` column, served from `sequence`; + /// `None` means the fragment has no version metadata and every row reads + /// as version 1. + pub(crate) fn with_row_created_at_version( + &mut self, + sequence: Option>, + ) -> &mut Self { self.with_row_created_at_version = true; - - // Load the version sequence if not already loaded - if self.created_at_sequence.is_none() - && let Some(meta) = &self.fragment.created_at_version_meta - && let Ok(sequence) = meta.load_sequence() - { - self.created_at_sequence = Some(Arc::new(sequence)); - } - // If no metadata or load fails, sequence remains None (will default to version 1) + self.created_at_sequence = sequence; // Add the version column to the output schema self.output_schema = self diff --git a/rust/lance/src/dataset/rowids.rs b/rust/lance/src/dataset/rowids.rs index ea9035746ae..71718232909 100644 --- a/rust/lance/src/dataset/rowids.rs +++ b/rust/lance/src/dataset/rowids.rs @@ -1,21 +1,31 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright The Lance Authors +mod spill; mod validate; use super::Dataset; use crate::io::deletion::read_dataset_deletion_file; -use crate::session::caches::{RowIdIndexKey, RowIdSequenceKey}; +use crate::session::caches::{RowIdIndexKey, RowIdSequenceKey, RowVersionSequenceKey}; use crate::{Error, Result}; use futures::{Stream, StreamExt, TryFutureExt, TryStreamExt}; use lance_core::utils::{address::RowAddress, deletion::DeletionVector}; use lance_select::{RowAddrSelection, RowAddrTreeMap}; use lance_table::{ - format::{Fragment, ROW_ID_FIELD_ID, RowIdMeta}, + format::{ + Fragment, ROW_CREATED_AT_VERSION_FIELD_ID, ROW_ID_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, RowDatasetVersionMeta, RowDatasetVersionSequence, + RowIdMeta, + }, rowids::{FragmentRowIdIndex, RowIdIndex, RowIdSequence, read_row_ids}, }; use std::sync::Arc; +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, +}; pub(super) use validate::validate_stable_row_ids; /// Load a row id sequence from the given dataset and fragment. @@ -33,20 +43,87 @@ pub async fn load_row_id_sequence( }; dataset .metadata_cache - .get_or_insert_with_key(key, || read_row_id_sequence(fragment)) + .get_or_insert_with_key(key, || read_row_id_sequence(dataset, fragment)) .await } /// Decode the row id sequence of `fragment`, bypassing every cache. -async fn read_row_id_sequence(fragment: &Fragment) -> Result { +async fn read_row_id_sequence(dataset: &Dataset, fragment: &Fragment) -> Result { match &fragment.row_id_meta { None => Err(Error::internal("Missing row id meta")), Some(RowIdMeta::Inline(data)) => read_row_ids(data), - Some(RowIdMeta::Column) => Err(Error::not_supported(format!( - "row ids of fragment {} are spilled to a data file column, which this build \ - cannot read", - fragment.id - ))), + Some(RowIdMeta::Column) => spill::read_spilled_row_ids(dataset, fragment).await, + } +} + +/// Which of a fragment's two per-row version sequences is meant. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RowVersionKind { + /// The dataset version each row first appeared at. + CreatedAt, + /// The dataset version each row was last written at. + LastUpdatedAt, +} + +impl RowVersionKind { + fn meta(self, fragment: &Fragment) -> Option<&RowDatasetVersionMeta> { + match self { + Self::CreatedAt => fragment.created_at_version_meta.as_ref(), + Self::LastUpdatedAt => fragment.last_updated_at_version_meta.as_ref(), + } + } + + /// The reserved field id of the hidden column a spilled sequence lives in. + pub fn field_id(self) -> i32 { + match self { + Self::CreatedAt => ROW_CREATED_AT_VERSION_FIELD_ID, + Self::LastUpdatedAt => ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, + } + } +} + +/// Load one of `fragment`'s per-row version sequences, wherever it is stored. +/// +/// `None` when the fragment carries no such metadata, which readers treat as +/// every row being at version 1. A sequence spilled to a data file is cached +/// per fragment and file; an inline one is decoded from the manifest bytes. +pub async fn load_row_version_sequence( + dataset: &Dataset, + fragment: &Fragment, + kind: RowVersionKind, +) -> Result>> { + let Some(meta) = kind.meta(fragment) else { + return Ok(None); + }; + match meta { + RowDatasetVersionMeta::Column => { + let data_file = fragment.row_lineage_file(kind.field_id())?.ok_or_else(|| { + Error::corrupt_file( + dataset.base.clone(), + format!( + "fragment {} marks its {kind:?} versions as spilled but none of its \ + data files carries field {}", + fragment.id, + kind.field_id() + ), + ) + })?; + let key = RowVersionSequenceKey { + fragment_id: fragment.id, + field_id: kind.field_id(), + data_file, + }; + dataset + .metadata_cache + .get_or_insert_with_key(key, || { + spill::read_spilled_versions(dataset, fragment, kind.field_id()) + }) + .await + .map(Some) + } + RowDatasetVersionMeta::Inline(_) => meta + .load_sequence() + .map(|sequence| Some(Arc::new(sequence))), } } @@ -252,7 +329,7 @@ async fn read_fragment_row_id_index( dataset: &Dataset, fragment: &Fragment, ) -> Result { - let row_id_sequence = Arc::new(read_row_id_sequence(fragment).await?); + let row_id_sequence = Arc::new(read_row_id_sequence(dataset, fragment).await?); let deletion_vector = match &fragment.deletion_file { None => Arc::new(DeletionVector::default()), Some(deletion_file) => { diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs new file mode 100644 index 00000000000..9fae144d33b --- /dev/null +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -0,0 +1,612 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Spilling a fragment's row lineage sequences -- its row ids and its +//! created-at and last-updated-at versions -- out of the manifest and into +//! hidden columns of a Lance data file. +//! +//! A row id sequence is run-encoded, so an appended fragment costs about 20 +//! bytes of manifest and never needs to leave it. A fragment whose rows came +//! from many places -- the output of compacting a shuffled table, for instance +//! -- has no runs to exploit and falls back to 8 bytes per row. The version +//! sequences are run-length encoded and degrade the same way once a fragment +//! interleaves rows written at many versions. Inline, that cost is paid again +//! in every manifest version, so the manifest grows with the table and every +//! commit rewrites all of it. +//! +//! Spilled, each sequence is an ordinary `UInt64` column carrying one of the +//! reserved field ids [`ROW_ID_FIELD_ID`], [`ROW_CREATED_AT_VERSION_FIELD_ID`] +//! or [`ROW_LAST_UPDATED_AT_VERSION_FIELD_ID`], written with the same encodings +//! and read with the same reader as user data. The file is one of the +//! fragment's own data files, found by that field id; the fragment's metadata +//! only marks the sequence as spilled. A fragment's spilled sequences written +//! here share one such file. + +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 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::version::ConcreteFileVersion; +use lance_file::versions; +use lance_file::writer::FileWriterOptions; +use lance_io::ReadBatchParams; +use lance_io::scheduler::{ScanScheduler, SchedulerConfig}; +use lance_table::format::{ + DataFile, Fragment, ROW_CREATED_AT_VERSION_FIELD_ID, ROW_ID_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, RowDatasetVersionMeta, RowDatasetVersionSequence, + RowIdMeta, +}; +use lance_table::rowids::version::write_dataset_versions; +use lance_table::rowids::{RowIdSequence, write_row_ids}; +use object_store::path::Path; + +use super::super::Dataset; +use crate::dataset::fragment::write::generate_random_filename; +use crate::{Error, Result}; + +/// Rows per batch handed to the file writer and read back from it. +const SPILL_BATCH_ROWS: usize = 64 * 1024; + +/// Encoded sequences at or below this size stay in the manifest, unless +/// [`INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY`] says otherwise. +/// +/// Matches the inline limit the format has always documented for the +/// lineage oneofs. A `Range` row id sequence -- every appended fragment -- +/// encodes to a few dozen bytes and is nowhere near it, and so does a +/// single-run version sequence. +pub const DEFAULT_INLINE_ROW_LINEAGE_MAX_BYTES: usize = 200 * 1024; + +/// Table config key that turns spilling on: `"true"` lets compaction move +/// oversized lineage sequences out of the manifest. Absent or anything else, +/// every sequence stays inline however large it grows, which is what every +/// released build does; a table that never sets it stays readable by them. +pub const SPILL_ROW_LINEAGE_CONFIG_KEY: &str = "lance.row_lineage.spill"; + +/// Table config key overriding [`DEFAULT_INLINE_ROW_LINEAGE_MAX_BYTES`], as a +/// byte count. Only consulted when [`SPILL_ROW_LINEAGE_CONFIG_KEY`] is on. +pub const INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY: &str = "lance.row_lineage.inline_max_bytes"; + +/// The largest encoded lineage sequence `dataset` keeps inline, or `None` when +/// the table does not spill at all. +/// +/// Spilling needs the table's opt-in, a build that understands the feature +/// flag (a build that does not would write a dataset it then refuses to +/// open), and v2 data files: the format allows the columns only there, since +/// a legacy v1 file has no `column_indices` to locate them by. +pub fn inline_row_lineage_max_bytes(dataset: &Dataset) -> Result> { + let config = dataset.config(); + let enabled = config + .get(SPILL_ROW_LINEAGE_CONFIG_KEY) + .is_some_and(|value| value.eq_ignore_ascii_case("true")); + if !enabled + || !lance_table::feature_flags::spilled_row_lineage_enabled() + || dataset.manifest.data_storage_format.lance_file_format() == ConcreteFileVersion::V1 + { + return Ok(None); + } + let Some(value) = config.get(INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY) else { + return Ok(Some(DEFAULT_INLINE_ROW_LINEAGE_MAX_BYTES)); + }; + value.parse().map(Some).map_err(|error| { + Error::invalid_input(format!( + "table config {INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY}={value:?} is not a byte \ + count: {error}" + )) + }) +} + +/// The per-row lineage of one fragment, in row offset order. +pub struct RowLineage { + pub row_ids: RowIdSequence, + pub created_at: RowDatasetVersionSequence, + pub last_updated_at: RowDatasetVersionSequence, +} + +/// Where each of a fragment's lineage sequences ended up, and the data file +/// the spilled ones were written to, which the fragment has to list among its +/// files. [`Self::apply`] does both. +pub struct PlacedRowLineage { + pub row_ids: RowIdMeta, + pub created_at: RowDatasetVersionMeta, + pub last_updated_at: RowDatasetVersionMeta, + pub file: Option, +} + +impl PlacedRowLineage { + /// Put the placement on `fragment`: its three lineage arms, and the + /// lineage file, if any, as one more of its data files. + pub fn apply(self, fragment: &mut Fragment) { + fragment.row_id_meta = Some(self.row_ids); + fragment.created_at_version_meta = Some(self.created_at); + fragment.last_updated_at_version_meta = Some(self.last_updated_at); + if let Some(file) = self.file { + fragment.files.push(file); + } + } +} + +/// Place each sequence of `lineage` either inline in the manifest or in a +/// hidden column of a new data file, as the table's spill policy (see +/// [`inline_row_lineage_max_bytes`]) and the sequence's encoded size call for. +/// Every sequence that spills goes into one file, which the caller adds to the +/// fragment's data files (see [`PlacedRowLineage::apply`]). +/// +/// Only correct for lineage that a commit conflict cannot change: row ids and +/// versions carried over from existing rows. Lineage assigned at commit time -- +/// an appended fragment's row ids, an inserted row's created-at version -- has +/// to stay inline, where the commit can still rewrite it. +pub async fn place_row_lineage( + dataset: &Dataset, + lineage: &RowLineage, +) -> Result { + let (can_spill, limit) = match inline_row_lineage_max_bytes(dataset)? { + Some(limit) => (true, limit), + None => (false, usize::MAX), + }; + let inline_row_ids = write_row_ids(&lineage.row_ids); + let inline_created_at = write_dataset_versions(&lineage.created_at); + let inline_last_updated_at = write_dataset_versions(&lineage.last_updated_at); + + // Materialized up front rather than streamed from the sequence iterators: + // `RowIdSequence::iter` returns a boxed `dyn DoubleEndedIterator`, which is + // not `Send`, so holding it across the write below would make this future + // non-`Send` and every caller of `compact_files` along with it -- + // including the Python bindings, which spawn that future. + let mut columns: Vec<(i32, &str, ArrayRef)> = Vec::with_capacity(3); + if can_spill && inline_row_ids.len() > limit { + let ids = UInt64Array::from(lineage.row_ids.iter().collect::>()); + columns.push((ROW_ID_FIELD_ID, ROW_ID, Arc::new(ids))); + } + if can_spill && inline_created_at.len() > limit { + let versions = UInt64Array::from(lineage.created_at.versions().collect::>()); + columns.push(( + ROW_CREATED_AT_VERSION_FIELD_ID, + ROW_CREATED_AT_VERSION, + Arc::new(versions), + )); + } + if can_spill && inline_last_updated_at.len() > limit { + let versions = UInt64Array::from(lineage.last_updated_at.versions().collect::>()); + columns.push(( + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION, + Arc::new(versions), + )); + } + + let spilled = if columns.is_empty() { + None + } else { + Some(write_lineage_file(dataset, &columns).await?) + }; + let holds = |field_id: i32| { + spilled + .as_ref() + .is_some_and(|file| file.fields.contains(&field_id)) + }; + + Ok(PlacedRowLineage { + row_ids: if holds(ROW_ID_FIELD_ID) { + RowIdMeta::Column + } else { + RowIdMeta::Inline(inline_row_ids.into()) + }, + created_at: if holds(ROW_CREATED_AT_VERSION_FIELD_ID) { + RowDatasetVersionMeta::Column + } else { + RowDatasetVersionMeta::Inline(inline_created_at.into()) + }, + last_updated_at: if holds(ROW_LAST_UPDATED_AT_VERSION_FIELD_ID) { + RowDatasetVersionMeta::Column + } else { + RowDatasetVersionMeta::Inline(inline_last_updated_at.into()) + }, + file: spilled, + }) +} + +/// Write `columns` as the hidden columns of one new data file and return the +/// [`DataFile`] that locates them, listing the columns' field ids in order. +async fn write_lineage_file( + dataset: &Dataset, + columns: &[(i32, &str, ArrayRef)], +) -> Result { + let file_version = dataset.manifest.data_storage_format.version; + let filename = format!("{}.lance", generate_random_filename()); + let full_path = dataset.data_dir().join(filename.as_str()); + + let arrow_schema = Arc::new(ArrowSchema::new( + columns + .iter() + .map(|(_, name, _)| ArrowField::new(*name, DataType::UInt64, false)) + .collect::>(), + )); + let schema = Schema::try_from(arrow_schema.as_ref())?; + let object_writer = dataset.object_store.create(&full_path).await?; + let mut writer = versions::create_writer( + file_version, + object_writer, + schema, + FileWriterOptions::default(), + )?; + + let num_rows = columns[0].2.len(); + for offset in (0..num_rows).step_by(SPILL_BATCH_ROWS) { + let len = SPILL_BATCH_ROWS.min(num_rows - offset); + let batch = RecordBatch::try_new( + arrow_schema.clone(), + columns + .iter() + .map(|(_, _, array)| array.slice(offset, len)) + .collect(), + )?; + writer.write_batch(&batch).await?; + } + let summary = writer.finish().await?; + + Ok(DataFile::new( + filename, + columns.iter().map(|(field_id, _, _)| *field_id).collect(), + (0..columns.len() as i32).collect(), + file_version, + std::num::NonZero::new(summary.size_bytes), + None, + )) +} + +/// Read back `fragment`'s row id sequence from the data file that carries it. +pub async fn read_spilled_row_ids(dataset: &Dataset, fragment: &Fragment) -> Result { + let ids = read_spilled_column(dataset, fragment, ROW_ID_FIELD_ID) + .boxed() + .await?; + Ok(RowIdSequence::from(ids.as_slice())) +} + +/// Read back one of `fragment`'s version sequences from the data file that +/// carries it; `field_id` says which of the two it is. +pub async fn read_spilled_versions( + dataset: &Dataset, + fragment: &Fragment, + field_id: i32, +) -> Result { + let versions = read_spilled_column(dataset, fragment, field_id) + .boxed() + .await?; + Ok(RowDatasetVersionSequence::from_versions(&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 +/// 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). +async fn read_spilled_column( + dataset: &Dataset, + fragment: &Fragment, + field_id: i32, +) -> Result> { + let data_file = fragment.row_lineage_file(field_id)?.ok_or_else(|| { + Error::corrupt_file( + dataset.base.clone(), + format!( + "fragment {} marks row lineage field {field_id} as spilled but none of its \ + data files carries it", + fragment.id + ), + ) + })?; + let column_index = data_file + .fields + .iter() + .position(|field| *field == field_id) + .and_then(|position| data_file.column_indices.get(position)) + .ok_or_else(|| { + Error::corrupt_file_named( + &data_file.path, + format!("spilled row lineage file does not carry field id {field_id}"), + ) + })?; + let column_index = u32::try_from(*column_index).map_err(|_| { + Error::corrupt_file_named( + &data_file.path, + format!("field id {field_id} has no column index in the spilled row lineage file"), + ) + })?; + + // 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 + .data_file_dir(data_file)? + .join(data_file.path.as_str()); + let object_store = dataset.object_store_for_data_file(data_file).await?; + let scheduler = ScanScheduler::new( + object_store.clone(), + SchedulerConfig::max_bandwidth(&object_store), + ); + 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( + &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], + }; + + let mut values: Vec = Vec::with_capacity(reader.num_rows() as usize); + let mut stream = reader + .read_stream_projected( + ReadBatchParams::RangeFull, + SPILL_BATCH_ROWS as u32, + 8, + projection, + FilterExpression::no_filter(), + ) + .await?; + while let Some(batch) = stream.try_next().await? { + let column = batch + .column(0) + .as_any() + .downcast_ref::() + .ok_or_else(|| { + Error::corrupt_file_named( + &data_file.path, + format!("spilled row lineage column {field_id} is not UInt64"), + ) + })?; + // A null has no row id or version to stand for; the format requires a + // value for every physical row. + if column.null_count() > 0 { + return Err(Error::corrupt_file_named( + &data_file.path, + format!("spilled row lineage column {field_id} holds nulls"), + )); + } + values.extend_from_slice(column.values()); + } + if let Some(physical_rows) = fragment.physical_rows + && values.len() != physical_rows + { + return Err(Error::corrupt_file_named( + &data_file.path, + format!( + "spilled row lineage column {field_id} holds {} values for a fragment of {} \ + physical rows", + values.len(), + physical_rows + ), + )); + } + + Ok(values) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::dataset::WriteParams; + use arrow_array::{Int32Array, RecordBatchIterator}; + use arrow_schema::Field; + use lance_core::utils::tempfile::TempStrDir; + use lance_file::version::LanceFileVersion; + + /// A sequence with no runs to exploit, which is what a globally shuffled + /// table produces and what forces the spill path. + fn scattered_row_ids(len: u64) -> RowIdSequence { + // A stride coprime with `len` visits every id exactly once in an order + // with no ascending run longer than one. + let ids: Vec = (0..len).map(|i| (i * 7919) % len).collect(); + RowIdSequence::from(ids.as_slice()) + } + + /// A version per row that alternates, so every row is its own run. + fn alternating_versions(len: u64, first: u64) -> RowDatasetVersionSequence { + let versions: Vec = (0..len).map(|i| first + i % 2).collect(); + RowDatasetVersionSequence::from_versions(&versions) + } + + fn test_schema() -> Arc { + Arc::new(ArrowSchema::new(vec![Field::new( + "i", + DataType::Int32, + false, + )])) + } + + async fn tiny_dataset(uri: &str) -> Dataset { + tiny_dataset_with_version(uri, LanceFileVersion::default()).await + } + + async fn tiny_dataset_with_version(uri: &str, version: LanceFileVersion) -> Dataset { + let schema = test_schema(); + let batch = RecordBatch::try_new( + schema.clone(), + vec![Arc::new(Int32Array::from(vec![1, 2, 3, 4]))], + ) + .unwrap(); + let reader = RecordBatchIterator::new(vec![Ok(batch)], schema); + Dataset::write( + reader, + uri, + Some(WriteParams { + enable_stable_row_ids: true, + data_storage_version: Some(version), + ..Default::default() + }), + ) + .await + .unwrap() + } + + fn versions_of(sequence: &RowDatasetVersionSequence) -> Vec { + sequence.versions().collect() + } + + #[tokio::test] + async fn spilled_lineage_shares_one_file_and_round_trips() { + let dir = TempStrDir::default(); + let mut dataset = tiny_dataset(dir.as_str()).await; + spill_everything(&mut dataset).await; + + let lineage = RowLineage { + row_ids: scattered_row_ids(20_000), + created_at: alternating_versions(20_000, 1), + last_updated_at: alternating_versions(20_000, 3), + }; + let placed = place_row_lineage(&dataset, &lineage).await.unwrap(); + assert!( + matches!(placed.row_ids, RowIdMeta::Column) + && matches!(placed.created_at, RowDatasetVersionMeta::Column) + && matches!(placed.last_updated_at, RowDatasetVersionMeta::Column), + "expected every sequence to spill" + ); + // One file carries all three, and it becomes one of the fragment's + // data files; the metadata only marks the sequences as spilled. + let mut fragment = Fragment::new(0); + fragment.physical_rows = Some(20_000); + placed.apply(&mut fragment); + assert_eq!(fragment.files.len(), 1); + assert_eq!( + fragment.files[0].fields.as_ref(), + [ + ROW_ID_FIELD_ID, + ROW_CREATED_AT_VERSION_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID + ] + ); + + let row_ids = read_spilled_row_ids(&dataset, &fragment).await.unwrap(); + assert_eq!( + row_ids.iter().collect::>(), + lineage.row_ids.iter().collect::>() + ); + let created_at = + read_spilled_versions(&dataset, &fragment, ROW_CREATED_AT_VERSION_FIELD_ID) + .await + .unwrap(); + assert_eq!(versions_of(&created_at), versions_of(&lineage.created_at)); + let last_updated_at = + read_spilled_versions(&dataset, &fragment, ROW_LAST_UPDATED_AT_VERSION_FIELD_ID) + .await + .unwrap(); + assert_eq!( + versions_of(&last_updated_at), + versions_of(&lineage.last_updated_at) + ); + } + + #[tokio::test] + async fn only_the_sequences_over_the_limit_spill() { + let dir = TempStrDir::default(); + let mut dataset = tiny_dataset(dir.as_str()).await; + // An appended fragment's row ids are a single `Range` and its versions + // a single run, so they encode to a few dozen bytes and must never + // leave the manifest, even next to a sequence that does. + let lineage = RowLineage { + row_ids: RowIdSequence::from(0..20_000), + created_at: alternating_versions(20_000, 1), + last_updated_at: RowDatasetVersionSequence::from_uniform_row_count(20_000, 1), + }; + + // A table that has not opted in never spills, whatever the size. + let placed = place_row_lineage(&dataset, &lineage).await.unwrap(); + assert!(matches!( + placed.created_at, + RowDatasetVersionMeta::Inline(_) + )); + assert!(placed.file.is_none()); + + dataset + .update_config([(SPILL_ROW_LINEAGE_CONFIG_KEY, "true")]) + .await + .unwrap(); + let placed = place_row_lineage(&dataset, &lineage).await.unwrap(); + assert!( + matches!(placed.row_ids, RowIdMeta::Inline(_)), + "a range sequence must stay inline, got {:?}", + placed.row_ids + ); + assert!( + matches!(placed.last_updated_at, RowDatasetVersionMeta::Inline(_)), + "a single-run sequence must stay inline, got {:?}", + placed.last_updated_at + ); + assert!( + matches!(placed.created_at, RowDatasetVersionMeta::Column), + "an alternating version sequence encodes past 200 KiB" + ); + assert_eq!( + placed.file.as_ref().unwrap().fields.as_ref(), + [ROW_CREATED_AT_VERSION_FIELD_ID] + ); + } + + /// The format allows the columns only in v2 files, so a legacy v1 table + /// keeps everything inline however it is configured. + #[tokio::test] + async fn a_legacy_v1_table_never_spills() { + let dir = TempStrDir::default(); + let mut dataset = tiny_dataset_with_version(dir.as_str(), LanceFileVersion::Legacy).await; + spill_everything(&mut dataset).await; + assert_eq!(inline_row_lineage_max_bytes(&dataset).unwrap(), None); + + let lineage = RowLineage { + row_ids: RowIdSequence::from(0..20_000), + created_at: alternating_versions(20_000, 1), + last_updated_at: RowDatasetVersionSequence::from_uniform_row_count(20_000, 1), + }; + let placed = place_row_lineage(&dataset, &lineage).await.unwrap(); + assert!(matches!(placed.row_ids, RowIdMeta::Inline(_))); + assert!(matches!( + placed.created_at, + RowDatasetVersionMeta::Inline(_) + )); + assert!(matches!( + placed.last_updated_at, + RowDatasetVersionMeta::Inline(_) + )); + assert!(placed.file.is_none()); + } + + /// 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. + async fn spill_everything(dataset: &mut Dataset) { + dataset + .update_config([ + (SPILL_ROW_LINEAGE_CONFIG_KEY, "true"), + (INLINE_ROW_LINEAGE_MAX_BYTES_CONFIG_KEY, "0"), + ]) + .await + .unwrap(); + } +} diff --git a/rust/lance/src/dataset/rowids/validate.rs b/rust/lance/src/dataset/rowids/validate.rs index 27bed6491ac..42d4044bd39 100644 --- a/rust/lance/src/dataset/rowids/validate.rs +++ b/rust/lance/src/dataset/rowids/validate.rs @@ -4,13 +4,12 @@ //! Integrity checks for the stable row id invariants that the row id index and the write //! paths rely on. Reached through [`Dataset::validate`]. -use super::load_row_id_sequence; +use super::{RowVersionKind, load_row_id_sequence, load_row_version_sequence}; use crate::dataset::Dataset; use crate::dataset::fragment::FileFragment; use crate::{Error, Result}; use futures::{StreamExt, TryStreamExt}; use lance_core::utils::deletion::DeletionVector; -use lance_table::format::RowDatasetVersionMeta; use lance_table::rowids::RowIdSequence; use roaring::RoaringTreemap; use std::sync::Arc; @@ -75,14 +74,20 @@ pub async fn validate_stable_row_ids(dataset: &Dataset) -> Result<()> { ))); } - for (name, meta) in [ - ("created_at", &metadata.created_at_version_meta), - ("last_updated_at", &metadata.last_updated_at_version_meta), + for (name, kind) in [ + ("created_at", RowVersionKind::CreatedAt), + ("last_updated_at", RowVersionKind::LastUpdatedAt), ] { - // Only inline version metadata can be read back; nothing writes the - // external form yet. - if let Some(meta @ RowDatasetVersionMeta::Inline(_)) = meta { - let versions = meta.load_sequence()?.len(); + // The external form is a valid encoding nothing writes and nothing + // can read yet, so it is the one placement not checked here. + if let Some(versions) = load_row_version_sequence(dataset, metadata, kind) + .await + .or_else(|error| match error { + Error::NotSupported { .. } => Ok(None), + other => Err(other), + })? + { + let versions = versions.len(); if versions != physical_rows { return Err(corrupt(format!( "Fragment {} has {} {} versions, but {} physical rows, in dataset {:?}", @@ -159,6 +164,7 @@ fn describe_first_live_slot(fragments: &[FragmentRowIds], row_id: u64) -> String #[cfg(test)] mod tests { use super::*; + use lance_table::format::RowDatasetVersionMeta; // Shared with the row id tests next door, which cover the same operations. use super::super::test::{compact, delete}; diff --git a/rust/lance/src/session/caches.rs b/rust/lance/src/session/caches.rs index 88b1556ad99..e5e0defc29c 100644 --- a/rust/lance/src/session/caches.rs +++ b/rust/lance/src/session/caches.rs @@ -19,7 +19,9 @@ use lance_core::{ }; use lance_select::RowAddrMask; use lance_table::{ - format::{DataFile, DeletionFile, DeletionFileType, Manifest, RowIdMeta}, + format::{ + DataFile, DeletionFile, DeletionFileType, Manifest, RowDatasetVersionSequence, RowIdMeta, + }, rowids::{RowIdIndex, RowIdSequence}, }; use object_store::path::Path; @@ -302,6 +304,45 @@ impl CacheKey for RowIdSequenceKey<'_> { } } +/// Cache key for one of a fragment's per-row version sequences that is spilled +/// to a data file column. +/// +/// Inline sequences are not cached: they decode straight from the manifest +/// bytes the fragment already holds. +#[derive(Debug)] +pub struct RowVersionSequenceKey<'a> { + pub fragment_id: u64, + /// Which sequence this is, by the reserved field id of its column. + pub field_id: i32, + /// The data file carrying the column, named freshly per rewrite, so its + /// path identifies the contents the way an inline sequence's digest does. + pub data_file: &'a DataFile, +} + +impl CacheKey for RowVersionSequenceKey<'_> { + type ValueType = RowDatasetVersionSequence; + fn key(&self) -> Cow<'_, str> { + Cow::Owned(format!( + "row_version_sequence/{}/{}", + self.fragment_id, self.field_id + )) + } + fn type_name() -> &'static str { + "RowDatasetVersionSequence" + } + + fn schema() -> CacheKeySchema { + CacheKeySchema::new("lance.dataset.row-version-sequence-key", 1) + } + + fn write_key(&self, builder: &mut KeyBuilder) { + builder.write_u64(self.fragment_id); + builder.write_u64(self.field_id as u64); + builder.write_str(&self.data_file.path); + builder.write_u64(self.data_file.base_id.map_or(u64::MAX, u64::from)); + } +} + impl DSMetadataCache { /// Create a file-specific metadata cache with the given prefix. /// This is used by file readers and other components that need file-level caching.