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
64 changes: 63 additions & 1 deletion rust/lance-core/src/datatypes/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -347,10 +347,21 @@ impl Schema {
}
}

// A negative id is reserved for system use; the only ones a schema may
// carry are the hidden row lineage columns a data file stores next to
// the user columns, as top-level fields under their own names. A
// nested field reusing one of these ids fails the duplicate check.
let row_lineage_ids = self
.fields
.iter()
.filter(|field| crate::row_lineage_field_id(&field.name) == Some(field.id))
.map(|field| field.id)
.collect::<HashSet<_>>();

// Check for duplicate field ids
let mut seen_ids = HashSet::new();
for field in self.fields_pre_order() {
if field.id < 0 {
if field.id < 0 && !row_lineage_ids.contains(&field.id) {
return Err(Error::schema(format!(
"Field {} has a negative id {}",
field.name, field.id
Expand Down Expand Up @@ -1920,6 +1931,57 @@ mod tests {
assert!(error.to_string().contains("Duplicate field id 0"));
}

#[test]
fn test_validate_admits_row_lineage_ids_only_at_top_level() {
let uint64_field = |name: &str, id: i32| {
let mut field = Field::new_arrow(name, DataType::UInt64, false).unwrap();
field.id = id;
field
};
let mut key = Field::new_arrow("i", DataType::Int32, false).unwrap();
key.id = 0;

// A hidden lineage column: top-level, under its own name's reserved id.
let schema = Schema {
fields: vec![key.clone(), uint64_field(ROW_ID, crate::ROW_ID_FIELD_ID)],
metadata: HashMap::new(),
};
schema.validate().unwrap();

// Another lineage column's id does not belong to this name.
let schema = Schema {
fields: vec![
key.clone(),
uint64_field(ROW_ID, crate::ROW_CREATED_AT_VERSION_FIELD_ID),
],
metadata: HashMap::new(),
};
let error = schema.validate().unwrap_err();
assert!(matches!(error, Error::Schema { .. }), "{error}");
assert!(error.to_string().contains("negative id"), "{error}");

// A data file stores lineage columns only at the top level.
let mut parent = Field::new_arrow(
"s",
DataType::Struct(ArrowFields::from(vec![ArrowField::new(
ROW_CREATED_AT_VERSION,
DataType::UInt64,
false,
)])),
true,
)
.unwrap();
parent.id = 1;
parent.children[0].id = crate::ROW_CREATED_AT_VERSION_FIELD_ID;
let schema = Schema {
fields: vec![key, parent],
metadata: HashMap::new(),
};
let error = schema.validate().unwrap_err();
assert!(matches!(error, Error::Schema { .. }), "{error}");
assert!(error.to_string().contains("negative id"), "{error}");
}

#[test]
fn test_resolve_quoted_fields() {
// Test that top-level fields with dots are rejected during validation
Expand Down
27 changes: 27 additions & 0 deletions rust/lance-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,33 @@ pub static ROW_LAST_UPDATED_AT_VERSION_FIELD: LazyLock<ArrowField> =
pub static ROW_CREATED_AT_VERSION_FIELD: LazyLock<ArrowField> =
LazyLock::new(|| ArrowField::new(ROW_CREATED_AT_VERSION, DataType::UInt64, true));

/// Field id of the hidden `_rowid` column 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` the tombstone for a field superseded by
/// a later data file, and `-3..=-5` the three row lineage columns.
pub const ROW_ID_FIELD_ID: i32 = -3;
/// Field id of the hidden `_row_created_at_version` column 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 a spilled
/// last-updated-at version sequence lives in.
pub const ROW_LAST_UPDATED_AT_VERSION_FIELD_ID: i32 = -5;

/// The reserved field id of a row lineage column, by column name.
///
/// These are the only negative field ids a data file schema may carry: a
/// lineage sequence spilled into the fragment's own data file is written as a
/// column under this id, next to the user columns.
pub fn row_lineage_field_id(column_name: &str) -> Option<i32> {
match column_name {
ROW_ID => Some(ROW_ID_FIELD_ID),
ROW_CREATED_AT_VERSION => Some(ROW_CREATED_AT_VERSION_FIELD_ID),
ROW_LAST_UPDATED_AT_VERSION => Some(ROW_LAST_UPDATED_AT_VERSION_FIELD_ID),
_ => None,
}
}

/// Check if a column name is a system column.
///
/// System columns are virtual columns that are computed at read time and don't
Expand Down
16 changes: 3 additions & 13 deletions rust/lance-table/src/format/row_ids.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,19 +12,9 @@ use serde::{Deserialize, Deserializer, Serialize, Serializer};

use super::pb;

/// 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;
pub use lance_core::{
ROW_CREATED_AT_VERSION_FIELD_ID, ROW_ID_FIELD_ID, ROW_LAST_UPDATED_AT_VERSION_FIELD_ID,
};

/// A reference to a part of a file, used by the fragment reuse index details.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
Expand Down
35 changes: 26 additions & 9 deletions rust/lance/src/dataset/fragment.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1754,7 +1754,10 @@ impl FileFragment {
/// Verifies:
/// * All field ids in the fragment are distinct
/// * Within each data file, field ids are in increasing order
/// * All data files exist and have the same length
/// * All data files holding user data exist and have the same length. A
/// file kept only for the spilled row lineage it carries is not opened;
/// [`Dataset::validate`] reads that lineage back, which checks its
/// length.
/// * Field ids are distinct between data files.
/// * Deletion file exists and has rowids in the correct range
/// * `Fragment.physical_rows` matches length of file
Expand Down Expand Up @@ -1821,14 +1824,28 @@ impl FileFragment {
data_file.validate(&self.dataset.data_file_dir(data_file)?)?;
}

// 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));
// A file that holds no field of the dataset schema is not opened when
// it holds no user field at all, or when the fragment keeps it for a
// spilled row lineage sequence it carries -- a data file whose user
// columns were all dropped or replaced after compaction wrote the
// lineage next to them. The sequences it carries are checked against
// `physical_rows` when they are validated. Any other file is opened,
// so a file listing only user fields the schema does not have is
// still reported.
let schema_field_ids = self
.dataset
.schema()
.fields_pre_order()
.map(|field| field.id)
.collect::<HashSet<_>>();
let spilled_field_ids = self.metadata.spilled_row_lineage_field_ids();
let user_data_files = self.metadata.files.iter().filter(|data_file| {
let fields = &data_file.fields;
let holds_schema_field = fields.iter().any(|id| schema_field_ids.contains(id));
let holds_user_field = fields.iter().any(|id| *id >= 0);
let holds_spilled_lineage = fields.iter().any(|id| spilled_field_ids.contains(id));
holds_schema_field || (holds_user_field && !holds_spilled_lineage)
});
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
Expand Down
Loading
Loading