Skip to content

Commit 44312bd

Browse files
committed
Fix V3 manifest read projection
Read manifest entries and manifest lists with the latest supported schema so V3-only fields are retained while older manifests resolve missing fields to null.
1 parent 82be040 commit 44312bd

2 files changed

Lines changed: 131 additions & 2 deletions

File tree

pyiceberg/manifest.py

Lines changed: 23 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,7 @@
5656
UNASSIGNED_SEQ = -1
5757
DEFAULT_BLOCK_SIZE = 67108864 # 64 * 1024 * 1024
5858
DEFAULT_READ_VERSION: Literal[2] = 2
59+
_LATEST_MANIFEST_VERSION: Literal[3] = 3
5960

6061
INITIAL_SEQUENCE_NUMBER = 0
6162

@@ -532,6 +533,22 @@ def equality_ids(self) -> list[int] | None:
532533
def sort_order_id(self) -> int | None:
533534
return self._data[15]
534535

536+
@property
537+
def first_row_id(self) -> int | None:
538+
return self._data[16] if len(self._data) > 16 else None
539+
540+
@property
541+
def referenced_data_file(self) -> str | None:
542+
return self._data[17] if len(self._data) > 17 else None
543+
544+
@property
545+
def content_offset(self) -> int | None:
546+
return self._data[18] if len(self._data) > 18 else None
547+
548+
@property
549+
def content_size_in_bytes(self) -> int | None:
550+
return self._data[19] if len(self._data) > 19 else None
551+
535552
# Spec ID should not be stored in the file
536553
_spec_id: int
537554

@@ -853,6 +870,10 @@ def partitions(self) -> list[PartitionFieldSummary] | None:
853870
def key_metadata(self) -> bytes | None:
854871
return self._data[14]
855872

873+
@property
874+
def first_row_id(self) -> int | None:
875+
return self._data[15] if len(self._data) > 15 else None
876+
856877
def has_added_files(self) -> bool:
857878
return self.added_files_count is None or self.added_files_count > 0
858879

