Skip to content

Commit 4df52c3

Browse files
committed
Support range-based reads for deletion vectors
1 parent d871cd2 commit 4df52c3

5 files changed

Lines changed: 228 additions & 7 deletions

File tree

‎pyiceberg/io/pyarrow.py‎

Lines changed: 3 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -144,11 +144,10 @@
144144
visit_with_partner,
145145
)
146146
from pyiceberg.table import DOWNCAST_NS_TIMESTAMP_TO_US_ON_WRITE, TableProperties
147-
from pyiceberg.table.deletion_vector import deletion_vectors_from_puffin_file
147+
from pyiceberg.table.deletion_vector import read_deletion_vector
148148
from pyiceberg.table.locations import load_location_provider
149149
from pyiceberg.table.metadata import TableMetadata
150150
from pyiceberg.table.name_mapping import NameMapping, apply_name_mapping
151-
from pyiceberg.table.puffin import PuffinFile
152151
from pyiceberg.transforms import IdentityTransform, TruncateTransform
153152
from pyiceberg.typedef import EMPTY_DICT, Properties, Record, TableVersion
154153
from pyiceberg.types import (
@@ -1139,10 +1138,8 @@ def _read_deletes(io: FileIO, data_file: DataFile) -> dict[str, pa.ChunkedArray]
11391138
for path in table.column("file_path").unique()
11401139
}
11411140
elif data_file.file_format == FileFormat.PUFFIN:
1142-
with io.new_input(data_file.file_path).open() as fi:
1143-
payload = fi.read()
1144-
1145-
return {dv.referenced_data_file: dv.to_vector() for dv in deletion_vectors_from_puffin_file(PuffinFile(payload))}
1141+
dv = read_deletion_vector(io, data_file)
1142+
return {dv.referenced_data_file: dv.to_vector()}
11461143
else:
11471144
raise ValueError(f"Delete file format not supported: {data_file.file_format}")
11481145

‎pyiceberg/manifest.py‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -532,6 +532,18 @@ def equality_ids(self) -> list[int] | None:
532532
def sort_order_id(self) -> int | None:
533533
return self._data[15]
534534

535+
@property
536+
def referenced_data_file(self) -> str | None:
537+
return self._data[17] if len(self._data) > 17 else None
538+
539+
@property
540+
def content_offset(self) -> int | None:
541+
return self._data[18] if len(self._data) > 18 else None
542+
543+
@property
544+
def content_size_in_bytes(self) -> int | None:
545+
return self._data[19] if len(self._data) > 19 else None
546+
535547
# Spec ID should not be stored in the file
536548
_spec_id: int
537549

‎pyiceberg/table/deletion_vector.py‎

Lines changed: 83 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,8 @@
1515
# specific language governing permissions and limitations
1616
# under the License.
1717
import math
18+
import struct
19+
import zlib
1820
from typing import TYPE_CHECKING
1921

2022
from pyroaring import BitMap, FrozenBitMap
@@ -24,9 +26,19 @@
2426
if TYPE_CHECKING:
2527
import pyarrow as pa
2628

29+
from pyiceberg.io import FileIO
30+
from pyiceberg.manifest import DataFile
31+
2732
EMPTY_BITMAP = FrozenBitMap()
2833
MAX_JAVA_SIGNED = int(math.pow(2, 31)) - 1
2934
PROPERTY_REFERENCED_DATA_FILE = "referenced-data-file"
35+
_MAX_DELETION_VECTOR_CONTENT_SIZE = 2**31 - 1
36+
_DV_BLOB_LENGTH = struct.Struct(">I")
37+
_DV_BLOB_MAGIC = struct.Struct("<I")
38+
_DV_BLOB_CRC = struct.Struct(">I")
39+
_DV_BLOB_MAGIC_NUMBER = 1681511377
40+
_ROARING_BITMAP_COUNT_SIZE_BYTES = 8
41+
_DV_BLOB_MIN_SIZE_BYTES = _DV_BLOB_LENGTH.size + _DV_BLOB_MAGIC.size + _ROARING_BITMAP_COUNT_SIZE_BYTES + _DV_BLOB_CRC.size
3042

3143

3244
class DeletionVector:
@@ -77,6 +89,77 @@ def to_vector(self) -> "pa.ChunkedArray":
7789
return self._bitmaps_to_chunked_array(self._bitmaps)
7890

7991

