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
111 changes: 63 additions & 48 deletions rust/lance/src/dataset/fragment.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand All @@ -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;
Expand Down Expand Up @@ -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!(
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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())
Expand All @@ -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
Expand Down Expand Up @@ -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<Arc<lance_table::rowids::version::RowDatasetVersionSequence>>,
) -> &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
Expand All @@ -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<Arc<lance_table::rowids::version::RowDatasetVersionSequence>>,
) -> &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
Expand Down
97 changes: 87 additions & 10 deletions rust/lance/src/dataset/rowids.rs
Original file line number Diff line number Diff line change
@@ -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.
Expand All @@ -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<RowIdSequence> {
async fn read_row_id_sequence(dataset: &Dataset, fragment: &Fragment) -> Result<RowIdSequence> {
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<Option<Arc<RowDatasetVersionSequence>>> {
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))),
}
}

Expand Down Expand Up @@ -252,7 +329,7 @@ async fn read_fragment_row_id_index(
dataset: &Dataset,
fragment: &Fragment,
) -> Result<FragmentRowIdIndex> {
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) => {
Expand Down
Loading
Loading