Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 22 additions & 5 deletions pyiceberg/io/pyarrow.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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}")

Expand Down Expand Up @@ -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(
Expand Down
24 changes: 20 additions & 4 deletions pyiceberg/table/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
Expand All @@ -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,
Expand All @@ -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
Expand Down
4 changes: 3 additions & 1 deletion pyiceberg/table/delete_file_index.py
Original file line number Diff line number Diff line change
Expand Up @@ -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())
Expand Down
19 changes: 19 additions & 0 deletions pyiceberg/table/deletion_vector.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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:
Expand Down Expand Up @@ -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]


Expand Down
41 changes: 41 additions & 0 deletions tests/catalog/test_scan_planning_models.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
78 changes: 78 additions & 0 deletions tests/io/test_pyarrow.py
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@
_ConvertToArrowSchema,
_determine_partitions,
_primitive_to_physical,
_read_all_delete_files,
_read_deletes,
_task_to_record_batches,
_to_requested_schema,
Expand Down Expand Up @@ -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: <length(4 BE)> <DV_MAGIC + 64-bit RoaringBitmapArray in the portable format shared by
Delta and Iceberg> <CRC-32(4 BE)>. 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
Expand Down
32 changes: 32 additions & 0 deletions tests/table/test_delete_file_index.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
33 changes: 32 additions & 1 deletion tests/table/test_deletion_vector.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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])]
Loading