Skip to content

Commit c108aa3

Browse files
committed
feat: support per-operation V2 writes
1 parent 63b60b1 commit c108aa3

6 files changed

Lines changed: 300 additions & 47 deletions

File tree

rust/lance/src/dataset/fragment.rs

Lines changed: 17 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -3537,8 +3537,8 @@ mod tests {
35373537
use super::*;
35383538
use crate::{
35393539
dataset::{
3540-
InsertBuilder,
3541-
transaction::{Operation, UpdateMode, UpdatedFragmentOffsets},
3540+
CommitBuilder, InsertBuilder,
3541+
transaction::{Operation, Transaction, UpdateMode, UpdatedFragmentOffsets},
35423542
},
35433543
session::Session,
35443544
utils::test::TestDatasetGenerator,
@@ -6511,26 +6511,24 @@ mod tests {
65116511
mismatched_writer.write_batch(&new_data).await.unwrap();
65126512
mismatched_writer.finish().await.unwrap();
65136513

6514-
let err = FileFragment::create_from_file("mismatched_file.lance", &dataset, 1, Some(128))
6515-
.await
6516-
.unwrap_err();
6517-
assert!(matches!(err, Error::InvalidInput { .. }));
6518-
assert!(err.to_string().contains("File version mismatch"));
6514+
let mismatched_frag =
6515+
FileFragment::create_from_file("mismatched_file.lance", &dataset, 0, Some(128))
6516+
.await
6517+
.unwrap();
6518+
assert_eq!(
6519+
mismatched_frag.files[0].file_version().unwrap(),
6520+
ConcreteFileVersion::V2_0
6521+
);
65196522

65206523
let op = Operation::Append {
6521-
fragments: vec![frag],
6524+
fragments: vec![mismatched_frag],
65226525
};
6523-
let dataset = Dataset::commit(
6524-
&dataset.uri,
6525-
op,
6526-
Some(dataset.version().version),
6527-
None,
6528-
None,
6529-
Default::default(),
6530-
false,
6531-
)
6532-
.await
6533-
.unwrap();
6526+
let transaction = Transaction::new_from_version(dataset.version().version, op);
6527+
let dataset = CommitBuilder::new(Arc::new(dataset))
6528+
.with_storage_format(LanceFileVersion::V2_0)
6529+
.execute(transaction)
6530+
.await
6531+
.unwrap();
65346532

65356533
assert_eq!(
65366534
dataset

rust/lance/src/dataset/tests/dataset_migrations.rs

Lines changed: 15 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -184,10 +184,11 @@ async fn test_fix_v0_8_0_broken_migration() {
184184
}
185185

186186
#[rstest]
187+
#[case::legacy(LanceFileVersion::Legacy)]
188+
#[case::v1_v2_mixed_rejected(LanceFileVersion::Stable)]
187189
#[tokio::test]
188190
async fn test_v0_8_14_invalid_index_fragment_bitmap(
189-
#[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
190-
data_storage_version: LanceFileVersion,
191+
#[case] data_storage_version: LanceFileVersion,
191192
) {
192193
// Old versions of lance could create an index whose fragment bitmap was
193194
// invalid because it did not include fragments that were part of the index
@@ -232,16 +233,25 @@ async fn test_v0_8_14_invalid_index_fragment_bitmap(
232233
let broken_version = dataset.version().version;
233234

234235
// Any transaction, no matter how simple, should trigger the fragment bitmap to be recalculated
235-
dataset
236+
let append_result = dataset
236237
.append(
237238
data,
238239
Some(WriteParams {
239240
data_storage_version: Some(data_storage_version),
240241
..Default::default()
241242
}),
242243
)
243-
.await
244-
.unwrap();
244+
.await;
245+
246+
if matches!(data_storage_version, LanceFileVersion::Stable) {
247+
let error = append_result.unwrap_err();
248+
assert!(
249+
error.to_string().contains("do not have a single version"),
250+
"{error}"
251+
);
252+
return;
253+
}
254+
append_result.unwrap();
245255

246256
for idx in dataset.load_indices().await.unwrap().iter() {
247257
// The corrupt fragment_bitmap does not contain 0 but the

rust/lance/src/dataset/versions/mod.rs

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -506,9 +506,23 @@ pub async fn create_fragment_from_file(
506506
fragment_id: usize,
507507
physical_rows: Option<usize>,
508508
) -> Result<Fragment> {
509-
if file_version != dataset_version {
509+
let same_family = matches!(
510+
(file_version, dataset_version),
511+
(ConcreteFileVersion::V1, ConcreteFileVersion::V1)
512+
| (
513+
ConcreteFileVersion::V2_0
514+
| ConcreteFileVersion::V2_1
515+
| ConcreteFileVersion::V2_2
516+
| ConcreteFileVersion::V2_3,
517+
ConcreteFileVersion::V2_0
518+
| ConcreteFileVersion::V2_1
519+
| ConcreteFileVersion::V2_2
520+
| ConcreteFileVersion::V2_3
521+
)
522+
);
523+
if !same_family {
510524
return Err(Error::invalid_input(format!(
511-
"File version mismatch. Dataset version: {:?} Fragment version: {:?}",
525+
"File version family mismatch. Dataset fallback: {:?} Fragment version: {:?}",
512526
dataset_version, file_version
513527
)));
514528
}

rust/lance/src/dataset/write.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -326,7 +326,9 @@ pub struct WriteParams {
326326
/// of lance.
327327
/// Lance file version 2.3 enables RLE v2 run length widths by default.
328328
///
329-
/// If not specified then the latest stable version will be used.
329+
/// For an existing dataset, an explicit version is the exact target for
330+
/// this operation; if omitted, the manifest storage version is used as the
331+
/// fallback. New datasets default to the latest stable version.
330332
pub data_storage_version: Option<LanceFileVersion>,
331333

332334
/// Experimental: if set to true, the writer will use stable row ids.

rust/lance/src/dataset/write/commit.rs

Lines changed: 113 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ use lance_file::version::{ConcreteFileVersion, LanceFileVersion};
99
use lance_io::object_store::{ObjectStore, ObjectStoreParams};
1010
use lance_select::RowAddrTreeMap;
1111
use lance_table::{
12-
format::{DataStorageFormat, is_detached_version},
12+
format::{DataFile, DataStorageFormat, Fragment, is_detached_version},
1313
io::commit::{CommitConfig, CommitHandler, ManifestNamingScheme},
1414
};
1515

@@ -59,6 +59,79 @@ pub struct CommitBuilder<'a> {
5959
/// Default timeout applied to [`CommitBuilder::execute`] when none is set.
6060
pub const DEFAULT_COMMIT_TIMEOUT: Duration = Duration::from_secs(1800);
6161

62+
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
63+
enum OperationVersionState {
64+
Versionless,
65+
NoFiles,
66+
ReferencesTarget,
67+
ReferencesOther,
68+
}
69+
70+
fn data_files_version_state<'a>(
71+
files: impl IntoIterator<Item = &'a DataFile>,
72+
version: ConcreteFileVersion,
73+
) -> Result<OperationVersionState> {
74+
let mut saw_file = false;
75+
for file in files {
76+
saw_file = true;
77+
if file.file_version()? == version {
78+
return Ok(OperationVersionState::ReferencesTarget);
79+
}
80+
}
81+
Ok(if saw_file {
82+
OperationVersionState::ReferencesOther
83+
} else {
84+
OperationVersionState::NoFiles
85+
})
86+
}
87+
88+
fn fragments_version_state<'a>(
89+
fragments: impl IntoIterator<Item = &'a Fragment>,
90+
version: ConcreteFileVersion,
91+
) -> Result<OperationVersionState> {
92+
data_files_version_state(
93+
fragments
94+
.into_iter()
95+
.flat_map(Fragment::referenced_lance_files),
96+
version,
97+
)
98+
}
99+
100+
fn operation_version_state(
101+
operation: &Operation,
102+
version: ConcreteFileVersion,
103+
) -> Result<OperationVersionState> {
104+
match operation {
105+
Operation::Append { fragments }
106+
| Operation::Overwrite { fragments, .. }
107+
| Operation::Merge { fragments, .. } => fragments_version_state(fragments, version),
108+
Operation::Rewrite { groups, .. } => fragments_version_state(
109+
groups.iter().flat_map(|group| group.new_fragments.iter()),
110+
version,
111+
),
112+
Operation::Update {
113+
updated_fragments,
114+
new_fragments,
115+
..
116+
} => fragments_version_state(
117+
updated_fragments.iter().chain(new_fragments.iter()),
118+
version,
119+
),
120+
Operation::DataReplacement { replacements } => data_files_version_state(
121+
replacements.iter().map(|replacement| &replacement.1),
122+
version,
123+
),
124+
Operation::DataOverlay { groups } => data_files_version_state(
125+
groups
126+
.iter()
127+
.flat_map(|group| group.overlays.iter())
128+
.map(|overlay| &overlay.data_file),
129+
version,
130+
),
131+
_ => Ok(OperationVersionState::Versionless),
132+
}
133+
}
134+
62135
impl<'a> CommitBuilder<'a> {
63136
pub fn new(dest: impl Into<WriteDestination<'a>>) -> Self {
64137
Self {
@@ -98,8 +171,9 @@ impl<'a> CommitBuilder<'a> {
98171
/// This is only needed when creating a new empty table. If any data files are
99172
/// passed, the storage format will be inferred from the data files.
100173
///
101-
/// All data files must use the same storage format as the existing dataset.
102-
/// If a different format is passed, an error will be returned.
174+
/// For an existing dataset, the manifest fallback remains unchanged. If
175+
/// prewritten fragments introduce another exact V2 version, the commit
176+
/// derives the required mixed-version capability from the final manifest.
103177
pub fn with_storage_format(mut self, storage_format: LanceFileVersion) -> Self {
104178
self.storage_format = Some(storage_format.resolve());
105179

@@ -409,20 +483,24 @@ impl<'a> CommitBuilder<'a> {
409483
} else {
410484
self.use_stable_row_ids.unwrap_or(false)
411485
};
412-
// Validate storage format matches existing dataset
413-
if let Some(ds) = dest.dataset()
486+
487+
if let Some(dataset) = dest.dataset()
414488
&& let Some(storage_format) = self.storage_format
489+
&& dataset.manifest.data_storage_format.lance_file_format() != storage_format
490+
&& !matches!(transaction.operation, Operation::Overwrite { .. })
491+
&& matches!(
492+
operation_version_state(&transaction.operation, storage_format)?,
493+
OperationVersionState::Versionless | OperationVersionState::ReferencesOther
494+
)
415495
{
416-
let passed_storage_format = DataStorageFormat::new(storage_format);
417-
if ds.manifest.data_storage_format != passed_storage_format
418-
&& !matches!(transaction.operation, Operation::Overwrite { .. })
419-
{
420-
return Err(Error::invalid_input_source(format!(
421-
"Storage format mismatch. Existing dataset uses {:?}, but new data uses {:?}",
422-
ds.manifest.data_storage_format,
423-
passed_storage_format
424-
).into()));
425-
}
496+
return Err(Error::invalid_input_source(
497+
format!(
498+
"Storage format mismatch. Existing dataset fallback is {:?}, but the commit requested {:?} without referencing data files in that version",
499+
dataset.manifest.data_storage_format,
500+
DataStorageFormat::new(storage_format)
501+
)
502+
.into(),
503+
));
426504
}
427505

428506
let manifest_config = ManifestWriteConfig {
@@ -652,6 +730,26 @@ mod tests {
652730
}
653731
}
654732

733+
#[test]
734+
fn empty_file_writes_do_not_require_a_target_version() {
735+
let operation = Operation::Update {
736+
updated_fragments: vec![],
737+
new_fragments: vec![],
738+
removed_fragment_ids: vec![],
739+
fields_modified: vec![],
740+
compacted_sstables: vec![],
741+
fields_for_preserving_frag_bitmap: vec![],
742+
update_mode: None,
743+
inserted_rows_filter: None,
744+
updated_fragment_offsets: None,
745+
};
746+
747+
assert_eq!(
748+
operation_version_state(&operation, ConcreteFileVersion::V2_2).unwrap(),
749+
OperationVersionState::NoFiles
750+
);
751+
}
752+
655753
#[derive(Debug)]
656754
struct SlowConflictingCommitHandler;
657755

0 commit comments

Comments
 (0)