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
60 changes: 59 additions & 1 deletion rust/lance-table/src/format/fragment.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,10 @@ 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::RowIdMeta;
use super::row_ids::{
ROW_CREATED_AT_VERSION_FIELD_ID, ROW_ID_FIELD_ID, ROW_LAST_UPDATED_AT_VERSION_FIELD_ID,
RowIdMeta,
};
use crate::format::pb;

use crate::rowids::version::{
Expand Down Expand Up @@ -600,6 +603,61 @@ impl Fragment {
Ok(file)
}

/// The reserved field ids of the row lineage sequences this fragment
/// marks as spilled: [`ROW_ID_FIELD_ID`] when its row ids are, and
/// likewise [`ROW_CREATED_AT_VERSION_FIELD_ID`] and
/// [`ROW_LAST_UPDATED_AT_VERSION_FIELD_ID`] for its versions.
///
/// These ids are never in the dataset schema, so code that decides whether
/// a data file is still needed by its schema fields has to keep a file
/// carrying one of them as well: that file is the sequence's only copy.
pub fn spilled_row_lineage_field_ids(&self) -> Vec<i32> {
let mut field_ids = Vec::new();
if matches!(self.row_id_meta, Some(RowIdMeta::Column)) {
field_ids.push(ROW_ID_FIELD_ID);
}
if matches!(
self.created_at_version_meta,
Some(RowDatasetVersionMeta::Column)
) {
field_ids.push(ROW_CREATED_AT_VERSION_FIELD_ID);
}
if matches!(
self.last_updated_at_version_meta,
Some(RowDatasetVersionMeta::Column)
) {
field_ids.push(ROW_LAST_UPDATED_AT_VERSION_FIELD_ID);
}
field_ids
}

/// Check that every sequence this fragment marks as spilled has exactly
/// one carrier among [`Self::files`], and that the carrier is a v2 file,
/// the only version that can hold the columns.
///
/// This reads metadata only. Committing a fragment that fails it would
/// publish lineage no reader can load, and once cleanup removed the
/// unreferenced carrier that lineage would be lost for good.
pub(crate) fn validate_row_lineage_carriers(&self) -> Result<()> {
for field_id in self.spilled_row_lineage_field_ids() {
let Some(file) = self.row_lineage_file(field_id)? else {
return Err(Error::internal(format!(
"cannot commit fragment {}: it marks row lineage field {} as spilled but \
none of its data files carries it",
self.id, field_id
)));
};
if file.file_version()? == ConcreteFileVersion::V1 {
return Err(Error::internal(format!(
"cannot commit fragment {}: its spilled row lineage field {} is carried by \
legacy v1 data file {}, which cannot hold row lineage columns",
self.id, field_id, file.path
)));
}
}
Ok(())
}

pub fn from_json(json: &str) -> Result<Self> {
let fragment: Self = serde_json::from_str(json)?;
Ok(fragment)
Expand Down
1 change: 1 addition & 0 deletions rust/lance-table/src/transaction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ pub use operation::{
TaggedRewriteAssembly, UpdateMode, UpdatedFragmentOffsets, reordered_sources,
};
pub use prepare::{FragReuseUpdate, PreparedIndices};
pub use row_version::has_writer_placed_lineage;
pub use update_map::{
UpdateMap, UpdateMapEntry, translate_config_updates, translate_schema_metadata_updates,
};
Expand Down
224 changes: 216 additions & 8 deletions rust/lance-table/src/transaction/manifest_build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1202,15 +1202,18 @@ impl Transaction {

// We might have removed all fields for certain data files, so
// we should remove the data files that are no longer relevant.
// A file carrying a spilled row lineage sequence stays: its
// reserved ids are never in the schema, and it is the only copy.
let remaining_field_ids = schema
.fields_pre_order()
.map(|f| f.id)
.collect::<HashSet<_>>();
for fragment in final_fragments.iter_mut() {
let spilled = fragment.spilled_row_lineage_field_ids();
fragment.files.retain(|file| {
file.fields
.iter()
.any(|field_id| remaining_field_ids.contains(field_id))
file.fields.iter().any(|field_id| {
remaining_field_ids.contains(field_id) || spilled.contains(field_id)
})
});
}

Expand Down Expand Up @@ -1367,14 +1370,22 @@ impl Transaction {
// the dataset schema: a file kept alive only by
// tombstones or by ids the schema no longer defines is
// unreachable to readers, uncollectable by cleanup,
// and reported corrupt by validate().
// and reported corrupt by validate(). The exception is
// a file carrying one of the fragment's spilled row
// lineage sequences: their reserved ids are never in
// the schema, the file holds their only copy, and the
// replaced user columns in it are dead space until
// compaction.
let live_ids = schema
.fields_pre_order()
.map(|field| field.id)
.collect::<HashSet<i32>>();
new_frag
.files
.retain(|file| file.fields.iter().any(|f| live_ids.contains(f)));
let spilled = new_frag.spilled_row_lineage_field_ids();
new_frag.files.retain(|file| {
file.fields
.iter()
.any(|f| live_ids.contains(f) || spilled.contains(f))
});
new_frag.files.push(new_file.clone());
}

Expand Down Expand Up @@ -1550,6 +1561,17 @@ impl Transaction {
}
}

// A spilled row lineage sequence lives only in its carrier file. An
// operation that dropped or duplicated that file would otherwise go
// unnoticed until the next read, by which time cleanup may have
// deleted the only copy. Every fragment is checked, not only the ones
// this operation touched: the check reads metadata only, skips
// fragments that spill nothing, and also refuses to build on a
// manifest that already lost a carrier.
for fragment in &final_fragments {
fragment.validate_row_lineage_carriers()?;
}

let user_requested_version = match (&config.storage_format, config.use_legacy_format) {
(Some(storage_format), _) => Some(storage_format.lance_file_format()),
(None, Some(true)) => Some(ConcreteFileVersion::V1),
Expand Down Expand Up @@ -1902,7 +1924,8 @@ mod tests {
use crate::format::overlay::OverlayCoverage;
use crate::format::pb;
use crate::format::{
DeletionFile, DeletionFileType, RowDatasetVersionMeta, RowDatasetVersionSequence, RowIdMeta,
DeletionFile, DeletionFileType, ROW_CREATED_AT_VERSION_FIELD_ID, ROW_ID_FIELD_ID,
RowDatasetVersionMeta, RowDatasetVersionSequence, RowIdMeta,
};
use crate::rowids::{RowIdSequence, write_row_ids};
use crate::transaction::test_support::{
Expand Down Expand Up @@ -4400,6 +4423,191 @@ mod tests {
assert_eq!(fragment.overlays[0].data_file.fields.as_ref(), &[5]);
}

/// A projection keeps only the data files that still share a field with
/// the schema. The reserved ids of spilled row lineage are never in it, so
/// without an exception the file carrying them, their only copy, would go.
#[test]
fn project_keeps_file_carrying_spilled_row_lineage() {
let mut fragment = Fragment::new(0);
fragment.physical_rows = Some(4);
fragment.row_id_meta = Some(RowIdMeta::Column);
fragment.created_at_version_meta = Some(RowDatasetVersionMeta::Column);
fragment.files = vec![
DataFile::new(
"d.lance",
vec![0, 1],
vec![0, 1],
ConcreteFileVersion::V2_0,
None,
None,
),
DataFile::new(
"l.lance",
vec![ROW_ID_FIELD_ID, ROW_CREATED_AT_VERSION_FIELD_ID],
vec![0, 1],
ConcreteFileVersion::V2_0,
None,
None,
),
];
let schema = ArrowSchema::new(vec![
ArrowField::new("a", DataType::Int32, false),
ArrowField::new("b", DataType::Int32, false),
]);
let mut manifest = Manifest::new(
LanceSchema::try_from(&schema).unwrap(),
Arc::new(vec![fragment]),
DataStorageFormat::new(ConcreteFileVersion::V2_0),
HashMap::new(),
);
manifest.reader_feature_flags = FLAG_STABLE_ROW_IDS;
manifest.writer_feature_flags = FLAG_STABLE_ROW_IDS;

// Drop field 1, as `drop_columns(["b"])` would.
let transaction = Transaction::new(
manifest.version,
Operation::Project {
schema: manifest.schema.project_by_ids(&[0], true),
preserves_nullability: true,
},
None,
);
let (result, _) = transaction
.build_manifest(Some(&manifest), vec![], "txn", &default_build_config())
.unwrap();

let fragment = &result.fragments[0];
let paths = fragment
.files
.iter()
.map(|file| file.path.as_str())
.collect::<Vec<_>>();
assert_eq!(paths, ["d.lance", "l.lance"]);
for field_id in [ROW_ID_FIELD_ID, ROW_CREATED_AT_VERSION_FIELD_ID] {
assert_eq!(
fragment
.row_lineage_file(field_id)
.unwrap()
.map(|file| file.path.as_str()),
Some("l.lance"),
"field {field_id}"
);
}
}

/// Replacing part of a wider file tombstones the replaced fields where
/// they live and drops any file left without a schema field. The file
/// carrying the fragment's spilled row ids may have none, yet it is their
/// only copy and must stay next to the new file that answers for the user
/// fields. The carrier is either a file of its own, as an update writes
/// it, or a file that also holds user columns, all of which are replaced.
#[rstest::rstest]
#[case::lineage_only_file(
vec![("v.lance", vec![3, 4, 5, 6]), ("l.lance", vec![ROW_ID_FIELD_ID])],
vec![3, 4],
vec![
("v.lance", vec![TOMBSTONE_FIELD_ID, TOMBSTONE_FIELD_ID, 5, 6]),
("l.lance", vec![ROW_ID_FIELD_ID]),
("v-new.lance", vec![3, 4]),
],
"l.lance"
)]
#[case::in_file_carrier(
vec![("v.lance", vec![3, 4, 5, 6, ROW_ID_FIELD_ID])],
vec![3, 4, 5, 6],
vec![
(
"v.lance",
vec![
TOMBSTONE_FIELD_ID,
TOMBSTONE_FIELD_ID,
TOMBSTONE_FIELD_ID,
TOMBSTONE_FIELD_ID,
ROW_ID_FIELD_ID,
],
),
("v-new.lance", vec![3, 4, 5, 6]),
],
"v.lance"
)]
fn data_replacement_keeps_lineage_carrier(
#[case] files: Vec<(&str, Vec<i32>)>,
#[case] replaced_fields: Vec<i32>,
#[case] expected_files: Vec<(&str, Vec<i32>)>,
#[case] carrier: &str,
) {
let mut fragment = Fragment::new(0);
fragment.row_id_meta = Some(RowIdMeta::Column);
fragment.files = files
.into_iter()
.map(|(path, fields)| {
let column_indices = (0..fields.len() as i32).collect();
DataFile::new(
path,
fields,
column_indices,
ConcreteFileVersion::V2_0,
None,
None,
)
})
.collect();

let fragment = replace_fields(fragment, replaced_fields, 1, 1).unwrap();
let files = fragment
.files
.iter()
.map(|file| (file.path.as_str(), file.fields.to_vec()))
.collect::<Vec<_>>();
assert_eq!(files, expected_files);
assert_eq!(
fragment
.row_lineage_file(ROW_ID_FIELD_ID)
.unwrap()
.map(|file| file.path.as_str()),
Some(carrier)
);
}

/// Whatever hands the commit a fragment whose spilled row ids have no v2
/// carrier -- an operation that dropped the file, or a writer that never
/// attached it -- the commit refuses it rather than publish lineage that
/// no reader can load.
#[rstest::rstest]
#[case::no_carrier(vec![0], ConcreteFileVersion::V2_0)]
#[case::v1_carrier(vec![0, ROW_ID_FIELD_ID], ConcreteFileVersion::V1)]
fn build_manifest_rejects_spilled_arm_without_carrier(
#[case] fields: Vec<i32>,
#[case] version: ConcreteFileVersion,
) {
let manifest = sample_manifest();
let mut fragment = Fragment::new(0);
fragment.physical_rows = Some(4);
fragment.row_id_meta = Some(RowIdMeta::Column);
// Column indices do not matter to the commit, only ids and version.
let data_file = DataFile::new("d.lance", fields, vec![], version, None, None);
fragment.files = vec![data_file];
let transaction = Transaction::new(
manifest.version,
Operation::Merge {
fragments: vec![fragment],
schema: manifest.schema.clone(),
preserves_nullability: true,
},
None,
);

let err = transaction
.build_manifest(Some(&manifest), vec![], "txn", &default_build_config())
.unwrap_err();
assert!(matches!(err, Error::Internal { .. }), "{err:?}");
let message = err.to_string();
assert!(
message.contains("fragment 0") && message.contains(&ROW_ID_FIELD_ID.to_string()),
"{message}"
);
}

#[test]
fn test_data_overlay_build_manifest_merges_duplicate_groups() {
// Two groups targeting the same fragment must both survive (a HashMap
Expand Down
Loading
Loading