diff --git a/pyiceberg/io/pyarrow.py b/pyiceberg/io/pyarrow.py index 2dcb8a5795..4a6fae6b31 100644 --- a/pyiceberg/io/pyarrow.py +++ b/pyiceberg/io/pyarrow.py @@ -149,11 +149,10 @@ visit_with_partner, ) from pyiceberg.table import DOWNCAST_NS_TIMESTAMP_TO_US_ON_WRITE, TableProperties -from pyiceberg.table.deletion_vector import deletion_vectors_from_puffin_file +from pyiceberg.table.deletion_vector import DeletionVector from pyiceberg.table.locations import load_location_provider from pyiceberg.table.metadata import TableMetadata from pyiceberg.table.name_mapping import NameMapping, apply_name_mapping -from pyiceberg.table.puffin import PuffinFile from pyiceberg.transforms import IdentityTransform, TruncateTransform from pyiceberg.typedef import EMPTY_DICT, Properties, Record, TableVersion from pyiceberg.types import ( @@ -1171,10 +1170,19 @@ def _read_deletes(io: FileIO, data_file: DataFile) -> dict[str, pa.ChunkedArray] for path in table.column("file_path").unique() } elif data_file.file_format == FileFormat.PUFFIN: + # Read the deletion vector blob directly using the offset and length recorded in the + # manifest, without parsing the Puffin footer. + referenced_data_file = data_file.referenced_data_file + offset = data_file.content_offset + length = data_file.content_size_in_bytes + if referenced_data_file is None or offset is None or length is None: + raise ValueError(f"Invalid deletion vector, missing referenced data file, offset, or length: {data_file.file_path}") + with io.new_input(data_file.file_path).open() as fi: - payload = fi.read() + fi.seek(offset) + blob = fi.read(length) - return {dv.referenced_data_file: dv.to_vector() for dv in deletion_vectors_from_puffin_file(PuffinFile(payload))} + return {referenced_data_file: DeletionVector.from_blob(blob, referenced_data_file).to_vector()} else: raise ValueError(f"Delete file format not supported: {data_file.file_format}") @@ -1743,7 +1751,16 @@ def _task_to_record_batches( def _read_all_delete_files(io: FileIO, tasks: Iterable[FileScanTask]) -> dict[str, list[ChunkedArray]]: deletes_per_file: dict[str, list[ChunkedArray]] = {} - unique_deletes = set(itertools.chain.from_iterable([task.delete_files for task in tasks])) + # Deduplicate on (file_path, content_offset): several deletion vectors can be packed into one + # Puffin file, sharing a file_path but differing by content_offset. DataFile equality keys only + # on file_path, so a plain set would collapse them and drop all but one vector. + unique_deletes = list( + { + (delete_file.file_path, delete_file.content_offset): delete_file + for task in tasks + for delete_file in task.delete_files + }.values() + ) if len(unique_deletes) > 0: executor = ExecutorFactory.get_or_create() deletes_per_files: Iterator[dict[str, ChunkedArray]] = executor.map( diff --git a/pyiceberg/table/__init__.py b/pyiceberg/table/__init__.py index 2c5c26800c..6854288263 100644 --- a/pyiceberg/table/__init__.py +++ b/pyiceberg/table/__init__.py @@ -2277,7 +2277,7 @@ def from_rest_response( delete_file = delete_files[idx] if isinstance(delete_file, RESTEqualityDeleteFile): raise NotImplementedError(f"PyIceberg does not yet support equality deletes: {delete_file.file_path}") - resolved_deletes.add(_rest_file_to_data_file(delete_file)) + resolved_deletes.add(_rest_file_to_data_file(delete_file, default_referenced_data_file=data_file.file_path)) return FileScanTask( data_file=data_file, @@ -2286,9 +2286,14 @@ def from_rest_response( ) -def _rest_file_to_data_file(rest_file: RESTContentFile) -> DataFile: - """Convert a REST content file to a manifest DataFile.""" - from pyiceberg.catalog.rest.scan_planning import RESTDataFile +def _rest_file_to_data_file(rest_file: RESTContentFile, default_referenced_data_file: str | None = None) -> DataFile: + """Convert a REST content file to a manifest DataFile. + + default_referenced_data_file supplies the referenced data file for a position delete when the + REST response omits it; the field is optional in the REST schema, but the offset-based deletion + vector read requires it. + """ + from pyiceberg.catalog.rest.scan_planning import RESTDataFile, RESTPositionDeleteFile if isinstance(rest_file, RESTDataFile): column_sizes = rest_file.column_sizes.to_dict() if rest_file.column_sizes else None @@ -2301,6 +2306,14 @@ def _rest_file_to_data_file(rest_file: RESTContentFile) -> DataFile: null_value_counts = None nan_value_counts = None + referenced_data_file = None + content_offset = None + content_size_in_bytes = None + if isinstance(rest_file, RESTPositionDeleteFile): + referenced_data_file = rest_file.referenced_data_file or default_referenced_data_file + content_offset = rest_file.content_offset + content_size_in_bytes = rest_file.content_size_in_bytes + data_file = DataFile.from_args( content=DataFileContent.from_rest_type(rest_file.content), file_path=rest_file.file_path, @@ -2314,6 +2327,9 @@ def _rest_file_to_data_file(rest_file: RESTContentFile) -> DataFile: nan_value_counts=nan_value_counts, split_offsets=rest_file.split_offsets, sort_order_id=rest_file.sort_order_id, + referenced_data_file=referenced_data_file, + content_offset=content_offset, + content_size_in_bytes=content_size_in_bytes, ) data_file.spec_id = rest_file.spec_id return data_file diff --git a/pyiceberg/table/delete_file_index.py b/pyiceberg/table/delete_file_index.py index 3f513aabe5..0348cbead9 100644 --- a/pyiceberg/table/delete_file_index.py +++ b/pyiceberg/table/delete_file_index.py @@ -115,7 +115,9 @@ def is_empty(self) -> bool: def add_delete_file(self, manifest_entry: ManifestEntry, partition_key: Record | None = None) -> None: delete_file = manifest_entry.data_file seq = manifest_entry.sequence_number or INITIAL_SEQUENCE_NUMBER - target_path = _referenced_data_file_path(delete_file) + # referenced_data_file identifies the target directly (always set for deletion vectors) and + # is authoritative; fall back to path bounds only when it is absent. + target_path = delete_file.referenced_data_file or _referenced_data_file_path(delete_file) if target_path: deletes = self._by_path.setdefault(target_path, PositionDeletes()) diff --git a/pyiceberg/table/deletion_vector.py b/pyiceberg/table/deletion_vector.py index 88fb3daf73..cbe0dab36d 100644 --- a/pyiceberg/table/deletion_vector.py +++ b/pyiceberg/table/deletion_vector.py @@ -15,6 +15,7 @@ # specific language governing permissions and limitations # under the License. import math +import zlib from typing import TYPE_CHECKING from pyroaring import BitMap, FrozenBitMap @@ -27,6 +28,7 @@ EMPTY_BITMAP = FrozenBitMap() MAX_JAVA_SIGNED = int(math.pow(2, 31)) - 1 PROPERTY_REFERENCED_DATA_FILE = "referenced-data-file" +DV_MAGIC = b"\xd1\xd3\x39\x64" class DeletionVector: @@ -76,10 +78,27 @@ def _bitmaps_to_chunked_array(bitmaps: list[BitMap]) -> "pa.ChunkedArray": def to_vector(self) -> "pa.ChunkedArray": return self._bitmaps_to_chunked_array(self._bitmaps) + @staticmethod + def from_blob(blob: bytes, referenced_data_file: str) -> "DeletionVector": + return DeletionVector( + referenced_data_file=referenced_data_file, + bitmaps=DeletionVector._deserialize_bitmap(_extract_vector_payload(blob)), + ) + def _extract_vector_payload(blob_payload: bytes) -> bytes: """Strip deletion-vector-v1 blob framing: length(4 big-endian) + DV magic(4) ... CRC(4 big-endian).""" length_prefix = int.from_bytes(blob_payload[0:4], "big") + magic = blob_payload[4:8] + if magic != DV_MAGIC: + raise ValueError(f"Invalid magic bytes for deletion vector: {magic!r}, expected {DV_MAGIC!r}") + + body = blob_payload[4 : 4 + length_prefix] # magic + serialized bitmap; the CRC covers this + stored_crc = int.from_bytes(blob_payload[4 + length_prefix : 8 + length_prefix], "big") + computed_crc = zlib.crc32(body) & 0xFFFFFFFF + if computed_crc != stored_crc: + raise ValueError(f"Invalid CRC for deletion vector: {computed_crc:#010x}, expected {stored_crc:#010x}") + return blob_payload[8 : 4 + length_prefix] diff --git a/tests/catalog/test_scan_planning_models.py b/tests/catalog/test_scan_planning_models.py index caf571d322..7844ac90a5 100644 --- a/tests/catalog/test_scan_planning_models.py +++ b/tests/catalog/test_scan_planning_models.py @@ -205,6 +205,47 @@ def test_equality_delete_file() -> None: assert equality_delete.equality_ids == [1, 2] +def test_rest_position_delete_file_to_data_file_propagates_deletion_vector_fields() -> None: + from pyiceberg.table import _rest_file_to_data_file + + rest_file = RESTPositionDeleteFile.model_validate( + { + **_rest_position_delete_file(file_path="s3://bucket/table/deletion_vector.puffin", file_format="puffin"), + "referenced-data-file": "s3://bucket/table/data/file.parquet", + } + ) + + data_file = _rest_file_to_data_file(rest_file) + + assert data_file.referenced_data_file == "s3://bucket/table/data/file.parquet" + assert data_file.content_offset == 100 + assert data_file.content_size_in_bytes == 200 + + +def test_from_rest_response_fills_missing_referenced_data_file() -> None: + from pyiceberg.table import FileScanTask + + # A valid REST position delete file may omit referenced-data-file; the containing task identifies + # the target data file, so the conversion must fill it in for the offset-based deletion vector read. + rest_task = RESTFileScanTask.model_validate( + { + "data-file": _rest_data_file(file_path="s3://bucket/table/data/file.parquet"), + "delete-file-references": [0], + } + ) + delete_file = TypeAdapter(RESTDeleteFile).validate_python( + _rest_position_delete_file(file_path="s3://bucket/table/deletion_vector.puffin", file_format="puffin") + ) + assert delete_file.referenced_data_file is None + + task = FileScanTask.from_rest_response(rest_task, [delete_file]) + + delete = next(iter(task.delete_files)) + assert delete.referenced_data_file == "s3://bucket/table/data/file.parquet" + assert delete.content_offset == 100 + assert delete.content_size_in_bytes == 200 + + def test_file_format_case_insensitive() -> None: for fmt in ["parquet", "PARQUET", "Parquet"]: data_file = _rest_data_file(file_format=fmt) diff --git a/tests/io/test_pyarrow.py b/tests/io/test_pyarrow.py index 4d5d4431cb..a3423be2d3 100644 --- a/tests/io/test_pyarrow.py +++ b/tests/io/test_pyarrow.py @@ -75,6 +75,7 @@ _ConvertToArrowSchema, _determine_partitions, _primitive_to_physical, + _read_all_delete_files, _read_deletes, _task_to_record_batches, _to_requested_schema, @@ -1840,6 +1841,83 @@ def test_read_deletes(deletes_file: str, request: pytest.FixtureRequest) -> None assert list(deletes.values())[0] == pa.chunked_array([[1, 3, 5]]) +def _delta_dv_entry(positions: list[int]) -> bytes: + """One deletion vector framed as Delta's DeletionVectorStore writes it in a .bin file. + + Layout: . The bitmap holds `positions` in a single bitmap keyed at 0. + """ + import zlib + + from pyroaring import BitMap + + from pyiceberg.table.deletion_vector import DV_MAGIC + + bitmap = (1).to_bytes(8, "little") + (0).to_bytes(4, "little") + BitMap(positions).serialize() + data = DV_MAGIC + bitmap + return len(data).to_bytes(4, "big") + data + (zlib.crc32(data) & 0xFFFFFFFF).to_bytes(4, "big") + + +def test_read_deletes_deletion_vector_in_delta_bin_file(tmp_path: Path) -> None: + entry = _delta_dv_entry([1, 3, 5]) + dv_path = f"{tmp_path}/deletion_vector.bin" + with open(dv_path, "wb") as f: + f.write(b"\x01" + entry) + + referenced_data_file = "s3://bucket/data.parquet" + data_file = DataFile.from_args( + file_path=dv_path, + file_format=FileFormat.PUFFIN, + content=DataFileContent.POSITION_DELETES, + referenced_data_file=referenced_data_file, + # Iceberg's content_offset points at the length prefix (past Delta's 1-byte version header), + # and content_size_in_bytes counts the length prefix and CRC that Delta's sizeInBytes omits. + content_offset=1, + content_size_in_bytes=len(entry), + ) + + deletes = _read_deletes(PyArrowFileIO(), data_file) + + assert deletes.keys() == {referenced_data_file} + assert deletes[referenced_data_file] == pa.chunked_array([[1, 3, 5]]) + + +def test_read_all_delete_files_reads_multiple_deletion_vectors_in_one_file(tmp_path: Path) -> None: + # Two deletion vectors packed into one file: same file_path, different content_offset. DataFile + # equality keys only on file_path, so deduplication must not collapse them into a single read. + entry_a = _delta_dv_entry([1, 3, 5]) + entry_b = _delta_dv_entry([2, 4]) + dv_path = f"{tmp_path}/deletion_vectors.bin" + with open(dv_path, "wb") as f: + f.write(b"\x01" + entry_a + entry_b) + + def _dv(referenced_data_file: str, offset: int, size: int) -> DataFile: + return DataFile.from_args( + file_path=dv_path, + file_format=FileFormat.PUFFIN, + content=DataFileContent.POSITION_DELETES, + referenced_data_file=referenced_data_file, + content_offset=offset, + content_size_in_bytes=size, + ) + + def _task(data_file_path: str, delete_file: DataFile) -> FileScanTask: + data_file = DataFile.from_args(file_path=data_file_path, file_format=FileFormat.PARQUET) + return FileScanTask(data_file=data_file, delete_files={delete_file}) + + deletes = _read_all_delete_files( + PyArrowFileIO(), + [ + _task("s3://bucket/a.parquet", _dv("s3://bucket/a.parquet", 1, len(entry_a))), + _task("s3://bucket/b.parquet", _dv("s3://bucket/b.parquet", 1 + len(entry_a), len(entry_b))), + ], + ) + + assert deletes.keys() == {"s3://bucket/a.parquet", "s3://bucket/b.parquet"} + assert deletes["s3://bucket/a.parquet"][0] == pa.chunked_array([[1, 3, 5]]) + assert deletes["s3://bucket/b.parquet"][0] == pa.chunked_array([[2, 4]]) + + def test_delete(deletes_file: str, request: pytest.FixtureRequest, table_schema_simple: Schema) -> None: # Determine file format from the file extension file_format = FileFormat.PARQUET if deletes_file.endswith(".parquet") else FileFormat.ORC diff --git a/tests/table/test_delete_file_index.py b/tests/table/test_delete_file_index.py index 09dd9ac81b..6d60dca83b 100644 --- a/tests/table/test_delete_file_index.py +++ b/tests/table/test_delete_file_index.py @@ -161,6 +161,38 @@ def test_dvs_treated_as_position_deletes() -> None: assert all(d.content == DataFileContent.POSITION_DELETES for d in result) +def test_deletion_vectors_matched_by_referenced_data_file() -> None: + # Two deletion vectors packed into one Puffin file (shared file_path), without path bounds, each + # referencing a different data file. They must be routed by referenced_data_file, not collapsed + # into one partition bucket and deduplicated away by their shared file_path. + index = DeleteFileIndex() + + def _dv(referenced_data_file: str, content_offset: int) -> ManifestEntry: + delete_file = DataFile.from_args( + content=DataFileContent.POSITION_DELETES, + file_path="s3://bucket/deletion_vectors.puffin", + file_format=FileFormat.PUFFIN, + partition=Record(), + record_count=10, + file_size_in_bytes=100, + referenced_data_file=referenced_data_file, + content_offset=content_offset, + content_size_in_bytes=40, + ) + return ManifestEntry.from_args(status=ManifestEntryStatus.ADDED, sequence_number=2, data_file=delete_file) + + index.add_delete_file(_dv("s3://bucket/a.parquet", 1)) + index.add_delete_file(_dv("s3://bucket/b.parquet", 41)) + + deletes_a = index.for_data_file(1, _create_data_file(file_path="s3://bucket/a.parquet")) + deletes_b = index.for_data_file(1, _create_data_file(file_path="s3://bucket/b.parquet")) + + assert {d.referenced_data_file for d in deletes_a} == {"s3://bucket/a.parquet"} + assert {d.content_offset for d in deletes_a} == {1} + assert {d.referenced_data_file for d in deletes_b} == {"s3://bucket/b.parquet"} + assert {d.content_offset for d in deletes_b} == {41} + + def test_cannot_add_after_indexing() -> None: group = PositionDeletes() group.add(_create_positional_delete(sequence_number=1).data_file, 1) diff --git a/tests/table/test_deletion_vector.py b/tests/table/test_deletion_vector.py index 788216f8b3..c7eeaa1916 100644 --- a/tests/table/test_deletion_vector.py +++ b/tests/table/test_deletion_vector.py @@ -14,12 +14,13 @@ # KIND, either express or implied. See the License for the # specific language governing permissions and limitations # under the License. +import zlib from os import path import pytest from pyroaring import BitMap -from pyiceberg.table.deletion_vector import DeletionVector +from pyiceberg.table.deletion_vector import DV_MAGIC, DeletionVector def _open_file(file: str) -> bytes: @@ -71,3 +72,33 @@ def test_map_high_vals() -> None: with pytest.raises(ValueError, match="Key 4022190063 is too large, max 2147483647 to maintain compatibility with Java impl"): _ = DeletionVector._deserialize_bitmap(puffin) + + +def test_from_blob_invalid_magic() -> None: + bitmap = _open_file("64map32bitvals.bin") + length_prefix = 4 + len(bitmap) + blob = length_prefix.to_bytes(4, "big") + b"\x00\x00\x00\x00" + bitmap + (0).to_bytes(4, "big") + + with pytest.raises(ValueError, match="Invalid magic bytes for deletion vector"): + _ = DeletionVector.from_blob(blob, "s3://bucket/data.parquet") + + +def test_from_blob_invalid_crc() -> None: + bitmap = _open_file("64map32bitvals.bin") + data = DV_MAGIC + bitmap + # A trailing CRC of zero does not match the CRC-32 of (magic + bitmap). + blob = len(data).to_bytes(4, "big") + data + (0).to_bytes(4, "big") + + with pytest.raises(ValueError, match="Invalid CRC for deletion vector"): + _ = DeletionVector.from_blob(blob, "s3://bucket/data.parquet") + + +def test_from_blob() -> None: + bitmap = _open_file("64map32bitvals.bin") + data = DV_MAGIC + bitmap + blob = len(data).to_bytes(4, "big") + data + (zlib.crc32(data) & 0xFFFFFFFF).to_bytes(4, "big") + + dv = DeletionVector.from_blob(blob, "s3://bucket/data.parquet") + + assert dv.referenced_data_file == "s3://bucket/data.parquet" + assert dv._bitmaps == [BitMap([0, 1, 2, 3, 4, 5, 6, 7, 8, 9])]