diff --git a/pyiceberg/manifest.py b/pyiceberg/manifest.py index 4a61cc5de2..7037c69437 100644 --- a/pyiceberg/manifest.py +++ b/pyiceberg/manifest.py @@ -19,7 +19,7 @@ import math import threading from abc import ABC, abstractmethod -from collections.abc import Iterator +from collections.abc import Callable, Iterator from copy import copy from enum import Enum from types import TracebackType @@ -879,13 +879,19 @@ def has_added_files(self) -> bool: def has_existing_files(self) -> bool: return self.existing_files_count is None or self.existing_files_count > 0 - def fetch_manifest_entry(self, io: FileIO, discard_deleted: bool = True) -> list[ManifestEntry]: + def fetch_manifest_entry( + self, + io: FileIO, + discard_deleted: bool = True, + entry_filter: Callable[[ManifestEntry], bool] | None = None, + ) -> list[ManifestEntry]: """ Read the manifest entries from the manifest file. Args: io: The FileIO to fetch the file. discard_deleted: Filter on live entries. + entry_filter: Optional predicate to filter manifest entries. Returns: An Iterator of manifest entries. @@ -897,11 +903,17 @@ def fetch_manifest_entry(self, io: FileIO, discard_deleted: bool = True) -> list read_types={-1: ManifestEntry, 2: DataFile}, read_enums={0: ManifestEntryStatus, 101: FileFormat, 134: DataFileContent}, ) as reader: - return [ + result = [] + + for entry in reader: + if discard_deleted and entry.status == ManifestEntryStatus.DELETED: + continue _inherit_from_manifest(entry, self) - for entry in reader - if not discard_deleted or entry.status != ManifestEntryStatus.DELETED - ] + + if entry_filter is None or entry_filter(entry): + result.append(entry) + + return result def __eq__(self, other: Any) -> bool: """Return the equality of two instances of the ManifestFile class.""" diff --git a/pyiceberg/table/__init__.py b/pyiceberg/table/__init__.py index 9624eac981..fca718f5ec 100644 --- a/pyiceberg/table/__init__.py +++ b/pyiceberg/table/__init__.py @@ -2325,11 +2325,9 @@ def _open_manifest( Returns: A list of ManifestEntry that matches the provided filters. """ - return [ - manifest_entry - for manifest_entry in manifest.fetch_manifest_entry(io, discard_deleted=True) - if partition_filter(manifest_entry.data_file) and metrics_evaluator(manifest_entry.data_file) - ] + return manifest.fetch_manifest_entry( + io, discard_deleted=True, entry_filter=lambda e: partition_filter(e.data_file) and metrics_evaluator(e.data_file) + ) def _min_sequence_number(manifests: list[ManifestFile]) -> int: diff --git a/tests/utils/test_manifest.py b/tests/utils/test_manifest.py index 6811ab2942..ff91224e4d 100644 --- a/tests/utils/test_manifest.py +++ b/tests/utils/test_manifest.py @@ -199,6 +199,34 @@ def test_read_manifest_entry(generated_manifest_entry_file: str) -> None: assert data_file.sort_order_id == 0 +def test_fetch_manifest_entry_with_filter(generated_manifest_entry_file: str) -> None: + manifest = ManifestFile.from_args( + manifest_path=generated_manifest_entry_file, + manifest_length=0, + partition_spec_id=0, + added_snapshot_id=0, + sequence_number=0, + partitions=[], + ) + + all_entries = manifest.fetch_manifest_entry(PyArrowFileIO()) + assert len(all_entries) == 2 + + # Entry 1 has tpep_pickup_day=1925 & entry 2 has tpep_pickup_day=None + matched = manifest.fetch_manifest_entry( + PyArrowFileIO(), + entry_filter=lambda e: e.data_file.partition[1] == 1925, + ) + assert len(matched) == 1 + assert matched[0].data_file.record_count == 19513 + + no_match = manifest.fetch_manifest_entry( + PyArrowFileIO(), + entry_filter=lambda e: e.data_file.partition[1] == 9999, + ) + assert len(no_match) == 0 + + def test_read_manifest_entry_v3_fields(tmp_path: Path) -> None: io = PyArrowFileIO()