Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions rust/lance-table/src/format/manifest.rs
Original file line number Diff line number Diff line change
Expand Up @@ -750,6 +750,13 @@ pub struct ManifestBuildConfig {
/// It bypasses the "cannot enable stable row ids on existing dataset" guard and
/// sets `manifest.next_row_id` to the provided value before activating the flag.
pub migration_next_row_id: Option<u64>,
/// Row lineage sequences of the current manifest's fragments that live
/// outside the manifest, read ahead of the build. An update that rewrites
/// rows needs the existing row ids and created-at versions to carry each
/// row's lineage over, and a partial column rewrite needs the existing
/// last-updated-at versions; the build cannot read a data file itself. Only
/// consulted for fragments whose sequences are spilled.
pub spilled_row_lineage: std::sync::Arc<crate::rowids::version::SpilledRowLineage>,
}

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
Expand Down
54 changes: 44 additions & 10 deletions rust/lance-table/src/rowids/version.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
//! update version for each row in a Lance dataset, enabling efficient
//! cross-version diff operations.

use std::collections::HashMap;
use std::{ops::Range, sync::Arc};

use lance_core::Error;
Expand All @@ -21,6 +22,25 @@ use crate::format::{Fragment, pb};
use crate::rowids::segment::U64Segment;
use crate::rowids::{RowIdSequence, read_row_ids};

/// One fragment's row lineage sequences that live outside the manifest, read
/// ahead of a commit by the caller.
///
/// Building a manifest is synchronous and has no object store, so it cannot
/// read a sequence spilled to a data file column (`RowIdMeta::Column`,
/// `RowDatasetVersionMeta::Column`). A caller that can do IO loads them first
/// and hands them over in [`ManifestBuildConfig`](crate::format::ManifestBuildConfig).
/// Only the spilled sequences need to be present; inline ones are decoded on
/// the spot.
#[derive(Debug, Clone, Default)]
pub struct LoadedRowLineage {
pub row_ids: Option<Arc<RowIdSequence>>,
pub created_at: Option<Arc<RowDatasetVersionSequence>>,
pub last_updated_at: Option<Arc<RowDatasetVersionSequence>>,
}

/// [`LoadedRowLineage`] per fragment id.
pub type SpilledRowLineage = HashMap<u64, LoadedRowLineage>;

/// A run of identical versions over a contiguous span of row positions.
///
/// Span is expressed as a U64Segment over row offsets (0..N within a fragment),
Expand Down Expand Up @@ -695,11 +715,15 @@ pub fn refresh_row_latest_update_meta_for_full_frag_rewrite_cols(
/// `updated_offsets` are local row offsets (within the fragment) that have been updated.
/// Existing version metadata is preserved and only the updated positions are set to `current_version`.
/// If no existing metadata is present, positions default to `prev_version`.
///
/// A fragment whose existing versions are spilled has them read from
/// `spilled`; the refreshed sequence is placed inline.
pub fn refresh_row_latest_update_meta_for_partial_frag_rewrite_cols(
fragment: &mut Fragment,
updated_offsets: &[usize],
current_version: u64,
prev_version: u64,
spilled: &SpilledRowLineage,
) -> Result<()> {
// Determine row count for fragment
let row_count_u64: u64 = if let Some(pr) = fragment.physical_rows {
Expand Down Expand Up @@ -727,17 +751,26 @@ 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<u64> = 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) {
let base_seq = 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() {
// live in a data file, which this path cannot read. Defaulting
// them would silently rewrite their lineage, so the caller has
// to have read them ahead of time.
let loaded = spilled
.get(&fragment.id)
.and_then(|lineage| lineage.last_updated_at.clone())
.ok_or_else(|| {
Error::not_supported(format!(
"fragment {} stores its last-updated-at versions outside the \
manifest and they were not loaded ahead of the commit",
fragment.id
))
})?;
Some(loaded.as_ref().clone())
} else {
meta.load_sequence().ok()
};
if let Some(base_seq) = base_seq {
base_versions.extend(base_seq.versions().take(row_count_u64 as usize));
base_versions.resize(row_count_u64 as usize, prev_version);
} else {
Expand Down Expand Up @@ -901,6 +934,7 @@ mod tests {
&[ROWS - 1],
3,
1,
&Default::default(),
)
.unwrap();

Expand Down
3 changes: 3 additions & 0 deletions rust/lance-table/src/transaction/manifest_build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -912,6 +912,7 @@ impl Transaction {
&offsets,
new_version,
prev_version,
&config.spilled_row_lineage,
)?;
}
}
Expand All @@ -934,6 +935,7 @@ impl Transaction {
existing_fragments,
new_fragments.as_mut_slice(),
new_version,
&config.spilled_row_lineage,
)?;
}

Expand Down Expand Up @@ -1498,6 +1500,7 @@ impl Transaction {
&covered_offsets,
new_version,
1,
&config.spilled_row_lineage,
)?;
}
}
Expand Down
155 changes: 96 additions & 59 deletions rust/lance-table/src/transaction/row_version.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,11 @@
//! fragment and offset it came from, which is what most of this module does.

use crate::format::{Fragment, RowDatasetVersionMeta, RowDatasetVersionSequence, RowIdMeta};
use crate::rowids::version::build_version_meta;
use crate::rowids::version::{SpilledRowLineage, build_version_meta};
use crate::rowids::{RowIdSequence, read_row_ids, write_row_ids};
use crate::transaction::Transaction;
use lance_core::{Error, Result};
use std::borrow::Cow;
use std::cmp::Ordering;
use std::collections::{HashMap, HashSet};

Expand Down Expand Up @@ -52,10 +53,15 @@ fn resolve_created_at_version(

/// For each new fragment produced by an update, set `created_at_version_meta`
/// (preserved from the original rows) and `last_updated_at_version_meta`.
///
/// `spilled` supplies the row ids and created-at versions of existing
/// fragments that keep them outside the manifest; see
/// [`SpilledRowLineage`].
pub(super) fn resolve_update_version_metadata(
existing_fragments: &[Fragment],
new_fragments: &mut [Fragment],
new_version: u64,
spilled: &SpilledRowLineage,
) -> Result<()> {
// Collect only the row IDs we actually need to resolve, those appearing in new_fragments
// with inline metadata. This bounds the lookup map to O(updated rows) instead of O(all dataset rows)
Expand Down Expand Up @@ -91,14 +97,20 @@ pub(super) fn resolve_update_version_metadata(
// 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
)));
}
Some(RowIdMeta::Inline(data)) => read_row_ids(data).ok().map(Cow::Owned),
Some(RowIdMeta::Column) => Some(Cow::Borrowed(
spilled
.get(&frag.id)
.and_then(|lineage| lineage.row_ids.as_deref())
.ok_or_else(|| {
Error::not_supported(format!(
"fragment {} stores its row ids outside the manifest and they \
were not loaded ahead of the commit; cannot resolve which rows \
this update rewrote",
frag.id
))
})?,
)),
None => None,
};
if let Some(seq) = seq {
Expand Down Expand Up @@ -138,15 +150,20 @@ pub(super) fn resolve_update_version_metadata(
};
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() {
// file, which this path cannot read. Defaulting them would silently
// rewrite their lineage, so the caller has to have read them.
let loaded = spilled
.get(&frag.id)
.and_then(|lineage| lineage.created_at.clone())
.ok_or_else(|| {
Error::not_supported(format!(
"fragment {} stores its created-at versions outside the manifest \
and they were not loaded ahead of the commit",
frag.id
))
})?;
version_cache.insert(frag.id, loaded.as_ref().clone());
} else if let Ok(seq) = meta.load_sequence() {
version_cache.insert(frag.id, seq);
}
}
Expand Down Expand Up @@ -388,59 +405,79 @@ mod tests {
}

#[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<u64>| {
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,
fn test_resolve_update_versions_reads_spilled_source_lineage_from_config() {
// The rewritten rows' sources are found by scanning existing row ids,
// which a spilled fragment does not expose here: the commit's caller
// reads them ahead of time. Without that, skipping the fragment would
// make its rows look freshly inserted, so it is an error instead.
let existing = vec![Fragment {
id: 1,
physical_rows: Some(50),
row_id_meta: Some(RowIdMeta::Column),
files: vec![],
overlays: vec![],
deletion_file: None,
last_updated_at_version_meta: Some(inline_versions(50, 2)),
created_at_version_meta: Some(RowDatasetVersionMeta::Column),
}];
let new_fragment = |ids: std::ops::Range<u64>| Fragment {
id: 2,
physical_rows: Some(2),
row_id_meta: Some(RowIdMeta::Inline(
write_row_ids(&RowIdSequence::from(ids)).into(),
)),
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();
let mut new_fragments = vec![new_fragment(10..12)];
let error =
resolve_update_version_metadata(&existing, &mut new_fragments, 9, &Default::default())
.unwrap_err();
assert!(matches!(error, Error::NotSupported { .. }), "{error:?}");

// A new fragment whose own row ids are spilled.
let mut new_fragments = vec![Fragment {
// A new fragment whose own row ids are spilled cannot be resolved.
let mut spilled_new = vec![Fragment {
row_id_meta: Some(RowIdMeta::Column),
..fragment(2, 50)
..new_fragment(0..2)
}];
let error = resolve_update_version_metadata(&[], &mut new_fragments, 9).unwrap_err();
let error = resolve_update_version_metadata(&[], &mut spilled_new, 9, &Default::default())
.unwrap_err();
assert!(matches!(error, Error::NotSupported { .. }), "{error:?}");

// Fragment 1 holds ids 100..150 and rows 10 and 11 were created at
// versions 3 and 4 respectively; that is what the rewrite must keep.
let mut created_at_versions = vec![2; 50];
created_at_versions[10] = 3;
created_at_versions[11] = 4;
let spilled = SpilledRowLineage::from([(
1,
crate::rowids::version::LoadedRowLineage {
row_ids: Some(Arc::new(RowIdSequence::from(100..150))),
created_at: Some(Arc::new(RowDatasetVersionSequence::from_versions(
&created_at_versions,
))),
last_updated_at: None,
},
)]);
let mut new_fragments = vec![new_fragment(110..112)];
resolve_update_version_metadata(&existing, &mut new_fragments, 9, &spilled).unwrap();
let versions = |meta: &Option<RowDatasetVersionMeta>| {
meta.as_ref()
.unwrap()
.load_sequence()
.unwrap()
.versions()
.collect::<Vec<_>>()
};
assert_eq!(versions(&new_fragments[0].created_at_version_meta), [3, 4]);
assert_eq!(
versions(&new_fragments[0].last_updated_at_version_meta),
[9, 9]
);
}

#[test]
Expand Down
1 change: 1 addition & 0 deletions rust/lance-table/src/transaction/test_support.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ pub fn default_build_config() -> ManifestBuildConfig {
storage_format: None,
disable_transaction_file: false,
migration_next_row_id: None,
spilled_row_lineage: Default::default(),
}
}

Expand Down
1 change: 1 addition & 0 deletions rust/lance/src/dataset.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4155,6 +4155,7 @@ impl ManifestWriteConfig {
storage_format: self.storage_format.clone(),
disable_transaction_file: self.disable_transaction_file,
migration_next_row_id: self.migration_next_row_id,
spilled_row_lineage: Default::default(),
}
}
}
Expand Down
Loading
Loading