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..33b20eb0c86 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(); + + 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 = async move { + while let Some(batch) = reader.next().await { + let batch = batch?; + 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(()); + } + } + } + Ok::<(), Error>(()) + }; + + let writers = receivers + .into_iter() + .zip(group_schemas) + .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(); + 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 { + 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) { + 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..e8c2f7e46bb --- /dev/null +++ b/rust/lance/src/dataset/rewrite_columns.rs @@ -0,0 +1,413 @@ +// 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(); + 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 + .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,