diff --git a/docs/src/format/table/index.md b/docs/src/format/table/index.md index f9da132cf3b..a20d4c564c2 100644 --- a/docs/src/format/table/index.md +++ b/docs/src/format/table/index.md @@ -119,6 +119,17 @@ Field ids might be replaced with `-2`, a tombstone value. In this case that column should be ignored. This used, for example, when rewriting a column: The old data file replaces the field id with `-2` to ignore the old data, and a new data file is appended to the fragment. +Every negative field id is reserved for system use and never names a field of the +dataset schema. A reader MUST skip any negative id when it projects the dataset schema +onto a data file, rather than treat it as a schema field or reject the file. Besides +`-1` (not yet assigned; only ever exists in memory and must not be written) and the +`-2` tombstone above, `-3`, `-4` and `-5` are the hidden `_rowid`, +`_row_created_at_version` and `_row_last_updated_at_version` columns that hold a +fragment's row lineage sequences when they are not stored in the manifest; see +[Row ID and Lineage](row_id_lineage.md). Such a column always lives in one of the +fragment's `files`, next to user columns or in a file holding nothing else, and at +most one file of a fragment may carry each of these ids. + ## Data Files Data files store column data for a fragment using the Lance file format. diff --git a/docs/src/format/table/row_id_lineage.md b/docs/src/format/table/row_id_lineage.md index 0d83431f972..2756d7757f9 100644 --- a/docs/src/format/table/row_id_lineage.md +++ b/docs/src/format/table/row_id_lineage.md @@ -185,14 +185,81 @@ The implementation selects the most compact encoding based on the value range, c -#### Inline and External Storage +#### Inline and Spilled Storage + +`DataFragment` defines inline and column alternatives for row ID sequences and row +version sequences. This allows small sequences to stay in the manifest (fewer +IOPS) while larger sequences (resulting from frequent updates) move outside the +manifest. + +Sequences small enough (~200KB encoded and under) are stored inline in the fragment +metadata to avoid additional I/O. An inline sequence is rewritten into every manifest +version. + +A larger sequence is **spilled to a hidden column of one of the fragment's data +files**. The column arm of the oneof (`column_row_ids`, `column_created_at_versions` +or `column_last_updated_at_versions`) is an empty `RowLineageColumn` marker: it carries +no file reference, because the file is one of the fragment's `files` and is found by a +reserved negative field id in that entry's `fields`. Exactly one file of the fragment +carries each reserved field; zero or more than one is corruption. The `column_indices` +entry paired with the field id locates the column the way it does for a user column, so +a lineage column +may share a file with the user columns or with the other lineage columns, or sit in a +file that holds nothing else. Reading it uses the ordinary data file reader and its +encodings. + +The marker is valid only in a fragment whose data files are Lance v2 files. A legacy +v1 data file has no `column_indices` to locate the column by, and a fragment cannot +mix v1 and v2 files, so a writer on a dataset that stores v1 files MUST leave every +sequence inline; a marker in a fragment with a v1 data file is corruption. + +The columns have this schema, with the field ids `-3`, `-4` and `-5` respectively. +The three names are reserved: a writer MUST reject a user column with any of them +(every Lance write path does), so a hidden column never collides with a field of the +dataset schema. + +```python +import pyarrow as pa + +row_lineage_columns = pa.schema([ + pa.field("_rowid", pa.uint64(), nullable=False), + pa.field("_row_created_at_version", pa.uint64(), nullable=False), + pa.field("_row_last_updated_at_version", pa.uint64(), nullable=False), +]) +``` -`DataFragment` defines inline and external metadata fields as valid wire alternatives for row ID sequences and row version sequences. -These fields do not currently imply a size-based switching threshold. -Current Lance writers store all three sequence types inline in the fragment metadata regardless of their encoded size and do not emit the external alternatives. +Each column present holds exactly `physical_rows` values, one per physical row in +physical row order, deleted rows included; the value at offset `i` is the row id or +version of the row at offset `i`, the same thing the inline encoding's `i`-th entry +would be. A null value, or a column whose length differs from `physical_rows`, is +corruption and MUST be rejected rather than read as a default. A file may carry any +subset of the three columns; a sequence whose arm is not the column marker is not read +from any file, whatever the file holds. + +Which sequences may leave the manifest follows from when their values are known. +A value the commit assigns -- an appended fragment's row ids, an inserted row's +created-at version, every row's last-updated-at version -- can change when a commit +conflict is retried, so it stays inline where the retry can rewrite it; those +sequences are single runs and cost a few bytes. A value carried over from existing +rows -- the row ids and created-at versions that compaction or a row rewrite +preserves -- is fixed before the commit and may be written to a data file. +A writer spills only on a table that opts in through the `lance.row_lineage.spill` +config key; a table that never sets it is unchanged. + +A writer that emits any column arm MUST set the spilled row lineage feature flag +(bit 11, value 2048) in both the reader and writer flag words. A reader without that +bit sees an unset oneof and would take the fragment to have no row IDs at all, on a +table whose manifest says every fragment has them. + +Field numbers 6, 8 and 10 of `DataFragment` (`external_row_ids`, +`external_last_updated_at_versions`, `external_created_at_versions`) once named an +opaque byte range in a file holding the same encoding as the inline arm. No Lance +writer ever emitted them; the column arms replace that design, and the numbers and +names are reserved. -Current Lance readers can load externally stored row ID sequences. -The format also permits external created-at and last-updated-at version sequences, but current Lance readers cannot load them; this is an implementation limitation, not an invalid encoding. +!!! note + Spilled row lineage sequences are not yet a released feature. A released build + treats bit 11 as an unknown feature flag and refuses the dataset.
DataFragment row_id_sequence field @@ -201,7 +268,7 @@ The format also permits external created-at and last-updated-at version sequence message DataFragment { oneof row_id_sequence { bytes inline_row_ids = 5; - ExternalFile external_row_ids = 6; + RowLineageColumn column_row_ids = 12; } } ``` @@ -288,7 +355,7 @@ RowDatasetVersionSequence { message DataFragment { oneof created_at_version_sequence { bytes inline_created_at_versions = 9; - ExternalFile external_created_at_versions = 10; + RowLineageColumn column_created_at_versions = 14; } } ``` @@ -331,7 +398,7 @@ New physical row (current): message DataFragment { oneof last_updated_at_version_sequence { bytes inline_last_updated_at_versions = 7; - ExternalFile external_last_updated_at_versions = 8; + RowLineageColumn column_last_updated_at_versions = 13; } } ``` diff --git a/docs/src/format/table/versioning.md b/docs/src/format/table/versioning.md index 4b12e2d2de6..16891717203 100644 --- a/docs/src/format/table/versioning.md +++ b/docs/src/format/table/versioning.md @@ -34,7 +34,8 @@ they should return an "unsupported" error on any read or write operation. | 256 | `FLAG_MIXED_DATA_FILE_VERSIONS` | Yes | Yes | The snapshot may reference recognized V2 data files with different exact versions. Both bits must be set and remain set on later versions. | | 512 | `FLAG_FRAG_REUSE_WITH_STABLE_ROW_IDS` | Yes | Yes | The table uses stable row IDs and carries a [Fragment Reuse Index](../index/system/frag_reuse.md). | | 1024 | `FLAG_FRAGMENT_REUSE_INDEX` | Yes | Yes | The fragment reuse index records tagged transitions (`IndexMetadata.index_version >= 1`). Readers must translate row addresses through them; writers must preserve them. An implementation without this flag would decode the details as the legacy format and silently drop the transitions when it next rewrites the fragment reuse index. See [FRI index versions](../index/system/frag_reuse.md#fri-index-versions). | +| 2048 | `FLAG_UNSTABLE_SPILLED_ROW_LINEAGE` | Yes | Yes | Some fragment stores its row ids or row version sequences as hidden columns of a data file rather than inline. A reader without this flag would see the fragment as having no row ids. Unstable: release builds reject it unless explicitly opted in. | -Flags with bit values 2048 and above are unknown; unknown flags cause implementations to reject the dataset with an "unsupported" error. The paired mixed-version reader and writer bits must either both be set or both be clear; a half-set manifest is invalid. +Flags with bit values 4096 and above are unknown; unknown flags cause implementations to reject the dataset with an "unsupported" error. The paired mixed-version reader and writer bits must either both be set or both be clear; a half-set manifest is invalid. diff --git a/protos/table.proto b/protos/table.proto index 785d3c71455..10823e47608 100644 --- a/protos/table.proto +++ b/protos/table.proto @@ -412,28 +412,46 @@ message DataFragment { oneof row_id_sequence { // Current Lance writers store row ids inline regardless of encoded size. bytes inline_row_ids = 5; - // Supported by current Lance readers, but not emitted by current Lance writers. - ExternalFile external_row_ids = 6; + /* The row ids are a hidden column of one of this fragment's `files`: the + * entry whose `fields` carries the reserved id -3. See RowLineageColumn. + * + * A writer MUST set the spilled row lineage feature flag (2048) when any + * fragment uses this arm or either of the other column arms below. A + * reader that does not understand the flag would see an unset oneof and + * treat the fragment as having no row ids at all. + */ + RowLineageColumn column_row_ids = 12; } // row_id_sequence oneof last_updated_at_version_sequence { // Current Lance writers store last-updated versions inline regardless of encoded size. bytes inline_last_updated_at_versions = 7; - /* Valid external alternative. Current Lance writers do not emit this field, - * and current Lance readers cannot load it. + /* The versions are a hidden column of one of this fragment's `files`: the + * entry whose `fields` carries the reserved id -5. Gated like + * `column_row_ids`. */ - ExternalFile external_last_updated_at_versions = 8; + RowLineageColumn column_last_updated_at_versions = 13; } // last_updated_at_version_sequence oneof created_at_version_sequence { // Current Lance writers store created-at versions inline regardless of encoded size. bytes inline_created_at_versions = 9; - /* Valid external alternative. Current Lance writers do not emit this field, - * and current Lance readers cannot load it. + /* The versions are a hidden column of one of this fragment's `files`: the + * entry whose `fields` carries the reserved id -4. Gated like + * `column_row_ids`. */ - ExternalFile external_created_at_versions = 10; + RowLineageColumn column_created_at_versions = 14; } // created_at_version_sequence + /* The `external_*` arms of the three sequence oneofs above: an opaque byte + * range in a file, holding the same encoding as the inline arm. No Lance + * writer ever emitted them; the column arms replace that design. Reserved so + * a reader never has to interpret them. + */ + reserved 6, 8, 10; + reserved "external_row_ids", "external_last_updated_at_versions", + "external_created_at_versions"; + /* Number of original rows in the fragment, this includes rows that are now marked with * deletion tombstones. To compute the current number of rows, subtract * `deletion_file.num_deleted_rows` from this value. @@ -451,6 +469,11 @@ message DataFile { * used for "tombstoned", meaning a field that is no longer in use. This is often * because the original field id was reassigned to a different data file. * + * Every negative value is reserved for system use and never names a field of the + * dataset schema. -3, -4 and -5 are the hidden row lineage columns (see + * RowLineageColumn); a reader projecting the dataset schema onto a file skips + * them, and any other negative value, rather than reject the file. + * * In Lance v1 IDs are assigned based on position in the file, offset by the max * existing field id in the table (if any already). So when a fragment is first created * with one file of N columns, the field ids will be 1, 2, ..., N. If a second fragment @@ -634,6 +657,26 @@ message DeletionFile { optional uint32 base_id = 7; } // DeletionFile +/* Marks a row lineage sequence that is stored as a hidden column of one of the + * fragment's `files` rather than inline: `_rowid` under the reserved field id + * -3, `_row_created_at_version` under -4, `_row_last_updated_at_version` under + * -5. Exactly one entry of `DataFragment.files` carries that id in its + * `fields`, and its `column_indices` locates the column the way it does for a + * user column, so the column may share a file with the user columns or with the + * other lineage columns, or sit in a file that holds nothing else. The column is + * a non-nullable uint64 with exactly `physical_rows` values, in physical row + * order; see the table format documentation for the schema. + * + * Carries no fields: the file is found by its field id, so the file's metadata + * lives in `files` alone. + * + * Valid only in a fragment whose data files are Lance v2 files. A legacy v1 + * data file has no `column_indices` to locate the column by, so a writer on a + * dataset that stores v1 files leaves every sequence inline. + */ +message RowLineageColumn {} + +// A byte range of a file. Used by the fragment reuse index details. message ExternalFile { // Path to the file, relative to the root of the table. string path = 1; diff --git a/python/src/rowids.rs b/python/src/rowids.rs index 0ad83c90579..4b6fd1de3d4 100644 --- a/python/src/rowids.rs +++ b/python/src/rowids.rs @@ -48,8 +48,8 @@ impl PyRowIdSequence { fn from_inline_metadata(metadata: PyRef<'_, PyRowIdMeta>) -> PyResult { match &metadata.0 { RowIdMeta::Inline(data) => read_row_ids(data).infer_error().map(Self), - RowIdMeta::External(_) => Err(PyNotImplementedError::new_err( - "Row ids stored in an external file cannot be read into a RowIdSequence", + RowIdMeta::Column => Err(PyNotImplementedError::new_err( + "Row ids stored outside the manifest cannot be read into a RowIdSequence", )), } } diff --git a/rust/lance-table/src/feature_flags.rs b/rust/lance-table/src/feature_flags.rs index 3c67753910d..019b332dfa6 100644 --- a/rust/lance-table/src/feature_flags.rs +++ b/rust/lance-table/src/feature_flags.rs @@ -64,8 +64,24 @@ pub const FLAG_FRAG_REUSE_WITH_STABLE_ROW_IDS: u64 = 1 << 9; /// preserves them during maintenance. Legacy-only FRI does not set this bit. /// Bit 9 is taken by the stable-row-id FRI compatibility flag. pub const FLAG_FRAGMENT_REUSE_INDEX: u64 = 1 << 10; +/// Some fragment stores a row lineage sequence -- its row ids, or its +/// created-at or last-updated-at versions -- as a hidden column of a data file +/// rather than inline in the manifest (`RowIdMeta::Column`, +/// `RowDatasetVersionMeta::Column`). +/// +/// A reader without this bit sees an unset `row_id_sequence` oneof and would +/// take the fragment to have no row ids at all, on a table whose manifest says +/// every fragment has them; a writer without it would carry the fragment +/// forward and drop the sequence. Both must refuse the table. +/// +/// Every earlier build has its unknown boundary at or below this bit, so each +/// already refuses such a dataset without a change of its own. Spilled row +/// lineage is not yet a released feature: this build understands the bit only +/// in debug builds or when [`ENABLE_UNSTABLE_SPILLED_ROW_LINEAGE_ENV`] is set, +/// mirroring [`FLAG_UNSTABLE_DATA_OVERLAY_FILES`]. +pub const FLAG_UNSTABLE_SPILLED_ROW_LINEAGE: u64 = 1 << 11; /// The first bit that is unknown as a feature flag -pub const FLAG_UNKNOWN: u64 = 1 << 11; +pub const FLAG_UNKNOWN: u64 = 1 << 12; const _: () = assert!(FLAG_COVERED_INDEX_METADATA < FLAG_UNKNOWN); // The fence needs a bit the current released build already refuses, which means @@ -77,6 +93,7 @@ const _: () = assert!(FLAG_MIXED_DATA_FILE_VERSIONS < FLAG_UNKNOWN); const _: () = assert!(FLAG_FRAG_REUSE_WITH_STABLE_ROW_IDS >= 1 << 8); const _: () = assert!(FLAG_FRAG_REUSE_WITH_STABLE_ROW_IDS < FLAG_UNKNOWN); const _: () = assert!(FLAG_FRAGMENT_REUSE_INDEX < FLAG_UNKNOWN); +const _: () = assert!(FLAG_UNSTABLE_SPILLED_ROW_LINEAGE < FLAG_UNKNOWN); pub(crate) const STICKY_PAIRED_FLAGS: u64 = FLAG_MIXED_DATA_FILE_VERSIONS; @@ -84,6 +101,12 @@ pub(crate) const STICKY_PAIRED_FLAGS: u64 = FLAG_MIXED_DATA_FILE_VERSIONS; /// overlay files before the feature is generally released. pub const ENABLE_UNSTABLE_DATA_OVERLAY_FILES_ENV: &str = "LANCE_ENABLE_UNSTABLE_DATA_OVERLAY_FILES"; +/// Environment variable that opts a release build into reading and writing row +/// lineage sequences spilled to hidden data file columns before the feature is +/// generally released. +pub const ENABLE_UNSTABLE_SPILLED_ROW_LINEAGE_ENV: &str = + "LANCE_ENABLE_UNSTABLE_SPILLED_ROW_LINEAGE"; + /// Set the reader and writer feature flags in the manifest based on the contents of the manifest. pub fn apply_feature_flags( manifest: &mut Manifest, @@ -153,6 +176,15 @@ pub fn apply_feature_flags( manifest.writer_feature_flags |= FLAG_UNSTABLE_DATA_OVERLAY_FILES; } + let has_spilled_row_lineage = manifest + .fragments + .iter() + .any(|frag| frag.has_spilled_row_lineage()); + if has_spilled_row_lineage { + manifest.reader_feature_flags |= FLAG_UNSTABLE_SPILLED_ROW_LINEAGE; + manifest.writer_feature_flags |= FLAG_UNSTABLE_SPILLED_ROW_LINEAGE; + } + if disable_transaction_file { manifest.writer_feature_flags |= FLAG_DISABLE_TRANSACTION_FILE; } @@ -188,6 +220,13 @@ fn data_overlay_files_enabled() -> bool { cfg!(debug_assertions) || std::env::var_os(ENABLE_UNSTABLE_DATA_OVERLAY_FILES_ENV).is_some() } +/// Whether this build understands row lineage sequences spilled to data file +/// columns: always in debug builds, and in release builds only when +/// [`ENABLE_UNSTABLE_SPILLED_ROW_LINEAGE_ENV`] is set, like data overlay files. +pub fn spilled_row_lineage_enabled() -> bool { + cfg!(debug_assertions) || std::env::var_os(ENABLE_UNSTABLE_SPILLED_ROW_LINEAGE_ENV).is_some() +} + /// Clear `flag` from `flags` when its gating feature is not enabled in this /// build; leave it set otherwise. One call per unstable flag, so support for /// several unstable features chains cleanly. @@ -200,7 +239,7 @@ fn mark_supported(flags: &mut u64, flag: u64, feature_enabled: bool) { /// The feature-flag bits this build understands, given whether overlay support /// is enabled. Split out from [`supported_flags`] so the policy is testable /// without toggling the build profile or environment. -fn supported_flags_when(overlay_enabled: bool) -> u64 { +fn supported_flags_when(overlay_enabled: bool, spilled_row_lineage_enabled: bool) -> u64 { let mut supported = FLAG_UNKNOWN - 1; mark_supported( &mut supported, @@ -212,11 +251,16 @@ fn supported_flags_when(overlay_enabled: bool) -> u64 { // Bit 10 now falls below the unknown boundary, so keep tagged FRI refused // until its reader/writer handling lands. mark_supported(&mut supported, FLAG_FRAGMENT_REUSE_INDEX, false); + mark_supported( + &mut supported, + FLAG_UNSTABLE_SPILLED_ROW_LINEAGE, + spilled_row_lineage_enabled, + ); supported } fn supported_flags() -> u64 { - supported_flags_when(data_overlay_files_enabled()) + supported_flags_when(data_overlay_files_enabled(), spilled_row_lineage_enabled()) } pub fn can_read_dataset(reader_flags: u64) -> bool { @@ -378,15 +422,31 @@ mod tests { fn test_data_overlay_flag_release_gating() { // Release default (overlays disabled): the overlay flag is treated as // unknown so the dataset is refused, while other known flags still pass. - let supported = supported_flags_when(false); + let supported = supported_flags_when(false, false); assert_eq!(supported & FLAG_UNSTABLE_DATA_OVERLAY_FILES, 0); assert_eq!(FLAG_DELETION_FILES & !supported, 0); assert_ne!(FLAG_UNSTABLE_DATA_OVERLAY_FILES & !supported, 0); // Enabled (debug or env opt-in): the overlay flag is understood. - let supported = supported_flags_when(true); + let supported = supported_flags_when(true, false); assert_eq!(FLAG_UNSTABLE_DATA_OVERLAY_FILES & !supported, 0); } + #[test] + fn test_spilled_row_lineage_flag_release_gating() { + // Every earlier build has its unknown boundary at or below this bit + // (256 for v11, 512 for the v12 and v13 pre-releases, 2048 on main + // before this flag), so each refuses the dataset without a change of + // its own. + assert_eq!(FLAG_UNSTABLE_SPILLED_ROW_LINEAGE, 2048); + // A build that has not opted in refuses the dataset; one that has + // understands it, and either way the other known flags still pass. + let supported = supported_flags_when(true, false); + assert_ne!(FLAG_UNSTABLE_SPILLED_ROW_LINEAGE & !supported, 0); + assert_eq!(FLAG_MIXED_DATA_FILE_VERSIONS & !supported, 0); + let supported = supported_flags_when(true, true); + assert_eq!(FLAG_UNSTABLE_SPILLED_ROW_LINEAGE & !supported, 0); + } + #[test] fn test_apply_feature_flags_sets_overlay_flag() { use crate::format::overlay::{DataOverlayFile, OverlayCoverage}; @@ -426,6 +486,73 @@ mod tests { ); } + #[test] + fn test_apply_feature_flags_sets_spilled_row_lineage_flag() { + use crate::format::{DataFile, DataStorageFormat, Fragment, ROW_ID_FIELD_ID, RowIdMeta}; + use crate::rowids::version::RowDatasetVersionMeta; + use arrow_schema::{Field as ArrowField, Schema as ArrowSchema}; + use lance_core::datatypes::Schema; + use std::collections::HashMap; + use std::sync::Arc; + + let arrow_schema = ArrowSchema::new(vec![ArrowField::new( + "id", + arrow_schema::DataType::Int64, + false, + )]); + let schema = Schema::try_from(&arrow_schema).unwrap(); + let spilled = + DataFile::new_legacy_from_fields("lineage.lance", vec![ROW_ID_FIELD_ID], None); + // Each of the three sequences on its own is enough to require the flag. + let mut spilled_fragments = Vec::new(); + for kind in 0..3 { + let mut fragment = Fragment::new(kind); + fragment.files.push(spilled.clone()); + // Every fragment of a stable-row-id table carries row ids; the + // version-only cases keep theirs inline. + fragment.row_id_meta = Some(RowIdMeta::Inline(vec![].into())); + match kind { + 0 => fragment.row_id_meta = Some(RowIdMeta::Column), + 1 => fragment.created_at_version_meta = Some(RowDatasetVersionMeta::Column), + _ => fragment.last_updated_at_version_meta = Some(RowDatasetVersionMeta::Column), + } + spilled_fragments.push(fragment); + } + for fragment in spilled_fragments { + let mut manifest = Manifest::new( + schema.clone(), + Arc::new(vec![fragment]), + DataStorageFormat::default(), + HashMap::new(), + ); + apply_feature_flags(&mut manifest, true, false).unwrap(); + assert_ne!( + manifest.reader_feature_flags & FLAG_UNSTABLE_SPILLED_ROW_LINEAGE, + 0 + ); + assert_ne!( + manifest.writer_feature_flags & FLAG_UNSTABLE_SPILLED_ROW_LINEAGE, + 0 + ); + } + + // An inline sequence does not, so a table whose sequences all fit in + // the manifest stays readable by builds without the flag. + let mut fragment = Fragment::new(0); + fragment.row_id_meta = Some(RowIdMeta::Inline(vec![].into())); + let mut manifest = Manifest::new( + schema, + Arc::new(vec![fragment]), + DataStorageFormat::default(), + HashMap::new(), + ); + apply_feature_flags(&mut manifest, true, false).unwrap(); + assert_eq!( + manifest.reader_feature_flags & FLAG_UNSTABLE_SPILLED_ROW_LINEAGE, + 0 + ); + } + #[test] fn test_write_check() { assert!(can_write_dataset(0)); diff --git a/rust/lance-table/src/format.rs b/rust/lance-table/src/format.rs index 1c9e0e37c8c..23959a0cb61 100644 --- a/rust/lance-table/src/format.rs +++ b/rust/lance-table/src/format.rs @@ -23,7 +23,10 @@ pub use manifest::{ SelfDescribingFileReader, WriterVersion, is_detached_version, populate_manifest_schema_dictionaries, }; -pub use row_ids::{ExternalFile, InlineRowIds, RowIdMeta}; +pub use row_ids::{ + ExternalFile, InlineRowIds, ROW_CREATED_AT_VERSION_FIELD_ID, ROW_ID_FIELD_ID, + ROW_LAST_UPDATED_AT_VERSION_FIELD_ID, RowIdMeta, +}; pub use transaction::{Transaction, operation_may_change_schema}; use lance_core::{Error, Result}; diff --git a/rust/lance-table/src/format/fragment.rs b/rust/lance-table/src/format/fragment.rs index c7b423a72cb..9d5763e88dd 100644 --- a/rust/lance-table/src/format/fragment.rs +++ b/rust/lance-table/src/format/fragment.rs @@ -13,7 +13,7 @@ use object_store::path::Path; use serde::{Deserialize, Deserializer, Serialize, Serializer}; use super::overlay::{DataOverlayFile, TOMBSTONE_FIELD_ID, sort_overlays_newest_last}; -use super::row_ids::{ExternalFile, RowIdMeta}; +use super::row_ids::RowIdMeta; use crate::format::pb; use crate::rowids::version::{ @@ -318,13 +318,9 @@ impl DataFileFieldInterner { pb::data_fragment::LastUpdatedAtVersionSequence::InlineLastUpdatedAtVersions(data) => { Ok(RowDatasetVersionMeta::Inline(cache.intern(data))) } - pb::data_fragment::LastUpdatedAtVersionSequence::ExternalLastUpdatedAtVersions( - file, - ) => Ok(RowDatasetVersionMeta::External(ExternalFile { - path: file.path, - offset: file.offset, - size: file.size, - })), + pb::data_fragment::LastUpdatedAtVersionSequence::ColumnLastUpdatedAtVersions(_) => { + Ok(RowDatasetVersionMeta::Column) + } } } @@ -337,12 +333,8 @@ impl DataFileFieldInterner { pb::data_fragment::CreatedAtVersionSequence::InlineCreatedAtVersions(data) => { Ok(RowDatasetVersionMeta::Inline(cache.intern(data))) } - pb::data_fragment::CreatedAtVersionSequence::ExternalCreatedAtVersions(file) => { - Ok(RowDatasetVersionMeta::External(ExternalFile { - path: file.path, - offset: file.offset, - size: file.size, - })) + pb::data_fragment::CreatedAtVersionSequence::ColumnCreatedAtVersions(_) => { + Ok(RowDatasetVersionMeta::Column) } } } @@ -571,6 +563,43 @@ impl Fragment { .chain(overlays.iter_mut().map(|overlay| &mut overlay.data_file)) } + /// Whether any of this fragment's row lineage sequences lives in a data + /// file column rather than inline. + pub fn has_spilled_row_lineage(&self) -> bool { + matches!(self.row_id_meta, Some(RowIdMeta::Column)) + || matches!( + self.created_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ) + || matches!( + self.last_updated_at_version_meta, + Some(RowDatasetVersionMeta::Column) + ) + } + + /// The data file holding the row lineage column with the reserved + /// `field_id`, which is the one entry of [`Self::files`] whose fields carry + /// it. `None` when no file does, which for a sequence whose metadata says + /// it is spilled is corruption; so is more than one file carrying the id, + /// which this reports as an error. + pub fn row_lineage_file(&self, field_id: i32) -> Result> { + let mut carriers = self + .files + .iter() + .filter(|file| file.fields.contains(&field_id)); + let file = carriers.next(); + if let Some(extra) = carriers.next() { + return Err(Error::corrupt_file_named( + &extra.path, + format!( + "fragment {} has more than one data file carrying row lineage field {}", + self.id, field_id + ), + )); + } + Ok(file) + } + pub fn from_json(json: &str) -> Result { let fragment: Self = serde_json::from_str(json)?; Ok(fragment) @@ -735,12 +764,8 @@ impl From<&Fragment> for pb::DataFragment { RowIdMeta::Inline(data) => { pb::data_fragment::RowIdSequence::InlineRowIds(data.bytes().clone()) } - RowIdMeta::External(file) => { - pb::data_fragment::RowIdSequence::ExternalRowIds(pb::ExternalFile { - path: file.path.clone(), - offset: file.offset, - size: file.size, - }) + RowIdMeta::Column => { + pb::data_fragment::RowIdSequence::ColumnRowIds(pb::RowLineageColumn {}) } }); let last_updated_at_version_sequence = diff --git a/rust/lance-table/src/format/row_ids.rs b/rust/lance-table/src/format/row_ids.rs index b70cf939759..a9a8a74e0ce 100644 --- a/rust/lance-table/src/format/row_ids.rs +++ b/rust/lance-table/src/format/row_ids.rs @@ -12,7 +12,21 @@ use serde::{Deserialize, Deserializer, Serialize, Serializer}; use super::pb; -/// A reference to a part of a file. +/// 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; + +/// A reference to a part of a file, used by the fragment reuse index details. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)] pub struct ExternalFile { pub path: String, @@ -138,7 +152,16 @@ impl DeepSizeOf for InlineRowIdsInner { #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)] pub enum RowIdMeta { Inline(InlineRowIds), - External(ExternalFile), + /// The sequence is spilled to a hidden [`ROW_ID_FIELD_ID`] column of one of + /// the fragment's data files, one row id per physical row, in offset order. + /// The file is the entry of [`Fragment::files`](super::Fragment::files) + /// whose fields carry that id; see + /// [`Fragment::row_lineage_file`](super::Fragment::row_lineage_file). + /// + /// An ordinary column rather than an opaque byte range: it carries the + /// file's encodings and page layout, so it can be read back a page at a + /// time instead of whole. + Column, } impl TryFrom for RowIdMeta { @@ -147,13 +170,7 @@ impl TryFrom for RowIdMeta { fn try_from(value: pb::data_fragment::RowIdSequence) -> Result { match value { pb::data_fragment::RowIdSequence::InlineRowIds(data) => Ok(Self::Inline(data.into())), - pb::data_fragment::RowIdSequence::ExternalRowIds(file) => { - Ok(Self::External(ExternalFile { - path: file.path.clone(), - offset: file.offset, - size: file.size, - })) - } + pb::data_fragment::RowIdSequence::ColumnRowIds(_) => Ok(Self::Column), } } } diff --git a/rust/lance-table/src/rowids/version.rs b/rust/lance-table/src/rowids/version.rs index 1e89db8326a..1b2ed5509d1 100644 --- a/rust/lance-table/src/rowids/version.rs +++ b/rust/lance-table/src/rowids/version.rs @@ -17,7 +17,7 @@ use serde::de::Deserializer; use serde::ser::Serializer; use serde::{Deserialize, Serialize}; -use crate::format::{ExternalFile, Fragment, pb}; +use crate::format::{Fragment, pb}; use crate::rowids::segment::U64Segment; use crate::rowids::{RowIdSequence, read_row_ids}; @@ -207,6 +207,28 @@ impl RowDatasetVersionSequence { Self { runs: vec![run] } } + /// Run-length encode one version per row, in row offset order. + pub fn from_versions(versions: &[u64]) -> Self { + let mut runs = Vec::new(); + let mut run_start = 0u64; + for (i, window) in versions.windows(2).enumerate() { + if window[0] != window[1] { + runs.push(RowDatasetVersionRun { + span: U64Segment::Range(run_start..i as u64 + 1), + version: window[0], + }); + run_start = i as u64 + 1; + } + } + if let Some(last) = versions.last() { + runs.push(RowDatasetVersionRun { + span: U64Segment::Range(run_start..versions.len() as u64), + version: *last, + }); + } + Self { runs } + } + /// Number of rows tracked by this sequence (sum of run lengths). pub fn len(&self) -> u64 { self.runs.iter().map(|s| s.len() as u64).sum() @@ -356,10 +378,25 @@ impl<'a> Iterator for VersionsIter<'a> { pub enum RowDatasetVersionMeta { /// Small sequences stored inline in the fragment metadata Inline(Arc<[u8]>), - /// Large sequences stored in external files - External(ExternalFile), + /// The sequence is spilled to a hidden column of one of the fragment's + /// data files, one version per physical row in offset order, at + /// [`ROW_CREATED_AT_VERSION_FIELD_ID`](crate::format::ROW_CREATED_AT_VERSION_FIELD_ID) + /// or + /// [`ROW_LAST_UPDATED_AT_VERSION_FIELD_ID`](crate::format::ROW_LAST_UPDATED_AT_VERSION_FIELD_ID) + /// depending on which sequence this is; see + /// [`Fragment::row_lineage_file`](crate::format::Fragment::row_lineage_file). + /// Reading it needs IO, so it goes through the dataset's loader rather + /// than [`Self::load_sequence`]. + Column, } +/// The JSON form of [`RowDatasetVersionMeta::Column`]: `{"column": {}}`, an +/// empty object under the arm's name, mirroring the empty `RowLineageColumn` +/// protobuf message. The name alone tells the arms apart; there is nothing to +/// carry. +#[derive(Serialize, Deserialize)] +struct ColumnMarker {} + // Custom Serialize: convert Arc<[u8]> to slice for transparent JSON output impl Serialize for RowDatasetVersionMeta { fn serialize(&self, serializer: S) -> std::result::Result { @@ -367,7 +404,7 @@ impl Serialize for RowDatasetVersionMeta { #[serde(untagged)] enum Helper<'a> { Inline { inline: &'a [u8] }, - External { external: &'a ExternalFile }, + Column { column: ColumnMarker }, } match self { @@ -375,7 +412,10 @@ impl Serialize for RowDatasetVersionMeta { inline: data.as_ref(), } .serialize(serializer), - Self::External(file) => Helper::External { external: file }.serialize(serializer), + Self::Column => Helper::Column { + column: ColumnMarker {}, + } + .serialize(serializer), } } } @@ -387,12 +427,15 @@ impl<'de> Deserialize<'de> for RowDatasetVersionMeta { #[serde(untagged)] enum Helper { Inline { inline: Vec }, - External { external: ExternalFile }, + Column { column: ColumnMarker }, } match Helper::deserialize(deserializer)? { Helper::Inline { inline } => Ok(Self::Inline(Arc::from(inline))), - Helper::External { external } => Ok(Self::External(external)), + Helper::Column { column } => { + let ColumnMarker {} = column; + Ok(Self::Column) + } } } } @@ -404,18 +447,18 @@ impl RowDatasetVersionMeta { Ok(Self::Inline(Arc::from(bytes))) } - /// Create external metadata reference - pub fn from_external_file(path: String, offset: u64, size: u64) -> Self { - Self::External(ExternalFile { path, offset, size }) - } - - /// Load the version sequence from this metadata + /// Decode the version sequence stored inline in this metadata. + /// + /// A sequence stored outside the manifest needs IO to read, which this + /// synchronous accessor cannot do; it is an error here and is loaded + /// through the dataset instead. pub fn load_sequence(&self) -> lance_core::Result { match self { Self::Inline(data) => read_dataset_versions(data), - Self::External(_file) => { - todo!("External file loading not yet implemented") - } + Self::Column => Err(Error::not_supported( + "row version sequence spilled to a data file column cannot be decoded from \ + fragment metadata alone", + )), } } } @@ -430,13 +473,9 @@ pub fn last_updated_at_version_meta_to_pb( data.to_vec(), ) } - RowDatasetVersionMeta::External(file) => { - pb::data_fragment::LastUpdatedAtVersionSequence::ExternalLastUpdatedAtVersions( - pb::ExternalFile { - path: file.path.clone(), - offset: file.offset, - size: file.size, - }, + RowDatasetVersionMeta::Column => { + pb::data_fragment::LastUpdatedAtVersionSequence::ColumnLastUpdatedAtVersions( + pb::RowLineageColumn {}, ) } }) @@ -450,13 +489,9 @@ pub fn created_at_version_meta_to_pb( RowDatasetVersionMeta::Inline(data) => { pb::data_fragment::CreatedAtVersionSequence::InlineCreatedAtVersions(data.to_vec()) } - RowDatasetVersionMeta::External(file) => { - pb::data_fragment::CreatedAtVersionSequence::ExternalCreatedAtVersions( - pb::ExternalFile { - path: file.path.clone(), - offset: file.offset, - size: file.size, - }, + RowDatasetVersionMeta::Column => { + pb::data_fragment::CreatedAtVersionSequence::ColumnCreatedAtVersions( + pb::RowLineageColumn {}, ) } }) @@ -637,8 +672,9 @@ pub fn refresh_row_latest_update_meta_for_full_frag_rewrite_cols( let sequence = read_row_ids(data).unwrap(); sequence.len() } - // Follow existing behavior: external sequence not yet supported here - crate::format::RowIdMeta::External(_file) => 0, + // Follow existing behavior: a sequence that is not inline needs IO + // to read, which this synchronous path cannot do. + crate::format::RowIdMeta::Column => 0, } } else { 0 @@ -674,9 +710,13 @@ pub fn refresh_row_latest_update_meta_for_partial_frag_rewrite_cols( let sequence = read_row_ids(data).unwrap(); sequence.len() } - crate::format::RowIdMeta::External(_file) => { - // Preserve original behavior for external sequences - todo!("External file loading not yet implemented") + // Reading these needs IO, which this synchronous path cannot do. + // Reachable only for a fragment that also has no `physical_rows`. + crate::format::RowIdMeta::Column => { + return Err(Error::not_supported( + "refreshing row update versions for a fragment whose row id \ + sequence is stored outside the manifest", + )); } } } else { @@ -687,6 +727,16 @@ pub fn refresh_row_latest_update_meta_for_partial_frag_rewrite_cols( // Build base version vector from existing meta or previous dataset version let mut base_versions: Vec = Vec::with_capacity(row_count_u64 as usize); if let Some(meta) = fragment.last_updated_at_version_meta.as_ref() { + if matches!(meta, RowDatasetVersionMeta::Column) { + // The existing versions of the rows this update leaves alone + // live in a data file, which this commit-time path cannot read. + // Defaulting them would silently rewrite their lineage. + return Err(Error::not_supported(format!( + "fragment {} stores its last-updated-at versions outside the manifest; \ + partially rewriting its columns is not supported yet", + fragment.id + ))); + } if let Ok(base_seq) = meta.load_sequence() { base_versions.extend(base_seq.versions().take(row_count_u64 as usize)); base_versions.resize(row_count_u64 as usize, prev_version); @@ -741,13 +791,9 @@ impl TryFrom for RowDatasetVers pb::data_fragment::LastUpdatedAtVersionSequence::InlineLastUpdatedAtVersions(data) => { Ok(Self::Inline(Arc::from(data))) } - pb::data_fragment::LastUpdatedAtVersionSequence::ExternalLastUpdatedAtVersions( - file, - ) => Ok(Self::External(ExternalFile { - path: file.path, - offset: file.offset, - size: file.size, - })), + pb::data_fragment::LastUpdatedAtVersionSequence::ColumnLastUpdatedAtVersions(_) => { + Ok(Self::Column) + } } } } @@ -760,12 +806,8 @@ impl TryFrom for RowDatasetVersionM pb::data_fragment::CreatedAtVersionSequence::InlineCreatedAtVersions(data) => { Ok(Self::Inline(Arc::from(data))) } - pb::data_fragment::CreatedAtVersionSequence::ExternalCreatedAtVersions(file) => { - Ok(Self::External(ExternalFile { - path: file.path, - offset: file.offset, - size: file.size, - })) + pb::data_fragment::CreatedAtVersionSequence::ColumnCreatedAtVersions(_) => { + Ok(Self::Column) } } } @@ -775,6 +817,43 @@ impl TryFrom for RowDatasetVersionM mod tests { use super::*; + #[test] + fn column_meta_json_is_an_empty_object_under_the_arm_name() { + let json = serde_json::to_string(&RowDatasetVersionMeta::Column).unwrap(); + assert_eq!(json, r#"{"column":{}}"#); + let parsed: RowDatasetVersionMeta = serde_json::from_str(&json).unwrap(); + assert_eq!(parsed, RowDatasetVersionMeta::Column); + let inline = RowDatasetVersionMeta::Inline(Arc::from(vec![1u8, 2, 3])); + let parsed: RowDatasetVersionMeta = + serde_json::from_str(&serde_json::to_string(&inline).unwrap()).unwrap(); + assert_eq!(parsed, inline); + } + + #[test] + fn from_versions_run_length_encodes() { + assert!( + RowDatasetVersionSequence::from_versions(&[]) + .runs + .is_empty() + ); + + let single = RowDatasetVersionSequence::from_versions(&[3, 3, 3]); + assert_eq!(single.runs.len(), 1); + assert_eq!(single.runs[0].version, 3); + assert_eq!(single.len(), 3); + + let alternating = RowDatasetVersionSequence::from_versions(&[1, 2, 1, 2]); + assert_eq!( + alternating + .runs + .iter() + .map(|run| run.version) + .collect::>(), + [1, 2, 1, 2] + ); + assert_eq!(alternating.versions().collect::>(), [1, 2, 1, 2]); + } + #[test] fn test_version_random_access() { let seq = RowDatasetVersionSequence { diff --git a/rust/lance-table/src/transaction/row_version.rs b/rust/lance-table/src/transaction/row_version.rs index 71c6229aa46..2ae085b9dc3 100644 --- a/rust/lance-table/src/transaction/row_version.rs +++ b/rust/lance-table/src/transaction/row_version.rs @@ -9,10 +9,7 @@ //! `created_at` correct across an update means tracing each new row back to the //! fragment and offset it came from, which is what most of this module does. -use crate::format::{ - Fragment, RowDatasetVersionMeta, RowDatasetVersionRun, RowDatasetVersionSequence, RowIdMeta, -}; -use crate::rowids::segment::U64Segment; +use crate::format::{Fragment, RowDatasetVersionMeta, RowDatasetVersionSequence, RowIdMeta}; use crate::rowids::version::build_version_meta; use crate::rowids::{RowIdSequence, read_row_ids, write_row_ids}; use crate::transaction::Transaction; @@ -89,9 +86,22 @@ pub(super) fn resolve_update_version_metadata( let mut sorted_frags: Vec<&Fragment> = existing_fragments.iter().collect(); sorted_frags.sort_by_key(|f| f.id); for frag in sorted_frags { - if let Some(RowIdMeta::Inline(data)) = &frag.row_id_meta - && let Ok(seq) = read_row_ids(data) - { + // Finding which existing row each rewritten row came from means + // scanning every existing fragment's row ids, and a spilled sequence + // cannot be read here. Skipping such a fragment would make its rows + // look freshly inserted and stamp them with a new created-at version. + let seq = match &frag.row_id_meta { + Some(RowIdMeta::Inline(data)) => read_row_ids(data).ok(), + Some(RowIdMeta::Column) => { + return Err(Error::not_supported(format!( + "fragment {} stores its row ids outside the manifest; updating rows \ + of a table with spilled row ids is not supported yet", + frag.id + ))); + } + None => None, + }; + if let Some(seq) = seq { // Range pre-filter: skip the per-row inner loop when the fragment's // bounding row-id range has no overlap with [needed_min, needed_max]. // row_id_range() returns None for empty sequences, which are also skipped. @@ -118,29 +128,40 @@ pub(super) fn resolve_update_version_metadata( // (a protobuf decode) for every single updated row, even when many rows originate // from the same fragment. let source_frag_ids: HashSet = row_id_to_source.values().map(|(f, _)| f.id).collect(); - let version_cache: HashMap = existing_fragments + let mut version_cache: HashMap = HashMap::new(); + for frag in existing_fragments .iter() .filter(|f| source_frag_ids.contains(&f.id)) - .filter_map(|frag| { - let seq = frag - .created_at_version_meta - .as_ref()? - .load_sequence() - .ok()?; - Some((frag.id, seq)) - }) - .collect(); + { + let Some(meta) = &frag.created_at_version_meta else { + continue; + }; + if matches!(meta, RowDatasetVersionMeta::Column) { + // The rewritten rows' original created-at versions live in a data + // file, which this commit-time path cannot read. Defaulting them + // would silently rewrite their lineage. + return Err(Error::not_supported(format!( + "fragment {} stores its created-at versions outside the manifest; \ + updating rows it holds is not supported yet", + frag.id + ))); + } + if let Ok(seq) = meta.load_sequence() { + version_cache.insert(frag.id, seq); + } + } for fragment in new_fragments.iter_mut() { let row_ids = match &fragment.row_id_meta { Some(RowIdMeta::Inline(data)) => read_row_ids(data).ok(), - Some(RowIdMeta::External(_)) => { - log::warn!( - "Fragment {} has external row ID metadata; \ - version tracking will use defaults", - fragment.id, - ); - None + Some(RowIdMeta::Column) => { + // Resolving the versions needs the row ids, which this + // commit-time path cannot read back from a data file. + return Err(Error::not_supported(format!( + "fragment {} stores its row ids outside the manifest; committing it \ + through an update is not supported yet", + fragment.id + ))); } None => None, }; @@ -165,8 +186,7 @@ pub(super) fn resolve_update_version_metadata( .collect(); debug_assert_eq!(created_at_versions.len(), physical_rows); - let runs = encode_version_runs(&created_at_versions); - let created_at_seq = RowDatasetVersionSequence { runs }; + let created_at_seq = RowDatasetVersionSequence::from_versions(&created_at_versions); fragment.created_at_version_meta = Some( RowDatasetVersionMeta::from_sequence(&created_at_seq).map_err(|e| { Error::internal(format!( @@ -186,31 +206,6 @@ pub(super) fn resolve_update_version_metadata( Ok(()) } -/// Run-length encode a sequence of per-row versions into [`RowDatasetVersionRun`]s. -fn encode_version_runs(versions: &[u64]) -> Vec { - if versions.is_empty() { - return Vec::new(); - } - let mut runs = Vec::new(); - let mut current_version = versions[0]; - let mut run_start = 0u64; - for (i, &version) in versions.iter().enumerate().skip(1) { - if version != current_version { - runs.push(RowDatasetVersionRun { - span: U64Segment::Range(run_start..i as u64), - version: current_version, - }); - current_version = version; - run_start = i as u64; - } - } - runs.push(RowDatasetVersionRun { - span: U64Segment::Range(run_start..versions.len() as u64), - version: current_version, - }); - runs -} - impl Transaction { /// collect the pure(the num of row IDs are equal to the physical rows) "rewrite rows" updated fragment ids pub(super) fn collect_pure_rewrite_row_update_frags_ids( @@ -230,7 +225,8 @@ impl Transaction { let sequence = read_row_ids(data)?; sequence.len() as u64 } - _ => 0, + // A spilled sequence always covers every physical row. + RowIdMeta::Column => physical_rows, }; // only filter the fragments that match: all the rows have row id, @@ -264,6 +260,8 @@ impl Transaction { let sequence = read_row_ids(data)?; sequence.len() as u64 } + // A spilled sequence always covers every physical row. + Some(RowIdMeta::Column) => physical_rows, _ => 0, }; @@ -321,6 +319,8 @@ impl Transaction { #[cfg(test)] mod tests { use super::*; + use crate::format::RowDatasetVersionRun; + use crate::rowids::segment::U64Segment; use crate::transaction::test_support::{ created_at_versions, default_build_config, last_updated_at_versions, make_stable_row_id_manifest, update_txn, @@ -357,6 +357,92 @@ mod tests { } } + fn inline_versions(rows: u64, version: u64) -> RowDatasetVersionMeta { + RowDatasetVersionMeta::from_sequence(&RowDatasetVersionSequence::from_uniform_row_count( + rows, version, + )) + .unwrap() + } + + #[test] + fn test_assign_row_ids_spilled_is_complete() { + // A spilled sequence covers every physical row, so the commit must not + // top it up with fresh ids the way it does for a partial inline one. + let spilled = RowIdMeta::Column; + let mut fragments = vec![Fragment { + id: 1, + physical_rows: Some(50), + row_id_meta: Some(spilled.clone()), + files: vec![], + overlays: vec![], + deletion_file: None, + last_updated_at_version_meta: None, + created_at_version_meta: None, + }]; + let mut next_row_id = 100; + + Transaction::assign_row_ids(&mut next_row_id, &mut fragments).unwrap(); + + assert_eq!(next_row_id, 100); + assert_eq!(fragments[0].row_id_meta, Some(spilled)); + } + + #[test] + fn test_resolve_update_versions_refuses_spilled_lineage() { + // Resolving lineage here means reading row ids and versions back, which + // this commit-time path cannot do for a spilled sequence. Skipping such + // a fragment would make its rows look freshly inserted, so it refuses. + let inline_ids = |range: std::ops::Range| { + Some(RowIdMeta::Inline( + write_row_ids(&RowIdSequence::from(range)).into(), + )) + }; + let fragment = |id: u64, rows: usize| Fragment { + id, + physical_rows: Some(rows), + row_id_meta: None, + files: vec![], + overlays: vec![], + deletion_file: None, + last_updated_at_version_meta: None, + created_at_version_meta: None, + }; + + // An existing fragment with spilled row ids, when the update rewrote rows. + let existing = vec![Fragment { + row_id_meta: Some(RowIdMeta::Column), + created_at_version_meta: Some(inline_versions(50, 2)), + ..fragment(1, 50) + }]; + let mut new_fragments = vec![Fragment { + row_id_meta: inline_ids(10..12), + ..fragment(2, 2) + }]; + let error = resolve_update_version_metadata(&existing, &mut new_fragments, 9).unwrap_err(); + assert!(matches!(error, Error::NotSupported { .. }), "{error:?}"); + + // An existing fragment with spilled created-at versions. + let existing = vec![Fragment { + row_id_meta: inline_ids(0..50), + created_at_version_meta: Some(RowDatasetVersionMeta::Column), + ..fragment(1, 50) + }]; + let mut new_fragments = vec![Fragment { + row_id_meta: inline_ids(10..12), + ..fragment(2, 2) + }]; + let error = resolve_update_version_metadata(&existing, &mut new_fragments, 9).unwrap_err(); + assert!(matches!(error, Error::NotSupported { .. }), "{error:?}"); + + // A new fragment whose own row ids are spilled. + let mut new_fragments = vec![Fragment { + row_id_meta: Some(RowIdMeta::Column), + ..fragment(2, 50) + }]; + let error = resolve_update_version_metadata(&[], &mut new_fragments, 9).unwrap_err(); + assert!(matches!(error, Error::NotSupported { .. }), "{error:?}"); + } + #[test] fn test_assign_row_ids_existing_complete() { // Test with fragment that already has complete row IDs @@ -1113,27 +1199,4 @@ mod tests { // Row 12 → frag A offset 2 → version 2; row 20 → frag B offset 0 → version 8 assert_eq!(created_at_versions(&result, 10), vec![2, 8]); } - - #[test] - fn test_encode_version_runs_empty() { - let runs = encode_version_runs(&[]); - assert!(runs.is_empty()); - } - - #[test] - fn test_encode_version_runs_single_run() { - let runs = encode_version_runs(&[3, 3, 3]); - assert_eq!(runs.len(), 1); - assert_eq!(runs[0].version, 3); - } - - #[test] - fn test_encode_version_runs_alternating() { - let runs = encode_version_runs(&[1, 2, 1, 2]); - assert_eq!(runs.len(), 4); - assert_eq!(runs[0].version, 1); - assert_eq!(runs[1].version, 2); - assert_eq!(runs[2].version, 1); - assert_eq!(runs[3].version, 2); - } } diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index 2a0bff656fa..7e6172ecad0 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -3451,12 +3451,6 @@ impl Dataset { let mut file_paths: Vec<(String, Path)> = Vec::new(); let mut blob_dirs = HashSet::new(); for fragment in self.manifest.fragments.iter() { - if let Some(RowIdMeta::External(external_file)) = &fragment.row_id_meta { - return Err(Error::internal(format!( - "External row_id_meta is not supported yet. external file path: {}", - external_file.path - ))); - } for data_file in fragment.referenced_lance_files() { let base_root = if let Some(base_id) = data_file.base_id { let base_path = diff --git a/rust/lance/src/dataset/fragment.rs b/rust/lance/src/dataset/fragment.rs index 8153f45bb4a..37876dc9d07 100644 --- a/rust/lance/src/dataset/fragment.rs +++ b/rust/lance/src/dataset/fragment.rs @@ -51,8 +51,7 @@ use lance_io::ReadBatchParams; use lance_io::scheduler::{FileScheduler, ScanScheduler, SchedulerConfig}; use lance_io::stream::RecordBatchStream; use lance_io::utils::CachedFileSize; -use lance_table::format::overlay::TOMBSTONE_FIELD_ID; -use lance_table::format::{DataFile, DeletionFile, Fragment}; +use lance_table::format::{DataFile, DeletionFile, Fragment, RowDatasetVersionMeta}; use lance_table::io::deletion::{deletion_file_path, write_deletion_file}; use lance_table::rowids::RowIdSequence; use lance_table::utils::stream::{ @@ -974,6 +973,28 @@ 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 (opened_files, deletion_vec, row_id_sequence) = join!(open_files, deletion_vec_load, row_id_load); let opened_files = opened_files?; @@ -1519,9 +1540,11 @@ impl FileFragment { for data_file in &self.metadata.files { let last = -1; for field_id in data_file.fields.iter() { - // A tombstone marks a field superseded by a later data file. - // It is not a field id: it has no ordering and can repeat. - if *field_id == TOMBSTONE_FIELD_ID { + // Negative ids are not schema fields: the tombstone marks a + // field superseded by a later data file, and the others are + // hidden system columns such as spilled row lineage. None has + // an ordering, and a tombstone can repeat. + if *field_id < 0 { continue; } if *field_id <= last { diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index 9632d338035..7a6eb120b31 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -2739,7 +2739,7 @@ async fn recalc_versions_for_rewritten_fragments( let row_count = if let Some(row_id_meta) = &frag.row_id_meta { match row_id_meta { RowIdMeta::Inline(data) => lance_table::rowids::read_row_ids(data)?.len(), - RowIdMeta::External(_file) => frag.physical_rows.unwrap_or(0) as u64, + RowIdMeta::Column => frag.physical_rows.unwrap_or(0) as u64, } } else { frag.physical_rows.unwrap_or(0) as u64 diff --git a/rust/lance/src/dataset/rowids.rs b/rust/lance/src/dataset/rowids.rs index b09e4933d5c..ea9035746ae 100644 --- a/rust/lance/src/dataset/rowids.rs +++ b/rust/lance/src/dataset/rowids.rs @@ -11,7 +11,7 @@ use futures::{Stream, StreamExt, TryFutureExt, TryStreamExt}; use lance_core::utils::{address::RowAddress, deletion::DeletionVector}; use lance_select::{RowAddrSelection, RowAddrTreeMap}; use lance_table::{ - format::{Fragment, RowIdMeta}, + format::{Fragment, ROW_ID_FIELD_ID, RowIdMeta}, rowids::{FragmentRowIdIndex, RowIdIndex, RowIdSequence, read_row_ids}, }; use std::sync::Arc; @@ -29,30 +29,24 @@ pub async fn load_row_id_sequence( let key = RowIdSequenceKey { fragment_id: fragment.id, row_id_meta, + lineage_file: fragment.row_lineage_file(ROW_ID_FIELD_ID)?, }; dataset .metadata_cache - .get_or_insert_with_key(key, || read_row_id_sequence(dataset, fragment)) + .get_or_insert_with_key(key, || read_row_id_sequence(fragment)) .await } /// Decode the row id sequence of `fragment`, bypassing every cache. -async fn read_row_id_sequence(dataset: &Dataset, fragment: &Fragment) -> Result { +async fn read_row_id_sequence(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::External(file_slice)) => { - let path = dataset.base.clone().join(file_slice.path.as_str()); - let range = - file_slice.offset as usize..(file_slice.offset as usize + file_slice.size as usize); - let data = dataset - .object_store - .open(&path) - .await? - .get_range(range) - .await?; - 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 + ))), } } @@ -258,7 +252,7 @@ async fn read_fragment_row_id_index( dataset: &Dataset, fragment: &Fragment, ) -> Result { - let row_id_sequence = Arc::new(read_row_id_sequence(dataset, fragment).await?); + let row_id_sequence = Arc::new(read_row_id_sequence(fragment).await?); let deletion_vector = match &fragment.deletion_file { None => Arc::new(DeletionVector::default()), Some(deletion_file) => { diff --git a/rust/lance/src/session/caches.rs b/rust/lance/src/session/caches.rs index bcb05d0291d..88b1556ad99 100644 --- a/rust/lance/src/session/caches.rs +++ b/rust/lance/src/session/caches.rs @@ -19,7 +19,7 @@ use lance_core::{ }; use lance_select::RowAddrMask; use lance_table::{ - format::{DeletionFile, DeletionFileType, Manifest, RowIdMeta}, + format::{DataFile, DeletionFile, DeletionFileType, Manifest, RowIdMeta}, rowids::{RowIdIndex, RowIdSequence}, }; use object_store::path::Path; @@ -259,6 +259,9 @@ pub struct RowIdSequenceKey<'a> { /// which those bytes memoize on first use — an array-encoded sequence is /// 8 bytes per row, too much to rehash on every lookup. pub row_id_meta: &'a RowIdMeta, + /// The data file the sequence is spilled to, when `row_id_meta` says it + /// is one; identifies the contents the way an inline digest does. + pub lineage_file: Option<&'a DataFile>, } impl CacheKey for RowIdSequenceKey<'_> { @@ -282,11 +285,18 @@ impl CacheKey for RowIdSequenceKey<'_> { builder.write_variant(0); builder.write_fixed_bytes(data.digest()); } - RowIdMeta::External(file) => { + // The sequence lives in one of the fragment's data files, which is + // named freshly per rewrite; the file identifies the contents the + // way the inline digest does. + RowIdMeta::Column => { builder.write_variant(1); - builder.write_str(&file.path); - builder.write_u64(file.offset); - builder.write_u64(file.size); + match self.lineage_file { + Some(file) => { + builder.write_str(&file.path); + builder.write_u64(file.base_id.map_or(u64::MAX, u64::from)); + } + None => builder.write_str(""), + } } } } @@ -304,7 +314,6 @@ impl DSMetadataCache { mod tests { use std::sync::Arc; - use lance_table::format::ExternalFile; use lance_table::rowids::write_row_ids; use super::*; @@ -353,6 +362,7 @@ mod tests { let key = RowIdSequenceKey { fragment_id: 0, row_id_meta: &first_generation, + lineage_file: None, }; cache .insert_with_key(&key, Arc::new(RowIdSequence::from(0..100))) @@ -365,52 +375,7 @@ mod tests { .get_with_key(&RowIdSequenceKey { fragment_id: 0, row_id_meta: &second_generation, - }) - .await - .is_none() - ); - } - - #[tokio::test] - async fn row_id_sequence_key_separates_external_slices() { - // External metadata is a read-only legacy shape, but the same slice of - // the same file is the only thing that may share a cache entry. - let cache = LanceCache::with_capacity(4096); - let external = |offset| { - RowIdMeta::External(ExternalFile { - path: "_row_ids/1.rowids".into(), - offset, - size: 16, - }) - }; - let first_slice = external(0); - cache - .insert_with_key( - &RowIdSequenceKey { - fragment_id: 0, - row_id_meta: &first_slice, - }, - Arc::new(RowIdSequence::from(0..100)), - ) - .await; - - let second_slice = external(16); - assert!( - cache - .get_with_key(&RowIdSequenceKey { - fragment_id: 0, - row_id_meta: &second_slice, - }) - .await - .is_none() - ); - // An inline sequence never aliases an external one. - let inline = RowIdMeta::Inline(write_row_ids(&(0..100).into()).into()); - assert!( - cache - .get_with_key(&RowIdSequenceKey { - fragment_id: 0, - row_id_meta: &inline, + lineage_file: None, }) .await .is_none()