Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
27 commits
Select commit Hold shift + click to select a range
cbb3377
feat: introduce exact file format identity
Xuanwo Jul 21, 2026
543fe0d
fix(python): preserve data storage version getter
Xuanwo Jul 21, 2026
81a459f
Merge remote-tracking branch 'origin/main' into xuanwo/exact-file-for…
Xuanwo Jul 27, 2026
b7e645b
refactor: name exact file version identity
Xuanwo Jul 27, 2026
e7e92ad
test: use exact version for overlay files
Xuanwo Jul 27, 2026
88346c8
test: lock exact file-version wire contracts
Xuanwo Jul 27, 2026
d4bc6fc
refactor: make v1 the canonical legacy file owner
Xuanwo Jul 27, 2026
4c78d37
refactor: make encoding mechanisms version-free
Xuanwo Jul 27, 2026
aec4d2e
refactor: add exact current-format writers
Xuanwo Jul 27, 2026
d05c05f
test: keep encoding strategy setup lint-clean
Xuanwo Jul 27, 2026
0a68138
Merge branch 'xuanwo/exact-version-stack-04-encoding-mechanisms' into…
Xuanwo Jul 27, 2026
faf3146
refactor: activate exact current-format writers
Xuanwo Jul 27, 2026
daced25
refactor: compose exact current-format readers
Xuanwo Jul 27, 2026
0011ec9
refactor: centralize exact dataset write policy
Xuanwo Jul 27, 2026
c77a4b8
Merge remote-tracking branch 'origin/main' into xuanwo/exact-file-for…
Xuanwo Jul 27, 2026
c6e695c
Merge branch 'xuanwo/exact-file-format-identity' into xuanwo/exact-ve…
Xuanwo Jul 27, 2026
a77e601
Merge branch 'xuanwo/exact-version-stack-02-fixtures' into xuanwo/exa…
Xuanwo Jul 27, 2026
15ffaf3
Merge branch 'xuanwo/exact-version-stack-03-v1' into xuanwo/exact-ver…
Xuanwo Jul 27, 2026
0ed1f34
Merge branch 'xuanwo/exact-version-stack-04-encoding-mechanisms' into…
Xuanwo Jul 27, 2026
5632498
Merge branch 'xuanwo/exact-version-stack-05-file-runtime' into xuanwo…
Xuanwo Jul 27, 2026
a5ba57d
Merge branch 'xuanwo/exact-version-stack-06-writers' into xuanwo/exac…
Xuanwo Jul 27, 2026
e0adecb
Merge branch 'xuanwo/exact-version-stack-07-readers' into xuanwo/exac…
Xuanwo Jul 27, 2026
92c7a22
Merge remote-tracking branch 'origin/main' into xuanwo/exact-version-…
Xuanwo Aug 5, 2026
02cda3d
Merge remote-tracking branch 'origin/main' into xuanwo/exact-version-…
Xuanwo Aug 5, 2026
1270f5f
refactor: align exact dataset write policy with main
Xuanwo Aug 5, 2026
80d255b
chore: add missing license header
Xuanwo Aug 5, 2026
cb601f7
refactor: resolve dataset write version once
Xuanwo Aug 5, 2026
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
1 change: 1 addition & 0 deletions rust/lance/src/dataset.rs
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,7 @@ pub mod transaction;
pub mod udtf;
pub mod updater;
mod utils;
pub(crate) mod versions;
pub mod write;

pub(crate) use take::row_offsets_to_row_addresses;
Expand Down
55 changes: 27 additions & 28 deletions rust/lance/src/dataset/fragment/write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,9 @@ use lance_datafusion::chunker::{break_stream, chunk_stream};
use lance_datafusion::utils::StreamingWriteSource;
use lance_file::version::{ConcreteFileVersion, LanceFileVersion};
use lance_file::versions::v1::writer::FileWriter as V1FileWriter;
use lance_file::writer::FileWriterOptions;
use lance_file::writer::FileWriter;
use lance_io::object_store::ObjectStore;
use lance_io::traits::Writer;
use lance_io::utils::CachedFileSize;
use lance_table::format::{DataFile, Fragment};
use lance_table::io::manifest::ManifestDescribing;
Expand All @@ -22,7 +23,7 @@ use uuid::Uuid;
use crate::Result;
use crate::dataset::builder::DatasetBuilder;
use crate::dataset::utils::SchemaAdapter;
use crate::dataset::write::{do_write_fragments, validate_and_resolve_target_bases_with_primary};
use crate::dataset::write::validate_and_resolve_target_bases_with_primary;
use crate::dataset::{DATA_DIR, Dataset, ReadParams, WriteMode, WriteParams};

