Skip to content

Commit 5b6d113

Browse files
committed
Support range-based reads for deletion vectors
# Conflicts: # pyiceberg/table/deletion_vector.py
1 parent 9d36e23 commit 5b6d113

7 files changed

Lines changed: 332 additions & 20 deletions

File tree

pyiceberg/io/pyarrow.py

Lines changed: 2 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_vectors
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,7 @@ 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+
return {dv.referenced_data_file: dv.to_vector() for dv in read_deletion_vectors(io, data_file)}
11461142
else:
11471143
raise ValueError(f"Delete file format not supported: {data_file.file_format}")
11481144

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: 100 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,9 @@
1515
# specific language governing permissions and limitations
1616
# under the License.
1717
import math
18-
from typing import TYPE_CHECKING
18+
import struct
19+
import zlib
20+
from typing import TYPE_CHECKING, cast
1921

2022
from pyroaring import BitMap, FrozenBitMap
2123

@@ -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,17 +89,99 @@ def to_vector(self) -> "pa.ChunkedArray":
7789
return self._bitmaps_to_chunked_array(self._bitmaps)
7890

7991

80-
def _extract_vector_payload(blob_payload: bytes) -> bytes:
81-
"""Strip deletion-vector-v1 blob framing: length(4 big-endian) + DV magic(4) ... CRC(4 big-endian)."""
82-
length_prefix = int.from_bytes(blob_payload[0:4], "big")
83-
return blob_payload[8 : 4 + length_prefix]
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(dv: "DataFile") -> None:
127+
content_offset = dv.content_offset
128+
content_size_in_bytes = dv.content_size_in_bytes
129+
referenced_data_file = dv.referenced_data_file
130+
131+
if content_offset is None:
132+
raise ValueError(f"Invalid deletion vector, content offset is missing: {dv.file_path}")
133+
if content_size_in_bytes is None:
134+
raise ValueError(f"Invalid deletion vector, content size is missing: {dv.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: {dv.file_path}")
143+
144+
145+
def has_deletion_vector_content_reference(dv: "DataFile") -> bool:
146+
"""Return whether a deletion vector is described by manifest content-range metadata."""
147+
return dv.content_offset is not None or dv.content_size_in_bytes is not None or dv.referenced_data_file is not None
148+
149+
150+
def _read_deletion_vector(io: "FileIO", dv: "DataFile") -> DeletionVector:
151+
_validate_deletion_vector_content(dv)
152+
153+
content_offset = cast(int, dv.content_offset)
154+
content_size_in_bytes = cast(int, dv.content_size_in_bytes)
155+
referenced_data_file = cast(str, dv.referenced_data_file)
156+
157+
with io.new_input(dv.file_path).open() as fi:
158+
fi.seek(content_offset)
159+
payload = fi.read(content_size_in_bytes)
160+
161+
if len(payload) != content_size_in_bytes:
162+
raise ValueError(f"Could not read deletion vector, expected {content_size_in_bytes} bytes, got {len(payload)}")
163+
164+
return DeletionVector(
165+
referenced_data_file=referenced_data_file,
166+
bitmaps=_deserialize_dv_blob(payload, dv.record_count),
167+
)
168+
169+
170+
def read_deletion_vectors(io: "FileIO", dv: "DataFile") -> list[DeletionVector]:
171+
"""Read deletion vectors from a delete file or its manifest content range."""
172+
if has_deletion_vector_content_reference(dv):
173+
return [_read_deletion_vector(io, dv)]
174+
175+
with io.new_input(dv.file_path).open() as fi:
176+
return deletion_vectors_from_puffin_file(PuffinFile(fi.read()))
84177

85178

86179
def deletion_vectors_from_puffin_file(puffin_file: PuffinFile) -> list[DeletionVector]:
180+
"""Read all deletion vectors stored in a Puffin file."""
87181
return [
88182
DeletionVector(
89183
referenced_data_file=blob.properties[PROPERTY_REFERENCED_DATA_FILE],
90-
bitmaps=DeletionVector._deserialize_bitmap(_extract_vector_payload(puffin_file.get_blob_payload(blob))),
184+
bitmaps=_deserialize_dv_blob(puffin_file.get_blob_payload(blob)),
91185
)
92186
for blob in puffin_file.footer.blobs
93187
]

tests/avro/test_file.py

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@
5151
LongType,
5252
NestedField,
5353
StringType,
54+
StructType,
5455
TimestampType,
5556
TimestamptzType,
5657
TimeType,
@@ -89,7 +90,16 @@ def test_missing_schema() -> None:
8990

