diff --git a/protos/transaction.proto b/protos/transaction.proto index bbde96fed5d..e2813c50ba1 100644 --- a/protos/transaction.proto +++ b/protos/transaction.proto @@ -332,6 +332,24 @@ message Transaction { repeated DataReplacementGroup replacements = 1; } + // Files produced for one fragment by a physical column-layout rewrite. + // The files are positionally aligned with the fragment and collectively + // provide exactly field_ids. + message OptimizeColumnsGroup { + uint64 fragment_id = 1; + repeated uint32 field_ids = 2; + repeated DataFile new_files = 3; + uint64 physical_rows = 4; + } + + // Reorganize selected fields without changing their logical values or the + // fragment's row layout. materialized_through_version is the immutable + // snapshot cutoff used when the output files were produced. + message OptimizeColumns { + uint64 materialized_through_version = 1; + repeated OptimizeColumnsGroup groups = 2; + } + // Overlay files to append to a single fragment, in order (the last entry is // newest). The overlays are appended to the fragment's existing `overlays` // list; they do not replace it, so overlays written by concurrent commits are @@ -394,6 +412,7 @@ message Transaction { Clone clone = 113; UpdateBases update_bases = 114; DataOverlay data_overlay = 115; + OptimizeColumns optimize_columns = 116; } // Fields 200/202 (`blob_append` / `blob_overwrite`) previously represented blob dataset ops. diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index fa2746cd313..f48319372d3 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -79,6 +79,7 @@ pub mod index; pub mod mem_wal; mod metadata; pub mod optimize; +pub mod optimize_columns; pub(crate) mod overlay; pub mod progress; pub mod refs; diff --git a/rust/lance/src/dataset/fragment.rs b/rust/lance/src/dataset/fragment.rs index 87c63a1549d..9983f16de02 100644 --- a/rust/lance/src/dataset/fragment.rs +++ b/rust/lance/src/dataset/fragment.rs @@ -23,9 +23,10 @@ use datafusion::scalar::ScalarValue; use futures::future::{BoxFuture, try_join_all}; use futures::{FutureExt, Stream, StreamExt, TryFutureExt, TryStreamExt, join, stream}; use lance_arrow::json::{convert_json_columns, has_json_fields, is_arrow_json_field}; -use lance_arrow::{RecordBatchExt, SchemaExt}; +use lance_arrow::{FieldExt, RecordBatchExt, SchemaExt}; use lance_core::datatypes::{ - BlobHandling, NullabilityComparison, OnMissing, OnTypeMismatch, SchemaCompareOptions, + BlobHandling, BlobV2Layout, NullabilityComparison, OnMissing, OnTypeMismatch, + SchemaCompareOptions, }; use lance_core::utils::address::RowAddress; use lance_core::utils::deletion::DeletionVector; @@ -746,6 +747,153 @@ fn relax_nullability(field: &ArrowField) -> ArrowField { ArrowField::new(field.name(), data_type, true).with_metadata(field.metadata().clone()) } +fn arrow_field_with_data_type(field: &ArrowField, data_type: DataType) -> ArrowField { + ArrowField::new(field.name(), data_type, field.is_nullable()) + .with_metadata(field.metadata().clone()) +} + +fn validate_blob_v2_write_field(field: &ArrowField) -> Result<()> { + if !field.is_blob_v2() { + return Err(Error::invalid_input(format!( + "field '{}' is not tagged as a blob v2 column", + field.name() + ))); + } + let DataType::Struct(children) = field.data_type() else { + return Err(Error::invalid_input(format!( + "blob v2 field '{}' has non-struct type {}", + field.name(), + field.data_type() + ))); + }; + match BlobV2Layout::classify(children) { + Some(BlobV2Layout::Logical | BlobV2Layout::Prepared) => Ok(()), + actual => Err(Error::invalid_input(format!( + "blob v2 field '{}' has {} layout; expected logical or prepared writer input", + field.name(), + actual + .map(|layout| layout.to_string()) + .unwrap_or_else(|| "unrecognized".to_string()) + ))), + } +} + +/// Normalize accepted blob writer layouts to the dataset field shape so the +/// regular schema comparison can continue to validate every non-blob detail. +fn normalize_blob_v2_for_comparison( + source: &ArrowField, + target: Option<&ArrowField>, +) -> Result { + if let Some(target) = target + && target.is_blob_v2() + { + validate_blob_v2_write_field(source)?; + return Ok(target.clone()); + } + + let data_type = match source.data_type() { + DataType::Struct(source_children) => { + let target_children = target.and_then(|target| match target.data_type() { + DataType::Struct(children) => Some(children), + _ => None, + }); + DataType::Struct( + source_children + .iter() + .map(|source_child| { + let target_child = target_children.and_then(|children| { + children + .iter() + .find(|target_child| target_child.name() == source_child.name()) + .map(AsRef::as_ref) + }); + normalize_blob_v2_for_comparison(source_child, target_child).map(Arc::new) + }) + .collect::>>()? + .into(), + ) + } + DataType::List(source_child) => { + let target_child = target.and_then(|target| match target.data_type() { + DataType::List(child) => Some(child.as_ref()), + _ => None, + }); + DataType::List(Arc::new(normalize_blob_v2_for_comparison( + source_child, + target_child, + )?)) + } + DataType::LargeList(source_child) => { + let target_child = target.and_then(|target| match target.data_type() { + DataType::LargeList(child) => Some(child.as_ref()), + _ => None, + }); + DataType::LargeList(Arc::new(normalize_blob_v2_for_comparison( + source_child, + target_child, + )?)) + } + other => other.clone(), + }; + Ok(arrow_field_with_data_type(source, data_type)) +} + +/// Build a name-ordered projection while retaining complete logical blob v2 +/// fields (`position` and `size`) needed to preserve external object slices. +fn blob_aware_projection_field(source: &ArrowField, target: &ArrowField) -> Result { + if target.is_blob_v2() { + validate_blob_v2_write_field(source)?; + return Ok(source.clone()); + } + + let data_type = match target.data_type() { + DataType::Struct(target_children) => { + let DataType::Struct(source_children) = source.data_type() else { + return Ok(target.clone()); + }; + DataType::Struct( + target_children + .iter() + .map(|target_child| { + let source_child = source_children + .iter() + .find(|source_child| source_child.name() == target_child.name()) + .ok_or_else(|| { + Error::invalid_input(format!( + "field '{}' is missing child '{}'", + target.name(), + target_child.name() + )) + })?; + blob_aware_projection_field(source_child, target_child).map(Arc::new) + }) + .collect::>>()? + .into(), + ) + } + DataType::List(target_child) => { + let DataType::List(source_child) = source.data_type() else { + return Ok(target.clone()); + }; + DataType::List(Arc::new(blob_aware_projection_field( + source_child, + target_child, + )?)) + } + DataType::LargeList(target_child) => { + let DataType::LargeList(source_child) = source.data_type() else { + return Ok(target.clone()); + }; + DataType::LargeList(Arc::new(blob_aware_projection_field( + source_child, + target_child, + )?)) + } + other => other.clone(), + }; + Ok(arrow_field_with_data_type(target, data_type)) +} + impl FileFragment { /// Creates a new FileFragment. pub fn new(dataset: Arc, metadata: Fragment) -> Self { @@ -2289,13 +2437,6 @@ impl FileFragment { metadata: schema.metadata.clone(), }; let batch_schema = ArrowSchema::from(&writer_schema); - let projection_schema = ArrowSchema::new( - batch_schema - .fields() - .iter() - .map(|field| relax_nullability(field)) - .collect::>(), - ); let file_version = self .dataset @@ -2351,7 +2492,24 @@ impl FileFragment { if let Some(duplicate) = duplicate_field_path(batch.schema_ref().fields(), "") { return Err(self.schema_mismatch(format!("column '{duplicate}' appears twice"))); } - LanceSchema::try_from(batch.schema_ref().as_ref()) + let normalized_fields = batch + .schema_ref() + .fields() + .iter() + .map(|source| { + let target = batch_schema + .fields() + .iter() + .find(|target| target.name() == source.name()) + .map(AsRef::as_ref); + normalize_blob_v2_for_comparison(source, target).map(Arc::new) + }) + .collect::>>()?; + let normalized_schema = ArrowSchema::new_with_metadata( + normalized_fields, + batch.schema_ref().metadata().clone(), + ); + LanceSchema::try_from(&normalized_schema) .and_then(|staged| { staged.check_compatible( &writer_schema, @@ -2363,6 +2521,28 @@ impl FileFragment { ) }) .map_err(|mismatch| self.schema_mismatch(mismatch))?; + + let projection_fields = batch_schema + .fields() + .iter() + .map(|target| { + let source = batch + .schema_ref() + .fields() + .iter() + .find(|source| source.name() == target.name()) + .ok_or_else(|| { + self.schema_mismatch(format!( + "column '{}' is missing", + target.name() + )) + })?; + blob_aware_projection_field(source, target) + .map(|field| Arc::new(relax_nullability(&field))) + .map_err(|error| self.schema_mismatch(error)) + }) + .collect::>>()?; + let projection_schema = ArrowSchema::new(projection_fields); let batch = batch .project_by_schema(&projection_schema) .map_err(|err| self.schema_mismatch(err))?; diff --git a/rust/lance/src/dataset/optimize_columns.rs b/rust/lance/src/dataset/optimize_columns.rs new file mode 100644 index 00000000000..33e85cbabba --- /dev/null +++ b/rust/lance/src/dataset/optimize_columns.rs @@ -0,0 +1,975 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Physical column-layout compaction without changing fragment row layout. + +use std::collections::{HashMap, HashSet}; +use std::sync::Arc; + +use arrow_array::RecordBatch; +use futures::{StreamExt, stream}; +use lance_core::datatypes::{BlobHandling, Schema}; +use lance_file::version::ConcreteFileVersion; +use lance_table::format::{DataFile, Fragment}; +use serde::{Deserialize, Serialize}; + +use super::transaction::{Operation, OptimizeColumnsGroup, Transaction}; +use super::{CommitBuilder, Dataset}; +use crate::{Error, Result}; + +/// One requested physical data-file grouping. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ColumnGroup { + /// Top-level logical fields to place in the same data file. + pub fields: Vec, +} + +/// Options for [`optimize_columns`]. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct OptimizeColumnsOptions { + /// Explicit target file groups. Fields absent from all groups are not rewritten. + pub groups: Vec, + /// Fragment ids to consider, or every current fragment when absent. + pub fragment_ids: Option>, + /// Maximum number of group files staged concurrently. + pub max_concurrency: Option, +} + +/// Metrics returned by [`optimize_columns`]. +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct OptimizeColumnsMetrics { + /// Number of selected fragments inspected by the planner. + pub fragments_examined: usize, + /// Number of fragments whose provider layout was rewritten. + pub fragments_rewritten: usize, + /// Number of new data files installed. + pub files_added: usize, + /// Number of old data files removed from the latest fragment descriptors. + pub files_removed: usize, + /// Number of mixed old files retained after selected fields were tombstoned. + pub mixed_files_retained: usize, + /// Bytes read while staging and committing the operation. + pub bytes_read: u64, + /// Bytes written while staging and committing the operation. + pub bytes_written: u64, +} + +/// Current physical column-layout statistics for one fragment. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct FragmentColumnLayoutStats { + /// Fragment id. + pub fragment_id: u64, + /// Data files with at least one live field provider. + pub live_file_count: usize, + /// Known byte size for each data-file descriptor. + pub file_sizes: Vec>, + /// Number of live physical fields in each data-file descriptor. + pub fields_per_file: Vec, + /// Number of overlay descriptors attached to the fragment. + pub overlay_count: usize, + /// Fraction of data-file field slots that are tombstoned. + pub tombstoned_field_ratio: f64, +} + +#[derive(Debug, Clone)] +struct PlannedGroup { + fragment: Fragment, + schema: Schema, + field_names: Vec, + field_ids: Vec, +} + +#[derive(Debug)] +struct StagedGroup { + fragment_id: u64, + field_ids: Vec, + physical_rows: u64, + data_file: DataFile, +} + +impl Dataset { + /// Reorganize selected top-level fields into explicit per-fragment data files. + /// + /// Logical values, fragment ids, physical row order, deletion files, stable + /// row ids, and row-version metadata are preserved. Fields not named by any + /// group keep their current providers. + /// + /// # Example + /// + /// ``` + /// # use lance::Dataset; + /// # use lance::dataset::optimize_columns::{ColumnGroup, OptimizeColumnsOptions}; + /// # async fn compact_features(dataset: &mut Dataset) -> lance::Result<()> { + /// let metrics = dataset + /// .optimize_columns(OptimizeColumnsOptions { + /// groups: vec![ColumnGroup { + /// fields: vec!["feature_a".into(), "feature_b".into()], + /// }], + /// fragment_ids: None, + /// max_concurrency: Some(8), + /// }) + /// .await?; + /// assert!(metrics.files_added <= metrics.fragments_examined); + /// # Ok(()) + /// # } + /// ``` + pub async fn optimize_columns( + &mut self, + options: OptimizeColumnsOptions, + ) -> Result { + optimize_columns(self, options).await + } + + /// Return descriptor-only physical column-layout statistics. + pub fn column_layout_stats(&self) -> Vec { + self.manifest + .fragments + .iter() + .map(|fragment| { + let fields_per_file = fragment + .files + .iter() + .map(|file| file.fields.iter().filter(|field| **field >= 0).count()) + .collect::>(); + let live_file_count = fields_per_file.iter().filter(|count| **count > 0).count(); + let file_sizes = fragment + .files + .iter() + .map(|file| file.file_size_bytes.get().map(|size| size.get())) + .collect(); + let field_slots = fragment + .files + .iter() + .map(|file| file.fields.len()) + .sum::(); + let tombstones = fragment + .files + .iter() + .flat_map(|file| file.fields.iter()) + .filter(|field| **field < 0) + .count(); + + FragmentColumnLayoutStats { + fragment_id: fragment.id, + live_file_count, + file_sizes, + fields_per_file, + overlay_count: fragment.overlays.len(), + tombstoned_field_ratio: if field_slots == 0 { + 0.0 + } else { + tombstones as f64 / field_slots as f64 + }, + } + }) + .collect() + } +} + +/// Reorganize selected top-level fields into explicit per-fragment data files. +pub async fn optimize_columns( + dataset: &mut Dataset, + options: OptimizeColumnsOptions, +) -> Result { + let (selected_fragments, resolved_groups, max_concurrency) = + validate_and_resolve(dataset, &options)?; + let mut metrics = OptimizeColumnsMetrics { + fragments_examined: selected_fragments.len(), + ..Default::default() + }; + + let mut plans = Vec::new(); + let live_schema_ids = dataset + .schema() + .fields_pre_order() + .map(|field| field.id) + .collect::>(); + for fragment in &selected_fragments { + let materialized_fields = fragment + .files + .iter() + .flat_map(|file| file.fields.iter().copied()) + .chain( + fragment + .overlays + .iter() + .flat_map(|overlay| overlay.data_file.fields.iter().copied()), + ) + .filter(|field| *field >= 0) + .map(|field| field as u32) + .collect::>(); + + let mut rewritten_fields = HashSet::new(); + for (schema, field_names, requested_ids) in &resolved_groups { + let field_ids = requested_ids + .iter() + .copied() + .filter(|field| materialized_fields.contains(field)) + .collect::>(); + if field_ids.is_empty() || group_is_already_optimized(fragment, &field_ids) { + continue; + } + rewritten_fields.extend(field_ids.iter().copied()); + plans.push(PlannedGroup { + fragment: fragment.clone(), + schema: schema.clone(), + field_names: field_names.clone(), + field_ids, + }); + } + + if !rewritten_fields.is_empty() { + metrics.fragments_rewritten += 1; + for file in &fragment.files { + let selected_in_file = file + .fields + .iter() + .any(|field| *field >= 0 && rewritten_fields.contains(&(*field as u32))); + if !selected_in_file { + continue; + } + let retains_live_field = file.fields.iter().any(|field| { + live_schema_ids.contains(field) && !rewritten_fields.contains(&(*field as u32)) + }); + if retains_live_field { + metrics.mixed_files_retained += 1; + } else { + metrics.files_removed += 1; + } + } + } + } + + if plans.is_empty() { + return Ok(metrics); + } + + let io_before = dataset.object_store.io_stats_snapshot(); + let dataset_snapshot = Arc::new(dataset.clone()); + let staged_results = stream::iter(plans.into_iter().map(|plan| { + let dataset = dataset_snapshot.clone(); + async move { stage_group(dataset, plan).await } + })) + .buffer_unordered(max_concurrency) + .collect::>() + .await; + + let mut staged = Vec::with_capacity(staged_results.len()); + let mut first_error = None; + for result in staged_results { + match result { + Ok(group) => staged.push(group), + Err(error) if first_error.is_none() => first_error = Some(error), + Err(_) => {} + } + } + if let Some(error) = first_error { + discard_staged_groups(dataset_snapshot.as_ref(), &staged).await; + return Err(error); + } + + metrics.files_added = staged.len(); + let mut groups_by_fragment: HashMap = HashMap::new(); + for staged_group in staged { + let group = groups_by_fragment + .entry(staged_group.fragment_id) + .or_insert_with(|| OptimizeColumnsGroup { + fragment_id: staged_group.fragment_id, + field_ids: Vec::new(), + new_files: Vec::new(), + physical_rows: staged_group.physical_rows, + }); + group.field_ids.extend(staged_group.field_ids); + group.new_files.push(staged_group.data_file); + } + + let materialized_through_version = dataset.manifest.version; + let transaction = Transaction::new( + materialized_through_version, + Operation::OptimizeColumns { + materialized_through_version, + groups: groups_by_fragment.into_values().collect(), + }, + None, + ); + let committed = CommitBuilder::new(dataset_snapshot) + .execute(transaction) + .await?; + *dataset = committed; + + let io_after = dataset.object_store.io_stats_snapshot(); + metrics.bytes_read = io_after.read_bytes.saturating_sub(io_before.read_bytes); + metrics.bytes_written = io_after + .written_bytes + .saturating_sub(io_before.written_bytes); + Ok(metrics) +} + +type ResolvedGroup = (Schema, Vec, Vec); + +fn validate_and_resolve( + dataset: &Dataset, + options: &OptimizeColumnsOptions, +) -> Result<(Vec, Vec, usize)> { + if options.groups.is_empty() { + return Err(Error::invalid_input( + "OptimizeColumnsOptions.groups must not be empty", + )); + } + let max_concurrency = options + .max_concurrency + .unwrap_or_else(|| dataset.object_store.io_parallelism()); + if max_concurrency == 0 { + return Err(Error::invalid_input( + "OptimizeColumnsOptions.max_concurrency must be greater than zero", + )); + } + + let mut named_fields = HashSet::new(); + let mut resolved_groups = Vec::with_capacity(options.groups.len()); + for group in &options.groups { + if group.fields.is_empty() { + return Err(Error::invalid_input( + "OptimizeColumns ColumnGroup.fields must not be empty", + )); + } + let mut fields = Vec::with_capacity(group.fields.len()); + for name in &group.fields { + if !named_fields.insert(name.as_str()) { + return Err(Error::invalid_input(format!( + "OptimizeColumns field '{name}' appears in more than one group" + ))); + } + let field = dataset + .schema() + .fields + .iter() + .find(|field| field.name == *name) + .ok_or_else(|| { + Error::invalid_input(format!( + "OptimizeColumns field '{name}' is not a top-level dataset field" + )) + })?; + fields.push(field.clone()); + } + let schema = Schema { + fields, + metadata: dataset.schema().metadata.clone(), + }; + let field_ids: Vec = schema + .fields_pre_order() + .map(|field| field.id as u32) + .collect(); + resolved_groups.push((schema, group.fields.clone(), field_ids)); + } + + let selected_fragments = if let Some(fragment_ids) = &options.fragment_ids { + let unique = fragment_ids.iter().copied().collect::>(); + if unique.len() != fragment_ids.len() { + return Err(Error::invalid_input( + "OptimizeColumnsOptions.fragment_ids contains duplicates", + )); + } + let known = dataset + .manifest + .fragments + .iter() + .map(|fragment| fragment.id) + .collect::>(); + if let Some(missing) = fragment_ids.iter().find(|id| !known.contains(id)) { + return Err(Error::invalid_input(format!( + "OptimizeColumns fragment id {missing} does not exist" + ))); + } + dataset + .manifest + .fragments + .iter() + .filter(|fragment| unique.contains(&fragment.id)) + .cloned() + .collect::>() + } else { + dataset.manifest.fragments.as_ref().clone() + }; + + for fragment in &selected_fragments { + for file in &fragment.files { + if file.file_version()? == ConcreteFileVersion::V1 + && resolved_groups.iter().any(|(_, _, fields)| { + file.fields + .iter() + .any(|field| *field >= 0 && fields.contains(&(*field as u32))) + }) + { + return Err(Error::not_supported(format!( + "OptimizeColumns does not support legacy fragment {}", + fragment.id + ))); + } + } + } + + Ok((selected_fragments, resolved_groups, max_concurrency)) +} + +fn group_is_already_optimized(fragment: &Fragment, field_ids: &[u32]) -> bool { + let target = field_ids.iter().copied().collect::>(); + let has_materializable_overlay = fragment.overlays.iter().any(|overlay| { + overlay + .data_file + .fields + .iter() + .any(|field| *field >= 0 && target.contains(&(*field as u32))) + }); + if has_materializable_overlay { + return false; + } + + let providers = fragment + .files + .iter() + .filter(|file| { + file.fields + .iter() + .any(|field| *field >= 0 && target.contains(&(*field as u32))) + }) + .collect::>(); + if providers.len() != 1 { + return false; + } + let provided = providers[0] + .fields + .iter() + .filter(|field| **field >= 0) + .map(|field| *field as u32) + .collect::>(); + provided == target +} + +async fn stage_group(dataset: Arc, plan: PlannedGroup) -> Result { + let physical_rows = plan.fragment.physical_rows.ok_or_else(|| { + Error::invalid_input(format!( + "OptimizeColumns target fragment {} has no physical row count", + plan.fragment.id + )) + })? as u64; + let mut scanner = dataset.scan(); + let has_legacy_blob = plan + .schema + .fields_pre_order() + .any(|field| field.is_blob() && !field.is_blob_v2()); + if has_legacy_blob { + scanner.blob_handling(BlobHandling::AllBinary); + } + let has_blob_v2 = plan + .schema + .fields_pre_order() + .any(|field| field.is_blob_v2()); + if has_blob_v2 { + scanner.with_row_address(); + } + let names = plan + .field_names + .iter() + .map(String::as_str) + .collect::>(); + scanner + .project(&names)? + .with_fragments(vec![plan.fragment.clone()]) + .with_row_id() + .include_deleted_rows(); + let field_names = plan.field_names.clone(); + let rewrite_dataset = dataset.clone(); + let rewrite_schema = plan.schema.clone(); + let data = scanner.try_into_stream().await?.then(move |batch| { + let field_names = field_names.clone(); + let rewrite_dataset = rewrite_dataset.clone(); + let rewrite_schema = rewrite_schema.clone(); + async move { + let batch = batch?; + let batch = if has_blob_v2 { + super::optimize::transform_blob_v2_batch( + &rewrite_dataset, + &rewrite_schema, + batch, + false, + ) + .await? + } else { + batch + }; + project_user_fields(&batch, &field_names) + } + }); + + let fragment = super::FileFragment::new(dataset, plan.fragment); + let super::transaction::DataReplacementGroup(_, mut data_file) = + fragment.write_column(data, &plan.schema).await?; + let selected = plan.field_ids.iter().copied().collect::>(); + data_file.fields = data_file + .fields + .iter() + .map(|field| { + if *field >= 0 && !selected.contains(&(*field as u32)) { + lance_table::format::overlay::TOMBSTONE_FIELD_ID + } else { + *field + } + }) + .collect::>() + .into(); + + Ok(StagedGroup { + fragment_id: fragment.id() as u64, + field_ids: plan.field_ids, + physical_rows, + data_file, + }) +} + +fn project_user_fields(batch: &RecordBatch, field_names: &[String]) -> Result { + let indices = field_names + .iter() + .map(|name| { + batch.schema().index_of(name).map_err(|error| { + Error::internal(format!( + "OptimizeColumns scan did not return projected field '{name}': {error}" + )) + }) + }) + .collect::>>()?; + Ok(batch.project(&indices)?) +} + +async fn discard_staged_groups(dataset: &Dataset, groups: &[StagedGroup]) { + for group in groups { + let path = dataset.data_dir().join(group.data_file.path.as_str()); + if let Some(stem) = path + .filename() + .and_then(|name| name.strip_suffix(".lance")) + .map(|stem| dataset.data_dir().join(stem)) + && let Err(error) = dataset.object_store.remove_dir_all(stem.clone()).await + { + log::warn!("failed to delete staged OptimizeColumns blob sidecars '{stem}': {error}"); + } + if let Err(error) = dataset.object_store.delete(&path).await { + log::warn!( + "failed to delete staged OptimizeColumns file '{}': {error}", + group.data_file.path + ); + } + } +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use arrow_array::{ + ArrayRef, Int32Array, LargeBinaryArray, RecordBatch, RecordBatchIterator, StructArray, + }; + use arrow_schema::{DataType, Field, Schema as ArrowSchema}; + use futures::stream; + use lance_arrow::BLOB_META_KEY; + use lance_file::version::LanceFileVersion; + use lance_index::{IndexType, scalar::ScalarIndexParams}; + + use super::*; + use crate::dataset::NewColumnTransform; + use crate::dataset::write::WriteParams; + use crate::index::DatasetIndexExt; + + async fn wide_dataset() -> Dataset { + let schema = Arc::new(ArrowSchema::new(vec![ + Field::new("a", DataType::Int32, false), + Field::new("b", DataType::Int32, false), + Field::new("c", DataType::Int32, false), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int32Array::from(vec![1, 2, 3, 4])), + Arc::new(Int32Array::from(vec![10, 20, 30, 40])), + Arc::new(Int32Array::from(vec![100, 200, 300, 400])), + ], + ) + .unwrap(); + Dataset::write( + RecordBatchIterator::new([Ok(batch)], schema), + "memory://", + Some(WriteParams { + enable_stable_row_ids: true, + ..Default::default() + }), + ) + .await + .unwrap() + } + + async fn split_column(dataset: Dataset, name: &str) -> Dataset { + let mut scanner = dataset.scan(); + scanner.project(&[name]).unwrap(); + let batch = scanner.try_into_batch().await.unwrap(); + let schema = Schema { + fields: vec![dataset.schema().field(name).unwrap().clone()], + metadata: Default::default(), + }; + let replacement = dataset.get_fragments()[0] + .write_column(stream::iter([Ok(batch)]), &schema) + .await + .unwrap(); + let read_version = dataset.manifest.version; + CommitBuilder::new(Arc::new(dataset)) + .execute(Transaction::new( + read_version, + Operation::DataReplacement { + replacements: vec![replacement], + }, + None, + )) + .await + .unwrap() + } + + async fn scan_all_blobs_binary(dataset: &Dataset) -> RecordBatch { + let mut scanner = dataset.scan(); + scanner.blob_handling(BlobHandling::AllBinary); + scanner.try_into_batch().await.unwrap() + } + + #[tokio::test] + async fn test_optimize_columns_preserves_values_and_row_layout() { + let dataset = split_column(wide_dataset().await, "b").await; + let mut dataset = split_column(dataset, "c").await; + dataset.delete("a = 2").await.unwrap(); + + let version_before = dataset.manifest.version; + let fragment_before = dataset.manifest.fragments[0].clone(); + let batch_before = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(dataset.column_layout_stats()[0].live_file_count, 3); + + let metrics = dataset + .optimize_columns(OptimizeColumnsOptions { + groups: vec![ColumnGroup { + fields: vec!["b".to_string(), "c".to_string()], + }], + fragment_ids: None, + max_concurrency: Some(2), + }) + .await + .unwrap(); + + assert_eq!(metrics.fragments_examined, 1); + assert_eq!(metrics.fragments_rewritten, 1); + assert_eq!(metrics.files_added, 1); + assert_eq!(metrics.files_removed, 2); + assert_eq!(metrics.mixed_files_retained, 0); + assert!(metrics.bytes_read > 0); + assert!(metrics.bytes_written > 0); + assert_eq!(dataset.column_layout_stats()[0].live_file_count, 2); + + let fragment_after = &dataset.manifest.fragments[0]; + assert_eq!(fragment_after.id, fragment_before.id); + assert_eq!(fragment_after.physical_rows, fragment_before.physical_rows); + assert_eq!(fragment_after.deletion_file, fragment_before.deletion_file); + assert_eq!(fragment_after.row_id_meta, fragment_before.row_id_meta); + assert_eq!( + fragment_after.created_at_version_meta, + fragment_before.created_at_version_meta + ); + assert_eq!( + fragment_after.last_updated_at_version_meta, + fragment_before.last_updated_at_version_meta + ); + assert_eq!(dataset.scan().try_into_batch().await.unwrap(), batch_before); + assert_eq!( + dataset + .checkout_version(version_before) + .await + .unwrap() + .scan() + .try_into_batch() + .await + .unwrap(), + batch_before + ); + + let optimized_version = dataset.manifest.version; + let no_op = dataset + .optimize_columns(OptimizeColumnsOptions { + groups: vec![ColumnGroup { + fields: vec!["b".to_string(), "c".to_string()], + }], + fragment_ids: None, + max_concurrency: None, + }) + .await + .unwrap(); + assert_eq!(no_op.fragments_examined, 1); + assert_eq!(no_op.fragments_rewritten, 0); + assert_eq!(dataset.manifest.version, optimized_version); + } + + #[tokio::test] + async fn test_optimize_columns_validates_scope_before_writing() { + let mut dataset = wide_dataset().await; + let version = dataset.manifest.version; + let error = dataset + .optimize_columns(OptimizeColumnsOptions { + groups: vec![ + ColumnGroup { + fields: vec!["b".to_string()], + }, + ColumnGroup { + fields: vec!["b".to_string()], + }, + ], + fragment_ids: None, + max_concurrency: Some(1), + }) + .await + .unwrap_err(); + assert!(error.to_string().contains("more than one group"), "{error}"); + assert_eq!(dataset.manifest.version, version); + } + + #[tokio::test] + async fn test_optimize_columns_keeps_fileless_fields_fileless() { + let mut dataset = wide_dataset().await; + dataset + .add_columns( + NewColumnTransform::AllNulls(Arc::new(ArrowSchema::new(vec![Field::new( + "d", + DataType::Int32, + true, + )]))), + None, + None, + ) + .await + .unwrap(); + let fileless_id = dataset.schema().field("d").unwrap().id; + + let metrics = dataset + .optimize_columns(OptimizeColumnsOptions { + groups: vec![ColumnGroup { + fields: vec!["b".to_string(), "d".to_string()], + }], + fragment_ids: None, + max_concurrency: Some(1), + }) + .await + .unwrap(); + + assert_eq!(metrics.files_added, 1); + assert!( + dataset.manifest.fragments[0] + .files + .iter() + .all(|file| !file.fields.contains(&fileless_id)) + ); + let batch = dataset.scan().try_into_batch().await.unwrap(); + assert_eq!(batch.column_by_name("d").unwrap().null_count(), 4); + } + + #[tokio::test] + async fn test_optimize_columns_prunes_index_coverage_without_replacing_index() { + let dataset = split_column(wide_dataset().await, "b").await; + let mut dataset = split_column(dataset, "c").await; + dataset + .create_index( + &["b"], + IndexType::BTree, + None, + &ScalarIndexParams::default(), + false, + ) + .await + .unwrap(); + let before = dataset.load_indices().await.unwrap(); + let before = before.iter().find(|index| index.name == "b_idx").unwrap(); + let index_uuid = before.uuid; + assert!(before.fragment_bitmap.as_ref().unwrap().contains(0)); + + dataset + .optimize_columns(OptimizeColumnsOptions { + groups: vec![ColumnGroup { + fields: vec!["b".to_string(), "c".to_string()], + }], + fragment_ids: None, + max_concurrency: Some(1), + }) + .await + .unwrap(); + + let after = dataset.load_indices().await.unwrap(); + let after = after.iter().find(|index| index.name == "b_idx").unwrap(); + assert_eq!(after.uuid, index_uuid); + assert!(!after.fragment_bitmap.as_ref().unwrap().contains(0)); + } + + #[tokio::test] + async fn test_optimize_columns_rewrites_legacy_blob() { + let expected = vec![Some(b"first".as_slice()), None, Some(b"third".as_slice())]; + let schema = Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int32, false), + Field::new("blob", DataType::LargeBinary, true) + .with_metadata([(BLOB_META_KEY.to_string(), "true".to_string())].into()), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int32Array::from(vec![0, 1, 2])), + Arc::new(LargeBinaryArray::from_iter(expected)), + ], + ) + .unwrap(); + let mut dataset = Dataset::write( + RecordBatchIterator::new([Ok(batch)], schema), + "memory://", + Some(WriteParams { + data_storage_version: Some(LanceFileVersion::V2_1), + ..Default::default() + }), + ) + .await + .unwrap(); + let before = scan_all_blobs_binary(&dataset).await; + + dataset + .optimize_columns(OptimizeColumnsOptions { + groups: vec![ColumnGroup { + fields: vec!["blob".to_string()], + }], + fragment_ids: None, + max_concurrency: Some(1), + }) + .await + .unwrap(); + + assert_eq!(scan_all_blobs_binary(&dataset).await, before); + assert_eq!(dataset.column_layout_stats()[0].live_file_count, 2); + } + + #[tokio::test] + async fn test_optimize_columns_rewrites_blob_v2_and_nested_blob_v2() { + let mut blob_builder = crate::BlobArrayBuilder::new(3); + blob_builder.push_bytes(b"top-0").unwrap(); + blob_builder.push_null().unwrap(); + blob_builder.push_bytes(b"top-2").unwrap(); + + let mut nested_blob_builder = crate::BlobArrayBuilder::new(3); + nested_blob_builder.push_bytes(b"nested-0").unwrap(); + nested_blob_builder.push_bytes(b"nested-1").unwrap(); + nested_blob_builder.push_null().unwrap(); + let nested_fields = vec![crate::blob_field("nested_blob", true)]; + let nested: ArrayRef = Arc::new( + StructArray::try_new( + nested_fields.clone().into(), + vec![nested_blob_builder.finish().unwrap()], + None, + ) + .unwrap(), + ); + + let schema = Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int32, false), + crate::blob_field("blob", true), + Field::new("info", DataType::Struct(nested_fields.into()), true), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int32Array::from(vec![0, 1, 2])), + blob_builder.finish().unwrap(), + nested, + ], + ) + .unwrap(); + let mut dataset = Dataset::write( + RecordBatchIterator::new([Ok(batch)], schema), + "memory://", + Some(WriteParams { + data_storage_version: Some(LanceFileVersion::V2_2), + ..Default::default() + }), + ) + .await + .unwrap(); + let before = scan_all_blobs_binary(&dataset).await; + + dataset + .optimize_columns(OptimizeColumnsOptions { + groups: vec![ColumnGroup { + fields: vec!["blob".to_string(), "info".to_string()], + }], + fragment_ids: None, + max_concurrency: Some(1), + }) + .await + .unwrap(); + + assert_eq!(scan_all_blobs_binary(&dataset).await, before); + assert_eq!(dataset.column_layout_stats()[0].live_file_count, 2); + } + + #[tokio::test] + async fn test_optimize_columns_preserves_external_blob_v2_reference() { + use lance_core::utils::tempfile::TempDir; + use lance_table::format::BasePath; + + let test_dir = TempDir::default(); + let external_dir = TempDir::default(); + let external_path = external_dir.std_path().join("external.bin"); + std::fs::write(&external_path, b"external-payload").unwrap(); + let external_uri = format!("file://{}", external_path.display()); + let external_base = format!("file://{}", external_dir.std_path().display()); + let test_uri = test_dir.path_str(); + + let mut blob_builder = crate::BlobArrayBuilder::new(2); + blob_builder.push_uri(external_uri).unwrap(); + blob_builder.push_bytes(b"inline-payload").unwrap(); + let schema = Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int32, false), + crate::blob_field("blob", true), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int32Array::from(vec![0, 1])), + blob_builder.finish().unwrap(), + ], + ) + .unwrap(); + let mut dataset = Dataset::write( + RecordBatchIterator::new([Ok(batch)], schema), + &test_uri, + Some(WriteParams { + data_storage_version: Some(LanceFileVersion::V2_2), + initial_bases: Some(vec![BasePath { + id: 1, + name: Some("external".to_string()), + path: external_base, + is_dataset_root: false, + }]), + ..Default::default() + }), + ) + .await + .unwrap(); + let before = scan_all_blobs_binary(&dataset).await; + + dataset + .optimize_columns(OptimizeColumnsOptions { + groups: vec![ColumnGroup { + fields: vec!["blob".to_string()], + }], + fragment_ids: None, + max_concurrency: Some(1), + }) + .await + .unwrap(); + + assert_eq!(scan_all_blobs_binary(&dataset).await, before); + } +} diff --git a/rust/lance/src/dataset/tests/dataset_overlay_index_masking.rs b/rust/lance/src/dataset/tests/dataset_overlay_index_masking.rs index a783ca09905..4dde7f94089 100644 --- a/rust/lance/src/dataset/tests/dataset_overlay_index_masking.rs +++ b/rust/lance/src/dataset/tests/dataset_overlay_index_masking.rs @@ -32,6 +32,7 @@ use lance_file::writer::FileWriterOptions; use crate::Dataset; use crate::dataset::optimize::{CompactionOptions, compact_files, remapping}; +use crate::dataset::optimize_columns::{ColumnGroup, OptimizeColumnsOptions}; use crate::dataset::transaction::{DataOverlayGroup, Operation}; use crate::dataset::{WriteDestination, WriteParams}; use crate::index::vector::VectorIndexParams; @@ -252,6 +253,60 @@ async fn test_overlay_stale_drop_and_new_match(#[values(false, true)] stable_row assert_eq!(ids_matching(&dataset, "age = 20").await, vec![2]); } +#[tokio::test] +async fn test_optimize_columns_prunes_legacy_index_without_fragment_bitmap() { + let mut dataset = create_base_dataset().await; + build_age_index(&mut dataset).await; + + let mut dataset = commit_overlay( + dataset, + "legacy_index_age_overlay", + 0, + &[1], + OverlayCoverage::dense(RoaringBitmap::from_iter([1])), + vec![i32_array([Some(999)])], + ) + .await; + assert_eq!(ids_matching(&dataset, "age = 10").await, Vec::::new()); + assert_eq!(ids_matching(&dataset, "age = 999").await, vec![1]); + + let mut legacy_indices = dataset.load_indices().await.unwrap().as_ref().clone(); + legacy_indices[0].fragment_bitmap = None; + let metadata_key = crate::session::index_caches::IndexMetadataKey { + version: dataset.version().version, + store_identity: &dataset.object_store.store_prefix, + }; + dataset + .index_cache + .insert_with_key(&metadata_key, Arc::new(legacy_indices)) + .await; + assert!( + dataset.load_indices().await.unwrap()[0] + .fragment_bitmap + .is_none() + ); + + dataset + .optimize_columns(OptimizeColumnsOptions { + groups: vec![ColumnGroup { + fields: vec!["age".to_string()], + }], + fragment_ids: Some(vec![0]), + max_concurrency: Some(1), + }) + .await + .unwrap(); + + let indices = dataset.load_indices().await.unwrap(); + let index = indices + .iter() + .find(|index| index.name == "age_idx") + .unwrap(); + assert!(!index.fragment_bitmap.as_ref().unwrap().contains(0)); + assert_eq!(ids_matching(&dataset, "age = 10").await, Vec::::new()); + assert_eq!(ids_matching(&dataset, "age = 999").await, vec![1]); +} + /// Row-level BTree precision: when one row in a covered fragment is stale, only that row is /// blocked from the index result and re-evaluated on the stale-Take path. Non-stale rows in /// the same fragment (including one that matches the predicate) remain on the indexed path. diff --git a/rust/lance/src/dataset/transaction.rs b/rust/lance/src/dataset/transaction.rs index 49c1ee6cbd7..aa0d15c974d 100644 --- a/rust/lance/src/dataset/transaction.rs +++ b/rust/lance/src/dataset/transaction.rs @@ -268,6 +268,19 @@ pub struct Transaction { #[derive(Debug, Clone, DeepSizeOf, PartialEq)] pub struct DataReplacementGroup(pub u64, pub DataFile); +/// Files produced for one fragment by an [`Operation::OptimizeColumns`]. +#[derive(Debug, Clone, DeepSizeOf, PartialEq)] +pub struct OptimizeColumnsGroup { + /// Fragment whose physical column providers are reorganized. + pub fragment_id: u64, + /// Live physical field ids provided by `new_files`. + pub field_ids: Vec, + /// Positionally aligned files replacing the selected field providers. + pub new_files: Vec, + /// Physical row count encoded in every new file. + pub physical_rows: u64, +} + /// Overlay files to append to a single fragment, in order (the last entry is /// newest). The overlays are appended to the fragment's existing `overlays` /// list rather than replacing it, so overlays written by concurrent commits are @@ -440,6 +453,15 @@ pub enum Operation { DataReplacement { replacements: Vec, }, + /// Reorganize selected physical fields without changing logical values, + /// fragment ids, row addresses, deletion files, or row-version metadata. + OptimizeColumns { + /// Immutable dataset version through which base and overlay values were + /// materialized into `groups`. + materialized_through_version: u64, + /// Per-fragment provider changes. + groups: Vec, + }, /// Attach overlay files to fragments, supplying new values for a subset of /// `(physical offset, field)` cells without rewriting the fragments' base /// data files. See [`DataOverlayFile`] and the Data Overlay Files @@ -598,6 +620,7 @@ impl std::fmt::Display for Operation { Self::Project { .. } => write!(f, "Project"), Self::UpdateConfig { .. } => write!(f, "UpdateConfig"), Self::DataReplacement { .. } => write!(f, "DataReplacement"), + Self::OptimizeColumns { .. } => write!(f, "OptimizeColumns"), Self::DataOverlay { .. } => write!(f, "DataOverlay"), Self::Clone { .. } => write!(f, "Clone"), Self::UpdateMemWalState { .. } => write!(f, "UpdateMemWalState"), @@ -801,6 +824,16 @@ impl PartialEq for Operation { Self::DataReplacement { replacements: a }, Self::DataReplacement { replacements: b }, ) => a.len() == b.len() && a.iter().all(|r| b.contains(r)), + ( + Self::OptimizeColumns { + materialized_through_version: a_version, + groups: a_groups, + }, + Self::OptimizeColumns { + materialized_through_version: b_version, + groups: b_groups, + }, + ) => a_version == b_version && a_groups == b_groups, // Handle all remaining combinations. // We spell out all combinations explicitly to prevent // us accidentally handling a new case in the wrong way. @@ -1464,6 +1497,7 @@ impl PartialEq for Operation { } (Self::DataOverlay { groups: a }, Self::DataOverlay { groups: b }) => compare_vec(a, b), (Self::DataOverlay { .. }, _) | (_, Self::DataOverlay { .. }) => false, + (Self::OptimizeColumns { .. }, _) | (_, Self::OptimizeColumns { .. }) => false, } } } @@ -1640,6 +1674,7 @@ impl Operation { Self::Project { .. } => "Project", Self::UpdateConfig { .. } => "UpdateConfig", Self::DataReplacement { .. } => "DataReplacement", + Self::OptimizeColumns { .. } => "OptimizeColumns", Self::DataOverlay { .. } => "DataOverlay", Self::UpdateMemWalState { .. } => "UpdateMemWalState", Self::Clone { .. } => "Clone", @@ -2195,6 +2230,43 @@ impl Transaction { Ok(()) } + /// Reapply OptimizeColumns index pruning after legacy coverage migration. + /// + /// `migrate_indices` runs after [`Self::build_manifest`] and can replace an + /// unknown or corrupt fragment bitmap with the coverage stored in the index + /// files. That recovered coverage predates the overlays materialized by this + /// operation, so target fragments must be removed again before commit. + pub(crate) fn prune_optimize_columns_indices_after_migration( + &self, + indices: &mut [IndexMetadata], + ) -> Result<()> { + let Operation::OptimizeColumns { groups, .. } = &self.operation else { + return Ok(()); + }; + + for group in groups { + let fragment_id = u32::try_from(group.fragment_id).map_err(|_| { + Error::invalid_input(format!( + "OptimizeColumns fragment id {} does not fit in index coverage", + group.fragment_id + )) + })?; + for index in indices.iter_mut() { + let covers_rewritten_field = group.field_ids.iter().any(|field_id| { + index + .fields + .iter() + .any(|index_field| u32::try_from(*index_field).ok() == Some(*field_id)) + }); + if covers_rewritten_field && let Some(fragment_bitmap) = &mut index.fragment_bitmap + { + fragment_bitmap.remove(fragment_id); + } + } + } + Ok(()) + } + /// Create a new manifest from the current manifest and the transaction. /// /// `current_manifest` should only be None if the dataset does not yet exist. @@ -2969,6 +3041,191 @@ impl Transaction { &replaced_fields, ); } + Operation::OptimizeColumns { + materialized_through_version, + groups, + } => { + if *materialized_through_version != self.read_version { + return Err(Error::invalid_input(format!( + "OptimizeColumns materialized_through_version {} must equal transaction read_version {}", + materialized_through_version, self.read_version + ))); + } + let existing_fragments = maybe_existing_fragments?; + let live_schema_ids = schema + .fields_pre_order() + .map(|field| field.id) + .collect::>(); + let existing_fragment_ids = existing_fragments + .iter() + .map(|fragment| fragment.id) + .collect::>(); + let mut groups_by_fragment: HashMap> = + HashMap::new(); + + for group in groups { + if !existing_fragment_ids.contains(&group.fragment_id) { + return Err(Error::invalid_input(format!( + "OptimizeColumns targets fragment {}, which does not exist", + group.fragment_id + ))); + } + if group.field_ids.is_empty() || group.new_files.is_empty() { + return Err(Error::invalid_input(format!( + "OptimizeColumns group for fragment {} must contain fields and files", + group.fragment_id + ))); + } + + let expected_fields = group.field_ids.iter().copied().collect::>(); + if expected_fields.len() != group.field_ids.len() { + return Err(Error::invalid_input(format!( + "OptimizeColumns group for fragment {} contains duplicate field ids", + group.fragment_id + ))); + } + if expected_fields + .iter() + .any(|field| !live_schema_ids.contains(&(*field as i32))) + { + return Err(Error::invalid_input(format!( + "OptimizeColumns group for fragment {} contains a field absent from the current schema", + group.fragment_id + ))); + } + + let mut provided_fields = HashSet::new(); + for file in &group.new_files { + if file.file_version()? == ConcreteFileVersion::V1 { + return Err(Error::not_supported(format!( + "OptimizeColumns cannot install legacy data file '{}'", + file.path + ))); + } + for field in file.fields.iter().copied().filter(|field| *field >= 0) { + let field = field as u32; + if !provided_fields.insert(field) { + return Err(Error::invalid_input(format!( + "OptimizeColumns files for fragment {} provide field {} more than once", + group.fragment_id, field + ))); + } + } + } + if provided_fields != expected_fields { + return Err(Error::invalid_input(format!( + "OptimizeColumns files for fragment {} must collectively provide exactly field_ids", + group.fragment_id + ))); + } + + let existing_groups = groups_by_fragment.entry(group.fragment_id).or_default(); + if existing_groups.iter().any(|existing| { + existing + .field_ids + .iter() + .any(|field| expected_fields.contains(field)) + }) { + return Err(Error::invalid_input(format!( + "OptimizeColumns groups for fragment {} overlap", + group.fragment_id + ))); + } + existing_groups.push(group); + } + + let mut pruned_index_inputs = Vec::with_capacity(groups_by_fragment.len()); + for fragment in existing_fragments { + let Some(fragment_groups) = groups_by_fragment.get(&fragment.id) else { + final_fragments.push(fragment.clone()); + continue; + }; + + let physical_rows = fragment.physical_rows.ok_or_else(|| { + Error::invalid_input(format!( + "OptimizeColumns target fragment {} has no physical row count", + fragment.id + )) + })? as u64; + if fragment_groups + .iter() + .any(|group| group.physical_rows != physical_rows) + { + return Err(Error::invalid_input(format!( + "OptimizeColumns output row count does not match fragment {} physical row count {}", + fragment.id, physical_rows + ))); + } + + let selected_fields = fragment_groups + .iter() + .flat_map(|group| group.field_ids.iter().copied()) + .collect::>(); + for file in &fragment.files { + if file.file_version()? == ConcreteFileVersion::V1 + && file.fields.iter().any(|field| { + *field >= 0 && selected_fields.contains(&(*field as u32)) + }) + { + return Err(Error::not_supported(format!( + "OptimizeColumns does not support legacy fragment {}", + fragment.id + ))); + } + } + + let mut optimized = fragment.clone(); + for file in &mut optimized.files { + file.fields = file + .fields + .iter() + .map(|field| { + if *field >= 0 && selected_fields.contains(&(*field as u32)) { + TOMBSTONE_FIELD_ID + } else { + *field + } + }) + .collect::>() + .into(); + } + optimized.files.retain(|file| { + file.fields + .iter() + .any(|field| live_schema_ids.contains(field)) + }); + optimized.files.extend( + fragment_groups + .iter() + .flat_map(|group| group.new_files.iter().cloned()), + ); + + let selected_fields_vec = selected_fields.iter().copied().collect::>(); + let (mut materialized, newer): (Vec<_>, Vec<_>) = + optimized.overlays.drain(..).partition(|overlay| { + overlay.committed_version <= *materialized_through_version + }); + lance_table::format::overlay::tombstone_overlay_fields( + &mut materialized, + &selected_fields_vec, + ); + materialized.extend(newer); + optimized.overlays = materialized; + + pruned_index_inputs.push((optimized.clone(), selected_fields_vec)); + final_fragments.push(optimized); + } + + // Apply against the latest manifest so indices created after the + // materialization snapshot are conservatively pruned as well. + for (fragment, fields) in pruned_index_inputs { + Self::prune_updated_fields_from_indices( + &mut final_indices, + std::slice::from_ref(&fragment), + &fields, + ); + } + } Operation::DataOverlay { groups } => { // Stamp each overlay with the version this commit is producing. // build_manifest re-runs on every retry with an updated @@ -3926,6 +4183,34 @@ impl TryFrom for DataReplacementGroup { } } +impl From<&OptimizeColumnsGroup> for pb::transaction::OptimizeColumnsGroup { + fn from(group: &OptimizeColumnsGroup) -> Self { + Self { + fragment_id: group.fragment_id, + field_ids: group.field_ids.clone(), + new_files: group.new_files.iter().map(pb::DataFile::from).collect(), + physical_rows: group.physical_rows, + } + } +} + +impl TryFrom for OptimizeColumnsGroup { + type Error = Error; + + fn try_from(group: pb::transaction::OptimizeColumnsGroup) -> Result { + Ok(Self { + fragment_id: group.fragment_id, + field_ids: group.field_ids, + new_files: group + .new_files + .into_iter() + .map(DataFile::try_from) + .collect::>>()?, + physical_rows: group.physical_rows, + }) + } +} + impl From<&DataOverlayGroup> for pb::transaction::DataOverlayGroup { fn from(group: &DataOverlayGroup) -> Self { Self { @@ -4254,6 +4539,18 @@ impl TryFrom for Transaction { .map(DataReplacementGroup::try_from) .collect::>>()?, }, + Some(pb::transaction::Operation::OptimizeColumns( + pb::transaction::OptimizeColumns { + materialized_through_version, + groups, + }, + )) => Operation::OptimizeColumns { + materialized_through_version, + groups: groups + .into_iter() + .map(OptimizeColumnsGroup::try_from) + .collect::>>()?, + }, Some(pb::transaction::Operation::UpdateMemWalState( pb::transaction::UpdateMemWalState { compacted_sstables, @@ -4570,6 +4867,16 @@ impl From<&Transaction> for pb::Transaction { .collect(), }) } + Operation::OptimizeColumns { + materialized_through_version, + groups, + } => pb::transaction::Operation::OptimizeColumns(pb::transaction::OptimizeColumns { + materialized_through_version: *materialized_through_version, + groups: groups + .iter() + .map(pb::transaction::OptimizeColumnsGroup::from) + .collect(), + }), Operation::DataOverlay { groups } => { pb::transaction::Operation::DataOverlay(pb::transaction::DataOverlay { groups: groups @@ -4738,6 +5045,25 @@ pub fn validate_operation(manifest: Option<&Manifest>, operation: &Operation) -> } Ok(()) } + Operation::OptimizeColumns { + materialized_through_version, + groups, + } => { + if *materialized_through_version == 0 + || *materialized_through_version > manifest.version + { + return Err(Error::invalid_input(format!( + "OptimizeColumns materialized_through_version {} must be between 1 and current dataset version {}", + materialized_through_version, manifest.version + ))); + } + if groups.is_empty() { + return Err(Error::invalid_input( + "OptimizeColumns must contain at least one group", + )); + } + Ok(()) + } _ => Ok(()), } } @@ -8792,6 +9118,146 @@ mod tests { assert_eq!(fragment.overlays[0].data_file.fields.as_ref(), &[5]); } + #[test] + fn test_optimize_columns_uses_fixed_overlay_cutoff_and_preserves_row_metadata() { + let mut fragment = Fragment::new(0); + fragment.files = vec![DataFile::new( + "wide.lance", + vec![4, 5], + vec![0, 1], + ConcreteFileVersion::V2_0, + None, + None, + )]; + fragment.physical_rows = Some(2); + fragment.row_id_meta = Some(RowIdMeta::Inline( + write_row_ids(&RowIdSequence::from([10_u64, 11].as_slice())).into(), + )); + fragment.created_at_version_meta = Some( + RowDatasetVersionMeta::from_sequence(&RowDatasetVersionSequence { + runs: vec![RowDatasetVersionRun { + span: U64Segment::Range(0..2), + version: 1, + }], + }) + .unwrap(), + ); + fragment.last_updated_at_version_meta = fragment.created_at_version_meta.clone(); + fragment.overlays = vec![ + overlay_with_field(5, 6), + DataOverlayFile { + data_file: DataFile::new( + "newer.lance", + vec![5], + vec![0], + ConcreteFileVersion::V2_0, + None, + None, + ), + coverage: OverlayCoverage::dense(roaring::RoaringBitmap::from_iter([1u32])), + committed_version: 7, + }, + ]; + let metadata_before = ( + fragment.row_id_meta.clone(), + fragment.created_at_version_meta.clone(), + fragment.last_updated_at_version_meta.clone(), + ); + + let schema = ArrowSchema::new(vec![ + ArrowField::new("x", DataType::Int32, true), + ArrowField::new("a", DataType::Int32, true), + ArrowField::new("v", DataType::Int32, true), + ArrowField::new("y", DataType::Int32, true), + ]); + let mut lance_schema = LanceSchema::try_from(&schema).unwrap(); + lance_schema.fields[0].id = 3; + lance_schema.fields[1].id = 4; + lance_schema.fields[2].id = 5; + lance_schema.fields[3].id = 6; + let mut manifest = Manifest::new( + lance_schema, + Arc::new(vec![fragment]), + lance_table::format::DataStorageFormat::new(ConcreteFileVersion::V2_0), + HashMap::new(), + ); + manifest.version = 7; + let txn = Transaction::new( + 6, + Operation::OptimizeColumns { + materialized_through_version: 6, + groups: vec![OptimizeColumnsGroup { + fragment_id: 0, + field_ids: vec![5], + new_files: vec![DataFile::new( + "optimized.lance", + vec![5], + vec![0], + ConcreteFileVersion::V2_0, + None, + None, + )], + physical_rows: 2, + }], + }, + None, + ); + let (result, _) = txn + .build_manifest( + Some(&manifest), + vec![], + "txn", + &ManifestWriteConfig::default(), + ) + .unwrap(); + let fragment = &result.fragments[0]; + + assert_eq!(fragment.overlays.len(), 1); + assert_eq!(fragment.overlays[0].committed_version, 7); + assert_eq!(fragment.overlays[0].data_file.fields.as_ref(), &[5]); + assert!( + fragment + .files + .iter() + .any(|file| file.path == "optimized.lance") + ); + assert_eq!( + ( + fragment.row_id_meta.clone(), + fragment.created_at_version_meta.clone(), + fragment.last_updated_at_version_meta.clone(), + ), + metadata_before + ); + } + + #[test] + fn test_optimize_columns_operation_roundtrips() { + let transaction = Transaction::new( + 4, + Operation::OptimizeColumns { + materialized_through_version: 4, + groups: vec![OptimizeColumnsGroup { + fragment_id: 3, + field_ids: vec![7, 8], + new_files: vec![DataFile::new( + "optimized.lance", + vec![7, 8], + vec![0, 1], + ConcreteFileVersion::V2_0, + None, + None, + )], + physical_rows: 10, + }], + }, + None, + ); + let encoded = pb::Transaction::from(&transaction); + let decoded = Transaction::try_from(encoded).unwrap(); + assert_eq!(decoded, transaction); + } + #[test] fn test_data_overlay_build_manifest_merges_duplicate_groups() { // Two groups targeting the same fragment must both survive (a HashMap diff --git a/rust/lance/src/io/commit.rs b/rust/lance/src/io/commit.rs index 129d0aa21a0..ced36fb7136 100644 --- a/rust/lance/src/io/commit.rs +++ b/rust/lance/src/io/commit.rs @@ -1086,6 +1086,7 @@ pub(crate) async fn do_commit_detached_transaction( // Runs after the coverage derivation and can replace a fragment bitmap // while keeping its UUID, so anything it narrowed loses its position. let recovered_coverage = migrate_indices(dataset, &mut indices).await?; + transaction.prune_optimize_columns_indices_after_migration(&mut indices)?; Transaction::withdraw_coverage_invalidated_after_build( &mut indices, &recovered_coverage, @@ -1448,6 +1449,7 @@ pub(crate) async fn commit_transaction( // Runs after the coverage derivation and can replace a fragment bitmap // while keeping its UUID, so anything it narrowed loses its position. let recovered_coverage = migrate_indices(&dataset, &mut indices).await?; + transaction.prune_optimize_columns_indices_after_migration(&mut indices)?; Transaction::withdraw_coverage_invalidated_after_build( &mut indices, &recovered_coverage, diff --git a/rust/lance/src/io/commit/conflict_resolver.rs b/rust/lance/src/io/commit/conflict_resolver.rs index 99b5188f91b..7efc9a96ba5 100644 --- a/rust/lance/src/io/commit/conflict_resolver.rs +++ b/rust/lance/src/io/commit/conflict_resolver.rs @@ -172,6 +172,23 @@ impl<'a> TransactionRebase<'a> { conflicting_mem_wal_compacted_sstables: Vec::new(), }) } + Operation::OptimizeColumns { groups, .. } => { + let modified_fragment_ids = groups + .iter() + .map(|group| group.fragment_id) + .collect::>(); + let initial_fragments = + initial_fragments_for_rebase(dataset, &transaction, &modified_fragment_ids) + .await; + Ok(Self { + transaction, + affected_rows, + initial_fragments, + modified_fragment_ids, + conflicting_frag_reuse_indices: Vec::new(), + conflicting_mem_wal_compacted_sstables: Vec::new(), + }) + } Operation::DataOverlay { groups } => { let modified_fragment_ids = groups.iter().map(|g| g.fragment_id).collect::>(); @@ -295,6 +312,9 @@ impl<'a> TransactionRebase<'a> { Operation::DataReplacement { .. } => { self.check_data_replacement_txn(other_transaction, other_version) } + Operation::OptimizeColumns { .. } => { + self.check_optimize_columns_txn(other_transaction, other_version) + } Operation::DataOverlay { .. } => { self.check_data_overlay_txn(other_transaction, other_version) } @@ -357,6 +377,17 @@ impl<'a> TransactionRebase<'a> { Ok(()) } } + Operation::OptimizeColumns { groups, .. } => { + if groups + .iter() + .map(|group| group.fragment_id) + .any(|id| self.modified_fragment_ids.contains(&id)) + { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } Operation::Update { updated_fragments, removed_fragment_ids, @@ -558,6 +589,17 @@ impl<'a> TransactionRebase<'a> { Ok(()) } } + Operation::OptimizeColumns { groups, .. } => { + if groups + .iter() + .map(|group| group.fragment_id) + .any(|id| self.modified_fragment_ids.contains(&id)) + { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } Operation::Update { updated_fragments, removed_fragment_ids, @@ -804,6 +846,24 @@ impl<'a> TransactionRebase<'a> { } Ok(()) } + Operation::OptimizeColumns { groups, .. } => { + for group in groups { + for index in new_indices.iter_mut().filter(|index| { + index + .fields + .iter() + .any(|field| group.field_ids.contains(&(*field as u32))) + }) { + let Some(bitmap) = index.fragment_bitmap.as_mut() else { + return Err( + self.retryable_conflict_err(other_transaction, other_version) + ); + }; + bitmap.remove(group.fragment_id as u32); + } + } + Ok(()) + } Operation::UpdateMemWalState { compacted_sstables: other_compacted_sstables, .. @@ -919,6 +979,23 @@ impl<'a> TransactionRebase<'a> { } Ok(()) } + Operation::OptimizeColumns { + groups: optimized_groups, + .. + } => { + if optimized_groups.iter().any(|optimized| { + groups.iter().any(|rewrite| { + rewrite + .old_fragments + .iter() + .any(|fragment| fragment.id == optimized.fragment_id) + }) + }) { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } Operation::Merge { .. } => { Err(self.retryable_conflict_err(other_transaction, other_version)) } @@ -1054,7 +1131,7 @@ impl<'a> TransactionRebase<'a> { Ok(()) } } - Operation::UpdateMemWalState { .. } => { + Operation::UpdateMemWalState { .. } | Operation::OptimizeColumns { .. } => { Err(self.incompatible_conflict_err(other_transaction, other_version)) } Operation::Append { .. } @@ -1098,6 +1175,7 @@ impl<'a> TransactionRebase<'a> { | Operation::UpdateConfig { .. } | Operation::Clone { .. } | Operation::DataReplacement { .. } + | Operation::OptimizeColumns { .. } | Operation::DataOverlay { .. } => Ok(()), } } @@ -1255,6 +1333,22 @@ impl<'a> TransactionRebase<'a> { } Ok(()) } + Operation::OptimizeColumns { groups, .. } => { + let overlaps = groups.iter().any(|group| { + replacements.iter().any(|replacement| { + replacement.0 == group.fragment_id + && replacement.1.fields.iter().any(|field| { + *field >= 0 + && group.field_ids.contains(&(*field as u32)) + }) + }) + }); + if overlaps { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } Operation::Overwrite { .. } | Operation::Restore { .. } | Operation::UpdateMemWalState { .. } => { @@ -1266,6 +1360,136 @@ impl<'a> TransactionRebase<'a> { } } + fn check_optimize_columns_txn( + &mut self, + other_transaction: &Transaction, + other_version: u64, + ) -> Result<()> { + let Operation::OptimizeColumns { groups, .. } = &self.transaction.operation else { + return Err(wrong_operation_err(&self.transaction.operation)); + }; + + let mut fields_by_fragment: HashMap> = HashMap::new(); + for group in groups { + fields_by_fragment + .entry(group.fragment_id) + .or_default() + .extend(group.field_ids.iter().copied()); + } + let touches_target = |fragment_id: u64| fields_by_fragment.contains_key(&fragment_id); + + match &other_transaction.operation { + Operation::Append { .. } + | Operation::CreateIndex { .. } + | Operation::ReserveFragments { .. } + | Operation::UpdateConfig { .. } + | Operation::Clone { .. } + | Operation::UpdateBases { .. } + | Operation::DataOverlay { .. } => Ok(()), + Operation::Project { schema, .. } => { + if fields_by_fragment.values().any(|fields| { + fields + .iter() + .any(|field| schema.field_by_id(*field as i32).is_none()) + }) { + Err(self.incompatible_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } + Operation::Delete { + updated_fragments, + deleted_fragment_ids, + .. + } => { + if updated_fragments + .iter() + .map(|fragment| fragment.id) + .chain(deleted_fragment_ids.iter().copied()) + .any(touches_target) + { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } + Operation::Update { + updated_fragments, + removed_fragment_ids, + .. + } => { + if updated_fragments + .iter() + .map(|fragment| fragment.id) + .chain(removed_fragment_ids.iter().copied()) + .any(touches_target) + { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } + Operation::Rewrite { groups, .. } => { + if groups + .iter() + .flat_map(|group| group.old_fragments.iter()) + .map(|fragment| fragment.id) + .any(touches_target) + { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } + Operation::DataReplacement { replacements } => { + let overlaps = replacements.iter().any(|replacement| { + fields_by_fragment + .get(&replacement.0) + .is_some_and(|fields| { + replacement + .1 + .fields + .iter() + .any(|field| *field >= 0 && fields.contains(&(*field as u32))) + }) + }); + if overlaps { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } + Operation::OptimizeColumns { + groups: other_groups, + .. + } => { + let overlaps = other_groups.iter().any(|other_group| { + fields_by_fragment + .get(&other_group.fragment_id) + .is_some_and(|fields| { + other_group + .field_ids + .iter() + .any(|field| fields.contains(field)) + }) + }); + if overlaps { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } + Operation::Merge { .. } => { + Err(self.retryable_conflict_err(other_transaction, other_version)) + } + Operation::Overwrite { .. } + | Operation::Restore { .. } + | Operation::UpdateMemWalState { .. } => { + Err(self.incompatible_conflict_err(other_transaction, other_version)) + } + } + } + /// Conflict checks for our DataOverlay transaction against a concurrent one. /// /// Overlays are intentionally permissive (see the Data Overlay Files spec): @@ -1295,6 +1519,7 @@ impl<'a> TransactionRebase<'a> { | Operation::UpdateBases { .. } | Operation::Clone { .. } | Operation::DataReplacement { .. } + | Operation::OptimizeColumns { .. } | Operation::DataOverlay { .. } => Ok(()), // A concurrent Delete only tombstones rows via a deletion vector, // which preserves physical offsets; the overlay value for a deleted @@ -1404,6 +1629,7 @@ impl<'a> TransactionRebase<'a> { | Operation::Rewrite { .. } | Operation::Merge { .. } | Operation::DataReplacement { .. } + | Operation::OptimizeColumns { .. } | Operation::DataOverlay { .. } => { Err(self.retryable_conflict_err(other_transaction, other_version)) } @@ -1437,6 +1663,9 @@ impl<'a> TransactionRebase<'a> { | Operation::Project { .. } | Operation::Clone { .. } | Operation::UpdateConfig { .. } => Ok(()), + Operation::OptimizeColumns { .. } => { + Err(self.incompatible_conflict_err(other_transaction, other_version)) + } Operation::UpdateMemWalState { .. } => { Err(self.incompatible_conflict_err(other_transaction, other_version)) } @@ -1457,6 +1686,7 @@ impl<'a> TransactionRebase<'a> { | Operation::CreateIndex { .. } | Operation::Rewrite { .. } | Operation::DataReplacement { .. } + | Operation::OptimizeColumns { .. } | Operation::DataOverlay { .. } | Operation::Merge { .. } | Operation::ReserveFragments { .. } @@ -1487,6 +1717,21 @@ impl<'a> TransactionRebase<'a> { | Operation::Clone { .. } | Operation::ReserveFragments { .. } | Operation::UpdateBases { .. } => Ok(()), + Operation::OptimizeColumns { groups, .. } => { + let Operation::Project { schema, .. } = &self.transaction.operation else { + return Err(wrong_operation_err(&self.transaction.operation)); + }; + if groups.iter().any(|group| { + group + .field_ids + .iter() + .any(|field| schema.field_by_id(*field as i32).is_none()) + }) { + Err(self.incompatible_conflict_err(other_transaction, other_version)) + } else { + Ok(()) + } + } Operation::Merge { .. } | Operation::Project { .. } => { // Need to recompute the schema Err(self.retryable_conflict_err(other_transaction, other_version)) @@ -1547,6 +1792,7 @@ impl<'a> TransactionRebase<'a> { | Operation::CreateIndex { .. } | Operation::Rewrite { .. } | Operation::DataReplacement { .. } + | Operation::OptimizeColumns { .. } | Operation::DataOverlay { .. } | Operation::Merge { .. } | Operation::Restore { .. } @@ -1624,6 +1870,7 @@ impl<'a> TransactionRebase<'a> { | Operation::Overwrite { .. } | Operation::Delete { .. } | Operation::DataReplacement { .. } + | Operation::OptimizeColumns { .. } | Operation::DataOverlay { .. } | Operation::Merge { .. } | Operation::Restore { .. } @@ -1720,6 +1967,7 @@ impl<'a> TransactionRebase<'a> { Operation::Append { .. } | Operation::Overwrite { .. } | Operation::DataReplacement { .. } + | Operation::OptimizeColumns { .. } | Operation::Merge { .. } | Operation::Restore { .. } | Operation::ReserveFragments { .. } @@ -2258,7 +2506,7 @@ mod tests { use lance_table::io::deletion::{deletion_file_path, read_deletion_file}; use super::*; - use crate::dataset::transaction::{DataReplacementGroup, RewriteGroup}; + use crate::dataset::transaction::{DataReplacementGroup, OptimizeColumnsGroup, RewriteGroup}; use crate::dataset::write::WriteMode; use crate::session::caches::DeletionFileKey; use crate::{ @@ -2293,6 +2541,73 @@ mod tests { .unwrap() } + fn optimize_operation(fragment_id: u64, field_ids: Vec) -> Operation { + Operation::OptimizeColumns { + materialized_through_version: 1, + groups: vec![OptimizeColumnsGroup { + fragment_id, + field_ids: field_ids.clone(), + new_files: vec![DataFile::new_legacy_from_fields( + format!("optimized-{fragment_id}.lance"), + field_ids.into_iter().map(|field| field as i32).collect(), + None, + )], + physical_rows: 5, + }], + } + } + + #[tokio::test] + async fn test_optimize_columns_conflict_rules() { + let dataset = test_dataset(10, 2).await; + let assert_result = |result: Result<()>, conflicts: bool| { + if conflicts { + assert!( + matches!(result, Err(Error::RetryableCommitConflict { .. })), + "expected retryable conflict, got {result:?}" + ); + } else { + assert!(result.is_ok(), "expected compatibility, got {result:?}"); + } + }; + + for (other, conflicts) in [ + (optimize_operation(0, vec![1]), true), + (optimize_operation(0, vec![0]), false), + (optimize_operation(1, vec![1]), false), + ( + Operation::Delete { + updated_fragments: vec![dataset.manifest.fragments[0].clone()], + deleted_fragment_ids: vec![], + predicate: "a = 1".to_string(), + }, + true, + ), + (Operation::DataOverlay { groups: vec![] }, false), + ] { + let ours = Transaction::new_from_version(1, optimize_operation(0, vec![1])); + let other = Transaction::new_from_version(1, other); + let mut rebase = TransactionRebase::try_new(&dataset, ours, None) + .await + .unwrap(); + assert_result(rebase.check_txn(&other, 2), conflicts); + } + + let delete = Transaction::new_from_version( + 1, + Operation::Delete { + updated_fragments: vec![dataset.manifest.fragments[0].clone()], + deleted_fragment_ids: vec![], + predicate: "a = 1".to_string(), + }, + ); + let optimize = Transaction::new_from_version(1, optimize_operation(0, vec![1])); + let mut rebase = TransactionRebase::try_new(&dataset, delete, None) + .await + .unwrap(); + assert_result(rebase.check_txn(&optimize, 2), true); + } + /// Helper function for tests to create UpdateConfig operations using old-style parameters #[cfg(test)] fn create_update_config_for_test( @@ -4498,6 +4813,9 @@ mod tests { Operation::DataReplacement { replacements } => { Box::new(replacements.iter().map(|r| r.0)) } + Operation::OptimizeColumns { groups, .. } => { + Box::new(groups.iter().map(|group| group.fragment_id)) + } Operation::DataOverlay { groups } => Box::new(groups.iter().map(|g| g.fragment_id)), } }