/// Generates a filename optimized for S3 throughput using a UUID-based approach.
Expand Down Expand Up @@ -113,7 +114,18 @@ impl<'a> FragmentCreateBuilder<'a> {
// the single-fragment create path must do the same or the raw UTF-8 string bytes
// would be written into a column whose schema declares JSONB, corrupting reads.
let stream = SchemaAdapter::new(stream.schema()).to_physical_stream(stream);
self.write_impl(stream, schema, id).await
let version = self
.write_params
.map(|params| params.storage_version_or_default())
.unwrap_or_else(|| ConcreteFileVersion::from(LanceFileVersion::Stable));
crate::dataset::versions::write_fragment(
version,
self,
stream,
schema,
id.unwrap_or_default(),
)
.await
}

/// Write multi fragment which separated by max_rows_per_file.
Expand All @@ -125,12 +137,16 @@ impl<'a> FragmentCreateBuilder<'a> {
self.write_fragments_v2_impl(stream, schema).await
}

async fn write_v2_impl(
pub(crate) async fn write_current_impl<F>(
&self,
create_writer: F,
stream: SendableRecordBatchStream,
schema: Schema,
id: u64,
) -> Result<Fragment> {
) -> Result<Fragment>
where
F: FnOnce(Box<dyn Writer>, Schema, String) -> Result<(FileWriter, DataFile)>,
{
let params = self.write_params.map(Cow::Borrowed).unwrap_or_default();
let progress = params.progress.as_ref();

Expand All @@ -147,16 +163,7 @@ impl<'a> FragmentCreateBuilder<'a> {
let mut fragment = Fragment::new(id);
let full_path = base_path.clone().join(DATA_DIR).join(filename.clone());
let obj_writer = object_store.create(&full_path).await?;
let file_version =
ConcreteFileVersion::from(params.data_storage_version.unwrap_or_default());
let mut writer = lance_file::versions::create_writer(
file_version,
obj_writer,
schema,
FileWriterOptions::default(),
)?;

let data_file = DataFile::new_unstarted(filename, file_version);
let (mut writer, data_file) = create_writer(obj_writer, schema, filename)?;
fragment.files.push(data_file);

progress.begin(&fragment).await?;
Expand Down Expand Up @@ -208,7 +215,7 @@ impl<'a> FragmentCreateBuilder<'a> {

Self::validate_schema(&schema, stream.schema().as_ref())?;

let version = params.data_storage_version.unwrap_or_default();
let version = params.storage_version_or_default();
let needs_existing_dataset = params.target_base_names_or_paths.is_some()
|| params.target_bases.is_some()
|| params.target_all_bases.is_some()
Expand Down Expand Up @@ -239,35 +246,27 @@ impl<'a> FragmentCreateBuilder<'a> {
} else {
None
};
do_write_fragments(
crate::dataset::versions::write_fragments_direct(
version,
existing_dataset.as_ref(),
object_store,
&base_path,
&schema,
stream,
params,
version,
target_bases_info,
Vec::new(),
)
.await
}

async fn write_impl(
pub(crate) async fn write_v1_impl(
&self,
stream: SendableRecordBatchStream,
schema: Schema,
id: Option<u64>,
id: u64,
) -> Result<Fragment> {
let id = id.unwrap_or_default();

let params = self.write_params.map(Cow::Borrowed).unwrap_or_default();

let storage_version = params.storage_version_or_default();

if storage_version != LanceFileVersion::Legacy {
return self.write_v2_impl(stream, schema, id).await;
}
let progress = params.progress.as_ref();

Self::validate_schema(&schema, stream.schema().as_ref())?;
Expand Down
1 change: 1 addition & 0 deletions rust/lance/src/dataset/optimize.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1804,6 +1804,7 @@ async fn rewrite_files(
}
} else {
let (frags, _) = write_fragments_internal(
dataset.manifest.data_storage_format.lance_file_format(),
Some(dataset.as_ref()),
dataset.object_store.clone(),
&dataset.base,
Expand Down
5 changes: 3 additions & 2 deletions rust/lance/src/dataset/updater.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,8 @@ use lance_table::utils::stream::ReadBatchFutStream;
use super::Dataset;
use super::fragment::FragmentReader;
use super::scanner::get_default_batch_size;
use super::write::{GenericWriter, cleanup_data_fragments, open_update_writer};
use super::versions;
use super::write::{GenericWriter, cleanup_data_fragments};
use crate::dataset::FileFragment;
use crate::dataset::utils::SchemaAdapter;

Expand Down Expand Up @@ -147,7 +148,7 @@ impl Updater {
.data_storage_format
.lance_file_version()?;

open_update_writer(self.dataset(), &schema, data_storage_version).await
versions::open_update_writer(data_storage_version.into(), self.dataset(), &schema).await
}

/// Update one batch.
Expand Down
Loading
Loading