amogh-jahagirdar commented on code in PR #3478:
URL: https://github.com/apache/iceberg-python/pull/3478#discussion_r3591046072
##########
pyiceberg/table/deletion_vector.py:
##########
@@ -77,17 +89,103 @@ def to_vector(self) -> "pa.ChunkedArray":
return self._bitmaps_to_chunked_array(self._bitmaps)
-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")
- return blob_payload[8 : 4 + length_prefix]
+def _deserialize_dv_blob(blob: bytes, record_count: int | None = None) ->
list[BitMap]:
+ # The DV blob encoding matches Iceberg Java's BitmapPositionDeleteIndex:
+ # 4-byte big-endian bitmap-data length, 4-byte little-endian magic number,
+ # portable Roaring bitmap data, and 4-byte big-endian CRC-32.
+ if len(blob) < _DV_BLOB_MIN_SIZE_BYTES:
+ raise ValueError(f"Invalid deletion vector blob length: {len(blob)}")
+
+ bitmap_data_length = _DV_BLOB_LENGTH.unpack_from(blob)[0]
+ expected_bitmap_data_length = len(blob) - _DV_BLOB_LENGTH.size -
_DV_BLOB_CRC.size
+ if bitmap_data_length != expected_bitmap_data_length:
+ raise ValueError(f"Invalid bitmap data length: {bitmap_data_length},
expected {expected_bitmap_data_length}")
+
+ bitmap_data_offset = _DV_BLOB_LENGTH.size
+ crc_offset = bitmap_data_offset + bitmap_data_length
+ bitmap_data = blob[bitmap_data_offset:crc_offset]
+
+ magic_number = _DV_BLOB_MAGIC.unpack_from(bitmap_data)[0]
+ if magic_number != _DV_BLOB_MAGIC_NUMBER:
+ raise ValueError(f"Invalid magic number: {magic_number}, expected
{_DV_BLOB_MAGIC_NUMBER}")
+
+ checksum = zlib.crc32(bitmap_data) & 0xFFFFFFFF
+ expected_checksum = _DV_BLOB_CRC.unpack_from(blob, crc_offset)[0]
+ if checksum != expected_checksum:
+ raise ValueError("Invalid CRC")
+
+ bitmaps =
DeletionVector._deserialize_bitmap(bitmap_data[_DV_BLOB_MAGIC.size :])
+ if record_count is not None:
+ cardinality = sum(len(bitmap) for bitmap in bitmaps)
+ if cardinality != record_count:
+ raise ValueError(f"Invalid cardinality: {cardinality}, expected
{record_count}")
+
+ return bitmaps
+
+
+def _validate_deletion_vector_content(data_file: "DataFile") -> None:
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+
+ if content_offset is None:
+ raise ValueError(f"Invalid deletion vector, content offset is missing:
{data_file.file_path}")
+ if content_size_in_bytes is None:
+ raise ValueError(f"Invalid deletion vector, content size is missing:
{data_file.file_path}")
+ if content_offset < 0:
+ raise ValueError(f"Invalid deletion vector, content offset cannot be
negative: {content_offset}")
+ if content_size_in_bytes < 0:
+ raise ValueError(f"Invalid deletion vector, content size cannot be
negative: {content_size_in_bytes}")
+ if content_size_in_bytes > _MAX_DELETION_VECTOR_CONTENT_SIZE:
+ raise ValueError(f"Cannot read deletion vector larger than 2GB:
{content_size_in_bytes}")
+ if referenced_data_file is None:
+ raise ValueError(f"Invalid deletion vector, referenced data file is
missing: {data_file.file_path}")
+
+
+def has_deletion_vector_content_reference(data_file: "DataFile") -> bool:
+ return (
+ data_file.content_offset is not None
+ or data_file.content_size_in_bytes is not None
+ or data_file.referenced_data_file is not None
+ )
+
+
+def _read_deletion_vector(io: "FileIO", data_file: "DataFile") ->
DeletionVector:
+ _validate_deletion_vector_content(data_file)
+
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+ assert content_offset is not None
+ assert content_size_in_bytes is not None
+ assert referenced_data_file is not None
+
+ with io.new_input(data_file.file_path).open() as fi:
+ fi.seek(content_offset)
+ payload = fi.read(content_size_in_bytes)
+
+ if len(payload) != content_size_in_bytes:
+ raise ValueError(f"Could not read deletion vector, expected
{content_size_in_bytes} bytes, got {len(payload)}")
+
+ return DeletionVector(
+ referenced_data_file=referenced_data_file,
+ bitmaps=_deserialize_dv_blob(payload, data_file.record_count),
+ )
+
+
+def read_deletion_vectors(io: "FileIO", data_file: "DataFile") ->
list[DeletionVector]:
Review Comment:
Docstrings for the public APIs looks to be the convention in the project so
would be nice to provide just a simple one here
##########
pyiceberg/table/deletion_vector.py:
##########
@@ -77,17 +89,103 @@ def to_vector(self) -> "pa.ChunkedArray":
return self._bitmaps_to_chunked_array(self._bitmaps)
-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")
- return blob_payload[8 : 4 + length_prefix]
+def _deserialize_dv_blob(blob: bytes, record_count: int | None = None) ->
list[BitMap]:
+ # The DV blob encoding matches Iceberg Java's BitmapPositionDeleteIndex:
+ # 4-byte big-endian bitmap-data length, 4-byte little-endian magic number,
+ # portable Roaring bitmap data, and 4-byte big-endian CRC-32.
+ if len(blob) < _DV_BLOB_MIN_SIZE_BYTES:
+ raise ValueError(f"Invalid deletion vector blob length: {len(blob)}")
+
+ bitmap_data_length = _DV_BLOB_LENGTH.unpack_from(blob)[0]
+ expected_bitmap_data_length = len(blob) - _DV_BLOB_LENGTH.size -
_DV_BLOB_CRC.size
+ if bitmap_data_length != expected_bitmap_data_length:
+ raise ValueError(f"Invalid bitmap data length: {bitmap_data_length},
expected {expected_bitmap_data_length}")
+
+ bitmap_data_offset = _DV_BLOB_LENGTH.size
+ crc_offset = bitmap_data_offset + bitmap_data_length
+ bitmap_data = blob[bitmap_data_offset:crc_offset]
+
+ magic_number = _DV_BLOB_MAGIC.unpack_from(bitmap_data)[0]
+ if magic_number != _DV_BLOB_MAGIC_NUMBER:
+ raise ValueError(f"Invalid magic number: {magic_number}, expected
{_DV_BLOB_MAGIC_NUMBER}")
+
+ checksum = zlib.crc32(bitmap_data) & 0xFFFFFFFF
+ expected_checksum = _DV_BLOB_CRC.unpack_from(blob, crc_offset)[0]
+ if checksum != expected_checksum:
+ raise ValueError("Invalid CRC")
+
+ bitmaps =
DeletionVector._deserialize_bitmap(bitmap_data[_DV_BLOB_MAGIC.size :])
+ if record_count is not None:
+ cardinality = sum(len(bitmap) for bitmap in bitmaps)
+ if cardinality != record_count:
+ raise ValueError(f"Invalid cardinality: {cardinality}, expected
{record_count}")
+
+ return bitmaps
+
+
+def _validate_deletion_vector_content(data_file: "DataFile") -> None:
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+
+ if content_offset is None:
+ raise ValueError(f"Invalid deletion vector, content offset is missing:
{data_file.file_path}")
+ if content_size_in_bytes is None:
+ raise ValueError(f"Invalid deletion vector, content size is missing:
{data_file.file_path}")
+ if content_offset < 0:
+ raise ValueError(f"Invalid deletion vector, content offset cannot be
negative: {content_offset}")
+ if content_size_in_bytes < 0:
+ raise ValueError(f"Invalid deletion vector, content size cannot be
negative: {content_size_in_bytes}")
+ if content_size_in_bytes > _MAX_DELETION_VECTOR_CONTENT_SIZE:
+ raise ValueError(f"Cannot read deletion vector larger than 2GB:
{content_size_in_bytes}")
+ if referenced_data_file is None:
+ raise ValueError(f"Invalid deletion vector, referenced data file is
missing: {data_file.file_path}")
+
+
+def has_deletion_vector_content_reference(data_file: "DataFile") -> bool:
Review Comment:
I know the python module is using a class called DataFile that encapsulates
both DataFile and DeleteFile and we don't need to change that in this PR but at
least in terms of naming could we call the parameter `dv`? Makes it a bit
easier to read the actual function logic as we know what exactly we're working
with in this method.
##########
pyiceberg/table/deletion_vector.py:
##########
@@ -77,17 +89,103 @@ def to_vector(self) -> "pa.ChunkedArray":
return self._bitmaps_to_chunked_array(self._bitmaps)
-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")
- return blob_payload[8 : 4 + length_prefix]
+def _deserialize_dv_blob(blob: bytes, record_count: int | None = None) ->
list[BitMap]:
+ # The DV blob encoding matches Iceberg Java's BitmapPositionDeleteIndex:
+ # 4-byte big-endian bitmap-data length, 4-byte little-endian magic number,
+ # portable Roaring bitmap data, and 4-byte big-endian CRC-32.
+ if len(blob) < _DV_BLOB_MIN_SIZE_BYTES:
+ raise ValueError(f"Invalid deletion vector blob length: {len(blob)}")
+
+ bitmap_data_length = _DV_BLOB_LENGTH.unpack_from(blob)[0]
+ expected_bitmap_data_length = len(blob) - _DV_BLOB_LENGTH.size -
_DV_BLOB_CRC.size
+ if bitmap_data_length != expected_bitmap_data_length:
+ raise ValueError(f"Invalid bitmap data length: {bitmap_data_length},
expected {expected_bitmap_data_length}")
+
+ bitmap_data_offset = _DV_BLOB_LENGTH.size
+ crc_offset = bitmap_data_offset + bitmap_data_length
+ bitmap_data = blob[bitmap_data_offset:crc_offset]
+
+ magic_number = _DV_BLOB_MAGIC.unpack_from(bitmap_data)[0]
+ if magic_number != _DV_BLOB_MAGIC_NUMBER:
+ raise ValueError(f"Invalid magic number: {magic_number}, expected
{_DV_BLOB_MAGIC_NUMBER}")
+
+ checksum = zlib.crc32(bitmap_data) & 0xFFFFFFFF
+ expected_checksum = _DV_BLOB_CRC.unpack_from(blob, crc_offset)[0]
+ if checksum != expected_checksum:
+ raise ValueError("Invalid CRC")
+
+ bitmaps =
DeletionVector._deserialize_bitmap(bitmap_data[_DV_BLOB_MAGIC.size :])
+ if record_count is not None:
+ cardinality = sum(len(bitmap) for bitmap in bitmaps)
+ if cardinality != record_count:
+ raise ValueError(f"Invalid cardinality: {cardinality}, expected
{record_count}")
+
+ return bitmaps
+
+
+def _validate_deletion_vector_content(data_file: "DataFile") -> None:
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+
+ if content_offset is None:
+ raise ValueError(f"Invalid deletion vector, content offset is missing:
{data_file.file_path}")
+ if content_size_in_bytes is None:
+ raise ValueError(f"Invalid deletion vector, content size is missing:
{data_file.file_path}")
+ if content_offset < 0:
+ raise ValueError(f"Invalid deletion vector, content offset cannot be
negative: {content_offset}")
+ if content_size_in_bytes < 0:
+ raise ValueError(f"Invalid deletion vector, content size cannot be
negative: {content_size_in_bytes}")
+ if content_size_in_bytes > _MAX_DELETION_VECTOR_CONTENT_SIZE:
+ raise ValueError(f"Cannot read deletion vector larger than 2GB:
{content_size_in_bytes}")
+ if referenced_data_file is None:
+ raise ValueError(f"Invalid deletion vector, referenced data file is
missing: {data_file.file_path}")
+
+
+def has_deletion_vector_content_reference(data_file: "DataFile") -> bool:
+ return (
+ data_file.content_offset is not None
+ or data_file.content_size_in_bytes is not None
+ or data_file.referenced_data_file is not None
+ )
+
+
+def _read_deletion_vector(io: "FileIO", data_file: "DataFile") ->
DeletionVector:
+ _validate_deletion_vector_content(data_file)
+
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+ assert content_offset is not None
+ assert content_size_in_bytes is not None
+ assert referenced_data_file is not None
+
+ with io.new_input(data_file.file_path).open() as fi:
+ fi.seek(content_offset)
+ payload = fi.read(content_size_in_bytes)
+
+ if len(payload) != content_size_in_bytes:
+ raise ValueError(f"Could not read deletion vector, expected
{content_size_in_bytes} bytes, got {len(payload)}")
+
+ return DeletionVector(
+ referenced_data_file=referenced_data_file,
+ bitmaps=_deserialize_dv_blob(payload, data_file.record_count),
+ )
+
+
+def read_deletion_vectors(io: "FileIO", data_file: "DataFile") ->
list[DeletionVector]:
Review Comment:
also same as above, we should call the param `dv` imo, at this point we know
we are trying to read the DV.
##########
pyiceberg/table/deletion_vector.py:
##########
@@ -77,17 +89,103 @@ def to_vector(self) -> "pa.ChunkedArray":
return self._bitmaps_to_chunked_array(self._bitmaps)
-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")
- return blob_payload[8 : 4 + length_prefix]
+def _deserialize_dv_blob(blob: bytes, record_count: int | None = None) ->
list[BitMap]:
+ # The DV blob encoding matches Iceberg Java's BitmapPositionDeleteIndex:
+ # 4-byte big-endian bitmap-data length, 4-byte little-endian magic number,
+ # portable Roaring bitmap data, and 4-byte big-endian CRC-32.
+ if len(blob) < _DV_BLOB_MIN_SIZE_BYTES:
+ raise ValueError(f"Invalid deletion vector blob length: {len(blob)}")
+
+ bitmap_data_length = _DV_BLOB_LENGTH.unpack_from(blob)[0]
+ expected_bitmap_data_length = len(blob) - _DV_BLOB_LENGTH.size -
_DV_BLOB_CRC.size
+ if bitmap_data_length != expected_bitmap_data_length:
+ raise ValueError(f"Invalid bitmap data length: {bitmap_data_length},
expected {expected_bitmap_data_length}")
+
+ bitmap_data_offset = _DV_BLOB_LENGTH.size
+ crc_offset = bitmap_data_offset + bitmap_data_length
+ bitmap_data = blob[bitmap_data_offset:crc_offset]
+
+ magic_number = _DV_BLOB_MAGIC.unpack_from(bitmap_data)[0]
+ if magic_number != _DV_BLOB_MAGIC_NUMBER:
+ raise ValueError(f"Invalid magic number: {magic_number}, expected
{_DV_BLOB_MAGIC_NUMBER}")
+
+ checksum = zlib.crc32(bitmap_data) & 0xFFFFFFFF
+ expected_checksum = _DV_BLOB_CRC.unpack_from(blob, crc_offset)[0]
+ if checksum != expected_checksum:
+ raise ValueError("Invalid CRC")
+
+ bitmaps =
DeletionVector._deserialize_bitmap(bitmap_data[_DV_BLOB_MAGIC.size :])
+ if record_count is not None:
+ cardinality = sum(len(bitmap) for bitmap in bitmaps)
+ if cardinality != record_count:
+ raise ValueError(f"Invalid cardinality: {cardinality}, expected
{record_count}")
+
+ return bitmaps
+
+
+def _validate_deletion_vector_content(data_file: "DataFile") -> None:
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+
+ if content_offset is None:
+ raise ValueError(f"Invalid deletion vector, content offset is missing:
{data_file.file_path}")
+ if content_size_in_bytes is None:
+ raise ValueError(f"Invalid deletion vector, content size is missing:
{data_file.file_path}")
+ if content_offset < 0:
+ raise ValueError(f"Invalid deletion vector, content offset cannot be
negative: {content_offset}")
+ if content_size_in_bytes < 0:
+ raise ValueError(f"Invalid deletion vector, content size cannot be
negative: {content_size_in_bytes}")
+ if content_size_in_bytes > _MAX_DELETION_VECTOR_CONTENT_SIZE:
+ raise ValueError(f"Cannot read deletion vector larger than 2GB:
{content_size_in_bytes}")
+ if referenced_data_file is None:
+ raise ValueError(f"Invalid deletion vector, referenced data file is
missing: {data_file.file_path}")
+
+
+def has_deletion_vector_content_reference(data_file: "DataFile") -> bool:
+ return (
+ data_file.content_offset is not None
+ or data_file.content_size_in_bytes is not None
+ or data_file.referenced_data_file is not None
+ )
+
+
+def _read_deletion_vector(io: "FileIO", data_file: "DataFile") ->
DeletionVector:
+ _validate_deletion_vector_content(data_file)
+
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+ assert content_offset is not None
+ assert content_size_in_bytes is not None
+ assert referenced_data_file is not None
+
Review Comment:
Do we need these 3 assertions given what _validate_deletion_vector_content
already checks?
##########
pyiceberg/table/deletion_vector.py:
##########
@@ -77,17 +89,103 @@ def to_vector(self) -> "pa.ChunkedArray":
return self._bitmaps_to_chunked_array(self._bitmaps)
-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")
- return blob_payload[8 : 4 + length_prefix]
+def _deserialize_dv_blob(blob: bytes, record_count: int | None = None) ->
list[BitMap]:
+ # The DV blob encoding matches Iceberg Java's BitmapPositionDeleteIndex:
+ # 4-byte big-endian bitmap-data length, 4-byte little-endian magic number,
+ # portable Roaring bitmap data, and 4-byte big-endian CRC-32.
+ if len(blob) < _DV_BLOB_MIN_SIZE_BYTES:
+ raise ValueError(f"Invalid deletion vector blob length: {len(blob)}")
+
+ bitmap_data_length = _DV_BLOB_LENGTH.unpack_from(blob)[0]
+ expected_bitmap_data_length = len(blob) - _DV_BLOB_LENGTH.size -
_DV_BLOB_CRC.size
+ if bitmap_data_length != expected_bitmap_data_length:
+ raise ValueError(f"Invalid bitmap data length: {bitmap_data_length},
expected {expected_bitmap_data_length}")
+
+ bitmap_data_offset = _DV_BLOB_LENGTH.size
+ crc_offset = bitmap_data_offset + bitmap_data_length
+ bitmap_data = blob[bitmap_data_offset:crc_offset]
+
+ magic_number = _DV_BLOB_MAGIC.unpack_from(bitmap_data)[0]
+ if magic_number != _DV_BLOB_MAGIC_NUMBER:
+ raise ValueError(f"Invalid magic number: {magic_number}, expected
{_DV_BLOB_MAGIC_NUMBER}")
+
+ checksum = zlib.crc32(bitmap_data) & 0xFFFFFFFF
+ expected_checksum = _DV_BLOB_CRC.unpack_from(blob, crc_offset)[0]
+ if checksum != expected_checksum:
+ raise ValueError("Invalid CRC")
+
+ bitmaps =
DeletionVector._deserialize_bitmap(bitmap_data[_DV_BLOB_MAGIC.size :])
+ if record_count is not None:
+ cardinality = sum(len(bitmap) for bitmap in bitmaps)
+ if cardinality != record_count:
+ raise ValueError(f"Invalid cardinality: {cardinality}, expected
{record_count}")
+
+ return bitmaps
+
+
+def _validate_deletion_vector_content(data_file: "DataFile") -> None:
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+
+ if content_offset is None:
+ raise ValueError(f"Invalid deletion vector, content offset is missing:
{data_file.file_path}")
+ if content_size_in_bytes is None:
+ raise ValueError(f"Invalid deletion vector, content size is missing:
{data_file.file_path}")
+ if content_offset < 0:
+ raise ValueError(f"Invalid deletion vector, content offset cannot be
negative: {content_offset}")
Review Comment:
I am fine with having a more defensive implementation (the spec requires
writers to produce the offset/size/refereenced file for DVs anyways) but just
mentioning i think we only need to do these checks once and in one place only
rather than in multiple places.
##########
pyiceberg/table/deletion_vector.py:
##########
@@ -77,17 +89,103 @@ def to_vector(self) -> "pa.ChunkedArray":
return self._bitmaps_to_chunked_array(self._bitmaps)
-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")
- return blob_payload[8 : 4 + length_prefix]
+def _deserialize_dv_blob(blob: bytes, record_count: int | None = None) ->
list[BitMap]:
+ # The DV blob encoding matches Iceberg Java's BitmapPositionDeleteIndex:
+ # 4-byte big-endian bitmap-data length, 4-byte little-endian magic number,
+ # portable Roaring bitmap data, and 4-byte big-endian CRC-32.
+ if len(blob) < _DV_BLOB_MIN_SIZE_BYTES:
+ raise ValueError(f"Invalid deletion vector blob length: {len(blob)}")
+
+ bitmap_data_length = _DV_BLOB_LENGTH.unpack_from(blob)[0]
+ expected_bitmap_data_length = len(blob) - _DV_BLOB_LENGTH.size -
_DV_BLOB_CRC.size
+ if bitmap_data_length != expected_bitmap_data_length:
+ raise ValueError(f"Invalid bitmap data length: {bitmap_data_length},
expected {expected_bitmap_data_length}")
+
+ bitmap_data_offset = _DV_BLOB_LENGTH.size
+ crc_offset = bitmap_data_offset + bitmap_data_length
+ bitmap_data = blob[bitmap_data_offset:crc_offset]
+
+ magic_number = _DV_BLOB_MAGIC.unpack_from(bitmap_data)[0]
+ if magic_number != _DV_BLOB_MAGIC_NUMBER:
+ raise ValueError(f"Invalid magic number: {magic_number}, expected
{_DV_BLOB_MAGIC_NUMBER}")
+
+ checksum = zlib.crc32(bitmap_data) & 0xFFFFFFFF
+ expected_checksum = _DV_BLOB_CRC.unpack_from(blob, crc_offset)[0]
+ if checksum != expected_checksum:
+ raise ValueError("Invalid CRC")
+
+ bitmaps =
DeletionVector._deserialize_bitmap(bitmap_data[_DV_BLOB_MAGIC.size :])
+ if record_count is not None:
+ cardinality = sum(len(bitmap) for bitmap in bitmaps)
+ if cardinality != record_count:
+ raise ValueError(f"Invalid cardinality: {cardinality}, expected
{record_count}")
+
+ return bitmaps
+
+
+def _validate_deletion_vector_content(data_file: "DataFile") -> None:
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+
+ if content_offset is None:
+ raise ValueError(f"Invalid deletion vector, content offset is missing:
{data_file.file_path}")
+ if content_size_in_bytes is None:
+ raise ValueError(f"Invalid deletion vector, content size is missing:
{data_file.file_path}")
+ if content_offset < 0:
+ raise ValueError(f"Invalid deletion vector, content offset cannot be
negative: {content_offset}")
+ if content_size_in_bytes < 0:
+ raise ValueError(f"Invalid deletion vector, content size cannot be
negative: {content_size_in_bytes}")
+ if content_size_in_bytes > _MAX_DELETION_VECTOR_CONTENT_SIZE:
+ raise ValueError(f"Cannot read deletion vector larger than 2GB:
{content_size_in_bytes}")
+ if referenced_data_file is None:
+ raise ValueError(f"Invalid deletion vector, referenced data file is
missing: {data_file.file_path}")
+
+
+def has_deletion_vector_content_reference(data_file: "DataFile") -> bool:
+ return (
+ data_file.content_offset is not None
+ or data_file.content_size_in_bytes is not None
+ or data_file.referenced_data_file is not None
+ )
+
+
+def _read_deletion_vector(io: "FileIO", data_file: "DataFile") ->
DeletionVector:
Review Comment:
`dv` for the param name?
##########
pyiceberg/table/deletion_vector.py:
##########
@@ -77,17 +89,103 @@ def to_vector(self) -> "pa.ChunkedArray":
return self._bitmaps_to_chunked_array(self._bitmaps)
-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")
- return blob_payload[8 : 4 + length_prefix]
+def _deserialize_dv_blob(blob: bytes, record_count: int | None = None) ->
list[BitMap]:
+ # The DV blob encoding matches Iceberg Java's BitmapPositionDeleteIndex:
+ # 4-byte big-endian bitmap-data length, 4-byte little-endian magic number,
+ # portable Roaring bitmap data, and 4-byte big-endian CRC-32.
+ if len(blob) < _DV_BLOB_MIN_SIZE_BYTES:
+ raise ValueError(f"Invalid deletion vector blob length: {len(blob)}")
+
+ bitmap_data_length = _DV_BLOB_LENGTH.unpack_from(blob)[0]
+ expected_bitmap_data_length = len(blob) - _DV_BLOB_LENGTH.size -
_DV_BLOB_CRC.size
+ if bitmap_data_length != expected_bitmap_data_length:
+ raise ValueError(f"Invalid bitmap data length: {bitmap_data_length},
expected {expected_bitmap_data_length}")
+
+ bitmap_data_offset = _DV_BLOB_LENGTH.size
+ crc_offset = bitmap_data_offset + bitmap_data_length
+ bitmap_data = blob[bitmap_data_offset:crc_offset]
+
+ magic_number = _DV_BLOB_MAGIC.unpack_from(bitmap_data)[0]
+ if magic_number != _DV_BLOB_MAGIC_NUMBER:
+ raise ValueError(f"Invalid magic number: {magic_number}, expected
{_DV_BLOB_MAGIC_NUMBER}")
+
+ checksum = zlib.crc32(bitmap_data) & 0xFFFFFFFF
+ expected_checksum = _DV_BLOB_CRC.unpack_from(blob, crc_offset)[0]
+ if checksum != expected_checksum:
+ raise ValueError("Invalid CRC")
+
+ bitmaps =
DeletionVector._deserialize_bitmap(bitmap_data[_DV_BLOB_MAGIC.size :])
+ if record_count is not None:
+ cardinality = sum(len(bitmap) for bitmap in bitmaps)
+ if cardinality != record_count:
+ raise ValueError(f"Invalid cardinality: {cardinality}, expected
{record_count}")
+
+ return bitmaps
+
+
+def _validate_deletion_vector_content(data_file: "DataFile") -> None:
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+
+ if content_offset is None:
+ raise ValueError(f"Invalid deletion vector, content offset is missing:
{data_file.file_path}")
+ if content_size_in_bytes is None:
+ raise ValueError(f"Invalid deletion vector, content size is missing:
{data_file.file_path}")
+ if content_offset < 0:
+ raise ValueError(f"Invalid deletion vector, content offset cannot be
negative: {content_offset}")
+ if content_size_in_bytes < 0:
+ raise ValueError(f"Invalid deletion vector, content size cannot be
negative: {content_size_in_bytes}")
+ if content_size_in_bytes > _MAX_DELETION_VECTOR_CONTENT_SIZE:
+ raise ValueError(f"Cannot read deletion vector larger than 2GB:
{content_size_in_bytes}")
+ if referenced_data_file is None:
+ raise ValueError(f"Invalid deletion vector, referenced data file is
missing: {data_file.file_path}")
+
+
+def has_deletion_vector_content_reference(data_file: "DataFile") -> bool:
+ return (
+ data_file.content_offset is not None
+ or data_file.content_size_in_bytes is not None
+ or data_file.referenced_data_file is not None
+ )
+
+
+def _read_deletion_vector(io: "FileIO", data_file: "DataFile") ->
DeletionVector:
+ _validate_deletion_vector_content(data_file)
+
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+ assert content_offset is not None
+ assert content_size_in_bytes is not None
+ assert referenced_data_file is not None
+
Review Comment:
Also I think we're needlessly doing these checks at this layer in the stack
and it
's not needed in validate_deletion_vector_content.Writers of DVs are
required to produce all 3 fields in metadata as per the spec. And at this layer
in the stack we know this is a DV so we can assume that all 3 fields are set.
##########
pyiceberg/table/deletion_vector.py:
##########
@@ -77,17 +89,103 @@ def to_vector(self) -> "pa.ChunkedArray":
return self._bitmaps_to_chunked_array(self._bitmaps)
-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")
- return blob_payload[8 : 4 + length_prefix]
+def _deserialize_dv_blob(blob: bytes, record_count: int | None = None) ->
list[BitMap]:
+ # The DV blob encoding matches Iceberg Java's BitmapPositionDeleteIndex:
+ # 4-byte big-endian bitmap-data length, 4-byte little-endian magic number,
+ # portable Roaring bitmap data, and 4-byte big-endian CRC-32.
+ if len(blob) < _DV_BLOB_MIN_SIZE_BYTES:
+ raise ValueError(f"Invalid deletion vector blob length: {len(blob)}")
+
+ bitmap_data_length = _DV_BLOB_LENGTH.unpack_from(blob)[0]
+ expected_bitmap_data_length = len(blob) - _DV_BLOB_LENGTH.size -
_DV_BLOB_CRC.size
+ if bitmap_data_length != expected_bitmap_data_length:
+ raise ValueError(f"Invalid bitmap data length: {bitmap_data_length},
expected {expected_bitmap_data_length}")
+
+ bitmap_data_offset = _DV_BLOB_LENGTH.size
+ crc_offset = bitmap_data_offset + bitmap_data_length
+ bitmap_data = blob[bitmap_data_offset:crc_offset]
+
+ magic_number = _DV_BLOB_MAGIC.unpack_from(bitmap_data)[0]
+ if magic_number != _DV_BLOB_MAGIC_NUMBER:
+ raise ValueError(f"Invalid magic number: {magic_number}, expected
{_DV_BLOB_MAGIC_NUMBER}")
+
+ checksum = zlib.crc32(bitmap_data) & 0xFFFFFFFF
+ expected_checksum = _DV_BLOB_CRC.unpack_from(blob, crc_offset)[0]
+ if checksum != expected_checksum:
+ raise ValueError("Invalid CRC")
+
+ bitmaps =
DeletionVector._deserialize_bitmap(bitmap_data[_DV_BLOB_MAGIC.size :])
+ if record_count is not None:
+ cardinality = sum(len(bitmap) for bitmap in bitmaps)
+ if cardinality != record_count:
+ raise ValueError(f"Invalid cardinality: {cardinality}, expected
{record_count}")
+
+ return bitmaps
+
+
+def _validate_deletion_vector_content(data_file: "DataFile") -> None:
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+
+ if content_offset is None:
+ raise ValueError(f"Invalid deletion vector, content offset is missing:
{data_file.file_path}")
+ if content_size_in_bytes is None:
+ raise ValueError(f"Invalid deletion vector, content size is missing:
{data_file.file_path}")
+ if content_offset < 0:
+ raise ValueError(f"Invalid deletion vector, content offset cannot be
negative: {content_offset}")
+ if content_size_in_bytes < 0:
+ raise ValueError(f"Invalid deletion vector, content size cannot be
negative: {content_size_in_bytes}")
+ if content_size_in_bytes > _MAX_DELETION_VECTOR_CONTENT_SIZE:
+ raise ValueError(f"Cannot read deletion vector larger than 2GB:
{content_size_in_bytes}")
+ if referenced_data_file is None:
+ raise ValueError(f"Invalid deletion vector, referenced data file is
missing: {data_file.file_path}")
+
+
+def has_deletion_vector_content_reference(data_file: "DataFile") -> bool:
Review Comment:
See my comment below, I think we can get ridi of this. At this layer in the
stack, I think we know it's a DV and we don't need to do this field check as
writers are required to produce this metadata as per the spec anyways.
##########
pyiceberg/table/deletion_vector.py:
##########
@@ -77,17 +89,103 @@ def to_vector(self) -> "pa.ChunkedArray":
return self._bitmaps_to_chunked_array(self._bitmaps)
-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")
- return blob_payload[8 : 4 + length_prefix]
+def _deserialize_dv_blob(blob: bytes, record_count: int | None = None) ->
list[BitMap]:
+ # The DV blob encoding matches Iceberg Java's BitmapPositionDeleteIndex:
+ # 4-byte big-endian bitmap-data length, 4-byte little-endian magic number,
+ # portable Roaring bitmap data, and 4-byte big-endian CRC-32.
+ if len(blob) < _DV_BLOB_MIN_SIZE_BYTES:
+ raise ValueError(f"Invalid deletion vector blob length: {len(blob)}")
+
+ bitmap_data_length = _DV_BLOB_LENGTH.unpack_from(blob)[0]
+ expected_bitmap_data_length = len(blob) - _DV_BLOB_LENGTH.size -
_DV_BLOB_CRC.size
+ if bitmap_data_length != expected_bitmap_data_length:
+ raise ValueError(f"Invalid bitmap data length: {bitmap_data_length},
expected {expected_bitmap_data_length}")
+
+ bitmap_data_offset = _DV_BLOB_LENGTH.size
+ crc_offset = bitmap_data_offset + bitmap_data_length
+ bitmap_data = blob[bitmap_data_offset:crc_offset]
+
+ magic_number = _DV_BLOB_MAGIC.unpack_from(bitmap_data)[0]
+ if magic_number != _DV_BLOB_MAGIC_NUMBER:
+ raise ValueError(f"Invalid magic number: {magic_number}, expected
{_DV_BLOB_MAGIC_NUMBER}")
+
+ checksum = zlib.crc32(bitmap_data) & 0xFFFFFFFF
+ expected_checksum = _DV_BLOB_CRC.unpack_from(blob, crc_offset)[0]
+ if checksum != expected_checksum:
+ raise ValueError("Invalid CRC")
+
+ bitmaps =
DeletionVector._deserialize_bitmap(bitmap_data[_DV_BLOB_MAGIC.size :])
+ if record_count is not None:
+ cardinality = sum(len(bitmap) for bitmap in bitmaps)
+ if cardinality != record_count:
+ raise ValueError(f"Invalid cardinality: {cardinality}, expected
{record_count}")
+
+ return bitmaps
+
+
+def _validate_deletion_vector_content(data_file: "DataFile") -> None:
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+
+ if content_offset is None:
+ raise ValueError(f"Invalid deletion vector, content offset is missing:
{data_file.file_path}")
+ if content_size_in_bytes is None:
+ raise ValueError(f"Invalid deletion vector, content size is missing:
{data_file.file_path}")
+ if content_offset < 0:
+ raise ValueError(f"Invalid deletion vector, content offset cannot be
negative: {content_offset}")
+ if content_size_in_bytes < 0:
+ raise ValueError(f"Invalid deletion vector, content size cannot be
negative: {content_size_in_bytes}")
+ if content_size_in_bytes > _MAX_DELETION_VECTOR_CONTENT_SIZE:
+ raise ValueError(f"Cannot read deletion vector larger than 2GB:
{content_size_in_bytes}")
+ if referenced_data_file is None:
+ raise ValueError(f"Invalid deletion vector, referenced data file is
missing: {data_file.file_path}")
+
+
+def has_deletion_vector_content_reference(data_file: "DataFile") -> bool:
+ return (
+ data_file.content_offset is not None
+ or data_file.content_size_in_bytes is not None
+ or data_file.referenced_data_file is not None
+ )
+
+
+def _read_deletion_vector(io: "FileIO", data_file: "DataFile") ->
DeletionVector:
+ _validate_deletion_vector_content(data_file)
+
+ content_offset = data_file.content_offset
+ content_size_in_bytes = data_file.content_size_in_bytes
+ referenced_data_file = data_file.referenced_data_file
+ assert content_offset is not None
+ assert content_size_in_bytes is not None
+ assert referenced_data_file is not None
+
+ with io.new_input(data_file.file_path).open() as fi:
+ fi.seek(content_offset)
+ payload = fi.read(content_size_in_bytes)
+
+ if len(payload) != content_size_in_bytes:
+ raise ValueError(f"Could not read deletion vector, expected
{content_size_in_bytes} bytes, got {len(payload)}")
+
+ return DeletionVector(
+ referenced_data_file=referenced_data_file,
+ bitmaps=_deserialize_dv_blob(payload, data_file.record_count),
+ )
+
+
+def read_deletion_vectors(io: "FileIO", data_file: "DataFile") ->
list[DeletionVector]:
+ if has_deletion_vector_content_reference(data_file):
+ return [_read_deletion_vector(io, data_file)]
+
+ with io.new_input(data_file.file_path).open() as fi:
+ return deletion_vectors_from_puffin_file(PuffinFile(fi.read()))
Review Comment:
I don't think we need both branches. In both cases we are reading a single
DV in a given byte range.
##########
pyiceberg/table/deletion_vector.py:
##########
@@ -77,17 +89,103 @@ def to_vector(self) -> "pa.ChunkedArray":
return self._bitmaps_to_chunked_array(self._bitmaps)
-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")
- return blob_payload[8 : 4 + length_prefix]
+def _deserialize_dv_blob(blob: bytes, record_count: int | None = None) ->
list[BitMap]:
+ # The DV blob encoding matches Iceberg Java's BitmapPositionDeleteIndex:
+ # 4-byte big-endian bitmap-data length, 4-byte little-endian magic number,
+ # portable Roaring bitmap data, and 4-byte big-endian CRC-32.
+ if len(blob) < _DV_BLOB_MIN_SIZE_BYTES:
+ raise ValueError(f"Invalid deletion vector blob length: {len(blob)}")
+
+ bitmap_data_length = _DV_BLOB_LENGTH.unpack_from(blob)[0]
+ expected_bitmap_data_length = len(blob) - _DV_BLOB_LENGTH.size -
_DV_BLOB_CRC.size
+ if bitmap_data_length != expected_bitmap_data_length:
+ raise ValueError(f"Invalid bitmap data length: {bitmap_data_length},
expected {expected_bitmap_data_length}")
+
+ bitmap_data_offset = _DV_BLOB_LENGTH.size
+ crc_offset = bitmap_data_offset + bitmap_data_length
+ bitmap_data = blob[bitmap_data_offset:crc_offset]
+
+ magic_number = _DV_BLOB_MAGIC.unpack_from(bitmap_data)[0]
+ if magic_number != _DV_BLOB_MAGIC_NUMBER:
+ raise ValueError(f"Invalid magic number: {magic_number}, expected
{_DV_BLOB_MAGIC_NUMBER}")
+
+ checksum = zlib.crc32(bitmap_data) & 0xFFFFFFFF
+ expected_checksum = _DV_BLOB_CRC.unpack_from(blob, crc_offset)[0]
+ if checksum != expected_checksum:
+ raise ValueError("Invalid CRC")
+
+ bitmaps =
DeletionVector._deserialize_bitmap(bitmap_data[_DV_BLOB_MAGIC.size :])
+ if record_count is not None:
+ cardinality = sum(len(bitmap) for bitmap in bitmaps)
+ if cardinality != record_count:
+ raise ValueError(f"Invalid cardinality: {cardinality}, expected
{record_count}")
Review Comment:
I think this is fine, again as we expect these two values to be the same but
just remember implementations can choose to be a bit more relaxed (or vice
versa more strict) than the actual spec. Is it worth failing the read of the DV
if there's a mismatch? On one hand it indicates something incorrect in the
metadata, on the other hand, we could be blocking a read of the data
unnecessarily (because it wouldn't affect correctness of the result anyways).
So in this case I'd probably bias to the latter of not doing this check. But
I'll leave it up to you cc @kevinjqliu @rambleraptor in case you folks have
opinions here.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]