@@ -873,7 +894,7 @@ def fetch_manifest_entry(self, io: FileIO, discard_deleted: bool = True) -> list
873894
input_file = io.new_input(self.manifest_path)
874895
with AvroFile[ManifestEntry](
875896
input_file,
876-
MANIFEST_ENTRY_SCHEMAS[DEFAULT_READ_VERSION],
897+
MANIFEST_ENTRY_SCHEMAS[_LATEST_MANIFEST_VERSION],
877898
read_types={-1: ManifestEntry, 2: DataFile},
878899
read_enums={0: ManifestEntryStatus, 101: FileFormat, 134: DataFileContent},
879900
) as reader:
@@ -996,7 +1017,7 @@ def read_manifest_list(input_file: InputFile) -> Iterator[ManifestFile]:
9961017
"""
9971018
with AvroFile[ManifestFile](
9981019
input_file,
999-
MANIFEST_LIST_FILE_SCHEMAS[DEFAULT_READ_VERSION],
1020+
MANIFEST_LIST_FILE_SCHEMAS[_LATEST_MANIFEST_VERSION],
10001021
read_types={-1: ManifestFile, 508: PartitionFieldSummary},
10011022
read_enums={517: ManifestContent},
10021023
) as reader:

tests/utils/test_manifest.py

Lines changed: 108 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,9 +25,12 @@
2525

2626
import pyiceberg.manifest as manifest_module
2727
from pyiceberg.avro.codecs import AvroCompressionCodec
28+
from pyiceberg.avro.file import AvroOutputFile
2829
from pyiceberg.io import load_file_io
2930
from pyiceberg.io.pyarrow import PyArrowFileIO
3031
from pyiceberg.manifest import (
32+
MANIFEST_ENTRY_SCHEMAS,
33+
MANIFEST_LIST_FILE_SCHEMAS,
3134
DataFile,
3235
DataFileContent,
3336
FileFormat,
@@ -92,6 +95,10 @@ def test_read_manifest_entry(generated_manifest_entry_file: str) -> None:
9295
assert repr(data_file.partition) == "Record[1, 1925]"
9396
assert data_file.record_count == 19513
9497
assert data_file.file_size_in_bytes == 388872
98+
assert data_file.first_row_id is None
99+
assert data_file.referenced_data_file is None
100+
assert data_file.content_offset is None
101+
assert data_file.content_size_in_bytes is None
95102
assert data_file.column_sizes == {
96103
1: 53,
97104
2: 98153,
@@ -192,6 +199,73 @@ def test_read_manifest_entry(generated_manifest_entry_file: str) -> None:
192199
assert data_file.sort_order_id == 0
193200

194201

202+
def test_read_manifest_entry_v3_fields(tmp_path: Path) -> None:
203+
io = PyArrowFileIO()
204+
205+
def write_and_read(file_name: str, data_file: DataFile) -> DataFile:
206+
manifest_path = str(tmp_path / file_name)
207+
entry = ManifestEntry.from_args(
208+
_table_format_version=3,
209+
status=ManifestEntryStatus.ADDED,
210+
snapshot_id=25,
211+
sequence_number=1,
212+
file_sequence_number=1,
213+
data_file=data_file,
214+
)
215+
with AvroOutputFile[ManifestEntry](
216+
output_file=io.new_output(manifest_path),
217+
file_schema=MANIFEST_ENTRY_SCHEMAS[3],
218+
record_schema=MANIFEST_ENTRY_SCHEMAS[3],
219+
schema_name="manifest_entry",
220+
metadata={"format-version": "3"},
221+
) as writer:
222+
writer.write_block([entry])
223+
224+
manifest = ManifestFile.from_args(
225+
manifest_path=manifest_path,
226+
manifest_length=0,
227+
partition_spec_id=0,
228+
added_snapshot_id=25,
229+
sequence_number=1,
230+
min_sequence_number=1,
231+
)
232+
return manifest.fetch_manifest_entry(io)[0].data_file
233+
234+
data_file = write_and_read(
235+
"data-manifest.avro",
236+
DataFile.from_args(
237+
_table_format_version=3,
238+
content=DataFileContent.DATA,
239+
file_path="s3://bucket/data.parquet",
240+
file_format=FileFormat.PARQUET,
241+
partition=Record(),
242+
record_count=10,
243+
file_size_in_bytes=1024,
244+
first_row_id=34,
245+
),
246+
)
247+
assert data_file.first_row_id == 34
248+
249+
delete_file = write_and_read(
250+
"delete-manifest.avro",
251+
DataFile.from_args(
252+
_table_format_version=3,
253+
content=DataFileContent.POSITION_DELETES,
254+
file_path="s3://bucket/deletes.puffin",
255+
file_format=FileFormat.PUFFIN,
256+
partition=Record(),
257+
record_count=3,
258+
file_size_in_bytes=47,
259+
referenced_data_file="s3://bucket/data.parquet",
260+
content_offset=1,
261+
content_size_in_bytes=46,
262+
),
263+
)
264+
assert delete_file.referenced_data_file == "s3://bucket/data.parquet"
265+
assert delete_file.content_offset == 1
266+
assert delete_file.content_size_in_bytes == 46
267+
268+
195269
def test_read_manifest_list(generated_manifest_file_file_v1: str) -> None:
196270
input_file = PyArrowFileIO().new_input(generated_manifest_file_file_v1)
197271
manifest_list = list(read_manifest_list(input_file))[0]
@@ -216,6 +290,40 @@ def test_read_manifest_list(generated_manifest_file_file_v1: str) -> None:
216290
assert manifest_list.added_rows_count == 237993
217291
assert manifest_list.existing_rows_count == 0
218292
assert manifest_list.deleted_rows_count == 0
293+
assert manifest_list.first_row_id is None
294+
295+
296+
def test_read_manifest_list_v3_fields(tmp_path: Path) -> None:
297+
io = PyArrowFileIO()
298+
path = str(tmp_path / "manifest-list.avro")
299+
manifest = ManifestFile.from_args(
300+
_table_format_version=3,
301+
manifest_path="s3://bucket/manifest.avro",
302+
manifest_length=1024,
303+
partition_spec_id=0,
304+
content=ManifestContent.DATA,
305+
sequence_number=1,
306+
min_sequence_number=1,
307+
added_snapshot_id=25,
308+
added_files_count=1,
309+
existing_files_count=0,
310+
deleted_files_count=0,
311+
added_rows_count=10,
312+
existing_rows_count=0,
313+
deleted_rows_count=0,
314+
first_row_id=34,
315+
)
316+
with AvroOutputFile[ManifestFile](
317+
output_file=io.new_output(path),
318+
file_schema=MANIFEST_LIST_FILE_SCHEMAS[3],
319+
record_schema=MANIFEST_LIST_FILE_SCHEMAS[3],
320+
schema_name="manifest_file",
321+
metadata={"format-version": "3"},
322+
) as writer:
323+
writer.write_block([manifest])
324+
325+
read_manifest = list(read_manifest_list(io.new_input(path)))[0]
326+
assert read_manifest.first_row_id == 34
219327

220328

221329
def test_read_manifest_v1(generated_manifest_file_file_v1: str) -> None:

0 commit comments

Comments
 (0)