diff --git a/rust/lance/src/dataset/fragment.rs b/rust/lance/src/dataset/fragment.rs index 35411512c5f..c7baba5fda4 100644 --- a/rust/lance/src/dataset/fragment.rs +++ b/rust/lance/src/dataset/fragment.rs @@ -3537,8 +3537,8 @@ mod tests { use super::*; use crate::{ dataset::{ - InsertBuilder, - transaction::{Operation, UpdateMode, UpdatedFragmentOffsets}, + CommitBuilder, InsertBuilder, + transaction::{Operation, Transaction, UpdateMode, UpdatedFragmentOffsets}, }, session::Session, utils::test::TestDatasetGenerator, @@ -6511,26 +6511,24 @@ mod tests { mismatched_writer.write_batch(&new_data).await.unwrap(); mismatched_writer.finish().await.unwrap(); - let err = FileFragment::create_from_file("mismatched_file.lance", &dataset, 1, Some(128)) - .await - .unwrap_err(); - assert!(matches!(err, Error::InvalidInput { .. })); - assert!(err.to_string().contains("File version mismatch")); + let mismatched_frag = + FileFragment::create_from_file("mismatched_file.lance", &dataset, 0, Some(128)) + .await + .unwrap(); + assert_eq!( + mismatched_frag.files[0].file_version().unwrap(), + ConcreteFileVersion::V2_0 + ); let op = Operation::Append { - fragments: vec![frag], + fragments: vec![mismatched_frag], }; - let dataset = Dataset::commit( - &dataset.uri, - op, - Some(dataset.version().version), - None, - None, - Default::default(), - false, - ) - .await - .unwrap(); + let transaction = Transaction::new_from_version(dataset.version().version, op); + let dataset = CommitBuilder::new(Arc::new(dataset)) + .with_storage_format(LanceFileVersion::V2_0) + .execute(transaction) + .await + .unwrap(); assert_eq!( dataset diff --git a/rust/lance/src/dataset/tests/dataset_migrations.rs b/rust/lance/src/dataset/tests/dataset_migrations.rs index a9f58eb7c64..e881e6b73bd 100644 --- a/rust/lance/src/dataset/tests/dataset_migrations.rs +++ b/rust/lance/src/dataset/tests/dataset_migrations.rs @@ -184,10 +184,11 @@ async fn test_fix_v0_8_0_broken_migration() { } #[rstest] +#[case::legacy(LanceFileVersion::Legacy)] +#[case::v1_v2_mixed_rejected(LanceFileVersion::Stable)] #[tokio::test] async fn test_v0_8_14_invalid_index_fragment_bitmap( - #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)] - data_storage_version: LanceFileVersion, + #[case] data_storage_version: LanceFileVersion, ) { // Old versions of lance could create an index whose fragment bitmap was // invalid because it did not include fragments that were part of the index @@ -232,7 +233,7 @@ async fn test_v0_8_14_invalid_index_fragment_bitmap( let broken_version = dataset.version().version; // Any transaction, no matter how simple, should trigger the fragment bitmap to be recalculated - dataset + let append_result = dataset .append( data, Some(WriteParams { @@ -240,8 +241,17 @@ async fn test_v0_8_14_invalid_index_fragment_bitmap( ..Default::default() }), ) - .await - .unwrap(); + .await; + + if matches!(data_storage_version, LanceFileVersion::Stable) { + let error = append_result.unwrap_err(); + assert!( + error.to_string().contains("do not have a single version"), + "{error}" + ); + return; + } + append_result.unwrap(); for idx in dataset.load_indices().await.unwrap().iter() { // The corrupt fragment_bitmap does not contain 0 but the diff --git a/rust/lance/src/dataset/versions/mod.rs b/rust/lance/src/dataset/versions/mod.rs index afb94303017..af6c22a3f6f 100644 --- a/rust/lance/src/dataset/versions/mod.rs +++ b/rust/lance/src/dataset/versions/mod.rs @@ -506,9 +506,23 @@ pub async fn create_fragment_from_file( fragment_id: usize, physical_rows: Option, ) -> Result { - if file_version != dataset_version { + let same_family = matches!( + (file_version, dataset_version), + (ConcreteFileVersion::V1, ConcreteFileVersion::V1) + | ( + ConcreteFileVersion::V2_0 + | ConcreteFileVersion::V2_1 + | ConcreteFileVersion::V2_2 + | ConcreteFileVersion::V2_3, + ConcreteFileVersion::V2_0 + | ConcreteFileVersion::V2_1 + | ConcreteFileVersion::V2_2 + | ConcreteFileVersion::V2_3 + ) + ); + if !same_family { return Err(Error::invalid_input(format!( - "File version mismatch. Dataset version: {:?} Fragment version: {:?}", + "File version family mismatch. Dataset fallback: {:?} Fragment version: {:?}", dataset_version, file_version ))); } diff --git a/rust/lance/src/dataset/write.rs b/rust/lance/src/dataset/write.rs index 81b0fa37109..ed006cbda14 100644 --- a/rust/lance/src/dataset/write.rs +++ b/rust/lance/src/dataset/write.rs @@ -326,7 +326,9 @@ pub struct WriteParams { /// of lance. /// Lance file version 2.3 enables RLE v2 run length widths by default. /// - /// If not specified then the latest stable version will be used. + /// For an existing dataset, an explicit version is the exact target for + /// this operation; if omitted, the manifest storage version is used as the + /// fallback. New datasets default to the latest stable version. pub data_storage_version: Option, /// Experimental: if set to true, the writer will use stable row ids. diff --git a/rust/lance/src/dataset/write/commit.rs b/rust/lance/src/dataset/write/commit.rs index 7ab5b17a9de..4e49f7cee4f 100644 --- a/rust/lance/src/dataset/write/commit.rs +++ b/rust/lance/src/dataset/write/commit.rs @@ -9,7 +9,7 @@ use lance_file::version::{ConcreteFileVersion, LanceFileVersion}; use lance_io::object_store::{ObjectStore, ObjectStoreParams}; use lance_select::RowAddrTreeMap; use lance_table::{ - format::{DataStorageFormat, is_detached_version}, + format::{DataFile, DataStorageFormat, Fragment, is_detached_version}, io::commit::{CommitConfig, CommitHandler, ManifestNamingScheme}, }; @@ -59,6 +59,79 @@ pub struct CommitBuilder<'a> { /// Default timeout applied to [`CommitBuilder::execute`] when none is set. pub const DEFAULT_COMMIT_TIMEOUT: Duration = Duration::from_secs(1800); +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum OperationVersionState { + Versionless, + NoFiles, + ReferencesTarget, + ReferencesOther, +} + +fn data_files_version_state<'a>( + files: impl IntoIterator, + version: ConcreteFileVersion, +) -> Result { + let mut saw_file = false; + for file in files { + saw_file = true; + if file.file_version()? == version { + return Ok(OperationVersionState::ReferencesTarget); + } + } + Ok(if saw_file { + OperationVersionState::ReferencesOther + } else { + OperationVersionState::NoFiles + }) +} + +fn fragments_version_state<'a>( + fragments: impl IntoIterator, + version: ConcreteFileVersion, +) -> Result { + data_files_version_state( + fragments + .into_iter() + .flat_map(Fragment::referenced_lance_files), + version, + ) +} + +fn operation_version_state( + operation: &Operation, + version: ConcreteFileVersion, +) -> Result { + match operation { + Operation::Append { fragments } + | Operation::Overwrite { fragments, .. } + | Operation::Merge { fragments, .. } => fragments_version_state(fragments, version), + Operation::Rewrite { groups, .. } => fragments_version_state( + groups.iter().flat_map(|group| group.new_fragments.iter()), + version, + ), + Operation::Update { + updated_fragments, + new_fragments, + .. + } => fragments_version_state( + updated_fragments.iter().chain(new_fragments.iter()), + version, + ), + Operation::DataReplacement { replacements } => data_files_version_state( + replacements.iter().map(|replacement| &replacement.1), + version, + ), + Operation::DataOverlay { groups } => data_files_version_state( + groups + .iter() + .flat_map(|group| group.overlays.iter()) + .map(|overlay| &overlay.data_file), + version, + ), + _ => Ok(OperationVersionState::Versionless), + } +} + impl<'a> CommitBuilder<'a> { pub fn new(dest: impl Into>) -> Self { Self { @@ -98,8 +171,9 @@ impl<'a> CommitBuilder<'a> { /// This is only needed when creating a new empty table. If any data files are /// passed, the storage format will be inferred from the data files. /// - /// All data files must use the same storage format as the existing dataset. - /// If a different format is passed, an error will be returned. + /// For an existing dataset, the manifest fallback remains unchanged. If + /// prewritten fragments introduce another exact V2 version, the commit + /// derives the required mixed-version capability from the final manifest. pub fn with_storage_format(mut self, storage_format: LanceFileVersion) -> Self { self.storage_format = Some(storage_format.resolve()); @@ -409,20 +483,24 @@ impl<'a> CommitBuilder<'a> { } else { self.use_stable_row_ids.unwrap_or(false) }; - // Validate storage format matches existing dataset - if let Some(ds) = dest.dataset() + + if let Some(dataset) = dest.dataset() && let Some(storage_format) = self.storage_format + && dataset.manifest.data_storage_format.lance_file_format() != storage_format + && !matches!(transaction.operation, Operation::Overwrite { .. }) + && matches!( + operation_version_state(&transaction.operation, storage_format)?, + OperationVersionState::Versionless | OperationVersionState::ReferencesOther + ) { - let passed_storage_format = DataStorageFormat::new(storage_format); - if ds.manifest.data_storage_format != passed_storage_format - && !matches!(transaction.operation, Operation::Overwrite { .. }) - { - return Err(Error::invalid_input_source(format!( - "Storage format mismatch. Existing dataset uses {:?}, but new data uses {:?}", - ds.manifest.data_storage_format, - passed_storage_format - ).into())); - } + return Err(Error::invalid_input_source( + format!( + "Storage format mismatch. Existing dataset fallback is {:?}, but the commit requested {:?} without referencing data files in that version", + dataset.manifest.data_storage_format, + DataStorageFormat::new(storage_format) + ) + .into(), + )); } let manifest_config = ManifestWriteConfig { @@ -652,6 +730,26 @@ mod tests { } } + #[test] + fn empty_file_writes_do_not_require_a_target_version() { + let operation = Operation::Update { + updated_fragments: vec![], + new_fragments: vec![], + removed_fragment_ids: vec![], + fields_modified: vec![], + compacted_sstables: vec![], + fields_for_preserving_frag_bitmap: vec![], + update_mode: None, + inserted_rows_filter: None, + updated_fragment_offsets: None, + }; + + assert_eq!( + operation_version_state(&operation, ConcreteFileVersion::V2_2).unwrap(), + OperationVersionState::NoFiles + ); + } + #[derive(Debug)] struct SlowConflictingCommitHandler; diff --git a/rust/lance/src/dataset/write/insert.rs b/rust/lance/src/dataset/write/insert.rs index b5dfd4b2953..76794eb291b 100644 --- a/rust/lance/src/dataset/write/insert.rs +++ b/rust/lance/src/dataset/write/insert.rs @@ -420,11 +420,10 @@ impl<'a> InsertBuilder<'a> { .map(LanceFileVersion::resolve) .unwrap_or_else(|| dataset.manifest.data_storage_format.lance_file_format()) } - (_, WriteDestination::Dataset(dataset)) => { - // If appending to an existing dataset, always use the dataset version - let m = dataset.manifest.as_ref(); - m.data_storage_format.lance_file_format() - } + (_, WriteDestination::Dataset(dataset)) => params + .data_storage_version + .map(LanceFileVersion::resolve) + .unwrap_or_else(|| dataset.manifest.data_storage_format.lance_file_format()), // Otherwise (no existing dataset) fallback to the default if the user didn't specify (_, WriteDestination::Uri(_)) => params.storage_version_or_default(), }; @@ -457,12 +456,144 @@ mod test { use arrow_array::{ArrayRef, BinaryArray, Int32Array, RecordBatchReader, StructArray}; use arrow_schema::{ArrowError, DataType, Field, Schema}; use lance_arrow::BLOB_META_KEY; + use lance_core::utils::tempfile::TempStrDir; + use lance_table::feature_flags::FLAG_MIXED_DATA_FILE_VERSIONS; use lance_table::io::commit::{RenameCommitHandler, commit_handler_from_url}; use crate::session::Session; use super::*; + fn int_batch(values: Vec) -> RecordBatch { + let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)])); + RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(values))]).unwrap() + } + + #[tokio::test] + async fn append_resolves_explicit_or_fallback_exact_version() { + let test_dir = TempStrDir::default(); + let create_params = WriteParams { + data_storage_version: Some(LanceFileVersion::V2_0), + ..Default::default() + }; + let dataset = InsertBuilder::new(test_dir.as_str()) + .with_params(&create_params) + .execute(vec![int_batch(vec![0])]) + .await + .unwrap(); + + let fallback_params = WriteParams { + mode: WriteMode::Append, + ..Default::default() + }; + let dataset = InsertBuilder::new(Arc::new(dataset)) + .with_params(&fallback_params) + .execute(vec![int_batch(vec![1])]) + .await + .unwrap(); + assert_eq!( + dataset.manifest.fragments[1].files[0] + .file_version() + .unwrap(), + ConcreteFileVersion::V2_0 + ); + assert_eq!( + dataset.manifest.reader_feature_flags & FLAG_MIXED_DATA_FILE_VERSIONS, + 0 + ); + + let explicit_params = WriteParams { + mode: WriteMode::Append, + data_storage_version: Some(LanceFileVersion::V2_1), + ..Default::default() + }; + let dataset = InsertBuilder::new(Arc::new(dataset)) + .with_params(&explicit_params) + .execute(vec![int_batch(vec![2])]) + .await + .unwrap(); + + assert_eq!( + dataset.manifest.data_storage_format.lance_file_format(), + ConcreteFileVersion::V2_0 + ); + assert_eq!( + dataset.manifest.fragments[2].files[0] + .file_version() + .unwrap(), + ConcreteFileVersion::V2_1 + ); + assert_ne!( + dataset.manifest.reader_feature_flags & FLAG_MIXED_DATA_FILE_VERSIONS, + 0 + ); + assert_ne!( + dataset.manifest.writer_feature_flags & FLAG_MIXED_DATA_FILE_VERSIONS, + 0 + ); + assert_eq!(dataset.count_rows(None).await.unwrap(), 3); + } + + #[tokio::test] + async fn concurrent_exact_version_appends_derive_capability_from_final_manifest() { + let test_dir = TempStrDir::default(); + let create_params = WriteParams { + data_storage_version: Some(LanceFileVersion::V2_0), + ..Default::default() + }; + let original = Arc::new( + InsertBuilder::new(test_dir.as_str()) + .with_params(&create_params) + .execute(vec![int_batch(vec![0])]) + .await + .unwrap(), + ); + let v2_1_params = WriteParams { + mode: WriteMode::Append, + data_storage_version: Some(LanceFileVersion::V2_1), + ..Default::default() + }; + let v2_2_params = WriteParams { + mode: WriteMode::Append, + data_storage_version: Some(LanceFileVersion::V2_2), + ..Default::default() + }; + let tx_v2_1 = InsertBuilder::new(original.clone()) + .with_params(&v2_1_params) + .execute_uncommitted(vec![int_batch(vec![1])]) + .await + .unwrap(); + let tx_v2_2 = InsertBuilder::new(original.clone()) + .with_params(&v2_2_params) + .execute_uncommitted(vec![int_batch(vec![2])]) + .await + .unwrap(); + let dataset = CommitBuilder::new(original.clone()) + .execute(tx_v2_1) + .await + .unwrap(); + let dataset = CommitBuilder::new(Arc::new(dataset)) + .execute(tx_v2_2) + .await + .unwrap(); + + let versions = dataset + .manifest + .fragments + .iter() + .map(|fragment| fragment.files[0].file_version().unwrap()) + .collect::>(); + assert_eq!( + versions, + vec![ + ConcreteFileVersion::V2_0, + ConcreteFileVersion::V2_1, + ConcreteFileVersion::V2_2 + ] + ); + assert_eq!(dataset.count_rows(None).await.unwrap(), 3); + } + #[tokio::test] async fn test_pass_session() { let session = Arc::new(Session::new(0, 0, Default::default()));