From b3ac0d8bf20ffb977f2351d760e18bde1429fdf6 Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Thu, 17 Sep 2026 17:36:15 +0800 Subject: [PATCH 1/2] feat(dataset): read and write spilled row lineage columns Read the hidden row lineage columns defined by the previous commit and add the primitives that write them; nothing calls the writer yet. The row id loader gains a Column arm that finds the one data file of the fragment carrying the reserved field id and reads the column back through the ordinary file reader, projected by field id, into the same per-fragment cache as inline sequences. A new load_row_version_sequence loads either version sequence wherever it is stored, and FileFragment::open loads the version sequences asynchronously alongside the row ids instead of the reader builder decoding them synchronously and silently falling back to version 1 on any failure. A null in a lineage column, a length other than physical_rows, or more than one file carrying the id is corruption. Dataset::validate checks spilled version sequences too, and a file holding only lineage columns is not opened by a dataset field. place_row_lineage encodes a fragment's three sequences and, for each one whose encoding exceeds the table's inline budget, writes it as a column of one new data file per fragment, which the caller adds to the fragment's files; the metadata only marks the sequence as spilled. Spilling is opt-in per table through lance.row_lineage.spill, with lance.row_lineage.inline_max_bytes overriding the 200 KiB default, and never happens in a build that does not understand the feature flag. Co-authored-by: Will Jones Co-Authored-By: Claude Fable 5.1 --- rust/lance/src/dataset/fragment.rs | 111 ++-- rust/lance/src/dataset/rowids.rs | 97 +++- rust/lance/src/dataset/rowids/spill.rs | 603 ++++++++++++++++++++++ rust/lance/src/dataset/rowids/validate.rs | 24 +- rust/lance/src/session/caches.rs | 43 +- 5 files changed, 810 insertions(+), 68 deletions(-) create mode 100644 rust/lance/src/dataset/rowids/spill.rs 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..ca19b65b055 --- /dev/null +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -0,0 +1,603 @@ +// 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::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).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).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. +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. From 14ec6a2807b00dbce0e0bc6c6bf680bd589ac166 Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Thu, 24 Sep 2026 02:57:04 +0800 Subject: [PATCH 2/2] fix(dataset): box the spilled lineage column read Reading a spilled row id or version sequence drives the full data file reader. Inlined into the row id index build, which take reaches and index builds and optimize_indices reach through take, it made those futures too deep for the trait solver to prove Send and Sync (E0275), so the crate no longer compiled on top of current main. Boxing the read where the two sequence loaders call it cuts the type depth for every caller. Co-Authored-By: Claude Opus 5.5 --- rust/lance/src/dataset/rowids/spill.rs | 15 ++++++++++++--- 1 file changed, 12 insertions(+), 3 deletions(-) diff --git a/rust/lance/src/dataset/rowids/spill.rs b/rust/lance/src/dataset/rowids/spill.rs index ca19b65b055..9fae144d33b 100644 --- a/rust/lance/src/dataset/rowids/spill.rs +++ b/rust/lance/src/dataset/rowids/spill.rs @@ -26,7 +26,7 @@ use std::sync::Arc; use arrow_array::{Array, ArrayRef, RecordBatch, UInt64Array}; use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema}; -use futures::TryStreamExt; +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}; @@ -261,7 +261,9 @@ async fn write_lineage_file( /// 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).await?; + let ids = read_spilled_column(dataset, fragment, ROW_ID_FIELD_ID) + .boxed() + .await?; Ok(RowIdSequence::from(ids.as_slice())) } @@ -272,12 +274,19 @@ pub async fn read_spilled_versions( fragment: &Fragment, field_id: i32, ) -> Result { - let versions = read_spilled_column(dataset, fragment, field_id).await?; + 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,