Skip to content

Commit 98d6432

Browse files
committed
Read deletion vectors by content offset without parsing the Puffin footer
When a scan applies deletes, _read_deletes loads the deletion vector that applies to each data file. For Puffin deletion vectors it read the entire file into memory and parsed the footer to locate and deserialize every blob, then returned the one for the referenced data file. Instead, read only the referenced blob with a single ranged read using content_offset and content_size_in_bytes from the manifest, and take the referenced data file from the manifest as well, matching the Java and Rust readers. Validate the blob's DV_MAGIC and CRC-32 while stripping the framing, so a bad pointer or corrupted body is caught now that the Puffin footer is no longer validated. This is both faster and more permissive: - Performance: a deletion vector read is now a single ranged read of one blob rather than loading the whole Puffin file and deserializing every blob it contains. That cuts I/O (notably against object storage, where only the blob's byte range is fetched), CPU, and memory, and scales with the referenced vector rather than the size of the shared container. - Compatibility: deletion vectors that are not fully-formed Puffin files -- for example Delta-compatible vectors that omit the footer -- become readable, since the footer is never consulted. Reading one blob per manifest entry surfaces gaps that reading every blob previously masked, so match and preserve each deletion vector by its target: - Route deletion vectors by referenced_data_file in DeleteFileIndex. Deletion vectors need not carry path bounds, so without this they fall into a shared partition bucket and collapse by file_path, giving every data file in the partition the same vector. - Deduplicate delete files on (file_path, content_offset) in _read_all_delete_files. DataFile equality keys only on file_path, so multiple deletion vectors packed into one Puffin file would otherwise collapse into a single read. - Fill referenced_data_file from the scan task's data file when converting REST position deletes. The field is optional in the REST schema, but the offset read requires it.
1 parent 068aae5 commit 98d6432

8 files changed

Lines changed: 247 additions & 11 deletions

File tree

‎pyiceberg/io/pyarrow.py‎

Lines changed: 22 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -149,11 +149,10 @@
149149
visit_with_partner,
150150
)
151151
from pyiceberg.table import DOWNCAST_NS_TIMESTAMP_TO_US_ON_WRITE, TableProperties
152-
from pyiceberg.table.deletion_vector import deletion_vectors_from_puffin_file
152+
from pyiceberg.table.deletion_vector import DeletionVector
153153
from pyiceberg.table.locations import load_location_provider
154154
from pyiceberg.table.metadata import TableMetadata
155155
from pyiceberg.table.name_mapping import NameMapping, apply_name_mapping
156-
from pyiceberg.table.puffin import PuffinFile
157156
from pyiceberg.transforms import IdentityTransform, TruncateTransform
158157
from pyiceberg.typedef import EMPTY_DICT, Properties, Record, TableVersion
159158
from pyiceberg.types import (
@@ -1171,10 +1170,19 @@ def _read_deletes(io: FileIO, data_file: DataFile) -> dict[str, pa.ChunkedArray]
11711170
for path in table.column("file_path").unique()
11721171
}
11731172
elif data_file.file_format == FileFormat.PUFFIN:
1173+
# Read the deletion vector blob directly using the offset and length recorded in the
1174+
# manifest, without parsing the Puffin footer.
1175+
referenced_data_file = data_file.referenced_data_file
1176+
offset = data_file.content_offset
1177+
length = data_file.content_size_in_bytes
1178+
if referenced_data_file is None or offset is None or length is None:
1179+
raise ValueError(f"Invalid deletion vector, missing referenced data file, offset, or length: {data_file.file_path}")
1180+
11741181
with io.new_input(data_file.file_path).open() as fi:
1175-
payload = fi.read()
1182+
fi.seek(offset)
1183+
blob = fi.read(length)
11761184

1177-
return {dv.referenced_data_file: dv.to_vector() for dv in deletion_vectors_from_puffin_file(PuffinFile(payload))}
1185+
return {referenced_data_file: DeletionVector.from_blob(blob, referenced_data_file).to_vector()}
11781186
else:
11791187
raise ValueError(f"Delete file format not supported: {data_file.file_format}")
11801188