92+
def _deserialize_dv_blob(blob: bytes, record_count: int | None = None) -> list[BitMap]:
93+
# The DV blob encoding matches Iceberg Java's BitmapPositionDeleteIndex:
94+
# 4-byte big-endian bitmap-data length, 4-byte little-endian magic number,
95+
# portable Roaring bitmap data, and 4-byte big-endian CRC-32.
96+
if len(blob) < _DV_BLOB_MIN_SIZE_BYTES:
97+
raise ValueError(f"Invalid deletion vector blob length: {len(blob)}")
98+
99+
bitmap_data_length = _DV_BLOB_LENGTH.unpack_from(blob)[0]
100+
expected_bitmap_data_length = len(blob) - _DV_BLOB_LENGTH.size - _DV_BLOB_CRC.size
101+
if bitmap_data_length != expected_bitmap_data_length:
102+
raise ValueError(f"Invalid bitmap data length: {bitmap_data_length}, expected {expected_bitmap_data_length}")
103+
104+
bitmap_data_offset = _DV_BLOB_LENGTH.size
105+
crc_offset = bitmap_data_offset + bitmap_data_length
106+
bitmap_data = blob[bitmap_data_offset:crc_offset]
107+
108+
magic_number = _DV_BLOB_MAGIC.unpack_from(bitmap_data)[0]
109+
if magic_number != _DV_BLOB_MAGIC_NUMBER:
110+
raise ValueError(f"Invalid magic number: {magic_number}, expected {_DV_BLOB_MAGIC_NUMBER}")
111+
112+
checksum = zlib.crc32(bitmap_data) & 0xFFFFFFFF
113+
expected_checksum = _DV_BLOB_CRC.unpack_from(blob, crc_offset)[0]
114+
if checksum != expected_checksum:
115+
raise ValueError("Invalid CRC")
116+
117+
bitmaps = DeletionVector._deserialize_bitmap(bitmap_data[_DV_BLOB_MAGIC.size :])
118+
if record_count is not None:
119+
cardinality = sum(len(bitmap) for bitmap in bitmaps)
120+
if cardinality != record_count:
121+
raise ValueError(f"Invalid cardinality: {cardinality}, expected {record_count}")
122+
123+
return bitmaps
124+
125+
126+
def _validate_deletion_vector_content(data_file: "DataFile") -> tuple[int, int, str]:
127+
content_offset = data_file.content_offset
128+
content_size_in_bytes = data_file.content_size_in_bytes
129+
referenced_data_file = data_file.referenced_data_file
130+
131+
if content_offset is None:
132+
raise ValueError(f"Invalid deletion vector, content offset is missing: {data_file.file_path}")
133+
if content_size_in_bytes is None:
134+
raise ValueError(f"Invalid deletion vector, content size is missing: {data_file.file_path}")
135+
if content_offset < 0:
136+
raise ValueError(f"Invalid deletion vector, content offset cannot be negative: {content_offset}")
137+
if content_size_in_bytes < 0:
138+
raise ValueError(f"Invalid deletion vector, content size cannot be negative: {content_size_in_bytes}")
139+
if content_size_in_bytes > _MAX_DELETION_VECTOR_CONTENT_SIZE:
140+
raise ValueError(f"Cannot read deletion vector larger than 2GB: {content_size_in_bytes}")
141+
if referenced_data_file is None:
142+
raise ValueError(f"Invalid deletion vector, referenced data file is missing: {data_file.file_path}")
143+
144+
return content_offset, content_size_in_bytes, referenced_data_file
145+
146+
147+
def read_deletion_vector(io: "FileIO", data_file: "DataFile") -> DeletionVector:
148+
content_offset, content_size_in_bytes, referenced_data_file = _validate_deletion_vector_content(data_file)
149+
150+
with io.new_input(data_file.file_path).open() as fi:
151+
fi.seek(content_offset)
152+
payload = fi.read(content_size_in_bytes)
153+
154+
if len(payload) != content_size_in_bytes:
155+
raise ValueError(f"Could not read deletion vector, expected {content_size_in_bytes} bytes, got {len(payload)}")
156+
157+
return DeletionVector(
158+
referenced_data_file=referenced_data_file,
159+
bitmaps=_deserialize_dv_blob(payload, data_file.record_count),
160+
)
161+
162+
80163
def deletion_vectors_from_puffin_file(puffin_file: PuffinFile) -> list[DeletionVector]:
81164
return [
82165
DeletionVector(

‎tests/io/test_pyarrow.py‎

Lines changed: 86 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,11 +15,14 @@
1515
# specific language governing permissions and limitations
1616
# under the License.
1717
# pylint: disable=protected-access,unused-argument,redefined-outer-name
18+
import json
1819
import logging
1920
import os
21+
import struct
2022
import tempfile
2123
import uuid
2224
import warnings
25+
import zlib
2326
from collections.abc import Iterator
2427
from datetime import date, datetime, timezone
2528
from pathlib import Path
@@ -34,6 +37,7 @@
3437
import pytest
3538
from packaging import version
3639
from pyarrow.fs import AwsDefaultS3RetryStrategy, FileType, LocalFileSystem, S3FileSystem
40+
from pyroaring import BitMap
3741

3842
from pyiceberg.exceptions import ResolveError
3943
from pyiceberg.expressions import (
@@ -89,8 +93,14 @@
8993
from pyiceberg.partitioning import PartitionField, PartitionSpec
9094
from pyiceberg.schema import Schema, make_compatible_name, visit
9195
from pyiceberg.table import FileScanTask, TableProperties
96+
from pyiceberg.table.deletion_vector import (
97+
_DV_BLOB_MAGIC_NUMBER,
98+
PROPERTY_REFERENCED_DATA_FILE,
99+
deletion_vectors_from_puffin_file,
100+
)
92101
from pyiceberg.table.metadata import TableMetadataV2
93102
from pyiceberg.table.name_mapping import create_mapping_from_schema
103+
from pyiceberg.table.puffin import MAGIC_BYTES, PuffinFile
94104
from pyiceberg.transforms import HourTransform, IdentityTransform
95105
from pyiceberg.typedef import UTF8, Properties, Record, TableVersion
96106
from pyiceberg.types import (
@@ -1820,6 +1830,82 @@ def test_read_deletes(deletes_file: str, request: pytest.FixtureRequest) -> None
18201830
assert list(deletes.values())[0] == pa.chunked_array([[1, 3, 5]])
18211831

18221832

1833+
def _deletion_vector_bitmap_payload() -> bytes:
1834+
return (1).to_bytes(8, byteorder="little") + (0).to_bytes(4, byteorder="little") + BitMap([1, 3, 5]).serialize()
1835+
1836+
1837+
def _deletion_vector_blob(bitmap_payload: bytes) -> bytes:
1838+
bitmap_data = struct.pack("<I", _DV_BLOB_MAGIC_NUMBER) + bitmap_payload
1839+
return struct.pack(">I", len(bitmap_data)) + bitmap_data + struct.pack(">I", zlib.crc32(bitmap_data) & 0xFFFFFFFF)
1840+
1841+
1842+
def test_deletion_vectors_from_puffin_file(tmp_path: Path) -> None:
1843+
referenced_data_file = f"{tmp_path}/data.parquet"
1844+
bitmap_payload = _deletion_vector_bitmap_payload()
1845+
footer_payload = json.dumps(
1846+
{
1847+
"blobs": [
1848+
{
1849+
"type": "deletion-vector-v1",
1850+
"fields": [2147483546],
1851+
"snapshot-id": 1,
1852+
"sequence-number": 1,
1853+
"offset": 0,
1854+
"length": len(bitmap_payload),
1855+
"properties": {PROPERTY_REFERENCED_DATA_FILE: referenced_data_file},
1856+
}
1857+
],
1858+
"properties": {},
1859+
}
1860+
).encode()
1861+
puffin_payload = (
1862+
MAGIC_BYTES
1863+
+ b"\x00\x00\x00\x00"
1864+
+ bitmap_payload
1865+
+ footer_payload
1866+
+ len(footer_payload).to_bytes(4, byteorder="little")
1867+
+ b"\x00\x00\x00\x00"
1868+
+ MAGIC_BYTES
1869+
)
1870+
delete_file_path = f"{tmp_path}/deletes.puffin"
1871+
1872+
with open(delete_file_path, "wb") as f:
1873+
f.write(puffin_payload)
1874+
1875+
with open(delete_file_path, "rb") as f:
1876+
deletion_vectors = deletion_vectors_from_puffin_file(PuffinFile(f.read()))
1877+
1878+
assert {dv.referenced_data_file: dv.to_vector() for dv in deletion_vectors} == {
1879+
referenced_data_file: pa.chunked_array([[1, 3, 5]])
1880+
}
1881+
1882+
1883+
def test_read_deletion_vector_blob_from_content_range(tmp_path: Path) -> None:
1884+
referenced_data_file = f"{tmp_path}/data.parquet"
1885+
dv_blob = _deletion_vector_blob(_deletion_vector_bitmap_payload())
1886+
prefix = b"\x01not-a-puffin-file"
1887+
delete_file_path = f"{tmp_path}/deletes.bin"
1888+
1889+
with open(delete_file_path, "wb") as f:
1890+
f.write(prefix + dv_blob + b"trailing-bytes")
1891+
1892+
deletes = _read_deletes(
1893+
PyArrowFileIO(),
1894+
DataFile.from_args(
1895+
_table_format_version=3,
1896+
content=DataFileContent.POSITION_DELETES,
1897+
file_path=delete_file_path,
1898+
file_format=FileFormat.PUFFIN,
1899+
record_count=3,
1900+
referenced_data_file=referenced_data_file,
1901+
content_offset=len(prefix),
1902+
content_size_in_bytes=len(dv_blob),
1903+
),
1904+
)
1905+
1906+
assert deletes == {referenced_data_file: pa.chunked_array([[1, 3, 5]])}
1907+
1908+
18231909
def test_delete(deletes_file: str, request: pytest.FixtureRequest, table_schema_simple: Schema) -> None:
18241910
# Determine file format from the file extension
18251911
file_format = FileFormat.PARQUET if deletes_file.endswith(".parquet") else FileFormat.ORC

‎tests/table/test_deletion_vector.py‎

Lines changed: 44 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,12 +14,14 @@
1414
# KIND, either express or implied. See the License for the
1515
# specific language governing permissions and limitations
1616
# under the License.
17+
import struct
18+
import zlib
1719
from os import path
1820

1921
import pytest
2022
from pyroaring import BitMap
2123

22-
from pyiceberg.table.deletion_vector import DeletionVector
24+
from pyiceberg.table.deletion_vector import DeletionVector, _DV_BLOB_MAGIC_NUMBER, _deserialize_dv_blob
2325

2426

2527
def _open_file(file: str) -> bytes:
@@ -28,6 +30,47 @@ def _open_file(file: str) -> bytes:
2830
return f.read()
2931

3032

33+
def _dv_blob(bitmap_payload: bytes) -> bytes:
34+
bitmap_data = struct.pack("<I", _DV_BLOB_MAGIC_NUMBER) + bitmap_payload
35+
return struct.pack(">I", len(bitmap_data)) + bitmap_data + struct.pack(">I", zlib.crc32(bitmap_data) & 0xFFFFFFFF)
36+
37+
38+
def _bitmap_payload() -> bytes:
39+
return (1).to_bytes(8, byteorder="little") + (0).to_bytes(4, byteorder="little") + BitMap([1, 3, 5]).serialize()
40+
41+
42+
def test_deserialize_deletion_vector_blob() -> None:
43+
actual = _deserialize_dv_blob(_dv_blob(_bitmap_payload()), record_count=3)
44+
45+
assert actual == [BitMap([1, 3, 5])]
46+
47+
48+
def test_deserialize_deletion_vector_blob_invalid_length() -> None:
49+
with pytest.raises(ValueError, match="Invalid bitmap data length"):
50+
_deserialize_dv_blob(_dv_blob(_bitmap_payload())[:-1])
51+
52+
53+
def test_deserialize_deletion_vector_blob_invalid_magic() -> None:
54+
bitmap_data = struct.pack("<I", _DV_BLOB_MAGIC_NUMBER + 1) + _bitmap_payload()
55+
blob = struct.pack(">I", len(bitmap_data)) + bitmap_data + struct.pack(">I", zlib.crc32(bitmap_data) & 0xFFFFFFFF)
56+
57+
with pytest.raises(ValueError, match="Invalid magic number"):
58+
_deserialize_dv_blob(blob)
59+
60+
61+
def test_deserialize_deletion_vector_blob_invalid_crc() -> None:
62+
blob = bytearray(_dv_blob(_bitmap_payload()))
63+
blob[-1] ^= 1
64+
65+
with pytest.raises(ValueError, match="Invalid CRC"):
66+
_deserialize_dv_blob(bytes(blob))
67+
68+
69+
def test_deserialize_deletion_vector_blob_invalid_cardinality() -> None:
70+
with pytest.raises(ValueError, match="Invalid cardinality"):
71+
_deserialize_dv_blob(_dv_blob(_bitmap_payload()), record_count=4)
72+
73+
3174
def test_map_empty() -> None:
3275
puffin = _open_file("64mapempty.bin")
3376

0 commit comments

Comments
 (0)