From a7443070b61f7970e2b9c9a97c12f72c62e818b6 Mon Sep 17 00:00:00 2001 From: dshepelev15 Date: Wed, 16 Sep 2026 13:00:32 +0000 Subject: [PATCH 1/2] feat: keep column groups through compaction and add rewrite_columns Compaction re-encodes every column of the fragments it rewrites, so migrating a table with a few multi-KB embedding columns to a newer file version rewrites the whole table for the sake of its narrow columns, and compaction folds any per-column file layout back into one wide file, which makes later single-column updates rewrite every column again. `CompactionOptions.column_groups` (config key `lance.compaction.column_groups`) writes each listed group of top-level columns to its own data file per compacted fragment and the remaining columns to one shared file. One scan feeds one writer per group with the same planned row counts, so the files of a fragment stay row-aligned; binary copy is disabled and `max_bytes_per_file` ignored while groups are set. `Dataset::rewrite_columns` / `FileFragment::rewrite_columns` rewrite only the named columns of every fragment into a new data file in the requested V2 version and tombstone them in the files they came from. The other columns' files are neither read nor written; rows, fragment ids, row addresses and indices are unchanged. Fragments already in the requested layout are skipped, and the fragment-level call returns metadata for a single distributed `Update` commit in `RewriteColumns` mode. --- docs/src/guide/read_and_write.md | 20 + python/python/lance/dataset.py | 62 ++- python/python/lance/fragment.py | 24 ++ python/python/lance/lance/__init__.pyi | 6 + python/python/lance/optimize.py | 10 + python/python/tests/test_optimize.py | 40 ++ python/python/tests/test_schema_evolution.py | 75 ++++ python/src/dataset.rs | 22 + python/src/dataset/optimize.rs | 9 +- python/src/fragment.rs | 20 + rust/lance/src/dataset.rs | 1 + rust/lance/src/dataset/optimize.rs | 330 ++++++++++++++- rust/lance/src/dataset/rewrite_columns.rs | 412 +++++++++++++++++++ rust/lance/src/dataset/versions/mod.rs | 24 +- 14 files changed, 1028 insertions(+), 27 deletions(-) create mode 100644 rust/lance/src/dataset/rewrite_columns.rs diff --git a/docs/src/guide/read_and_write.md b/docs/src/guide/read_and_write.md index 40062b6ce33..2c6ca70207d 100644 --- a/docs/src/guide/read_and_write.md +++ b/docs/src/guide/read_and_write.md @@ -85,6 +85,26 @@ the target version, actual version, and file path. A persistent compaction targe can be set through `lance.compaction.data_storage_version` in the table config; an explicit operation target takes precedence. +Compaction re-encodes every column of the fragments it rewrites. To migrate +only some columns, `rewrite_columns` (Python: `LanceDataset.rewrite_columns`, +Rust: `Dataset::rewrite_columns`) reads just the named top-level columns of each +fragment, writes them to one new data file per fragment in the requested V2 +version, and tombstones them in the files they came from. The files holding the +other columns are neither read nor written, and rows, fragment ids, row +addresses and indices are unchanged. Fragments already in the requested layout +are skipped, so an interrupted rewrite can be rerun; `LanceFragment.rewrite_columns` +does the same for one fragment so the work can be spread over workers and +committed as a single `Update` in `rewrite_columns` mode. + +A compaction folds those files back into one per fragment unless it is told to +keep them apart: the `column_groups` compaction option (config key +`lance.compaction.column_groups`, groups separated by `;` and columns by `,`) +writes each listed group of columns to its own data file per fragment and the +remaining columns to one shared file. Setting the config key once makes every +later compaction preserve the layout that `rewrite_columns` produced. Binary +copy is disabled when groups are set, and `max_bytes_per_file` is ignored so +all groups split at the same rows. + ### Upgrading clients before mixed-version writes Before writing files that differ from the dataset default, upgrade every reader diff --git a/python/python/lance/dataset.py b/python/python/lance/dataset.py index 662e9307f59..ca30dc0c55f 100644 --- a/python/python/lance/dataset.py +++ b/python/python/lance/dataset.py @@ -2868,6 +2868,52 @@ def drop_columns(self, columns: List[str]): # Indices might have changed self._list_indices_res = None + def rewrite_columns( + self, + columns: List[str], + *, + data_storage_version: Optional[str] = None, + ): + """Rewrite columns into their own data files, leaving other files alone. + + Every fragment gets one new data file holding exactly ``columns``, and + those columns are tombstoned in the files they came from. The files + holding the other columns are not read or written, so this migrates + narrow columns to a newer data file version without re-encoding the + wide columns that dominate a table's size. Rows, fragment ids, row + addresses and indices are unchanged. + + Fragments whose ``columns`` already sit alone in a file of the requested + version are skipped, so an interrupted rewrite can be rerun. Fragments + are rewritten one after another; to spread the work over many workers, + call :meth:`lance.fragment.LanceFragment.rewrite_columns` per fragment + and commit the returned metadata in one + :class:`LanceOperation.Update` with ``update_mode="rewrite_columns"``. + + A later ``compact_files`` folds the new files back into one file per + fragment unless its ``column_groups`` lists the same columns. + + Parameters + ---------- + columns : list of str + Top-level column names to rewrite. + data_storage_version : str, optional + Data file version for the new files, such as ``"2.2"`` or + ``"stable"``. Defaults to the dataset's default write version and + never changes it. Must be a V2 version. + + Examples + -------- + >>> import lance + >>> import pyarrow as pa + >>> table = pa.table({"a": [1, 2, 3], "b": ["x", "y", "z"]}) + >>> dataset = lance.write_dataset(table, "example", data_storage_version="2.0") + >>> dataset.rewrite_columns(["b"], data_storage_version="2.2") + >>> [f.fields for f in dataset.get_fragments()[0].data_files()] + [[0, -2], [1]] + """ + self._ds.rewrite_columns(columns, data_storage_version) + def delete( self, predicate: Union[str, Expression], @@ -7511,6 +7557,7 @@ def compact_files( max_source_bytes: Optional[int] = None, excluded_fragment_ids: Optional[list[int]] = None, data_storage_version: Optional[str] = None, + column_groups: Optional[list[list[str]]] = None, ) -> CompactionMetrics: """Compacts small files in the dataset, reducing total number of files. @@ -7545,7 +7592,9 @@ def compact_files( ``lance.compaction.max_source_fragments``, ``lance.compaction.max_source_rows``, ``lance.compaction.max_source_bytes``, - ``lance.compaction.data_storage_version``. + ``lance.compaction.data_storage_version``, + ``lance.compaction.column_groups`` (groups separated by ``;``, + columns by ``,``, e.g. ``"embedding;caption,tags"``). Parameters ---------- @@ -7626,6 +7675,16 @@ def compact_files( Uses the compaction config target when set, otherwise the dataset's default write version. Does not change that default or the versions of unselected files. V1/V2 cross-family targets are rejected. + column_groups: list[list[str]], optional + Top-level columns to keep in their own data files. Each inner list + becomes one data file per compacted fragment holding exactly those + columns; every column not listed goes to one shared file. This is + the layout :meth:`LanceDataset.rewrite_columns` produces, so a + compaction configured with the same groups preserves it instead of + folding wide columns back next to narrow ones. Binary copy is + disabled when set, and ``max_bytes_per_file`` is ignored so every + group splits at the same rows. Uses the manifest config value when + not specified. Returns ------- @@ -7654,6 +7713,7 @@ def compact_files( max_source_bytes=max_source_bytes, excluded_fragment_ids=excluded_fragment_ids, data_storage_version=data_storage_version, + column_groups=column_groups, ).items() if v is not None } diff --git a/python/python/lance/fragment.py b/python/python/lance/fragment.py index 2be76ae2d4c..5033edbdaf6 100644 --- a/python/python/lance/fragment.py +++ b/python/python/lance/fragment.py @@ -1015,6 +1015,30 @@ def update_columns( return metadata, fields_modified, matched_offsets return metadata, fields_modified + def rewrite_columns( + self, + columns: List[str], + *, + data_storage_version: Optional[str] = None, + ) -> Optional[FragmentMetadata]: + """Rewrite columns of this fragment into one new data file. + + .. warning:: + + Internal API. This method is not intended to be used by end users. + + The per-fragment half of + :meth:`lance.dataset.LanceDataset.rewrite_columns`, for spreading a + rewrite over many workers. The new file is written but nothing is + committed: collect the returned metadata from every fragment and commit + it in one :class:`lance.dataset.LanceOperation.Update` with + ``update_mode="rewrite_columns"`` and no ``fields_modified``. + + Returns ``None`` when ``columns`` already sit alone in a file of the + requested version. + """ + return self._fragment.rewrite_columns(columns, data_storage_version) + def merge_columns( self, value_func: ( diff --git a/python/python/lance/lance/__init__.pyi b/python/python/lance/lance/__init__.pyi index d2ce9591a5a..9086e0d293a 100644 --- a/python/python/lance/lance/__init__.pyi +++ b/python/python/lance/lance/__init__.pyi @@ -636,6 +636,9 @@ class _Dataset: def validate(self): ... def migrate_manifest_paths_v2(self): ... def drop_columns(self, columns: List[str]): ... + def rewrite_columns( + self, columns: List[str], data_storage_version: Optional[str] = None + ): ... def add_columns_from_reader( self, reader: pa.RecordBatchReader, batch_size: Optional[int] = None ): ... @@ -759,6 +762,9 @@ class _Fragment: read_columns: Optional[List[str]], batch_size: Optional[int], ) -> Tuple[FragmentMetadata, LanceSchema]: ... + def rewrite_columns( + self, columns: List[str], data_storage_version: Optional[str] = None + ) -> Optional[FragmentMetadata]: ... def delete(self, predicate: str) -> Optional[_Fragment]: ... def delete_rows(self, offsets: List[int]) -> Optional[_Fragment]: ... def schema(self) -> pa.Schema: ... diff --git a/python/python/lance/optimize.py b/python/python/lance/optimize.py index ab18255c853..407a13a18b6 100644 --- a/python/python/lance/optimize.py +++ b/python/python/lance/optimize.py @@ -131,3 +131,13 @@ class CompactionOptions(TypedDict, total=False): input versions and no overlays; TryBinaryCopy reencodes ineligible inputs, while ForceBinaryCopy reports an error. """ + column_groups: Optional[list[list[str]]] + """ + Top-level columns to keep in their own data files. Each inner list becomes + one data file per compacted fragment holding exactly those columns; every + column not listed goes to one shared file. This is the layout + ``LanceDataset.rewrite_columns`` produces, so a compaction configured with + the same groups preserves it. Binary copy is disabled when set, and + ``max_bytes_per_file`` is ignored so every group splits at the same rows. + (default: None, one file per fragment) + """ diff --git a/python/python/tests/test_optimize.py b/python/python/tests/test_optimize.py index 34199956c1e..adde41cead7 100644 --- a/python/python/tests/test_optimize.py +++ b/python/python/tests/test_optimize.py @@ -909,3 +909,43 @@ def test_remap_row_addrs(tmp_path: Path): pa.array([old[i] for i in sample], pa.uint64()) ).to_pylist() assert remapped == [new[i] for i in sample] + + +def test_compact_files_column_groups(tmp_path: Path): + data = pa.table({"a": range(8), "b": [str(i) for i in range(8)], "c": range(8)}) + dataset = lance.write_dataset( + data, tmp_path / "dataset", max_rows_per_file=2, data_storage_version="2.0" + ) + dataset.delete("a = 5") + expected = dataset.to_table() + + with pytest.raises(OSError, match="more than once"): + dataset.optimize.compact_files(column_groups=[["c"], ["c"]]) + with pytest.raises(OSError, match="not a top-level column"): + dataset.optimize.compact_files(column_groups=[["missing"]]) + with pytest.raises(OSError, match="binary copy is not supported"): + dataset.optimize.compact_files( + column_groups=[["c"]], compaction_mode="force_binary_copy" + ) + + metrics = dataset.optimize.compact_files( + target_rows_per_fragment=100, + column_groups=[["c"]], + data_storage_version="2.2", + ) + assert metrics.fragments_added == 1 + assert metrics.files_added == 2 + (fragment,) = dataset.get_fragments() + files = fragment.data_files() + assert [f.fields for f in files] == [[0, 1], [2]] + assert all((f.file_major_version, f.file_minor_version) == (2, 2) for f in files) + assert dataset.to_table() == expected + assert dataset.data_storage_version == "2.0" + + # The persisted config makes later compactions keep the layout. + dataset.update_config({"lance.compaction.column_groups": "c"}) + dataset.insert(pa.table({"a": [8], "b": ["8"], "c": [8]})) + dataset.optimize.compact_files(target_rows_per_fragment=100) + (fragment,) = dataset.get_fragments() + assert [f.fields for f in fragment.data_files()] == [[0, 1], [2]] + assert dataset.count_rows() == 8 diff --git a/python/python/tests/test_schema_evolution.py b/python/python/tests/test_schema_evolution.py index abd89e0cded..709be633e32 100644 --- a/python/python/tests/test_schema_evolution.py +++ b/python/python/tests/test_schema_evolution.py @@ -627,3 +627,78 @@ def test_project_nullability_assertion_round_trips(tmp_path: Path): ) with pytest.raises(Exception, match="preempted"): lance.LanceDataset.commit(tmp_path, relax, read_version=written_at) + + +def test_rewrite_columns(tmp_path: Path): + table = pa.table({"a": range(8), "b": [str(i) for i in range(8)], "c": range(8)}) + dataset = lance.write_dataset( + table, tmp_path, max_rows_per_file=4, data_storage_version="2.0" + ) + dataset.delete("a = 3") + dataset.create_scalar_index("a", "BTREE") + expected = dataset.to_table() + version = dataset.version + + dataset.rewrite_columns(["c"], data_storage_version="2.2") + + assert dataset.version == version + 1 + for fragment in dataset.get_fragments(): + files = fragment.data_files() + assert [f.fields for f in files] == [[0, 1, -2], [2]] + assert (files[0].file_major_version, files[0].file_minor_version) == (2, 0) + assert (files[1].file_major_version, files[1].file_minor_version) == (2, 2) + assert dataset.to_table() == expected + assert dataset.data_storage_version == "2.0" + # The values did not change, so the index still covers every fragment. + (index,) = dataset.describe_indices() + (segment,) = index.segments + assert set(segment.fragment_ids) == {0, 1} + assert dataset.to_table(filter="a = 5").num_rows == 1 + + # Already in the requested layout: no new version. + dataset.rewrite_columns(["c"], data_storage_version="2.2") + assert dataset.version == version + 1 + + with pytest.raises(ValueError, match="not a top-level column"): + dataset.rewrite_columns(["missing"]) + + # A compaction that knows the group keeps `c` apart while moving the rest. + dataset.optimize.compact_files( + target_rows_per_fragment=100, + column_groups=[["c"]], + data_storage_version="2.2", + ) + (fragment,) = dataset.get_fragments() + files = fragment.data_files() + assert [f.fields for f in files] == [[0, 1], [2]] + assert all((f.file_major_version, f.file_minor_version) == (2, 2) for f in files) + assert dataset.to_table() == expected + + +def test_rewrite_columns_per_fragment_commit(tmp_path: Path): + table = pa.table({"a": range(6), "b": [str(i) for i in range(6)]}) + dataset = lance.write_dataset( + table, tmp_path, max_rows_per_file=3, data_storage_version="2.0" + ) + expected = dataset.to_table() + + # The distributed shape: rewrite each fragment on its own, commit once. + updated = [ + fragment.rewrite_columns(["b"], data_storage_version="2.1") + for fragment in dataset.get_fragments() + ] + assert all(metadata is not None for metadata in updated) + operation = lance.LanceOperation.Update( + updated_fragments=updated, update_mode="rewrite_columns" + ) + dataset = lance.LanceDataset.commit( + dataset, operation, read_version=dataset.version + ) + + for fragment in dataset.get_fragments(): + files = fragment.data_files() + assert [f.fields for f in files] == [[0, -2], [1]] + assert (files[1].file_major_version, files[1].file_minor_version) == (2, 1) + # Nothing left to do for this fragment. + assert fragment.rewrite_columns(["b"], data_storage_version="2.1") is None + assert dataset.to_table() == expected diff --git a/python/src/dataset.rs b/python/src/dataset.rs index 6b442fb3027..eb46a40f418 100644 --- a/python/src/dataset.rs +++ b/python/src/dataset.rs @@ -3272,6 +3272,28 @@ impl Dataset { Ok(()) } + #[pyo3(signature = (columns, data_storage_version = None))] + fn rewrite_columns( + &mut self, + columns: Vec, + data_storage_version: Option<&str>, + ) -> PyResult<()> { + let version = data_storage_version + .map(str::parse) + .transpose() + .infer_error()?; + let mut new_self = self.ds.as_ref().clone(); + let new_self = rt() + .spawn(None, async move { + let columns: Vec<&str> = columns.iter().map(String::as_str).collect(); + new_self.rewrite_columns(&columns, version).await?; + Ok(new_self) + })? + .infer_error()?; + self.ds = Arc::new(new_self); + Ok(()) + } + #[pyo3(signature = (reader, batch_size = None))] fn add_columns_from_reader( &mut self, diff --git a/python/src/dataset/optimize.rs b/python/src/dataset/optimize.rs index e17023b7450..cefb519a57a 100644 --- a/python/src/dataset/optimize.rs +++ b/python/src/dataset/optimize.rs @@ -92,6 +92,11 @@ fn parse_compaction_options( opts.data_storage_version = Some(version.parse().infer_error()?); } } + "column_groups" => { + opts.column_groups = value + .extract::>>>()? + .unwrap_or_default(); + } _ => { return Err(PyValueError::new_err(format!( "Invalid compaction option: {}", @@ -129,8 +134,8 @@ pub struct PyCompactionMetrics { /// int : The number of files that have been removed, including deletion files. #[pyo3(get)] pub files_removed: usize, - /// int : The number of files that have been added, which is always equal to the - /// number of fragments. + /// int : The number of files that have been added: one per new fragment, or + /// one per column group and fragment when ``column_groups`` is set. #[pyo3(get)] pub files_added: usize, } diff --git a/python/src/fragment.rs b/python/src/fragment.rs index ae7135443e0..7b40bc9bc55 100644 --- a/python/src/fragment.rs +++ b/python/src/fragment.rs @@ -380,6 +380,26 @@ impl FileFragment { Ok((PyLance(fragment), LanceSchema(schema))) } + #[pyo3(signature = (columns, data_storage_version = None))] + fn rewrite_columns( + &self, + columns: Vec, + data_storage_version: Option<&str>, + ) -> PyResult>> { + let version = data_storage_version + .map(str::parse) + .transpose() + .infer_error()?; + let fragment = self.fragment.clone(); + let updated = rt() + .spawn(None, async move { + let columns: Vec<&str> = columns.iter().map(String::as_str).collect(); + fragment.rewrite_columns(&columns, version).await + })? + .infer_error()?; + Ok(updated.map(PyLance)) + } + fn merge( &mut self, reader: PyArrowType, diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index ab415549d20..0ea97ce79d9 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -82,6 +82,7 @@ pub mod optimize; pub(crate) mod overlay; pub mod progress; pub mod refs; +pub mod rewrite_columns; pub mod rowids; pub mod scanner; mod schema_evolution; diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index 9632d338035..889aa90cffa 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -122,7 +122,7 @@ use lance_arrow::{list::ListArrayExt, r#struct::StructArrayExt}; use lance_core::Error; use lance_core::datatypes::{ BLOB_V2_LOGICAL_FIELDS, BLOB_V2_LOGICAL_TYPE, BlobHandling, BlobKind, BlobV2Layout, - Field as LanceField, + Field as LanceField, Schema as LanceSchema, }; use lance_core::utils::tokio::get_num_compute_intensive_cpus; use lance_core::utils::tracing::{DATASET_COMPACTING_EVENT, TRACE_DATASET_EVENTS}; @@ -322,6 +322,19 @@ pub struct CompactionOptions { /// }; /// ``` pub data_storage_version: Option, + /// Top-level columns to keep in their own data files. + /// + /// Each inner list becomes one data file per compacted fragment holding + /// exactly those columns; every column not listed goes to one shared file. + /// Empty (the default) writes every column to a single file per fragment. + /// + /// This is the layout [`Dataset::rewrite_columns`] produces, so a + /// compaction configured with the same groups preserves that split instead + /// of folding wide columns back into the same file as narrow ones. Binary + /// copy is disabled when groups are set, and `max_bytes_per_file` is + /// ignored so every group's files split at the same rows. + #[serde(default)] + pub column_groups: Vec>, /// Transaction properties to store with this commit. /// /// These key-value pairs are stored in the transaction file @@ -356,6 +369,7 @@ impl Default for CompactionOptions { excluded_fragment_ids: Vec::new(), max_overlays_per_fragment: Some(10), data_storage_version: None, + column_groups: Vec::new(), transaction_properties: None, } } @@ -386,6 +400,7 @@ impl CompactionOptions { /// - `lance.compaction.max_source_bytes` /// - `lance.compaction.max_overlays_per_fragment` /// - `lance.compaction.data_storage_version` + /// - `lance.compaction.column_groups` (groups separated by `;`, columns by `,`) pub fn from_dataset_config(config: &HashMap) -> Result { let mut opts = Self::default(); opts.apply_dataset_config(config)?; @@ -488,6 +503,20 @@ impl CompactionOptions { )) })?); } + "column_groups" => { + self.column_groups = value + .split(';') + .map(|group| { + group + .split(',') + .map(str::trim) + .filter(|column| !column.is_empty()) + .map(str::to_owned) + .collect::>() + }) + .filter(|group| !group.is_empty()) + .collect(); + } "binary_copy_read_batch_bytes" => { self.binary_copy_read_batch_bytes = Some(value.parse().map_err(|_| { Error::invalid_input(format!( @@ -562,6 +591,16 @@ impl CompactionOptions { ))); } } + + self.column_groups.retain(|group| !group.is_empty()); + let mut seen = HashSet::new(); + for column in self.column_groups.iter().flatten() { + if !seen.insert(column.as_str()) { + return Err(Error::invalid_input(format!( + "CompactionOptions::column_groups lists column \"{column}\" more than once" + ))); + } + } Ok(()) } @@ -647,6 +686,11 @@ pub(super) async fn can_use_binary_copy_current( return Ok(false); } + if !options.column_groups.is_empty() { + log::debug!("Binary copy disabled: column_groups splits each fragment across files"); + return Ok(false); + } + let has_blob_columns = dataset .schema() .fields_pre_order() @@ -811,6 +855,8 @@ impl CompactionPlanner for DefaultCompactionPlanner { dataset.manifest.data_storage_format.lance_file_format(), write_version, )?; + // Fail on an unknown column here, before tasks are distributed. + column_group_schemas(dataset.schema(), &self.options.column_groups)?; if self.options.defer_index_remap && dataset.manifest.uses_stable_row_ids() { return Err(Error::invalid_input( "defer_index_remap=true is not supported on datasets with stable row IDs: \ @@ -2580,19 +2626,33 @@ async fn rewrite_files( row_ids_rx = Some(rx); } } else { - let (frags, _) = write_fragments_internal_with_file_row_counts( - write_version, - Some(dataset.as_ref()), - dataset.object_store.clone(), - &dataset.base, - dataset.schema().clone(), - reader.expect("reader must be prepared for non-binary-copy path"), - params, - None, - Some(file_row_counts), - ) - .await?; - new_fragments = frags; + let reader = reader.expect("reader must be prepared for non-binary-copy path"); + new_fragments = if options.column_groups.is_empty() { + let (frags, _) = write_fragments_internal_with_file_row_counts( + write_version, + Some(dataset.as_ref()), + dataset.object_store.clone(), + &dataset.base, + dataset.schema().clone(), + reader, + params, + None, + Some(file_row_counts), + ) + .await?; + frags + } else { + let group_schemas = column_group_schemas(dataset.schema(), &options.column_groups)?; + write_column_group_fragments( + write_version, + dataset.as_ref(), + group_schemas, + reader, + params, + file_row_counts, + ) + .await? + }; } log::info!("Compaction task {}: file written", task_id); @@ -2660,6 +2720,135 @@ async fn rewrite_files( }) } +/// The schema of each data file a compacted fragment gets when column groups +/// are configured: the columns no group claims first (when any), then one +/// schema per group. Columns keep the dataset's order within each file. +fn column_group_schemas(schema: &LanceSchema, groups: &[Vec]) -> Result> { + let names = || schema.fields.iter().map(|field| field.name.as_str()); + for column in groups.iter().flatten() { + if !names().any(|name| name == column) { + return Err(Error::invalid_input(format!( + "column_groups names \"{column}\", which is not a top-level column of the dataset" + ))); + } + } + let claimed = |name: &str| groups.iter().flatten().any(|column| column == name); + let rest: Vec<&str> = names().filter(|name| !claimed(name)).collect(); + std::iter::once(rest) + .chain(groups.iter().map(|group| { + names() + .filter(|name| group.iter().any(|column| column == name)) + .collect() + })) + .filter(|columns: &Vec<&str>| !columns.is_empty()) + .map(|columns| schema.project(&columns)) + .collect() +} + +/// Write one data file per column group for every output fragment. A single +/// scan feeds one writer per group, and every writer closes its files at the +/// same planned row counts, so the files of one fragment stay row-aligned. +async fn write_column_group_fragments( + write_version: ConcreteFileVersion, + dataset: &Dataset, + group_schemas: Vec, + mut reader: SendableRecordBatchStream, + mut params: WriteParams, + file_row_counts: Vec, +) -> Result> { + use futures::SinkExt; + + // Only the planned row counts may close a file: a byte-driven split in one + // group would misalign it with the others. + params.max_bytes_per_file = usize::MAX; + versions::validate_write_schema(write_version, dataset.schema())?; + + let arrow_schema = reader.schema(); + let projections = group_schemas + .iter() + .map(|schema| { + schema + .fields + .iter() + .map(|field| arrow_schema.index_of(&field.name)) + .collect::, _>>() + }) + .collect::, _>>()?; + let (mut senders, receivers): (Vec<_>, Vec<_>) = (0..group_schemas.len()) + .map(|_| futures::channel::mpsc::channel::>(1)) + .unzip(); + + // Owns the senders so they drop, and the writers see end-of-stream, as + // soon as the scan is exhausted. + let forward_projections = projections.clone(); + let forward = async move { + while let Some(batch) = reader.next().await { + let batch = batch?; + for (sender, projection) in senders.iter_mut().zip(&forward_projections) { + if sender.send(Ok(batch.project(projection)?)).await.is_err() { + // That writer failed; its own error is what the join reports. + return Ok(()); + } + } + } + Ok::<(), Error>(()) + }; + + let writers = receivers + .into_iter() + .zip(group_schemas) + .zip(&projections) + .map(|((receiver, schema), projection)| { + let arrow_schema = Arc::new(arrow_schema.project(projection)?); + let stream: SendableRecordBatchStream = + Box::pin(RecordBatchStreamAdapter::new(arrow_schema, receiver)); + let params = params.clone(); + let file_row_counts = file_row_counts.clone(); + Ok(async move { + let mut seed_writers = + versions::create_seed_writers(write_version, Some(dataset), ¶ms).await?; + seed_writers.retain(|writer| schema.field(writer.column_name()).is_some()); + versions::write_fragments_direct( + write_version, + Some(dataset), + dataset.object_store.clone(), + &dataset.base, + &schema, + stream, + params, + None, + seed_writers, + Some(file_row_counts), + ) + .await + }) + }) + .collect::>>()?; + let (_, per_group) = futures::try_join!(forward, futures::future::try_join_all(writers))?; + + let mut per_group = per_group.into_iter(); + let mut fragments = per_group.next().unwrap_or_default(); + for group in per_group { + if group.len() != fragments.len() { + return Err(Error::internal(format!( + "column groups wrote {} and {} fragments for the same rows", + fragments.len(), + group.len() + ))); + } + for (fragment, other) in fragments.iter_mut().zip(group) { + if fragment.physical_rows != other.physical_rows { + return Err(Error::internal(format!( + "column groups wrote {:?} and {:?} rows for the same fragment", + fragment.physical_rows, other.physical_rows + ))); + } + fragment.files.extend(other.files); + } + } + Ok(fragments) +} + async fn rechunk_stable_row_ids( dataset: &Dataset, new_fragments: &mut [Fragment], @@ -3370,6 +3559,119 @@ mod tests { } } + #[rstest] + #[tokio::test] + async fn compaction_column_groups_write_one_file_per_group( + #[values(None, Some(LanceFileVersion::V2_2))] target: Option, + #[values( + CompactionMode::Reencode, + CompactionMode::TryBinaryCopy, + CompactionMode::ForceBinaryCopy + )] + mode: CompactionMode, + ) { + let batch = arrow_array::record_batch!( + ("a", Int32, [1, 2, 3, 4]), + ("b", Utf8, ["w", "x", "y", "z"]), + ("c", Int64, [10, 20, 30, 40]) + ) + .unwrap(); + let mut dataset = Dataset::write( + RecordBatchIterator::new([Ok(batch.clone())], batch.schema()), + "memory://", + Some(WriteParams { + max_rows_per_file: 2, + data_storage_version: Some(LanceFileVersion::V2_0), + ..Default::default() + }), + ) + .await + .unwrap(); + dataset.delete("a = 2").await.unwrap(); + let before = dataset.scan().try_into_batch().await.unwrap(); + + let options = CompactionOptions { + target_rows_per_fragment: 8, + column_groups: vec![vec!["c".to_string()], vec!["a".to_string()]], + data_storage_version: target, + compaction_mode: Some(mode), + ..Default::default() + }; + let result = compact_files(&mut dataset, options, None).await; + if mode == CompactionMode::ForceBinaryCopy { + // Groups split every fragment across files, which binary copy + // cannot produce. + assert!(matches!(result.unwrap_err(), Error::NotSupported { .. })); + return; + } + let metrics = result.unwrap(); + assert_eq!(metrics.fragments_added, 1); + assert_eq!(metrics.files_added, 3); + + let expected = target.unwrap_or(LanceFileVersion::V2_0).resolve(); + let fragment = &dataset.manifest.fragments[0]; + // Unclaimed columns first, then the groups in the order given. + assert_eq!( + fragment + .files + .iter() + .map(|file| file.fields.to_vec()) + .collect::>(), + vec![vec![1], vec![2], vec![0]] + ); + assert!( + fragment + .files + .iter() + .all(|file| file.file_version().unwrap() == expected) + ); + assert_eq!(dataset.scan().try_into_batch().await.unwrap(), before); + dataset.validate().await.unwrap(); + } + + #[tokio::test] + async fn compaction_column_groups_validation() { + let mut duplicated = CompactionOptions { + column_groups: vec![vec!["a".to_string()], vec!["a".to_string()]], + ..Default::default() + }; + let err = duplicated.validate().unwrap_err(); + assert!(err.to_string().contains("more than once"), "{err}"); + + let mut config = HashMap::new(); + config.insert( + "lance.compaction.column_groups".to_string(), + " c ; a, b ;".to_string(), + ); + let options = CompactionOptions::from_dataset_config(&config).unwrap(); + assert_eq!( + options.column_groups, + vec![ + vec!["c".to_string()], + vec!["a".to_string(), "b".to_string()] + ] + ); + + let batch = arrow_array::record_batch!(("a", Int32, [1, 2])).unwrap(); + let dataset = Dataset::write( + RecordBatchIterator::new([Ok(batch.clone())], batch.schema()), + "memory://", + None, + ) + .await + .unwrap(); + let err = plan_compaction( + &dataset, + &CompactionOptions { + column_groups: vec![vec!["missing".to_string()]], + ..Default::default() + }, + ) + .await + .unwrap_err(); + assert!(matches!(err, Error::InvalidInput { .. }), "{err}"); + } + #[rstest] #[tokio::test] async fn test_compact_empty( diff --git a/rust/lance/src/dataset/rewrite_columns.rs b/rust/lance/src/dataset/rewrite_columns.rs new file mode 100644 index 00000000000..b624471a0b7 --- /dev/null +++ b/rust/lance/src/dataset/rewrite_columns.rs @@ -0,0 +1,412 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Rewrite a few columns of every fragment into their own data files while +//! leaving the files that hold the other columns untouched. +//! +//! Compaction rewrites whole fragments, so migrating a table to a newer data +//! file version through compaction re-encodes every column, including the +//! wide ones that dominate its size. `rewrite_columns` instead re-reads only +//! the named columns, writes them to one new file per fragment in the requested +//! version, and tombstones them where they used to live. Rows, fragment ids, +//! row addresses and indices are unchanged. +//! +//! The resulting layout survives compaction only when the compaction is told +//! about it: set [`CompactionOptions::column_groups`] (or the +//! `lance.compaction.column_groups` dataset config key) to the same groups. +//! +//! [`CompactionOptions::column_groups`]: super::optimize::CompactionOptions::column_groups + +use std::collections::HashSet; + +use lance_core::datatypes::Schema; +use lance_file::version::{ConcreteFileVersion, LanceFileVersion}; +use lance_table::format::{Fragment, overlay::TOMBSTONE_FIELD_ID}; + +use super::fragment::FileFragment; +use super::transaction::{Operation, Transaction, UpdateMode}; +use super::{Dataset, versions, write::cleanup_data_fragments}; +use crate::{Error, Result}; + +/// The exact V2 version the new files get: `requested`, or the dataset's +/// default write version, which is never changed. +fn write_version( + dataset: &Dataset, + requested: Option, +) -> Result { + let default_version = dataset.manifest.data_storage_format.lance_file_format(); + let version = requested + .map(LanceFileVersion::resolve) + .unwrap_or(default_version); + versions::validate_write_version(default_version, version)?; + if version == ConcreteFileVersion::V1 { + return Err(Error::not_supported( + "rewrite_columns requires the V2 file format: a V1 file cannot tombstone \ + single fields, so compact the dataset to V2 first", + )); + } + Ok(version) +} + +/// The schema of the new data file: `columns` in dataset order, all top-level. +fn rewrite_schema(dataset: &Dataset, columns: &[&str]) -> Result { + if columns.is_empty() { + return Err(Error::invalid_input( + "rewrite_columns needs at least one column", + )); + } + let schema = dataset.schema(); + for column in columns { + if !schema.fields.iter().any(|field| &field.name == column) { + return Err(Error::invalid_input(format!( + "Column \"{column}\" is not a top-level column of the dataset" + ))); + } + } + let ordered = schema + .fields + .iter() + .map(|field| field.name.as_str()) + .filter(|name| columns.contains(name)) + .collect::>(); + let schema = schema.project(&ordered)?; + if let Some(blob) = schema.fields_pre_order().find(|field| field.is_blob()) { + return Err(Error::not_supported(format!( + "rewrite_columns cannot rewrite blob column \"{}\"", + blob.name + ))); + } + Ok(schema) +} + +impl FileFragment { + /// Rewrite `columns` of this fragment into one new data file. + /// + /// The returned metadata has the new file appended and the rewritten + /// fields tombstoned in the files they came from; a file left holding only + /// tombstones is dropped. Nothing is committed: pass the metadata of every + /// rewritten fragment to [`Operation::Update`] with + /// [`UpdateMode::RewriteColumns`], which is how [`Dataset::rewrite_columns`] + /// commits and how a distributed caller commits from many workers. + /// + /// Returns `None` when the columns already sit alone in a file of the + /// requested version, so a retried or resumed rewrite skips finished work. + pub async fn rewrite_columns( + &self, + columns: &[&str], + data_storage_version: Option, + ) -> Result> { + let dataset = self.dataset(); + let write_schema = rewrite_schema(dataset, columns)?; + let write_version = write_version(dataset, data_storage_version)?; + let field_ids: HashSet = write_schema.field_ids().into_iter().collect(); + + let already_split = self.metadata().files.iter().any(|file| { + file.fields + .iter() + .copied() + .filter(|id| *id != TOMBSTONE_FIELD_ID) + .collect::>() + == field_ids + && file + .file_version() + .is_ok_and(|version| version == write_version) + }); + if already_split { + return Ok(None); + } + + let mut updater = self + .updater_with_version( + Some(columns), + Some((write_schema, dataset.schema().clone())), + None, + None, + write_version, + ) + .await?; + let written: Result = async { + while let Some(batch) = updater.next().await?.cloned() { + updater.update(batch).await?; + } + updater.finish().await + } + .await; + let mut fragment = match written { + Ok(fragment) => fragment, + Err(err) => { + updater.cleanup_unfinished_writer().await; + return Err(err); + } + }; + // No live rows, so no file was written and nothing to tombstone. + if fragment.files.len() == self.metadata().files.len() { + return Ok(None); + } + + let new_file = fragment.files.len() - 1; + for file in &mut fragment.files[..new_file] { + file.fields = file + .fields + .iter() + .map(|id| { + if field_ids.contains(id) { + TOMBSTONE_FIELD_ID + } else { + *id + } + }) + .collect::>() + .into(); + } + // A file holding only tombstones is unreachable to readers. + fragment + .files + .retain(|file| file.fields.iter().any(|id| *id != TOMBSTONE_FIELD_ID)); + Ok(Some(fragment)) + } +} + +impl Dataset { + /// Rewrite `columns` of every fragment into one new data file per fragment, + /// optionally in a different data file version, without touching the files + /// that hold the other columns. + /// + /// Fragments whose `columns` already sit alone in a file of the requested + /// version are skipped, so an interrupted rewrite can be rerun. Nothing is + /// committed when every fragment is skipped. + /// + /// Fragments are rewritten one after another. To spread the work over many + /// machines, call [`FileFragment::rewrite_columns`] per fragment and commit + /// the returned metadata in one [`Operation::Update`] with + /// [`UpdateMode::RewriteColumns`]. + /// + /// A later compaction folds the new files back into one file per fragment + /// unless its `column_groups` lists the same columns. + pub async fn rewrite_columns( + &mut self, + columns: &[&str], + data_storage_version: Option, + ) -> Result<()> { + let fragments = self.get_fragments(); + let mut updated_fragments = Vec::with_capacity(fragments.len()); + for fragment in &fragments { + match fragment + .rewrite_columns(columns, data_storage_version) + .await + { + Ok(Some(updated)) => updated_fragments.push(updated), + Ok(None) => {} + Err(err) => { + self.cleanup_rewritten_files(&updated_fragments).await; + return Err(err); + } + } + } + if updated_fragments.is_empty() { + return Ok(()); + } + + let transaction = Transaction::new( + self.manifest.version, + Operation::Update { + removed_fragment_ids: Vec::new(), + updated_fragments, + new_fragments: Vec::new(), + // The values are unchanged, so indices over these fields stay + // valid and overlays on them keep shadowing the same cells. + fields_modified: Vec::new(), + compacted_sstables: Vec::new(), + fields_for_preserving_frag_bitmap: Vec::new(), + update_mode: Some(UpdateMode::RewriteColumns), + inserted_rows_filter: None, + updated_fragment_offsets: None, + }, + None, + ); + self.apply_commit(transaction, &Default::default(), &Default::default()) + .await + } + + /// Delete the files a failed rewrite wrote: the last file of every + /// rewritten fragment is the only one that did not exist before. + async fn cleanup_rewritten_files(&self, updated_fragments: &[Fragment]) { + let new_files = updated_fragments + .iter() + .filter_map(|fragment| fragment.files.last()) + .map(|file| Fragment { + files: vec![file.clone()], + ..Fragment::new(0) + }) + .collect::>(); + cleanup_data_fragments(&self.object_store, &self.base, None, &new_files).await; + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::dataset::WriteParams; + use crate::dataset::optimize::{CompactionOptions, compact_files}; + use crate::index::DatasetIndexExt; + use arrow_array::{Int32Array, RecordBatch, RecordBatchIterator, StringArray}; + use arrow_schema::{DataType, Field, Schema as ArrowSchema}; + use lance_core::utils::tempfile::TempStrDir; + use lance_index::IndexType; + use lance_index::scalar::ScalarIndexParams; + use std::sync::Arc; + + fn batch(start: i32, rows: i32) -> RecordBatch { + let schema = Arc::new(ArrowSchema::new(vec![ + Field::new("a", DataType::Int32, false), + Field::new("b", DataType::Int32, true), + Field::new("c", DataType::Utf8, true), + ])); + RecordBatch::try_new( + schema, + vec![ + Arc::new(Int32Array::from_iter_values(start..start + rows)), + Arc::new(Int32Array::from_iter_values( + (start..start + rows).map(|v| v * 10), + )), + Arc::new(StringArray::from_iter_values( + (start..start + rows).map(|v| format!("row-{v}")), + )), + ], + ) + .unwrap() + } + + async fn write(uri: &str, version: LanceFileVersion) -> Dataset { + let data = batch(0, 8); + let reader = RecordBatchIterator::new([Ok(data.clone())], data.schema()); + Dataset::write( + reader, + uri, + Some(WriteParams { + max_rows_per_file: 4, + data_storage_version: Some(version), + ..Default::default() + }), + ) + .await + .unwrap() + } + + fn file_versions(fragment: &Fragment) -> Vec<(Vec, ConcreteFileVersion)> { + fragment + .files + .iter() + .map(|file| (file.fields.to_vec(), file.file_version().unwrap())) + .collect() + } + + #[tokio::test] + async fn rewrite_migrates_only_the_named_column() { + let mut dataset = write("memory://", LanceFileVersion::V2_0).await; + dataset.delete("a = 3").await.unwrap(); + let before = dataset.scan().try_into_batch().await.unwrap(); + let version = dataset.manifest.version; + + dataset + .rewrite_columns(&["c"], Some(LanceFileVersion::V2_2)) + .await + .unwrap(); + + assert_eq!(dataset.manifest.version, version + 1); + assert_eq!(dataset.manifest.fragments.len(), 2); + for fragment in dataset.manifest.fragments.iter() { + assert_eq!( + file_versions(fragment), + vec![ + (vec![0, 1, TOMBSTONE_FIELD_ID], ConcreteFileVersion::V2_0), + (vec![2], ConcreteFileVersion::V2_2), + ] + ); + } + assert_eq!(dataset.scan().try_into_batch().await.unwrap(), before); + assert_eq!( + dataset.manifest.data_storage_format.lance_file_format(), + ConcreteFileVersion::V2_0 + ); + dataset.validate().await.unwrap(); + + // Already in the requested layout: nothing to do, nothing committed. + dataset + .rewrite_columns(&["c"], Some(LanceFileVersion::V2_2)) + .await + .unwrap(); + assert_eq!(dataset.manifest.version, version + 1); + + // A compaction that knows the group keeps `c` in its own file while + // moving everything to the target version. + compact_files( + &mut dataset, + CompactionOptions { + target_rows_per_fragment: 100, + column_groups: vec![vec!["c".to_string()]], + data_storage_version: Some(LanceFileVersion::V2_2), + ..Default::default() + }, + None, + ) + .await + .unwrap(); + assert_eq!(dataset.manifest.fragments.len(), 1); + assert_eq!( + file_versions(&dataset.manifest.fragments[0]), + vec![ + (vec![0, 1], ConcreteFileVersion::V2_2), + (vec![2], ConcreteFileVersion::V2_2), + ] + ); + assert_eq!(dataset.scan().try_into_batch().await.unwrap(), before); + dataset.validate().await.unwrap(); + } + + #[tokio::test] + async fn rewrite_keeps_scalar_index() { + let mut dataset = write("memory://", LanceFileVersion::V2_0).await; + dataset + .create_index( + &["a"], + IndexType::BTree, + Some("a_idx".into()), + &ScalarIndexParams::default(), + false, + ) + .await + .unwrap(); + let before = dataset.load_indices().await.unwrap()[0].clone(); + + dataset.rewrite_columns(&["a"], None).await.unwrap(); + + let after = dataset.load_indices().await.unwrap()[0].clone(); + assert_eq!(after.uuid, before.uuid); + assert_eq!(after.fragment_bitmap, before.fragment_bitmap); + let filtered = dataset + .scan() + .filter("a = 5") + .unwrap() + .try_into_batch() + .await + .unwrap(); + assert_eq!(filtered.num_rows(), 1); + dataset.validate().await.unwrap(); + } + + #[tokio::test] + async fn rewrite_rejects_bad_input() { + let mut dataset = write("memory://", LanceFileVersion::V2_0).await; + let err = dataset + .rewrite_columns(&["missing"], None) + .await + .unwrap_err(); + assert!(matches!(err, Error::InvalidInput { .. }), "{err}"); + + let legacy_dir = TempStrDir::default(); + let mut legacy = write(&legacy_dir, LanceFileVersion::Legacy).await; + let err = legacy.rewrite_columns(&["c"], None).await.unwrap_err(); + assert!(matches!(err, Error::NotSupported { .. }), "{err}"); + } +} diff --git a/rust/lance/src/dataset/versions/mod.rs b/rust/lance/src/dataset/versions/mod.rs index be99a42bbe4..cc1ff03ab9c 100644 --- a/rust/lance/src/dataset/versions/mod.rs +++ b/rust/lance/src/dataset/versions/mod.rs @@ -106,7 +106,7 @@ pub fn schema_compare_options(version: ConcreteFileVersion) -> SchemaCompareOpti } } -async fn create_seed_writers( +pub(super) async fn create_seed_writers( version: ConcreteFileVersion, dataset: Option<&Dataset>, params: &WriteParams, @@ -146,21 +146,13 @@ pub async fn write_fragments( target_bases_info: Option>, file_row_counts: 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)?; - } - } + validate_write_schema(version, &schema)?; let seed_writers = create_seed_writers(version, dataset, ¶ms).await?; let fragments = write_fragments_direct( version, @@ -178,6 +170,18 @@ pub async fn write_fragments( Ok((fragments, schema)) } +/// Reject a schema whose blob columns the file `version` cannot store. +pub(super) fn validate_write_schema(version: ConcreteFileVersion, schema: &Schema) -> Result<()> { + match version { + ConcreteFileVersion::V1 | ConcreteFileVersion::V2_0 | ConcreteFileVersion::V2_1 => { + write::validate_legacy_blob_write_schema(schema, &format!("{version:?}")) + } + ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 => { + write::validate_blob_v2_write_schema(schema) + } + } +} + #[allow(clippy::too_many_arguments)] pub async fn write_fragments_direct( version: ConcreteFileVersion, From 7a9863983f28e4bd1361b21725609e0dc09d81db Mon Sep 17 00:00:00 2001 From: dshepelev15 Date: Wed, 16 Sep 2026 14:21:48 +0000 Subject: [PATCH 2/2] refactor: tighten column group alignment check and column lookup Fold the two internal errors on misaligned column group output into one, hand the forwarder its projections by value instead of cloning them, and look up an unknown rewrite column with find instead of a loop. --- rust/lance/src/dataset/optimize.rs | 40 +++++++++++------------ rust/lance/src/dataset/rewrite_columns.rs | 13 ++++---- 2 files changed, 27 insertions(+), 26 deletions(-) diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index 889aa90cffa..33b20eb0c86 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -2778,13 +2778,17 @@ async fn write_column_group_fragments( .map(|_| futures::channel::mpsc::channel::>(1)) .unzip(); + let group_arrow_schemas = projections + .iter() + .map(|projection| arrow_schema.project(projection).map(Arc::new)) + .collect::, _>>()?; + // Owns the senders so they drop, and the writers see end-of-stream, as // soon as the scan is exhausted. - let forward_projections = projections.clone(); let forward = async move { while let Some(batch) = reader.next().await { let batch = batch?; - for (sender, projection) in senders.iter_mut().zip(&forward_projections) { + for (sender, projection) in senders.iter_mut().zip(&projections) { if sender.send(Ok(batch.project(projection)?)).await.is_err() { // That writer failed; its own error is what the join reports. return Ok(()); @@ -2797,14 +2801,13 @@ async fn write_column_group_fragments( let writers = receivers .into_iter() .zip(group_schemas) - .zip(&projections) - .map(|((receiver, schema), projection)| { - let arrow_schema = Arc::new(arrow_schema.project(projection)?); + .zip(group_arrow_schemas) + .map(|((receiver, schema), arrow_schema)| { let stream: SendableRecordBatchStream = Box::pin(RecordBatchStreamAdapter::new(arrow_schema, receiver)); let params = params.clone(); let file_row_counts = file_row_counts.clone(); - Ok(async move { + async move { let mut seed_writers = versions::create_seed_writers(write_version, Some(dataset), ¶ms).await?; seed_writers.retain(|writer| schema.field(writer.column_name()).is_some()); @@ -2821,28 +2824,25 @@ async fn write_column_group_fragments( Some(file_row_counts), ) .await - }) + } }) - .collect::>>()?; + .collect::>(); let (_, per_group) = futures::try_join!(forward, futures::future::try_join_all(writers))?; let mut per_group = per_group.into_iter(); let mut fragments = per_group.next().unwrap_or_default(); for group in per_group { - if group.len() != fragments.len() { - return Err(Error::internal(format!( - "column groups wrote {} and {} fragments for the same rows", - fragments.len(), - group.len() - ))); + let aligned = group.len() == fragments.len() + && fragments + .iter() + .zip(&group) + .all(|(fragment, other)| fragment.physical_rows == other.physical_rows); + if !aligned { + return Err(Error::internal( + "column group writers did not split the rows at the same fragments", + )); } for (fragment, other) in fragments.iter_mut().zip(group) { - if fragment.physical_rows != other.physical_rows { - return Err(Error::internal(format!( - "column groups wrote {:?} and {:?} rows for the same fragment", - fragment.physical_rows, other.physical_rows - ))); - } fragment.files.extend(other.files); } } diff --git a/rust/lance/src/dataset/rewrite_columns.rs b/rust/lance/src/dataset/rewrite_columns.rs index b624471a0b7..e8c2f7e46bb 100644 --- a/rust/lance/src/dataset/rewrite_columns.rs +++ b/rust/lance/src/dataset/rewrite_columns.rs @@ -56,12 +56,13 @@ fn rewrite_schema(dataset: &Dataset, columns: &[&str]) -> Result { )); } let schema = dataset.schema(); - for column in columns { - if !schema.fields.iter().any(|field| &field.name == column) { - return Err(Error::invalid_input(format!( - "Column \"{column}\" is not a top-level column of the dataset" - ))); - } + if let Some(column) = columns + .iter() + .find(|column| !schema.fields.iter().any(|field| &field.name == *column)) + { + return Err(Error::invalid_input(format!( + "Column \"{column}\" is not a top-level column of the dataset" + ))); } let ordered = schema .fields