diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index 784ee4e1952..3fbb9e87c6e 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -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; diff --git a/rust/lance/src/dataset/fragment/write.rs b/rust/lance/src/dataset/fragment/write.rs index b2b2180a9e1..27723f1877f 100644 --- a/rust/lance/src/dataset/fragment/write.rs +++ b/rust/lance/src/dataset/fragment/write.rs @@ -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; @@ -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. @@ -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. @@ -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( &self, + create_writer: F, stream: SendableRecordBatchStream, schema: Schema, id: u64, - ) -> Result { + ) -> Result + where + F: FnOnce(Box, Schema, String) -> Result<(FileWriter, DataFile)>, + { let params = self.write_params.map(Cow::Borrowed).unwrap_or_default(); let progress = params.progress.as_ref(); @@ -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?; @@ -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() @@ -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, + id: u64, ) -> Result { - 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())?; diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index b6dce301d21..120a61c2fe8 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -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, diff --git a/rust/lance/src/dataset/updater.rs b/rust/lance/src/dataset/updater.rs index 24b0cbbe8dc..4314dd05b1f 100644 --- a/rust/lance/src/dataset/updater.rs +++ b/rust/lance/src/dataset/updater.rs @@ -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; @@ -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. diff --git a/rust/lance/src/dataset/versions/mod.rs b/rust/lance/src/dataset/versions/mod.rs new file mode 100644 index 00000000000..671e757eb10 --- /dev/null +++ b/rust/lance/src/dataset/versions/mod.rs @@ -0,0 +1,248 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Dataset write policies that differ across exact Lance file versions. +//! +//! File grammar belongs to `lance_file::versions`. This module contains only +//! operation-level dataset choices whose behavior actually differs by version. + +use std::sync::Arc; + +use datafusion::execution::SendableRecordBatchStream; +use futures::{StreamExt, TryStreamExt}; +use lance_core::{ + Result, + datatypes::{Schema, SchemaCompareOptions}, +}; +use lance_datafusion::chunker::{break_stream, chunk_stream}; +use lance_file::{ + version::ConcreteFileVersion, + versions as file_versions, + writer::{FileWriter, FileWriterOptions}, +}; +use lance_index::scalar::seed::IndexSeedWriter; +use lance_io::object_store::ObjectStore; +use lance_io::traits::Writer as ObjectWriter; +use lance_table::format::{DataFile, Fragment}; +use object_store::path::Path; + +use super::Dataset; +use super::fragment::write::FragmentCreateBuilder; +use super::utils::SchemaAdapter; +use super::write::{self, TargetBaseInfo, WriteParams, WriterOptions}; + +pub fn schema_compare_options(version: ConcreteFileVersion) -> SchemaCompareOptions { + match version { + ConcreteFileVersion::V1 => SchemaCompareOptions { + compare_dictionary: true, + ..Default::default() + }, + ConcreteFileVersion::V2_0 + | ConcreteFileVersion::V2_1 + | ConcreteFileVersion::V2_2 + | ConcreteFileVersion::V2_3 => SchemaCompareOptions::default(), + } +} + +async fn create_seed_writers( + version: ConcreteFileVersion, + dataset: Option<&Dataset>, + params: &WriteParams, +) -> Result>> { + match version { + ConcreteFileVersion::V1 => Ok(Vec::new()), + ConcreteFileVersion::V2_0 + | ConcreteFileVersion::V2_1 + | ConcreteFileVersion::V2_2 + | ConcreteFileVersion::V2_3 => write::create_seed_writers_current(dataset, params).await, + } +} + +fn create_current_file_writer( + version: ConcreteFileVersion, + object_writer: Box, + schema: Schema, + filename: String, + base_id: Option, +) -> Result<(FileWriter, DataFile)> { + let writer = + file_versions::create_writer(version, object_writer, schema, FileWriterOptions::default())?; + let mut data_file = DataFile::new_unstarted(filename, version); + data_file.base_id = base_id; + Ok((writer, data_file)) +} + +#[allow(clippy::too_many_arguments)] +pub async fn write_fragments( + version: ConcreteFileVersion, + dataset: Option<&Dataset>, + object_store: Arc, + base_dir: &Path, + normalized_schema: Schema, + data: SendableRecordBatchStream, + params: WriteParams, + target_bases_info: Option>, +) -> Result<(Vec, Schema)> { + let version_name = format!("{version:?}"); + let schema = write::prepare_write_schema( + dataset, + normalized_schema, + ¶ms, + schema_compare_options(version), + )?; + match version { + ConcreteFileVersion::V1 | ConcreteFileVersion::V2_0 | ConcreteFileVersion::V2_1 => { + write::validate_legacy_blob_write_schema(&schema, &version_name)?; + } + ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 => { + write::validate_blob_v2_write_schema(&schema)?; + } + } + let seed_writers = create_seed_writers(version, dataset, ¶ms).await?; + let fragments = write_fragments_direct( + version, + dataset, + object_store, + base_dir, + &schema, + data, + params, + target_bases_info, + seed_writers, + ) + .await?; + Ok((fragments, schema)) +} + +#[allow(clippy::too_many_arguments)] +pub async fn write_fragments_direct( + version: ConcreteFileVersion, + dataset: Option<&Dataset>, + object_store: Arc, + base_dir: &Path, + schema: &Schema, + data: SendableRecordBatchStream, + params: WriteParams, + target_bases_info: Option>, + seed_writers: Vec>, +) -> Result> { + let adapter = SchemaAdapter::new(data.schema()); + let data = adapter.to_physical_stream(data); + let buffered_reader = match version { + ConcreteFileVersion::V1 => chunk_stream(data, params.max_rows_per_group), + ConcreteFileVersion::V2_0 + | ConcreteFileVersion::V2_1 + | ConcreteFileVersion::V2_2 + | ConcreteFileVersion::V2_3 => break_stream(data, params.max_rows_per_file) + .map_ok(|batch| vec![batch]) + .boxed(), + }; + let external_base_resolver = match version { + ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 => { + write::blob_v2_external_base_resolver(dataset, ¶ms, schema).await? + } + ConcreteFileVersion::V1 | ConcreteFileVersion::V2_0 | ConcreteFileVersion::V2_1 => None, + }; + write::do_write_fragments_impl( + dataset, + object_store, + base_dir, + schema, + buffered_reader, + params, + move |object_store, schema, base_dir, options| async move { + open_writer(version, &object_store, &schema, &base_dir, options).await + }, + external_base_resolver, + target_bases_info, + seed_writers, + ) + .await +} + +pub async fn write_fragment( + version: ConcreteFileVersion, + builder: &FragmentCreateBuilder<'_>, + stream: SendableRecordBatchStream, + schema: Schema, + id: u64, +) -> Result { + match version { + ConcreteFileVersion::V1 => builder.write_v1_impl(stream, schema, id).await, + ConcreteFileVersion::V2_0 + | ConcreteFileVersion::V2_1 + | ConcreteFileVersion::V2_2 + | ConcreteFileVersion::V2_3 => { + builder + .write_current_impl( + move |object_writer, schema, filename| { + create_current_file_writer(version, object_writer, schema, filename, None) + }, + stream, + schema, + id, + ) + .await + } + } +} + +pub async fn open_writer( + version: ConcreteFileVersion, + object_store: &ObjectStore, + schema: &Schema, + base_dir: &Path, + options: WriterOptions, +) -> Result> { + match version { + ConcreteFileVersion::V1 => { + write::open_v1_writer(object_store, schema, base_dir, options).await + } + ConcreteFileVersion::V2_0 | ConcreteFileVersion::V2_1 => { + write::open_current_writer( + move |object_writer, schema, filename, base_id| { + create_current_file_writer(version, object_writer, schema, filename, base_id) + }, + object_store, + schema, + base_dir, + options, + ) + .await + } + ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 => { + write::open_current_blob_v2_writer( + move |object_writer, schema, filename, base_id| { + create_current_file_writer(version, object_writer, schema, filename, base_id) + }, + object_store, + schema, + base_dir, + options, + ) + .await + } + } +} + +pub async fn open_update_writer( + version: ConcreteFileVersion, + dataset: &Dataset, + schema: &Schema, +) -> Result> { + let external_base_resolver = match version { + ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 => { + write::blob_v2_external_base_resolver(Some(dataset), &WriteParams::default(), schema) + .await? + } + ConcreteFileVersion::V1 | ConcreteFileVersion::V2_0 | ConcreteFileVersion::V2_1 => None, + }; + open_writer( + version, + &dataset.object_store, + schema, + &dataset.base, + WriterOptions::update(dataset.session.store_registry(), external_base_resolver), + ) + .await +} diff --git a/rust/lance/src/dataset/write.rs b/rust/lance/src/dataset/write.rs index 18d28bd2d39..3547c088352 100644 --- a/rust/lance/src/dataset/write.rs +++ b/rust/lance/src/dataset/write.rs @@ -5,35 +5,34 @@ use arrow_array::RecordBatch; use bytes::Bytes; use chrono::TimeDelta; use datafusion::physical_plan::SendableRecordBatchStream; -use futures::{StreamExt, TryStreamExt}; +use futures::StreamExt; use lance_arrow::{ ARROW_EXT_NAME_KEY, BLOB_DEDICATED_SIZE_THRESHOLD_META_KEY, BLOB_INLINE_SIZE_THRESHOLD_META_KEY, BLOB_META_KEY, BLOB_PACK_FILE_SIZE_THRESHOLD_META_KEY, BLOB_V2_EXT_NAME, }; -use lance_core::datatypes::{ - NullabilityComparison, OnMissing, OnTypeMismatch, SchemaCompareOptions, -}; +use lance_core::datatypes::{NullabilityComparison, OnMissing, OnTypeMismatch}; use lance_core::utils::tracing::{ AUDIT_MODE_CREATE, AUDIT_MODE_DELETE, AUDIT_TYPE_DATA, TRACE_FILE_AUDIT, }; use lance_core::{Error, Result, datatypes::Schema}; -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, ManifestProvider as V1ManifestProvider, }; -use lance_file::writer::{self as current_writer, FileWriterOptions}; +use lance_file::writer::{self as current_writer}; use lance_io::object_store::{ ObjectStore, ObjectStoreParams, ObjectStoreRegistry, parse_base_scoped_key, }; +use lance_io::traits::Writer; use lance_table::format::{BasePath, DataFile, Fragment, IndexMetadata}; use lance_table::io::commit::{CommitHandler, commit_handler_from_url}; use lance_table::io::manifest::ManifestDescribing; use object_store::path::Path; use std::borrow::Cow; use std::collections::{BTreeSet, HashMap, HashSet}; +use std::future::Future; use std::num::NonZero; use std::sync::Arc; use std::sync::atomic::AtomicUsize; @@ -55,6 +54,7 @@ use super::fragment::write::generate_random_filename; use super::progress::{NoopFragmentWriteProgress, WriteFragmentProgress}; use super::transaction::Transaction; use super::utils::SchemaAdapter; +use super::versions; mod commit; pub mod delete; @@ -455,8 +455,8 @@ impl WriteParams { } } - pub fn storage_version_or_default(&self) -> LanceFileVersion { - self.data_storage_version.unwrap_or_default() + pub fn storage_version_or_default(&self) -> ConcreteFileVersion { + self.data_storage_version.unwrap_or_default().into() } pub fn store_registry(&self) -> Arc { @@ -596,40 +596,22 @@ pub async fn write_fragments( } #[allow(clippy::too_many_arguments)] -pub async fn do_write_fragments( +pub(super) async fn do_write_fragments_impl( dataset: Option<&Dataset>, object_store: Arc, base_dir: &Path, schema: &Schema, - data: SendableRecordBatchStream, + mut buffered_reader: futures::stream::BoxStream<'static, Result>>, params: WriteParams, - storage_version: LanceFileVersion, + open_writer: OpenWriter, + external_base_resolver: Option>, target_bases_info: Option>, mut seed_writers: Vec>, -) -> Result> { - let adapter = SchemaAdapter::new(data.schema()); - let data = adapter.to_physical_stream(data); - - let mut buffered_reader = if storage_version == LanceFileVersion::Legacy { - // In v1 we split the stream into row group sized batches - chunk_stream(data, params.max_rows_per_group) - } else { - // In v2 we don't care about group size but we do want to break - // the stream on file boundaries - break_stream(data, params.max_rows_per_file) - .map_ok(|batch| vec![batch]) - .boxed() - }; - - let external_base_resolver = if storage_version >= LanceFileVersion::V2_2 - && schema.fields_pre_order().any(|field| field.is_blob_v2()) - { - Some(Arc::new( - build_external_base_resolver(dataset, ¶ms).await?, - )) - } else { - None - }; +) -> Result> +where + OpenWriter: Fn(Arc, Schema, Path, WriterOptions) -> OpenWriterFuture + Send + Sync, + OpenWriterFuture: Future>> + Send, +{ let source_store_registry = dataset .map(|ds| ds.session.store_registry()) .unwrap_or_else(|| params.store_registry()); @@ -641,7 +623,7 @@ pub async fn do_write_fragments( object_store.clone(), base_dir, schema, - storage_version, + open_writer, target_bases_info, external_base_resolver, params.allow_external_blob_outside_bases, @@ -1265,6 +1247,20 @@ async fn build_external_base_resolver( Ok(ExternalBaseResolver::new(candidates, store_registry)) } +pub(super) async fn blob_v2_external_base_resolver( + dataset: Option<&Dataset>, + params: &WriteParams, + schema: &Schema, +) -> Result>> { + if schema.fields_pre_order().any(|field| field.is_blob_v2()) { + Ok(Some(Arc::new( + build_external_base_resolver(dataset, params).await?, + ))) + } else { + Ok(None) + } +} + /// Writes the given data to the dataset and returns fragments. /// /// NOTE: the fragments have not yet been assigned an ID. That must be done @@ -1274,8 +1270,13 @@ async fn build_external_base_resolver( /// This is a private variant that takes a `SendableRecordBatchStream` instead /// of a reader. We don't expose the stream at our interface because it is a /// DataFusion type. +/// +/// The caller must resolve `storage_version` once for the operation. Operations +/// that also select a commit format must reuse the same value when committing. +#[allow(clippy::too_many_arguments)] #[instrument(level = "debug", skip_all)] pub async fn write_fragments_internal( + storage_version: ConcreteFileVersion, dataset: Option<&Dataset>, object_store: Arc, base_dir: &Path, @@ -1303,103 +1304,73 @@ pub async fn write_fragments_internal( validate_external_blob_write_params(¶ms)?; let normalized_converted_schema = prepared_to_logical_blob_schema(&converted_schema)?; - let (schema, storage_version) = if let Some(dataset) = dataset { - match params.mode { - WriteMode::Append | WriteMode::Create => { - // Append mode, so we need to check compatibility - normalized_converted_schema.check_compatible( - dataset.schema(), - &SchemaCompareOptions { - // We don't care if the user claims their data is nullable / non-nullable. We will - // verify against the actual data. - compare_nullability: NullabilityComparison::Ignore, - allow_missing_if_nullable: true, - ignore_field_order: true, - compare_dictionary: dataset.is_legacy_storage(), - ..Default::default() - }, - )?; - validate_blob_threshold_metadata_for_append( - &normalized_converted_schema, - dataset.schema(), - )?; - let write_schema = dataset.schema().project_by_schema( - &normalized_converted_schema, - OnMissing::Error, - OnTypeMismatch::Error, - )?; - // Use the storage version from the dataset, ignoring any version from the user. - let data_storage_version = dataset - .manifest() - .data_storage_format - .lance_file_version()?; - (write_schema, data_storage_version) - } - WriteMode::Overwrite => { - // Overwrite, use the schema from the data. If the user specified - // a storage version use that. Otherwise use the version from the - // dataset. - let data_storage_version = params.data_storage_version.unwrap_or( - dataset - .manifest() - .data_storage_format - .lance_file_version()?, - ); - (normalized_converted_schema, data_storage_version) - } - } + versions::write_fragments( + storage_version, + dataset, + object_store, + base_dir, + normalized_converted_schema, + data, + params, + target_bases_info, + ) + .await +} + +pub(super) fn prepare_write_schema( + dataset: Option<&Dataset>, + normalized_converted_schema: Schema, + params: &WriteParams, + mut schema_compare_options: lance_core::datatypes::SchemaCompareOptions, +) -> Result { + let schema = if let Some(dataset) = dataset + && matches!(params.mode, WriteMode::Append | WriteMode::Create) + { + schema_compare_options.compare_nullability = NullabilityComparison::Ignore; + schema_compare_options.allow_missing_if_nullable = true; + schema_compare_options.ignore_field_order = true; + normalized_converted_schema.check_compatible(dataset.schema(), &schema_compare_options)?; + validate_blob_threshold_metadata_for_append( + &normalized_converted_schema, + dataset.schema(), + )?; + dataset.schema().project_by_schema( + &normalized_converted_schema, + OnMissing::Error, + OnTypeMismatch::Error, + )? } else { - // Brand new dataset, use the schema from the data and the storage version - // from the user or the default. - ( - normalized_converted_schema, - params.storage_version_or_default(), - ) + normalized_converted_schema }; + Ok(schema) +} - if storage_version < LanceFileVersion::V2_2 && schema.fields_pre_order().any(|f| f.is_blob_v2()) - { +pub(super) fn validate_legacy_blob_write_schema( + schema: &Schema, + version_debug: &str, +) -> Result<()> { + if schema.fields_pre_order().any(|field| field.is_blob_v2()) { return Err(Error::invalid_input(format!( - "Blob v2 requires file version >= 2.2 (got {:?})", - storage_version + "Blob v2 requires file version >= 2.2 (got {version_debug})" ))); } + Ok(()) +} - if storage_version >= LanceFileVersion::V2_2 - && let Some(blob_field_path) = legacy_blob_field_path(&schema) - { +pub(super) fn validate_blob_v2_write_schema(schema: &Schema) -> Result<()> { + if let Some(blob_field_path) = legacy_blob_field_path(schema) { return Err(Error::invalid_input(format!( "Legacy blob columns (field metadata key {BLOB_META_KEY:?}) are not supported for file version >= 2.2. Found legacy blob field: {blob_field_path}. Use the blob v2 extension type (ARROW:extension:name = \"lance.blob.v2\") and the new blob APIs (e.g. lance::blob::blob_field / lance::blob::BlobArrayBuilder)." ))); } - - let seed_writers = create_seed_writers(dataset, ¶ms, storage_version).await?; - - let fragments = do_write_fragments( - dataset, - object_store, - base_dir, - &schema, - data, - params, - storage_version, - target_bases_info, - seed_writers, - ) - .await?; - - Ok((fragments, schema)) + Ok(()) } -async fn create_seed_writers( +pub(crate) async fn create_seed_writers_current( dataset: Option<&Dataset>, params: &WriteParams, - storage_version: LanceFileVersion, ) -> Result>> { - // Seeds only make sense when appending to an existing dataset with V2 files. - if storage_version == LanceFileVersion::Legacy { - return Ok(Vec::new()); - } + // Seeds only make sense when appending to an existing dataset. if !matches!(params.mode, WriteMode::Append) { return Ok(Vec::new()); } @@ -1515,9 +1486,7 @@ where struct V2WriterAdapter { writer: current_writer::FileWriter, - version: ConcreteFileVersion, - path: String, - base_id: Option, + data_file: Option, preprocessor: Option, } @@ -1537,7 +1506,10 @@ impl GenericWriter for V2WriterAdapter { Ok(()) } fn data_file_path(&self) -> (&str, Option) { - (&self.path, self.base_id) + self.data_file + .as_ref() + .map(|data_file| (data_file.path.as_str(), data_file.base_id)) + .unwrap_or(("", None)) } async fn tell(&mut self) -> Result { Ok(self.writer.tell().await?) @@ -1559,14 +1531,13 @@ impl GenericWriter for V2WriterAdapter { .map(|(_, column_index)| *column_index as i32) .collect::>(); let write_summary = self.writer.finish().await?; - let data_file = DataFile::new( - std::mem::take(&mut self.path), - field_ids, - column_indices, - self.version, - NonZero::new(write_summary.size_bytes), - self.base_id, - ); + let mut data_file = self + .data_file + .take() + .ok_or_else(|| Error::internal("current writer was already finished"))?; + data_file.fields = field_ids.into(); + data_file.column_indices = column_indices.into(); + data_file.file_size_bytes = NonZero::new(write_summary.size_bytes).into(); Ok((write_summary.num_rows as u32, data_file)) } @@ -1579,60 +1550,8 @@ impl GenericWriter for V2WriterAdapter { } } -pub async fn open_writer( - object_store: &ObjectStore, - schema: &Schema, - base_dir: &Path, - storage_version: LanceFileVersion, -) -> Result> { - open_writer_with_options( - object_store, - schema, - base_dir, - storage_version, - WriterOptions { - add_data_dir: true, - ..Default::default() - }, - ) - .await -} - -pub(super) async fn open_update_writer( - dataset: &Dataset, - schema: &Schema, - storage_version: LanceFileVersion, -) -> Result> { - // add_columns / alter_columns reuse the normal writer stack, but they do not - // flow through WriteParams. Rebuild the external base resolver here so blob - // v2 reference columns can resolve dataset-registered external URIs. - let external_base_resolver = if storage_version >= LanceFileVersion::V2_2 - && schema.fields_pre_order().any(|f| f.is_blob_v2()) - { - Some(Arc::new( - build_external_base_resolver(Some(dataset), &WriteParams::default()).await?, - )) - } else { - None - }; - - open_writer_with_options( - &dataset.object_store, - schema, - &dataset.base, - storage_version, - WriterOptions { - add_data_dir: true, - external_base_resolver, - source_store_registry: dataset.session.store_registry(), - ..Default::default() - }, - ) - .await -} - #[derive(Default)] -struct WriterOptions { +pub(crate) struct WriterOptions { add_data_dir: bool, base_id: Option, external_base_resolver: Option>, @@ -1643,13 +1562,92 @@ struct WriterOptions { blob_pack_file_size_threshold: Option, } -async fn open_writer_with_options( +impl WriterOptions { + pub(super) fn update( + source_store_registry: Arc, + external_base_resolver: Option>, + ) -> Self { + Self { + add_data_dir: true, + external_base_resolver, + source_store_registry, + ..Default::default() + } + } +} + +pub(crate) async fn open_v1_writer( object_store: &ObjectStore, schema: &Schema, base_dir: &Path, - storage_version: LanceFileVersion, options: WriterOptions, ) -> Result> { + let WriterOptions { + add_data_dir, + base_id, + .. + } = options; + let (_data_file_key, filename, _data_dir, full_path) = + prepare_data_file_path(base_dir, add_data_dir); + Ok(Box::new(V1WriterAdapter { + writer: V1FileWriter::::try_new( + object_store, + &full_path, + schema.clone(), + &Default::default(), + ) + .await?, + path: filename, + base_id, + })) +} + +pub(in crate::dataset) async fn open_current_writer( + create_file_writer: F, + object_store: &ObjectStore, + schema: &Schema, + base_dir: &Path, + options: WriterOptions, +) -> Result> +where + F: FnOnce( + Box, + Schema, + String, + Option, + ) -> Result<(current_writer::FileWriter, DataFile)>, +{ + let WriterOptions { + add_data_dir, + base_id, + .. + } = options; + let (_data_file_key, filename, _data_dir, full_path) = + prepare_data_file_path(base_dir, add_data_dir); + let writer = object_store.create(&full_path).await?; + let (file_writer, data_file) = create_file_writer(writer, schema.clone(), filename, base_id)?; + Ok(Box::new(V2WriterAdapter { + writer: file_writer, + data_file: Some(data_file), + preprocessor: None, + })) +} + +pub(in crate::dataset) async fn open_current_blob_v2_writer( + create_file_writer: F, + object_store: &ObjectStore, + schema: &Schema, + base_dir: &Path, + options: WriterOptions, +) -> Result> +where + F: FnOnce( + Box, + Schema, + String, + Option, + ) -> Result<(current_writer::FileWriter, DataFile)>, +{ let WriterOptions { add_data_dir, base_id, @@ -1660,66 +1658,39 @@ async fn open_writer_with_options( source_store_params, blob_pack_file_size_threshold, } = options; + let (data_file_key, filename, data_dir, full_path) = + prepare_data_file_path(base_dir, add_data_dir); + let writer = object_store.create(&full_path).await?; + let (file_writer, data_file) = create_file_writer(writer, schema.clone(), filename, base_id)?; + let preprocessor = BlobPreprocessor::new( + object_store.clone(), + data_dir, + data_file_key, + schema, + external_base_resolver, + allow_external_blob_outside_bases, + external_blob_mode, + source_store_registry, + source_store_params, + blob_pack_file_size_threshold, + )?; + Ok(Box::new(V2WriterAdapter { + writer: file_writer, + data_file: Some(data_file), + preprocessor: Some(preprocessor), + })) +} +fn prepare_data_file_path(base_dir: &Path, add_data_dir: bool) -> (String, String, Path, Path) { let data_file_key = generate_random_filename(); let filename = format!("{}.lance", data_file_key); - let data_dir = if add_data_dir { base_dir.clone().join(DATA_DIR) } else { base_dir.clone() }; - let full_path = data_dir.clone().join(filename.as_str()); - - let writer = if storage_version == LanceFileVersion::Legacy { - Box::new(V1WriterAdapter { - writer: V1FileWriter::::try_new( - object_store, - &full_path, - schema.clone(), - &Default::default(), - ) - .await?, - path: filename, - base_id, - }) - } else { - let writer = object_store.create(&full_path).await?; - let enable_blob_v2 = storage_version >= LanceFileVersion::V2_2; - let version = ConcreteFileVersion::from(storage_version); - let file_writer = lance_file::versions::create_writer( - version, - writer, - schema.clone(), - FileWriterOptions::default(), - )?; - let preprocessor = if enable_blob_v2 { - Some(BlobPreprocessor::new( - object_store.clone(), - data_dir.clone(), - data_file_key.clone(), - schema, - external_base_resolver, - allow_external_blob_outside_bases, - external_blob_mode, - source_store_registry, - source_store_params, - blob_pack_file_size_threshold, - )?) - } else { - None - }; - let writer_adapter = V2WriterAdapter { - writer: file_writer, - version, - path: filename, - base_id, - preprocessor, - }; - Box::new(writer_adapter) as Box - }; - Ok(writer) + (data_file_key, filename, data_dir, full_path) } /// Reserved base id that refers to the dataset's primary storage in @@ -1744,13 +1715,13 @@ pub struct TargetBaseInfo { pub is_dataset_root: bool, } -struct WriterGenerator { +struct WriterGenerator { /// Default object store (used when no target bases specified) object_store: Arc, /// Default base directory (used when no target bases specified) base_dir: Path, schema: Schema, - storage_version: LanceFileVersion, + open_writer: OpenWriter, /// Target base information (if writing to specific bases) target_bases_info: Option>, external_base_resolver: Option>, @@ -1763,13 +1734,17 @@ struct WriterGenerator { next_base_index: AtomicUsize, } -impl WriterGenerator { +impl WriterGenerator +where + OpenWriter: Fn(Arc, Schema, Path, WriterOptions) -> OpenWriterFuture + Send + Sync, + OpenWriterFuture: Future>> + Send, +{ #[allow(clippy::too_many_arguments)] pub fn new( object_store: Arc, base_dir: &Path, schema: &Schema, - storage_version: LanceFileVersion, + open_writer: OpenWriter, target_bases_info: Option>, external_base_resolver: Option>, allow_external_blob_outside_bases: bool, @@ -1782,7 +1757,7 @@ impl WriterGenerator { object_store, base_dir: base_dir.clone(), schema: schema.clone(), - storage_version, + open_writer, target_bases_info, external_base_resolver, allow_external_blob_outside_bases, @@ -1810,11 +1785,10 @@ impl WriterGenerator { let fragment = Fragment::new(0); let writer = if let Some(base_info) = self.select_target_base() { - open_writer_with_options( - &base_info.object_store, - &self.schema, - &base_info.base_dir, - self.storage_version, + (self.open_writer)( + base_info.object_store.clone(), + self.schema.clone(), + base_info.base_dir.clone(), WriterOptions { add_data_dir: base_info.is_dataset_root, // Primary-storage slots stamp no base id, like a write @@ -1830,11 +1804,10 @@ impl WriterGenerator { ) .await? } else { - open_writer_with_options( - &self.object_store, - &self.schema, - &self.base_dir, - self.storage_version, + (self.open_writer)( + self.object_store.clone(), + self.schema.clone(), + self.base_dir.clone(), WriterOptions { add_data_dir: true, base_id: None, @@ -1895,6 +1868,7 @@ mod tests { use datafusion::{error::DataFusionError, physical_plan::stream::RecordBatchStreamAdapter}; use datafusion_physical_plan::RecordBatchStream; use futures::TryStreamExt; + use lance_datafusion::chunker::chunk_stream; use lance_datagen::{BatchCount, RowCount, array, gen_batch}; use lance_file::version::ConcreteFileVersion; use lance_file::versions::v1::reader::FileReader as V1FileReader; @@ -1902,6 +1876,32 @@ mod tests { use lance_io::traits::Reader; use lance_table::format::BasePath; + async fn open_v2_1_test_writer( + object_store: Arc, + schema: Schema, + base_dir: Path, + options: WriterOptions, + ) -> Result> { + open_current_writer( + |object_writer, schema, filename, base_id| { + let writer = lance_file::versions::v2_1::create_writer( + object_writer, + schema, + lance_file::writer::FileWriterOptions::default(), + )? + .into(); + let mut data_file = DataFile::new_unstarted(filename, ConcreteFileVersion::V2_1); + data_file.base_id = base_id; + Ok((writer, data_file)) + }, + &object_store, + &schema, + &base_dir, + options, + ) + .await + } + #[test] fn test_auto_cleanup_disabled_by_default() { // Auto-cleanup must be off by default: the cleanup hook is expensive on @@ -2017,6 +2017,7 @@ mod tests { let object_store = Arc::new(ObjectStore::memory()); write_fragments_internal( + write_params.storage_version_or_default(), None, object_store, &Path::from("test"), @@ -2070,6 +2071,7 @@ mod tests { let object_store = Arc::new(ObjectStore::memory()); write_fragments_internal( + write_params.storage_version_or_default(), None, object_store, &Path::from("test"), @@ -2131,6 +2133,7 @@ mod tests { let object_store = Arc::new(ObjectStore::memory()); write_fragments_internal( + write_params.storage_version_or_default(), None, object_store, &Path::from("test"), @@ -2244,6 +2247,7 @@ mod tests { let object_store = Arc::new(ObjectStore::memory()); let (fragments, _) = write_fragments_internal( + ConcreteFileVersion::from(version), None, object_store, &Path::from("test"), @@ -2319,6 +2323,7 @@ mod tests { let object_store = Arc::new(ObjectStore::memory()); let base_path = Path::from("test"); let (fragments, _) = write_fragments_internal( + ConcreteFileVersion::V1, None, object_store.clone(), &base_path, @@ -2528,7 +2533,7 @@ mod tests { object_store.clone(), &base_dir, &schema, - LanceFileVersion::Stable, + open_v2_1_test_writer, Some(target_bases), None, false, @@ -2574,11 +2579,11 @@ mod tests { let object_store = Arc::new(ObjectStore::memory()); let base_dir = Path::from("test/bucket2"); - let mut inner_writer = open_writer_with_options( + let mut inner_writer = versions::open_writer( + ConcreteFileVersion::from(LanceFileVersion::Stable), &object_store, &schema, &base_dir, - LanceFileVersion::Stable, WriterOptions { add_data_dir: false, // Don't add /data ..Default::default() @@ -2646,7 +2651,7 @@ mod tests { Arc::new(ObjectStore::memory()), &Path::from("default"), &schema, - LanceFileVersion::Stable, + open_v2_1_test_writer, Some(target_bases), None, false, @@ -3480,6 +3485,7 @@ mod tests { }; let result = write_fragments_internal( + write_params.storage_version_or_default(), None, object_store, &Path::from("test_empty"), @@ -3704,6 +3710,7 @@ mod tests { // Attempt to write data - should fail with IO error due to disk full let result = write_fragments_internal( + write_params.storage_version_or_default(), None, object_store, &Path::from("test_disk_full"), @@ -3885,14 +3892,14 @@ mod tests { futures::stream::iter(items), )); - let result = do_write_fragments( + let result = versions::write_fragments_direct( + ConcreteFileVersion::V2_1, None, object_store.clone(), &base_dir, &schema, stream, WriteParams::default(), - LanceFileVersion::V2_1, None, Vec::new(), ) @@ -3943,7 +3950,8 @@ mod tests { futures::stream::iter(items), )); - let result = do_write_fragments( + let result = versions::write_fragments_direct( + ConcreteFileVersion::V2_1, None, object_store.clone(), &base_dir, @@ -3953,7 +3961,6 @@ mod tests { max_rows_per_file: 3, ..Default::default() }, - LanceFileVersion::V2_1, None, Vec::new(), ) @@ -4171,7 +4178,8 @@ mod tests { is_dataset_root: true, }]; - let result = do_write_fragments( + let result = versions::write_fragments_direct( + ConcreteFileVersion::V2_1, None, object_store.clone(), &base_dir, @@ -4181,7 +4189,6 @@ mod tests { max_rows_per_file: 3, ..Default::default() }, - LanceFileVersion::V2_1, Some(target_bases), vec![], ) diff --git a/rust/lance/src/dataset/write/commit.rs b/rust/lance/src/dataset/write/commit.rs index 75fe8ba6440..fd57f68f9e7 100644 --- a/rust/lance/src/dataset/write/commit.rs +++ b/rust/lance/src/dataset/write/commit.rs @@ -40,7 +40,7 @@ pub struct CommitBuilder<'a> { dest: WriteDestination<'a>, use_stable_row_ids: Option, enable_v2_manifest_paths: bool, - storage_format: Option, + storage_format: Option, commit_handler: Option>, store_params: Option, object_store: Option>, @@ -98,6 +98,11 @@ impl<'a> CommitBuilder<'a> { /// All data files must use the same storage format as the existing dataset. /// If a different format is passed, an error will be returned. pub fn with_storage_format(mut self, storage_format: LanceFileVersion) -> Self { + self.storage_format = Some(storage_format.into()); + self + } + + pub(crate) fn with_exact_storage_format(mut self, storage_format: ConcreteFileVersion) -> Self { self.storage_format = Some(storage_format); self } @@ -389,8 +394,7 @@ impl<'a> CommitBuilder<'a> { if let Some(ds) = dest.dataset() && let Some(storage_format) = self.storage_format { - let passed_storage_format = - DataStorageFormat::new(ConcreteFileVersion::from(storage_format)); + let passed_storage_format = DataStorageFormat::new(storage_format); if ds.manifest.data_storage_format != passed_storage_format && !matches!(transaction.operation, Operation::Overwrite { .. }) { @@ -404,10 +408,7 @@ impl<'a> CommitBuilder<'a> { let manifest_config = ManifestWriteConfig { use_stable_row_ids, - storage_format: self - .storage_format - .map(ConcreteFileVersion::from) - .map(DataStorageFormat::new), + storage_format: self.storage_format.map(DataStorageFormat::new), ..Default::default() }; diff --git a/rust/lance/src/dataset/write/insert.rs b/rust/lance/src/dataset/write/insert.rs index 02759c9d4fd..c0c40b8bb53 100644 --- a/rust/lance/src/dataset/write/insert.rs +++ b/rust/lance/src/dataset/write/insert.rs @@ -7,10 +7,12 @@ use std::sync::Arc; use arrow_array::{RecordBatch, RecordBatchIterator}; use datafusion::execution::SendableRecordBatchStream; use humantime::format_duration; -use lance_core::datatypes::{NullabilityComparison, Schema, SchemaCompareOptions}; +use lance_core::datatypes::{NullabilityComparison, Schema}; use lance_core::is_system_column; use lance_core::utils::tracing::{DATASET_WRITING_EVENT, TRACE_DATASET_EVENTS}; use lance_datafusion::utils::StreamingWriteSource; +use lance_file::version::ConcreteFileVersion; +#[cfg(test)] use lance_file::version::LanceFileVersion; use lance_io::object_store::ObjectStore; use lance_table::feature_flags::can_write_dataset; @@ -136,7 +138,7 @@ impl<'a> InsertBuilder<'a> { async fn do_commit(context: &WriteContext<'_>, transaction: Transaction) -> Result { let mut commit_builder = CommitBuilder::new(context.dest.clone()) .use_stable_row_ids(context.params.enable_stable_row_ids) - .with_storage_format(context.storage_version) + .with_exact_storage_format(context.storage_version) .enable_v2_manifest_paths(context.params.enable_v2_manifest_paths) .with_commit_handler(context.commit_handler.clone()) .with_object_store(context.object_store.clone()) @@ -214,6 +216,7 @@ impl<'a> InsertBuilder<'a> { .await?; let (written_fragments, written_schema) = write_fragments_internal( + context.storage_version, context.dest.dataset(), context.object_store.clone(), &context.base_path, @@ -316,13 +319,11 @@ impl<'a> InsertBuilder<'a> { context.params.enable_stable_row_ids = dataset.manifest.uses_stable_row_ids(); } - let schema_cmp_opts = SchemaCompareOptions { - compare_dictionary: dataset.manifest.should_use_legacy_format(), - compare_nullability: NullabilityComparison::Ignore, - allow_missing_if_nullable: true, - ignore_field_order: true, - ..Default::default() - }; + let version = dataset.manifest.data_storage_format.lance_file_format(); + let mut schema_cmp_opts = crate::dataset::versions::schema_compare_options(version); + schema_cmp_opts.compare_nullability = NullabilityComparison::Ignore; + schema_cmp_opts.allow_missing_if_nullable = true; + schema_cmp_opts.ignore_field_order = true; let normalized_data_schema = prepared_to_logical_blob_schema(data_schema)?; normalized_data_schema.check_compatible(dataset.schema(), &schema_cmp_opts)?; @@ -416,15 +417,15 @@ impl<'a> InsertBuilder<'a> { (WriteMode::Overwrite, WriteDestination::Dataset(dataset)) => { // If overwriting an existing dataset, allow the user to specify but use // the existing version if they don't - params.data_storage_version.map(Ok).unwrap_or_else(|| { - let m = dataset.manifest.as_ref(); - m.data_storage_format.lance_file_version() - })? + params + .data_storage_version + .map(ConcreteFileVersion::from) + .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_version()? + m.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(), @@ -448,7 +449,7 @@ struct WriteContext<'a> { object_store: Arc, base_path: Path, commit_handler: Arc, - storage_version: LanceFileVersion, + storage_version: ConcreteFileVersion, } #[cfg(test)] diff --git a/rust/lance/src/dataset/write/merge_insert.rs b/rust/lance/src/dataset/write/merge_insert.rs index a325d5de037..5fefa8603e8 100644 --- a/rust/lance/src/dataset/write/merge_insert.rs +++ b/rust/lance/src/dataset/write/merge_insert.rs @@ -57,7 +57,7 @@ use crate::{ dataset::{ fragment::{FileFragment, FragReadConfig}, transaction::{Operation, Transaction}, - write::{merge_insert::logical_plan::MergeInsertPlanner, open_writer}, + write::merge_insert::logical_plan::MergeInsertPlanner, }, index::DatasetIndexInternalExt, io::exec::{ @@ -1260,11 +1260,15 @@ impl MergeInsertJob { .manifest() .data_storage_format .lance_file_version()?; - let mut writer = open_writer( + let mut writer = crate::dataset::versions::open_writer( + data_storage_version.into(), &dataset.object_store, &write_schema, &dataset.base, - data_storage_version, + super::WriterOptions { + add_data_dir: true, + ..Default::default() + }, ) .await?; @@ -1486,6 +1490,7 @@ impl MergeInsertJob { )?; let (fragments, _) = write_fragments_internal( + dataset.manifest.data_storage_format.lance_file_format(), Some(dataset.as_ref()), dataset.object_store.clone(), &dataset.base, @@ -2267,6 +2272,10 @@ impl MergeInsertJob { } else { let cleanup_bases = target_bases_info.clone(); let (mut new_fragments, _) = write_fragments_internal( + self.dataset + .manifest + .data_storage_format + .lance_file_format(), Some(&self.dataset), self.dataset.object_store.clone(), &self.dataset.base, diff --git a/rust/lance/src/dataset/write/merge_insert/exec/write.rs b/rust/lance/src/dataset/write/merge_insert/exec/write.rs index 7f17da26965..b4bb8112ae7 100644 --- a/rust/lance/src/dataset/write/merge_insert/exec/write.rs +++ b/rust/lance/src/dataset/write/merge_insert/exec/write.rs @@ -992,6 +992,7 @@ impl ExecutionPlan for FullSchemaMergeInsertExec { // Keep a copy so failures after the write can clean up routed files. let cleanup_bases = target_bases_info.clone(); let (mut new_fragments, _) = write_fragments_internal( + dataset.manifest.data_storage_format.lance_file_format(), Some(&dataset), dataset.object_store.clone(), &dataset.base, diff --git a/rust/lance/src/dataset/write/update.rs b/rust/lance/src/dataset/write/update.rs index 68bc1c2f7cc..075ab0916b4 100644 --- a/rust/lance/src/dataset/write/update.rs +++ b/rust/lance/src/dataset/write/update.rs @@ -320,18 +320,17 @@ impl UpdateJob { }); let stream = RecordBatchStreamAdapter::new(schema, stream); - let version = self - .dataset - .manifest() - .data_storage_format - .lance_file_version()?; let (mut new_fragments, _) = write_fragments_internal( + self.dataset + .manifest + .data_storage_format + .lance_file_format(), Some(&self.dataset), self.dataset.object_store.clone(), &self.dataset.base, self.dataset.schema().clone(), Box::pin(stream), - WriteParams::with_storage_version(version), + WriteParams::default(), None, // TODO: support multiple bases for update ) .await?;