diff --git a/mkdocs/docs/api.md b/mkdocs/docs/api.md index 0a2a72b6f0..fd2f5a2624 100644 --- a/mkdocs/docs/api.md +++ b/mkdocs/docs/api.md @@ -1485,7 +1485,7 @@ table.manage_snapshots().remove_branch("dev").commit() ## Table Maintenance -PyIceberg provides table maintenance operations through the `table.maintenance` API. This provides a clean interface for performing maintenance tasks like snapshot expiration. +PyIceberg provides table maintenance operations through the `table.maintenance` API. This provides a clean interface for performing maintenance tasks like snapshot expiration and manifest rewriting. ### Snapshot Expiration @@ -1527,6 +1527,45 @@ def cleanup_old_snapshots(table_name: str, snapshot_ids: list[int]): cleanup_old_snapshots("analytics.user_events", [12345, 67890, 11111]) ``` +### Manifest Rewriting + +Rewrite the current snapshot's data manifests without changing any data. Live entries are regrouped into new manifests sized by the `commit.manifest.target-size-bytes` table property and written as `EXISTING` entries that keep their sequence numbers, which keeps scan planning fast on tables that accumulate many small manifests through frequent appends. Delete manifests are kept as-is, and the result is committed as a `replace` snapshot: + +```python +table.maintenance.rewrite_manifests().commit() +``` + +Nothing is committed when no manifests would be merged, so the call is safe to schedule periodically. To find out in advance whether a rewrite would do anything: + +```python +if table.maintenance.rewrite_manifests().rewrites_needed(): + table.maintenance.rewrite_manifests().commit() +``` + +The `replace` snapshot records what the operation did in its summary: + +```python +table.current_snapshot().summary +# {'operation': 'replace', 'manifests-created': '1', 'manifests-kept': '0', +# 'manifests-replaced': '3', 'entries-processed': '3', ...} +``` + +By default every data manifest is a candidate. Pass a predicate to `rewrite_if` to choose which ones to rewrite; manifests that do not match are kept exactly as they are. This is the way to target a subset, such as manifests written by an older library version whose contents need re-encoding: + +```python +# Rewrite only the manifests smaller than 1 MiB, leaving larger ones untouched +table.maintenance.rewrite_manifests().rewrite_if(lambda manifest: manifest.manifest_length < 1024 * 1024).commit() +``` + +A predicate also turns off the shortcut that leaves a lone manifest alone: a manifest that matches is rewritten even when there is nothing to merge it with. `rewrite_if(lambda manifest: True)` therefore rewrites every data manifest, which is not the same as passing no predicate at all. + + + +!!! note "V3 tables" + Rewriting manifests on V3 tables raises `NotImplementedError`, because the `first-row-id` of rewritten manifests has to be preserved and that support is still pending ([#3621](https://github.com/apache/iceberg-python/issues/3621)). + + + ## Views If PyIceberg is unable to automatically determine view support on your REST Catalog, you can manually specify, `"view-endpoints-supported": "true"`: diff --git a/pyiceberg/table/maintenance.py b/pyiceberg/table/maintenance.py index 0fcda35ae9..d0f4e0d0ea 100644 --- a/pyiceberg/table/maintenance.py +++ b/pyiceberg/table/maintenance.py @@ -24,7 +24,7 @@ if TYPE_CHECKING: from pyiceberg.table import Table - from pyiceberg.table.update.snapshot import ExpireSnapshots + from pyiceberg.table.update.snapshot import ExpireSnapshots, RewriteManifests class MaintenanceTable: @@ -43,3 +43,17 @@ def expire_snapshots(self) -> ExpireSnapshots: from pyiceberg.table.update.snapshot import ExpireSnapshots return ExpireSnapshots(transaction=Transaction(self.tbl, autocommit=True)) + + def rewrite_manifests(self) -> RewriteManifests: + """Return a RewriteManifests operation that merges the current snapshot's data manifests. + + Entries are rewritten as EXISTING, keeping their sequence numbers; delete + manifests are kept as-is. The result is committed as a `replace` snapshot. + + Returns: + RewriteManifests operation; call commit() to execute it. + """ + from pyiceberg.table import Transaction + from pyiceberg.table.update.snapshot import RewriteManifests + + return RewriteManifests(transaction=Transaction(self.tbl, autocommit=True), io=self.tbl.io) diff --git a/pyiceberg/table/snapshots.py b/pyiceberg/table/snapshots.py index 5e9e519a01..14d15b2c66 100644 --- a/pyiceberg/table/snapshots.py +++ b/pyiceberg/table/snapshots.py @@ -351,7 +351,7 @@ def _partition_summary(self, update_metrics: UpdateMetrics) -> str: def update_snapshot_summaries(summary: Summary, previous_summary: Mapping[str, str] | None = None) -> Summary: - if summary.operation not in {Operation.APPEND, Operation.OVERWRITE, Operation.DELETE}: + if summary.operation not in {Operation.APPEND, Operation.OVERWRITE, Operation.DELETE, Operation.REPLACE}: raise ValueError(f"Operation not implemented: {summary.operation}") if not previous_summary: diff --git a/pyiceberg/table/update/snapshot.py b/pyiceberg/table/update/snapshot.py index 57215dca04..12bb38afb9 100644 --- a/pyiceberg/table/update/snapshot.py +++ b/pyiceberg/table/update/snapshot.py @@ -1312,3 +1312,185 @@ def older_than(self, dt: datetime) -> ExpireSnapshots: if snapshot.timestamp_ms < expire_from and snapshot.snapshot_id not in protected_ids: self._snapshot_ids_to_expire.add(snapshot.snapshot_id) return self + + +class RewriteManifests(_SnapshotProducer["RewriteManifests"]): + """Rewrite the current snapshot's data manifests without changing data. + + Live entries from the rewritten data manifests are regrouped into new + manifests sized by `commit.manifest.target-size-bytes`, written as EXISTING + entries that keep their sequence numbers. Entries with status DELETED are + dropped, matching the reference implementation, which rewrites live entries + only. Delete manifests are kept as-is. The result is committed as a + `replace` snapshot; if no manifests need merging, no snapshot is committed. + """ + + _computed_manifests: list[ManifestFile] | None + _manifest_predicate: Callable[[ManifestFile], bool] | None + + _rewritten_count: int + _created_count: int + _kept_count: int + _entries_processed: int + + def __init__( + self, + transaction: Transaction, + io: FileIO, + commit_uuid: uuid.UUID | None = None, + snapshot_properties: dict[str, str] = EMPTY_DICT, + branch: str | None = MAIN_BRANCH, + ) -> None: + super().__init__(Operation.REPLACE, transaction, io, commit_uuid, snapshot_properties, branch) + if transaction.table_metadata.format_version >= 3: + raise NotImplementedError( + "Rewriting manifests is not yet supported for V3 tables: " + "the first-row-id of rewritten manifests must be preserved, " + "see: https://github.com/apache/iceberg-python/issues/3621" + ) + self._manifest_predicate = None + self._rewritten_count = 0 + self._created_count = 0 + self._kept_count = 0 + self._entries_processed = 0 + self._computed_manifests = None + + def rewrite_if(self, predicate: Callable[[ManifestFile], bool]) -> RewriteManifests: + """Filter which manifests should be rewritten. + + Passing a predicate also disables the optimization that keeps single-manifest + groups as-is, allowing single manifests to be rewritten when they match the predicate. + + Args: + predicate: A function that takes a ManifestFile and returns True if it should be rewritten. + + Returns: + This RewriteManifests instance for method chaining. + """ + self._manifest_predicate = predicate + return self + + def _deleted_entries(self) -> list[ManifestEntry]: + return [] + + def _target_size_bytes(self) -> int: + from pyiceberg.table import TableProperties + + return property_as_int( # type: ignore + self._transaction.table_metadata.properties, + TableProperties.MANIFEST_TARGET_SIZE_BYTES, + TableProperties.MANIFEST_TARGET_SIZE_BYTES_DEFAULT, + ) + + def _group_by_target_size(self, manifests: list[ManifestFile]) -> list[list[ManifestFile]]: + """Pack manifests into groups whose source sizes add up to roughly the target size.""" + target_size = self._target_size_bytes() + groups: list[list[ManifestFile]] = [] + current_group: list[ManifestFile] = [] + current_size = 0 + for manifest in manifests: + if current_group and current_size + manifest.manifest_length > target_size: + groups.append(current_group) + current_group = [] + current_size = 0 + current_group.append(manifest) + current_size += manifest.manifest_length + if current_group: + groups.append(current_group) + return groups + + def _existing_manifests(self) -> list[ManifestFile]: + if self._computed_manifests is not None: + return self._computed_manifests + + snapshot = self._transaction.table_metadata.snapshot_by_name(self._target_branch or MAIN_BRANCH) + if snapshot is None: + self._computed_manifests = [] + return self._computed_manifests + + data_manifests_by_spec: defaultdict[int, list[ManifestFile]] = defaultdict(list) + kept_manifests: list[ManifestFile] = [] + for manifest in snapshot.manifests(self._io): + if manifest.content == ManifestContent.DATA and ( + self._manifest_predicate is None or self._manifest_predicate(manifest) + ): + data_manifests_by_spec[manifest.partition_spec_id].append(manifest) + else: + kept_manifests.append(manifest) + + new_manifests: list[ManifestFile] = [] + for spec_id, manifests in data_manifests_by_spec.items(): + for group in self._group_by_target_size(manifests): + if len(group) == 1 and self._manifest_predicate is None: + # nothing to merge and no predicate specified; keep the manifest as-is + kept_manifests.append(group[0]) + continue + + entries = (entry for manifest in group for entry in manifest.fetch_manifest_entry(self._io, discard_deleted=True)) + first_entry = next(entries, None) + if first_entry is None: + kept_manifests.extend(group) + continue + + with self.new_manifest_writer(self.spec(spec_id)) as writer: + for entry in itertools.chain([first_entry], entries): + writer.existing(entry) + self._entries_processed += 1 + + new_manifests.append(writer.to_manifest_file()) + self._created_count += 1 + self._rewritten_count += len(group) + + self._kept_count = len(kept_manifests) + self.snapshot_properties = { + **self.snapshot_properties, + "manifests-created": str(self._created_count), + "manifests-kept": str(self._kept_count), + "manifests-replaced": str(self._rewritten_count), + "entries-processed": str(self._entries_processed), + } + self._computed_manifests = new_manifests + kept_manifests + return self._computed_manifests + + def _refresh_for_retry(self) -> None: + """Reset state for a retry attempt, discarding the plan built from the stale branch head.""" + super()._refresh_for_retry() + # The plan and its counters depend on the branch head, which changes on retry. Keeping them + # would rebuild the replacement snapshot from the manifests of the superseded snapshot, + # dropping whatever was committed concurrently. + self._computed_manifests = None + self._rewritten_count = 0 + self._created_count = 0 + self._kept_count = 0 + self._entries_processed = 0 + + def _commit(self) -> UpdatesAndRequirements: + self._existing_manifests() + if self._created_count == 0: + # nothing was merged; committing would only produce a pointless replace snapshot + return (), () + return super()._commit() + + def rewrites_needed(self) -> bool: + """Return whether committing would rewrite any manifest. + + This mirrors the grouping that `_existing_manifests` performs, so a snapshot whose data + manifests all end up kept as-is reports False. A predicate that matches only manifests + holding no live entries still reports True, because deciding that requires reading them. + """ + snapshot = self._transaction.table_metadata.snapshot_by_name(self._target_branch or MAIN_BRANCH) + if snapshot is None: + return False + + data_manifests_by_spec: defaultdict[int, list[ManifestFile]] = defaultdict(list) + for manifest in snapshot.manifests(self._io): + if manifest.content == ManifestContent.DATA and ( + self._manifest_predicate is None or self._manifest_predicate(manifest) + ): + data_manifests_by_spec[manifest.partition_spec_id].append(manifest) + + return any( + len(group) > 1 or self._manifest_predicate is not None + for manifests in data_manifests_by_spec.values() + for group in self._group_by_target_size(manifests) + ) diff --git a/tests/integration/test_rewrite_manifests_interop.py b/tests/integration/test_rewrite_manifests_interop.py new file mode 100644 index 0000000000..bb44c3bfd8 --- /dev/null +++ b/tests/integration/test_rewrite_manifests_interop.py @@ -0,0 +1,74 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +# pylint:disable=redefined-outer-name +from typing import TYPE_CHECKING + +import pyarrow as pa +import pytest + +from pyiceberg.catalog import Catalog +from pyiceberg.exceptions import NoSuchTableError +from pyiceberg.manifest import ManifestContent + +if TYPE_CHECKING: + from pyspark.sql import SparkSession + + +@pytest.mark.integration +def test_spark_reads_table_after_rewrite_manifests(session_catalog: Catalog, spark: "SparkSession") -> None: + identifier = "default.test_rewrite_manifests_interop" + try: + session_catalog.drop_table(identifier) + except NoSuchTableError: + pass + + table = session_catalog.create_table(identifier, schema=pa.schema([pa.field("id", pa.int64())])) + for i in range(3): + table.append(pa.table({"id": pa.array([i * 3 + 1, i * 3 + 2, i * 3 + 3], type=pa.int64())})) + + table = session_catalog.load_table(identifier) + snapshot = table.current_snapshot() + assert snapshot is not None + assert len([m for m in snapshot.manifests(table.io) if m.content == ManifestContent.DATA]) == 3 + + table.maintenance.rewrite_manifests().commit() + + table = session_catalog.load_table(identifier) + snapshot = table.current_snapshot() + assert snapshot is not None + assert len([m for m in snapshot.manifests(table.io) if m.content == ManifestContent.DATA]) == 1 + + # Spark must read the rewritten table with the same data + spark_rows = spark.table(f"{identifier}").collect() + assert sorted(row.id for row in spark_rows) == list(range(1, 10)) + + # Spark must see the replace snapshot and the preserved data files + snapshots = spark.sql(f"SELECT operation FROM {identifier}.snapshots ORDER BY committed_at").collect() + assert [row.operation for row in snapshots] == ["append", "append", "append", "replace"] + + files = spark.sql(f"SELECT file_path FROM {identifier}.files").collect() + assert len(files) == 3 + + # Spark sees the same manifest consolidation + manifests = spark.sql(f"SELECT path FROM {identifier}.manifests").collect() + assert len(manifests) == 1 + + # time travel to the pre-rewrite snapshot still works from Spark + previous_snapshot_id = snapshot.parent_snapshot_id + assert previous_snapshot_id is not None + previous_rows = spark.sql(f"SELECT id FROM {identifier} VERSION AS OF {previous_snapshot_id}").collect() + assert sorted(row.id for row in previous_rows) == list(range(1, 10)) diff --git a/tests/table/test_commit_retry.py b/tests/table/test_commit_retry.py index ce5dca96aa..98d15d4f9b 100644 --- a/tests/table/test_commit_retry.py +++ b/tests/table/test_commit_retry.py @@ -302,6 +302,40 @@ def test_delete_files_refresh_clears_compute_deletes_cache(catalog: Catalog) -> assert "_compute_deletes" not in producer.__dict__ +def test_rewrite_manifests_refresh_clears_computed_manifests(catalog: Catalog) -> None: + """Verify that _refresh_for_retry discards the plan RewriteManifests computed for the old head.""" + catalog.create_namespace("default") + schema = _test_schema() + table = catalog.create_table("default.rewrite_cache_test", schema=schema) + + import pyarrow as pa + + table.append(pa.table({"x": [1, 2, 3]})) + table.append(pa.table({"x": [4, 5, 6]})) + table = catalog.load_table("default.rewrite_cache_test") + + from pyiceberg.table.update.snapshot import RewriteManifests + + tx = Transaction(table, autocommit=False) + producer = RewriteManifests(transaction=tx, io=table.io) + + # Plan against the current head, populating the cache and the counters + planned = producer._existing_manifests() + + assert planned + assert producer._created_count == 1 + assert producer._rewritten_count == 2 + + producer._refresh_for_retry() + + # Reusing either would rebuild the replacement snapshot from the superseded head + assert producer._computed_manifests is None + assert producer._created_count == 0 + assert producer._rewritten_count == 0 + assert producer._kept_count == 0 + assert producer._entries_processed == 0 + + def test_concurrent_overwrite_overwrite_raises_validation_exception(catalog: Catalog) -> None: """Concurrent overwrites on the same data should fail with ValidationException.""" catalog.create_namespace("default") diff --git a/tests/table/test_rewrite_manifests.py b/tests/table/test_rewrite_manifests.py new file mode 100644 index 0000000000..8c73ad2f08 --- /dev/null +++ b/tests/table/test_rewrite_manifests.py @@ -0,0 +1,305 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +from pathlib import Path + +import pyarrow as pa +import pytest + +from pyiceberg.catalog import Catalog +from pyiceberg.catalog.memory import InMemoryCatalog +from pyiceberg.manifest import ManifestContent, ManifestFile +from pyiceberg.table import Table +from pyiceberg.table.snapshots import Operation + + +@pytest.fixture +def catalog(tmp_path: Path) -> Catalog: + catalog = InMemoryCatalog("test.rewrite_manifests", warehouse=f"file://{tmp_path}") + catalog.create_namespace("default") + return catalog + + +def _arrow_table(offset: int = 0) -> pa.Table: + return pa.table({"id": pa.array([offset + 1, offset + 2, offset + 3], type=pa.int64())}) + + +def _create_table_with_appends(catalog: Catalog, appends: int = 3) -> Table: + table = catalog.create_table("default.test_rewrite", schema=pa.schema([pa.field("id", pa.int64())])) + for i in range(appends): + table.append(_arrow_table(offset=i * 3)) + return table + + +def _data_manifests(table: Table) -> list[ManifestFile]: + snapshot = table.current_snapshot() + assert snapshot is not None + return [m for m in snapshot.manifests(table.io) if m.content == ManifestContent.DATA] + + +def test_rewrite_manifests_merges_data_manifests(catalog: Catalog) -> None: + table = _create_table_with_appends(catalog, appends=3) + assert len(_data_manifests(table)) == 3 + rows_before = table.scan().to_arrow().sort_by("id") + + table.maintenance.rewrite_manifests().commit() + + table = catalog.load_table("default.test_rewrite") + manifests = _data_manifests(table) + assert len(manifests) == 1 + # entries are rewritten as EXISTING + assert manifests[0].existing_files_count == 3 + assert manifests[0].added_files_count == 0 + + # data is unchanged + assert table.scan().to_arrow().sort_by("id") == rows_before + + snapshot = table.current_snapshot() + assert snapshot is not None + assert snapshot.summary is not None + assert snapshot.summary.operation == Operation.REPLACE + assert snapshot.summary["manifests-created"] == "1" + assert snapshot.summary["manifests-replaced"] == "3" + assert snapshot.summary["entries-processed"] == "3" + # totals carry over unchanged + assert snapshot.summary["total-data-files"] == "3" + assert snapshot.summary["total-records"] == "9" + + +def _sequence_numbers_by_file(table: Table) -> dict[str, int]: + result: dict[str, int] = {} + for manifest in _data_manifests(table): + for entry in manifest.fetch_manifest_entry(table.io, discard_deleted=True): + assert entry.sequence_number is not None + result[entry.data_file.file_path] = entry.sequence_number + return result + + +def test_rewrite_manifests_preserves_sequence_numbers(catalog: Catalog) -> None: + table = _create_table_with_appends(catalog, appends=3) + entries_before = _sequence_numbers_by_file(table) + + table.maintenance.rewrite_manifests().commit() + + table = catalog.load_table("default.test_rewrite") + entries_after = _sequence_numbers_by_file(table) + assert entries_after == entries_before + # the merged manifest keeps the min sequence number of its entries + assert _data_manifests(table)[0].min_sequence_number == min(entries_before.values()) + + +def test_rewrite_manifests_single_manifest_is_noop(catalog: Catalog) -> None: + table = _create_table_with_appends(catalog, appends=1) + snapshot_before = table.current_snapshot() + assert snapshot_before is not None + manifest_path_before = _data_manifests(table)[0].manifest_path + + table.maintenance.rewrite_manifests().commit() + + table = catalog.load_table("default.test_rewrite") + # nothing to merge: no new snapshot is committed and the manifest is untouched + snapshot = table.current_snapshot() + assert snapshot is not None + assert snapshot.snapshot_id == snapshot_before.snapshot_id + assert _data_manifests(table)[0].manifest_path == manifest_path_before + + +def test_rewrite_manifests_respects_target_size(catalog: Catalog) -> None: + table = _create_table_with_appends(catalog, appends=4) + max_manifest_length = max(m.manifest_length for m in _data_manifests(table)) + + # allow two source manifests per group (2x fits, 3x exceeds), robust to small size variations + with table.transaction() as tx: + tx.set_properties({"commit.manifest.target-size-bytes": str(int(max_manifest_length * 2.5))}) + + table = catalog.load_table("default.test_rewrite") + table.maintenance.rewrite_manifests().commit() + + table = catalog.load_table("default.test_rewrite") + manifests = _data_manifests(table) + assert len(manifests) == 2 + assert all(m.existing_files_count == 2 for m in manifests) + + +def test_rewrites_needed(catalog: Catalog) -> None: + table = _create_table_with_appends(catalog, appends=1) + assert table.maintenance.rewrite_manifests().rewrites_needed() is False + + table.append(_arrow_table(offset=3)) + table = catalog.load_table("default.test_rewrite") + assert table.maintenance.rewrite_manifests().rewrites_needed() is True + + +def test_rewrite_manifests_with_predicate_selective(catalog: Catalog) -> None: + table = _create_table_with_appends(catalog, appends=3) + manifests_before = _data_manifests(table) + assert len(manifests_before) == 3 + + # 1. Extract all manifest paths and underlying data file paths before rewrite + paths_before = [m.manifest_path for m in manifests_before] + target_path = paths_before[0] + kept_paths_before = set(paths_before[1:]) + rows_before = table.scan().to_arrow().sort_by("id") + data_files_before = [ + entry.data_file.file_path for m in manifests_before for entry in m.fetch_manifest_entry(table.io, discard_deleted=True) + ] + + # 2. Execute selective rewrite + table.maintenance.rewrite_manifests().rewrite_if(lambda m: m.manifest_path == target_path).commit() + + # 3. Reload table and extract new paths for comprehensive verification + table = catalog.load_table("default.test_rewrite") + manifests_after = _data_manifests(table) + paths_after = {m.manifest_path for m in manifests_after} + + assert len(manifests_after) == 3 + # Manifest paths not rewritten remain unchanged + assert kept_paths_before.issubset(paths_after) + # Old path of rewritten target manifest is gone + assert target_path not in paths_after + # Exactly one new manifest was created + new_manifest_paths = paths_after - kept_paths_before + assert len(new_manifest_paths) == 1 + + # Verify table data and underlying data files are completely preserved + assert table.scan().to_arrow().sort_by("id") == rows_before + data_files_after = [ + entry.data_file.file_path for m in manifests_after for entry in m.fetch_manifest_entry(table.io, discard_deleted=True) + ] + assert set(data_files_before) == set(data_files_after) + + snapshot = table.current_snapshot() + assert snapshot is not None + assert snapshot.summary is not None + assert snapshot.summary["manifests-created"] == "1" + assert snapshot.summary["manifests-replaced"] == "1" + assert snapshot.summary["manifests-kept"] == "2" + assert snapshot.summary["entries-processed"] == "1" + + +def test_rewrite_manifests_single_manifest_with_predicate(catalog: Catalog) -> None: + table = _create_table_with_appends(catalog, appends=1) + snapshot_before = table.current_snapshot() + assert snapshot_before is not None + manifest_path_before = _data_manifests(table)[0].manifest_path + + # with predicate, even single manifest should be rewritten + table.maintenance.rewrite_manifests().rewrite_if(lambda m: True).commit() + + table = catalog.load_table("default.test_rewrite") + manifests_after = _data_manifests(table) + assert len(manifests_after) == 1 + assert manifests_after[0].manifest_path != manifest_path_before + assert manifests_after[0].existing_files_count == 1 + assert manifests_after[0].added_files_count == 0 + + snapshot = table.current_snapshot() + assert snapshot is not None + assert snapshot.snapshot_id != snapshot_before.snapshot_id + assert snapshot.summary is not None + assert snapshot.summary["manifests-created"] == "1" + assert snapshot.summary["manifests-replaced"] == "1" + assert snapshot.summary["manifests-kept"] == "0" + assert snapshot.summary["entries-processed"] == "1" + + +def test_rewrites_needed_with_predicate(catalog: Catalog) -> None: + table = _create_table_with_appends(catalog, appends=1) + # single manifest without predicate: False + assert table.maintenance.rewrite_manifests().rewrites_needed() is False + # single manifest with matching predicate: True + assert table.maintenance.rewrite_manifests().rewrite_if(lambda m: True).rewrites_needed() is True + # single manifest with non-matching predicate: False + assert table.maintenance.rewrite_manifests().rewrite_if(lambda m: False).rewrites_needed() is False + + +def test_rewrite_manifests_with_fully_deleted_manifest(catalog: Catalog) -> None: + table = catalog.create_table("default.test_fully_deleted", schema=pa.schema([pa.field("id", pa.int64())])) + table.append(_arrow_table(offset=0)) # manifest 1: id 1, 2, 3 + table.append(_arrow_table(offset=3)) # manifest 2: id 4, 5, 6 + table.delete("id <= 3") # deletes id 1, 2, 3 + + snapshot_before = table.current_snapshot() + manifests_before = _data_manifests(table) + assert len(manifests_before) == 2 + paths_before = [m.manifest_path for m in manifests_before] + + # Target rewrite for manifest containing only deleted entries + table.maintenance.rewrite_manifests().rewrite_if(lambda m: (m.deleted_files_count or 0) > 0).commit() + + table = catalog.load_table("default.test_fully_deleted") + assert table.current_snapshot() == snapshot_before + + # Verify both manifest paths remain unchanged + manifests_after = _data_manifests(table) + assert len(manifests_after) == 2 + assert [m.manifest_path for m in manifests_after] == paths_before + + # Verify the second manifest with live data (ids 4, 5, 6) is intact and readable + assert table.scan().to_arrow().sort_by("id") == _arrow_table(offset=3) + live_entries_manifest2 = manifests_after[1].fetch_manifest_entry(table.io, discard_deleted=True) + assert len(live_entries_manifest2) == 1 + assert live_entries_manifest2[0].data_file.record_count == 3 + + +def test_rewrite_manifests_predicate_matching_nothing(catalog: Catalog) -> None: + table = _create_table_with_appends(catalog, appends=2) + snapshot_before = table.current_snapshot() + table.maintenance.rewrite_manifests().rewrite_if(lambda m: False).commit() + table = catalog.load_table("default.test_rewrite") + assert table.current_snapshot() == snapshot_before + + +def test_rewrite_manifests_merges_live_and_fully_deleted_manifests(catalog: Catalog) -> None: + table = catalog.create_table("default.test_merge_deleted", schema=pa.schema([pa.field("id", pa.int64())])) + table.append(_arrow_table(offset=0)) # manifest 1: id 1, 2, 3 + table.append(_arrow_table(offset=3)) # manifest 2: id 4, 5, 6 + table.delete("id <= 3") # fully deletes manifest 1 + + manifests_before = _data_manifests(table) + assert len(manifests_before) == 2 + + # Plain rewrite_manifests() without predicate should merge live and fully-deleted manifests into 1 + assert table.maintenance.rewrite_manifests().rewrites_needed() is True + table.maintenance.rewrite_manifests().commit() + + table = catalog.load_table("default.test_merge_deleted") + manifests_after = _data_manifests(table) + assert len(manifests_after) == 1 + assert manifests_after[0].existing_files_count == 1 + assert manifests_after[0].added_files_count == 0 + + # Verify table data is fully preserved and matches remaining live records + assert table.scan().to_arrow().sort_by("id") == _arrow_table(offset=3) + + +def test_rewrites_needed_is_false_when_every_group_is_a_single_manifest(catalog: Catalog) -> None: + table = _create_table_with_appends(catalog, appends=3) + # A tiny target size puts every manifest in its own group, so nothing can be merged + with table.transaction() as tx: + tx.set_properties({"commit.manifest.target-size-bytes": "1"}) + + table = catalog.load_table("default.test_rewrite") + assert len(_data_manifests(table)) == 3 + + snapshot_before = table.current_snapshot() + assert table.maintenance.rewrite_manifests().rewrites_needed() is False + + table.maintenance.rewrite_manifests().commit() + + table = catalog.load_table("default.test_rewrite") + assert table.current_snapshot() == snapshot_before + assert len(_data_manifests(table)) == 3 diff --git a/tests/table/test_snapshots.py b/tests/table/test_snapshots.py index 5f1680ed59..9e23669cd0 100644 --- a/tests/table/test_snapshots.py +++ b/tests/table/test_snapshots.py @@ -399,10 +399,23 @@ def test_merge_snapshot_summaries_overwrite_summary() -> None: assert actual.additional_properties == expected -def test_invalid_operation() -> None: - with pytest.raises(ValueError) as e: - update_snapshot_summaries(summary=Summary(Operation.REPLACE)) - assert "Operation not implemented: Operation.REPLACE" in str(e.value) +def test_replace_operation_carries_totals() -> None: + actual = update_snapshot_summaries( + summary=Summary(Operation.REPLACE), + previous_summary={ + "total-data-files": "3", + "total-delete-files": "0", + "total-records": "9", + "total-files-size": "1234", + "total-position-deletes": "0", + "total-equality-deletes": "0", + }, + ) + + # a replace operation does not change any of the totals + assert actual["total-data-files"] == "3" + assert actual["total-records"] == "9" + assert actual["total-files-size"] == "1234" def test_invalid_type() -> None: