Skip to content
Draft
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
37 changes: 18 additions & 19 deletions rust/lance/src/dataset/fragment.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -6511,26 +6511,25 @@ 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)
.activate_mixed_file_versions()
.execute(transaction)
.await
.unwrap();

assert_eq!(
dataset
Expand Down
18 changes: 16 additions & 2 deletions rust/lance/src/dataset/versions/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -494,9 +494,23 @@ pub async fn create_fragment_from_file(
fragment_id: usize,
physical_rows: Option<usize>,
) -> Result<Fragment> {
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
)));
}
Expand Down
4 changes: 3 additions & 1 deletion rust/lance/src/dataset/write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<LanceFileVersion>,

/// Experimental: if set to true, the writer will use stable row ids.
Expand Down
29 changes: 25 additions & 4 deletions rust/lance/src/dataset/write/commit.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ use lance_file::version::{ConcreteFileVersion, LanceFileVersion};
use lance_io::object_store::{ObjectStore, ObjectStoreParams};
use lance_select::RowAddrTreeMap;
use lance_table::{
feature_flags::FLAG_MIXED_DATA_FILE_VERSIONS,
format::{DataStorageFormat, is_detached_version},
io::commit::{CommitConfig, CommitHandler, ManifestNamingScheme},
};
Expand Down Expand Up @@ -54,6 +55,7 @@ pub struct CommitBuilder<'a> {
timeout: Option<Duration>,
/// When `Some`, this commit is the second step of `migrate_to_stable_row_ids`.
migration_next_row_id: Option<u64>,
activate_mixed_file_versions: bool,
}

/// Default timeout applied to [`CommitBuilder::execute`] when none is set.
Expand All @@ -78,6 +80,7 @@ impl<'a> CommitBuilder<'a> {
transaction_properties: None,
timeout: Some(DEFAULT_COMMIT_TIMEOUT),
migration_next_row_id: None,
activate_mixed_file_versions: false,
}
}

Expand All @@ -98,8 +101,10 @@ 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, a format different from its manifest fallback
/// requires mixed data-file-version capability. Use
/// [`Self::activate_mixed_file_versions`] for prewritten fragments that
/// introduce the first non-fallback V2 file.
pub fn with_storage_format(mut self, storage_format: LanceFileVersion) -> Self {
self.storage_format = Some(storage_format.resolve());

Expand All @@ -111,6 +116,16 @@ impl<'a> CommitBuilder<'a> {
self
}

/// Activate mixed exact V2 data file versions for this commit.
///
/// This is required when committing prewritten fragments whose exact V2
/// version differs from the existing dataset's manifest fallback. The
/// activation is sticky and cannot be removed by later commits.
pub fn activate_mixed_file_versions(mut self) -> Self {
self.activate_mixed_file_versions = true;
self
}

/// Pass an object store to use.
pub fn with_object_store(mut self, object_store: Arc<ObjectStore>) -> Self {
self.object_store = Some(object_store);
Expand Down Expand Up @@ -292,7 +307,8 @@ impl<'a> CommitBuilder<'a> {
}
}

async fn execute_inner(self, transaction: Transaction) -> Result<Dataset> {
async fn execute_inner(self, mut transaction: Transaction) -> Result<Dataset> {
transaction.activate_mixed_file_versions |= self.activate_mixed_file_versions;
let session = self
.session
.or_else(|| self.dest.dataset().map(|ds| ds.session.clone()))
Expand Down Expand Up @@ -414,11 +430,16 @@ impl<'a> CommitBuilder<'a> {
&& let Some(storage_format) = self.storage_format
{
let passed_storage_format = DataStorageFormat::new(storage_format);
let mixed_enabled = ds.manifest.reader_feature_flags & FLAG_MIXED_DATA_FILE_VERSIONS
!= 0
&& ds.manifest.writer_feature_flags & FLAG_MIXED_DATA_FILE_VERSIONS != 0;
if ds.manifest.data_storage_format != passed_storage_format
&& !matches!(transaction.operation, Operation::Overwrite { .. })
&& !mixed_enabled
&& !transaction.activate_mixed_file_versions
{
return Err(Error::invalid_input_source(format!(
"Storage format mismatch. Existing dataset uses {:?}, but new data uses {:?}",
"Storage format mismatch. Existing dataset fallback is {:?}, but new data uses {:?}; activate mixed data-file versions for prewritten fragments",
ds.manifest.data_storage_format,
passed_storage_format
).into()));
Expand Down
162 changes: 154 additions & 8 deletions rust/lance/src/dataset/write/insert.rs
Original file line number Diff line number Diff line change
Expand Up @@ -276,16 +276,21 @@ impl<'a> InsertBuilder<'a> {
WriteMode::Append => Operation::Append { fragments },
};

let transaction = TransactionBuilder::new(
let mut transaction_builder = TransactionBuilder::new(
context
.dest
.dataset()
.map(|ds| ds.manifest.version)
.unwrap_or(0),
operation,
)
.transaction_properties(context.params.transaction_properties.clone())
.build();
.transaction_properties(context.params.transaction_properties.clone());

if context.activate_mixed_file_versions {
transaction_builder = transaction_builder.activate_mixed_file_versions();
}

let transaction = transaction_builder.build();

Ok(transaction)
}
Expand Down Expand Up @@ -420,22 +425,27 @@ 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(),
};

let activate_mixed_file_versions = matches!(params.mode, WriteMode::Append)
&& dest.dataset().is_some_and(|dataset| {
storage_version != dataset.manifest.data_storage_format.lance_file_format()
});

Ok(WriteContext {
params,
dest,
object_store,
base_path,
commit_handler,
storage_version,
activate_mixed_file_versions,
})
}
}
Expand All @@ -448,6 +458,7 @@ struct WriteContext<'a> {
base_path: Path,
commit_handler: Arc<dyn CommitHandler>,
storage_version: ConcreteFileVersion,
activate_mixed_file_versions: bool,
}

#[cfg(test)]
Expand All @@ -457,12 +468,147 @@ 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<i32>) -> 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_preserve_activation_intent() {
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();
assert!(tx_v2_1.activate_mixed_file_versions);
assert!(tx_v2_2.activate_mixed_file_versions);

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::<Vec<_>>();
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()));
Expand Down
Loading