@@ -1743,7 +1751,16 @@ def _task_to_record_batches(
17431751

17441752
def _read_all_delete_files(io: FileIO, tasks: Iterable[FileScanTask]) -> dict[str, list[ChunkedArray]]:
17451753
deletes_per_file: dict[str, list[ChunkedArray]] = {}
1746-
unique_deletes = set(itertools.chain.from_iterable([task.delete_files for task in tasks]))
1754+
# Deduplicate on (file_path, content_offset): several deletion vectors can be packed into one
1755+
# Puffin file, sharing a file_path but differing by content_offset. DataFile equality keys only
1756+
# on file_path, so a plain set would collapse them and drop all but one vector.
1757+
unique_deletes = list(
1758+
{
1759+
(delete_file.file_path, delete_file.content_offset): delete_file
1760+
for task in tasks
1761+
for delete_file in task.delete_files
1762+
}.values()
1763+
)
17471764
if len(unique_deletes) > 0:
17481765
executor = ExecutorFactory.get_or_create()
17491766
deletes_per_files: Iterator[dict[str, ChunkedArray]] = executor.map(

‎pyiceberg/table/__init__.py‎

Lines changed: 20 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2277,7 +2277,7 @@ def from_rest_response(
22772277
delete_file = delete_files[idx]
22782278
if isinstance(delete_file, RESTEqualityDeleteFile):
22792279
raise NotImplementedError(f"PyIceberg does not yet support equality deletes: {delete_file.file_path}")
2280-
resolved_deletes.add(_rest_file_to_data_file(delete_file))
2280+
resolved_deletes.add(_rest_file_to_data_file(delete_file, default_referenced_data_file=data_file.file_path))
22812281

22822282
return FileScanTask(
22832283
data_file=data_file,
@@ -2286,9 +2286,14 @@ def from_rest_response(
22862286
)
22872287

22882288

2289-
def _rest_file_to_data_file(rest_file: RESTContentFile) -> DataFile:
2290-
"""Convert a REST content file to a manifest DataFile."""
2291-
from pyiceberg.catalog.rest.scan_planning import RESTDataFile
2289+
def _rest_file_to_data_file(rest_file: RESTContentFile, default_referenced_data_file: str | None = None) -> DataFile:
2290+
"""Convert a REST content file to a manifest DataFile.
2291+
2292+
default_referenced_data_file supplies the referenced data file for a position delete when the
2293+
REST response omits it; the field is optional in the REST schema, but the offset-based deletion
2294+
vector read requires it.
2295+
"""
2296+
from pyiceberg.catalog.rest.scan_planning import RESTDataFile, RESTPositionDeleteFile
22922297

22932298
if isinstance(rest_file, RESTDataFile):
22942299
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:
23012306
null_value_counts = None
23022307
nan_value_counts = None
23032308

2309+
referenced_data_file = None
2310+
content_offset = None
2311+
content_size_in_bytes = None
2312+
if isinstance(rest_file, RESTPositionDeleteFile):
2313+
referenced_data_file = rest_file.referenced_data_file or default_referenced_data_file
2314+
content_offset = rest_file.content_offset
2315+
content_size_in_bytes = rest_file.content_size_in_bytes
2316+
23042317
data_file = DataFile.from_args(
23052318
content=DataFileContent.from_rest_type(rest_file.content),
23062319
file_path=rest_file.file_path,
@@ -2314,6 +2327,9 @@ def _rest_file_to_data_file(rest_file: RESTContentFile) -> DataFile:
23142327
nan_value_counts=nan_value_counts,
23152328
split_offsets=rest_file.split_offsets,
23162329
sort_order_id=rest_file.sort_order_id,
2330+
referenced_data_file=referenced_data_file,
2331+
content_offset=content_offset,
2332+
content_size_in_bytes=content_size_in_bytes,
23172333
)
23182334
data_file.spec_id = rest_file.spec_id
23192335
return data_file

‎pyiceberg/table/delete_file_index.py‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -115,7 +115,9 @@ def is_empty(self) -> bool:
115115
def add_delete_file(self, manifest_entry: ManifestEntry, partition_key: Record | None = None) -> None:
116116
delete_file = manifest_entry.data_file
117117
seq = manifest_entry.sequence_number or INITIAL_SEQUENCE_NUMBER
118-
target_path = _referenced_data_file_path(delete_file)
118+
# referenced_data_file identifies the target directly (always set for deletion vectors) and
119+
# is authoritative; fall back to path bounds only when it is absent.
120+
target_path = delete_file.referenced_data_file or _referenced_data_file_path(delete_file)
119121

120122
if target_path:
121123
deletes = self._by_path.setdefault(target_path, PositionDeletes())

‎pyiceberg/table/deletion_vector.py‎

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

2021
from pyroaring import BitMap, FrozenBitMap
@@ -27,6 +28,7 @@
2728
EMPTY_BITMAP = FrozenBitMap()
2829
MAX_JAVA_SIGNED = int(math.pow(2, 31)) - 1
2930
PROPERTY_REFERENCED_DATA_FILE = "referenced-data-file"
31+
DV_MAGIC = b"\xd1\xd3\x39\x64"
3032

3133

3234
class DeletionVector:
@@ -76,10 +78,27 @@ def _bitmaps_to_chunked_array(bitmaps: list[BitMap]) -> "pa.ChunkedArray":
7678
def to_vector(self) -> "pa.ChunkedArray":
7779
return self._bitmaps_to_chunked_array(self._bitmaps)
7880

81+
@staticmethod
82+
def from_blob(blob: bytes, referenced_data_file: str) -> "DeletionVector":
83+
return DeletionVector(
84+
referenced_data_file=referenced_data_file,
85+
bitmaps=DeletionVector._deserialize_bitmap(_extract_vector_payload(blob)),
86+
)
87+
7988

8089
def _extract_vector_payload(blob_payload: bytes) -> bytes:
8190
"""Strip deletion-vector-v1 blob framing: length(4 big-endian) + DV magic(4) ... CRC(4 big-endian)."""
8291
length_prefix = int.from_bytes(blob_payload[0:4], "big")
92+
magic = blob_payload[4:8]
93+
if magic != DV_MAGIC:
94+
raise ValueError(f"Invalid magic bytes for deletion vector: {magic!r}, expected {DV_MAGIC!r}")
95+
96+
body = blob_payload[4 : 4 + length_prefix] # magic + serialized bitmap; the CRC covers this
97+
stored_crc = int.from_bytes(blob_payload[4 + length_prefix : 8 + length_prefix], "big")
98+
computed_crc = zlib.crc32(body) & 0xFFFFFFFF
99+
if computed_crc != stored_crc:
100+
raise ValueError(f"Invalid CRC for deletion vector: {computed_crc:#010x}, expected {stored_crc:#010x}")
101+
83102
return blob_payload[8 : 4 + length_prefix]
84103

85104

‎tests/catalog/test_scan_planning_models.py‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -205,6 +205,47 @@ def test_equality_delete_file() -> None:
205205
assert equality_delete.equality_ids == [1, 2]
206206

207207

208+
def test_rest_position_delete_file_to_data_file_propagates_deletion_vector_fields() -> None:
209+
from pyiceberg.table import _rest_file_to_data_file
210+
211+
rest_file = RESTPositionDeleteFile.model_validate(
212+
{
213+
**_rest_position_delete_file(file_path="s3://bucket/table/deletion_vector.puffin", file_format="puffin"),
214+
"referenced-data-file": "s3://bucket/table/data/file.parquet",
215+
}
216+
)
217+
218+
data_file = _rest_file_to_data_file(rest_file)
219+
220+
assert data_file.referenced_data_file == "s3://bucket/table/data/file.parquet"
221+
assert data_file.content_offset == 100
222+
assert data_file.content_size_in_bytes == 200
223+
224+
225+
def test_from_rest_response_fills_missing_referenced_data_file() -> None:
226+
from pyiceberg.table import FileScanTask
227+
228+
# A valid REST position delete file may omit referenced-data-file; the containing task identifies
229+
# the target data file, so the conversion must fill it in for the offset-based deletion vector read.
230+
rest_task = RESTFileScanTask.model_validate(
231+
{
232+
"data-file": _rest_data_file(file_path="s3://bucket/table/data/file.parquet"),
233+
"delete-file-references": [0],
234+
}
235+
)
236+
delete_file = TypeAdapter(RESTDeleteFile).validate_python(
237+
_rest_position_delete_file(file_path="s3://bucket/table/deletion_vector.puffin", file_format="puffin")
238+
)
239+
assert delete_file.referenced_data_file is None
240+
241+
task = FileScanTask.from_rest_response(rest_task, [delete_file])
242+
243+
delete = next(iter(task.delete_files))
244+
assert delete.referenced_data_file == "s3://bucket/table/data/file.parquet"
245+
assert delete.content_offset == 100
246+
assert delete.content_size_in_bytes == 200
247+
248+
208249
def test_file_format_case_insensitive() -> None:
209250
for fmt in ["parquet", "PARQUET", "Parquet"]:
210251
data_file = _rest_data_file(file_format=fmt)

‎tests/io/test_pyarrow.py‎

Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,7 @@
7575
_ConvertToArrowSchema,
7676
_determine_partitions,
7777
_primitive_to_physical,
78+
_read_all_delete_files,
7879
_read_deletes,
7980
_task_to_record_batches,
8081
_to_requested_schema,
@@ -1840,6 +1841,83 @@ def test_read_deletes(deletes_file: str, request: pytest.FixtureRequest) -> None
18401841
assert list(deletes.values())[0] == pa.chunked_array([[1, 3, 5]])
18411842

18421843

1844+
def _delta_dv_entry(positions: list[int]) -> bytes:
1845+
"""One deletion vector framed as Delta's DeletionVectorStore writes it in a .bin file.
1846+
1847+
Layout: <length(4 BE)> <DV_MAGIC + 64-bit RoaringBitmapArray in the portable format shared by
1848+
Delta and Iceberg> <CRC-32(4 BE)>. The bitmap holds `positions` in a single bitmap keyed at 0.
1849+
"""
1850+
import zlib
1851+
1852+
from pyroaring import BitMap
1853+
1854+
from pyiceberg.table.deletion_vector import DV_MAGIC
1855+
1856+
bitmap = (1).to_bytes(8, "little") + (0).to_bytes(4, "little") + BitMap(positions).serialize()
1857+
data = DV_MAGIC + bitmap
1858+
return len(data).to_bytes(4, "big") + data + (zlib.crc32(data) & 0xFFFFFFFF).to_bytes(4, "big")
1859+
1860+
1861+
def test_read_deletes_deletion_vector_in_delta_bin_file(tmp_path: Path) -> None:
1862+
entry = _delta_dv_entry([1, 3, 5])
1863+
dv_path = f"{tmp_path}/deletion_vector.bin"
1864+
with open(dv_path, "wb") as f:
1865+
f.write(b"\x01" + entry)
1866+
1867+
referenced_data_file = "s3://bucket/data.parquet"
1868+
data_file = DataFile.from_args(
1869+
file_path=dv_path,
1870+
file_format=FileFormat.PUFFIN,
1871+
content=DataFileContent.POSITION_DELETES,
1872+
referenced_data_file=referenced_data_file,
1873+
# Iceberg's content_offset points at the length prefix (past Delta's 1-byte version header),
1874+
# and content_size_in_bytes counts the length prefix and CRC that Delta's sizeInBytes omits.
1875+
content_offset=1,
1876+
content_size_in_bytes=len(entry),
1877+
)
1878+
1879+
deletes = _read_deletes(PyArrowFileIO(), data_file)
1880+
1881+
assert deletes.keys() == {referenced_data_file}
1882+
assert deletes[referenced_data_file] == pa.chunked_array([[1, 3, 5]])
1883+
1884+
1885+
def test_read_all_delete_files_reads_multiple_deletion_vectors_in_one_file(tmp_path: Path) -> None:
1886+
# Two deletion vectors packed into one file: same file_path, different content_offset. DataFile
1887+
# equality keys only on file_path, so deduplication must not collapse them into a single read.
1888+
entry_a = _delta_dv_entry([1, 3, 5])
1889+
entry_b = _delta_dv_entry([2, 4])
1890+
dv_path = f"{tmp_path}/deletion_vectors.bin"
1891+
with open(dv_path, "wb") as f:
1892+
f.write(b"\x01" + entry_a + entry_b)
1893+
1894+
def _dv(referenced_data_file: str, offset: int, size: int) -> DataFile:
1895+
return DataFile.from_args(
1896+
file_path=dv_path,
1897+
file_format=FileFormat.PUFFIN,
1898+
content=DataFileContent.POSITION_DELETES,
1899+
referenced_data_file=referenced_data_file,
1900+
content_offset=offset,
1901+
content_size_in_bytes=size,
1902+
)
1903+
1904+
def _task(data_file_path: str, delete_file: DataFile) -> FileScanTask:
1905+
data_file = DataFile.from_args(file_path=data_file_path, file_format=FileFormat.PARQUET)
1906+
return FileScanTask(data_file=data_file, delete_files={delete_file})
1907+
1908+
deletes = _read_all_delete_files(
1909+
PyArrowFileIO(),
1910+
[
1911+
_task("s3://bucket/a.parquet", _dv("s3://bucket/a.parquet", 1, len(entry_a))),
1912+
_task("s3://bucket/b.parquet", _dv("s3://bucket/b.parquet", 1 + len(entry_a), len(entry_b))),
1913+
],
1914+
)
1915+
1916+
assert deletes.keys() == {"s3://bucket/a.parquet", "s3://bucket/b.parquet"}
1917+
assert deletes["s3://bucket/a.parquet"][0] == pa.chunked_array([[1, 3, 5]])
1918+
assert deletes["s3://bucket/b.parquet"][0] == pa.chunked_array([[2, 4]])
1919+
1920+
18431921
def test_delete(deletes_file: str, request: pytest.FixtureRequest, table_schema_simple: Schema) -> None:
18441922
# Determine file format from the file extension
18451923
file_format = FileFormat.PARQUET if deletes_file.endswith(".parquet") else FileFormat.ORC

‎tests/table/test_delete_file_index.py‎

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -161,6 +161,38 @@ def test_dvs_treated_as_position_deletes() -> None:
161161
assert all(d.content == DataFileContent.POSITION_DELETES for d in result)
162162

163163

164+
def test_deletion_vectors_matched_by_referenced_data_file() -> None:
165+
# Two deletion vectors packed into one Puffin file (shared file_path), without path bounds, each
166+
# referencing a different data file. They must be routed by referenced_data_file, not collapsed
167+
# into one partition bucket and deduplicated away by their shared file_path.
168+
index = DeleteFileIndex()
169+
170+
def _dv(referenced_data_file: str, content_offset: int) -> ManifestEntry:
171+
delete_file = DataFile.from_args(
172+
content=DataFileContent.POSITION_DELETES,
173+
file_path="s3://bucket/deletion_vectors.puffin",
174+
file_format=FileFormat.PUFFIN,
175+
partition=Record(),
176+
record_count=10,
177+
file_size_in_bytes=100,
178+
referenced_data_file=referenced_data_file,
179+
content_offset=content_offset,
180+
content_size_in_bytes=40,
181+
)
182+
return ManifestEntry.from_args(status=ManifestEntryStatus.ADDED, sequence_number=2, data_file=delete_file)
183+
184+
index.add_delete_file(_dv("s3://bucket/a.parquet", 1))
185+
index.add_delete_file(_dv("s3://bucket/b.parquet", 41))
186+
187+
deletes_a = index.for_data_file(1, _create_data_file(file_path="s3://bucket/a.parquet"))
188+
deletes_b = index.for_data_file(1, _create_data_file(file_path="s3://bucket/b.parquet"))
189+
190+
assert {d.referenced_data_file for d in deletes_a} == {"s3://bucket/a.parquet"}
191+
assert {d.content_offset for d in deletes_a} == {1}
192+
assert {d.referenced_data_file for d in deletes_b} == {"s3://bucket/b.parquet"}
193+
assert {d.content_offset for d in deletes_b} == {41}
194+
195+
164196
def test_cannot_add_after_indexing() -> None:
165197
group = PositionDeletes()
166198
group.add(_create_positional_delete(sequence_number=1).data_file, 1)

‎tests/table/test_deletion_vector.py‎

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

1920
import pytest
2021
from pyroaring import BitMap
2122

22-
from pyiceberg.table.deletion_vector import DeletionVector
23+
from pyiceberg.table.deletion_vector import DV_MAGIC, DeletionVector
2324

2425

2526
def _open_file(file: str) -> bytes:
@@ -71,3 +72,33 @@ def test_map_high_vals() -> None:
7172

7273
with pytest.raises(ValueError, match="Key 4022190063 is too large, max 2147483647 to maintain compatibility with Java impl"):
7374
_ = DeletionVector._deserialize_bitmap(puffin)
75+
76+
77+
def test_from_blob_invalid_magic() -> None:
78+
bitmap = _open_file("64map32bitvals.bin")
79+
length_prefix = 4 + len(bitmap)
80+
blob = length_prefix.to_bytes(4, "big") + b"\x00\x00\x00\x00" + bitmap + (0).to_bytes(4, "big")
81+
82+
with pytest.raises(ValueError, match="Invalid magic bytes for deletion vector"):
83+
_ = DeletionVector.from_blob(blob, "s3://bucket/data.parquet")
84+
85+
86+
def test_from_blob_invalid_crc() -> None:
87+
bitmap = _open_file("64map32bitvals.bin")
88+
data = DV_MAGIC + bitmap
89+
# A trailing CRC of zero does not match the CRC-32 of (magic + bitmap).
90+
blob = len(data).to_bytes(4, "big") + data + (0).to_bytes(4, "big")
91+
92+
with pytest.raises(ValueError, match="Invalid CRC for deletion vector"):
93+
_ = DeletionVector.from_blob(blob, "s3://bucket/data.parquet")
94+
95+
96+
def test_from_blob() -> None:
97+
bitmap = _open_file("64map32bitvals.bin")
98+
data = DV_MAGIC + bitmap
99+
blob = len(data).to_bytes(4, "big") + data + (zlib.crc32(data) & 0xFFFFFFFF).to_bytes(4, "big")
100+
101+
dv = DeletionVector.from_blob(blob, "s3://bucket/data.parquet")
102+
103+
assert dv.referenced_data_file == "s3://bucket/data.parquet"
104+
assert dv._bitmaps == [BitMap([0, 1, 2, 3, 4, 5, 6, 7, 8, 9])]

0 commit comments

Comments
 (0)