9091
# helper function to serialize our objects to dicts to enable
9192
# direct comparison with the dicts returned by fastavro
92-
def todict(obj: Any) -> Any:
93+
def todict(obj: Any, struct: StructType | None = None) -> Any:
94+
if struct is not None and isinstance(obj, Record):
95+
return {
96+
field.name: todict(
97+
getattr(obj, field.name),
98+
field.field_type if isinstance(field.field_type, StructType) else None,
99+
)
100+
for field in struct.fields
101+
if hasattr(obj, field.name)
102+
}
93103
if isinstance(obj, dict):
94104
data = []
95105
for k, v in obj.items():
@@ -157,7 +167,7 @@ def test_write_manifest_entry_with_iceberg_read_with_fastavro_v1() -> None:
157167

158168
fa_entry = next(it)
159169

160-
v2_entry = todict(entry)
170+
v2_entry = todict(entry, MANIFEST_ENTRY_SCHEMAS[2].as_struct())
161171

162172
# These are not written in V1
163173
del v2_entry["sequence_number"]
@@ -222,7 +232,7 @@ def test_write_manifest_entry_with_iceberg_read_with_fastavro_v2() -> None:
222232

223233
fa_entry = next(it)
224234

225-
assert todict(entry) == fa_entry
235+
assert todict(entry, MANIFEST_ENTRY_SCHEMAS[2].as_struct()) == fa_entry
226236

227237

228238
@pytest.mark.parametrize("format_version", [1, 2])
@@ -260,7 +270,7 @@ def test_write_manifest_entry_with_fastavro_read_with_iceberg(format_version: Ta
260270
schema = AvroSchemaConversion().iceberg_to_avro(MANIFEST_ENTRY_SCHEMAS[format_version], schema_name="manifest_entry")
261271

262272
with open(tmp_avro_file, "wb") as out:
263-
writer(out, schema, [todict(entry)])
273+
writer(out, schema, [todict(entry, MANIFEST_ENTRY_SCHEMAS[format_version].as_struct())])
264274

265275
# Read as V2
266276
with avro.AvroFile[ManifestEntry](

tests/integration/test_rest_manifest.py

Lines changed: 18 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -28,17 +28,28 @@
2828
from pyiceberg.avro.codecs import AvroCompressionCodec
2929
from pyiceberg.catalog import Catalog, load_catalog
3030
from pyiceberg.io.pyarrow import PyArrowFileIO
31-
from pyiceberg.manifest import DataFile, write_manifest
31+
from pyiceberg.manifest import MANIFEST_ENTRY_SCHEMAS, DataFile, write_manifest
3232
from pyiceberg.table import Table
3333
from pyiceberg.typedef import Record
34+
from pyiceberg.types import StructType
3435
from pyiceberg.utils.lazydict import LazyDict
3536

3637

3738
# helper function to serialize our objects to dicts to enable
3839
# direct comparison with the dicts returned by fastavro
39-
def todict(obj: Any, spec_keys: list[str]) -> Any:
40+
def todict(obj: Any, spec_keys: list[str], struct: StructType | None = None) -> Any:
4041
if type(obj) is Record:
4142
return {key: obj[pos] for key, pos in zip(spec_keys, range(len(obj)), strict=True)}
43+
if struct is not None and isinstance(obj, Record):
44+
return {
45+
field.name: todict(
46+
getattr(obj, field.name),
47+
spec_keys,
48+
field.field_type if isinstance(field.field_type, StructType) else None,
49+
)
50+
for field in struct.fields
51+
if hasattr(obj, field.name)
52+
}
4253
if isinstance(obj, dict) or isinstance(obj, LazyDict):
4354
data = []
4455
for k, v in obj.items():
@@ -111,7 +122,11 @@ def test_write_sample_manifest(table_test_all_types: Table, compression: AvroCom
111122
)
112123
wrapped_entry_v2 = copy(entry)
113124
wrapped_entry_v2.data_file = wrapped_data_file_v2_debug
114-
wrapped_entry_v2_dict = todict(wrapped_entry_v2, [field.name for field in test_spec.fields])
125+
wrapped_entry_v2_dict = todict(
126+
wrapped_entry_v2,
127+
[field.name for field in test_spec.fields],
128+
MANIFEST_ENTRY_SCHEMAS[2].as_struct(),
129+
)
115130

116131
with TemporaryDirectory() as tmpdir:
117132
tmp_avro_file = tmpdir + "/test_write_manifest.avro"

0 commit comments

Comments
 (0)