Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 23 additions & 2 deletions pyiceberg/manifest.py
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@
UNASSIGNED_SEQ = -1
DEFAULT_BLOCK_SIZE = 67108864 # 64 * 1024 * 1024
DEFAULT_READ_VERSION: Literal[2] = 2
_LATEST_MANIFEST_READ_VERSION: Literal[3] = 3

INITIAL_SEQUENCE_NUMBER = 0

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

@property
def first_row_id(self) -> int | None:
return self._data[16] if len(self._data) > 16 else None

@property
def referenced_data_file(self) -> str | None:
return self._data[17] if len(self._data) > 17 else None

@property
def content_offset(self) -> int | None:
return self._data[18] if len(self._data) > 18 else None

@property
def content_size_in_bytes(self) -> int | None:
return self._data[19] if len(self._data) > 19 else None

# Spec ID should not be stored in the file
_spec_id: int

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

@property
def first_row_id(self) -> int | None:
return self._data[15] if len(self._data) > 15 else None

def has_added_files(self) -> bool:
return self.added_files_count is None or self.added_files_count > 0

Expand All @@ -873,7 +894,7 @@ def fetch_manifest_entry(self, io: FileIO, discard_deleted: bool = True) -> list
input_file = io.new_input(self.manifest_path)
with AvroFile[ManifestEntry](
input_file,
MANIFEST_ENTRY_SCHEMAS[DEFAULT_READ_VERSION],
MANIFEST_ENTRY_SCHEMAS[_LATEST_MANIFEST_READ_VERSION],
read_types={-1: ManifestEntry, 2: DataFile},
read_enums={0: ManifestEntryStatus, 101: FileFormat, 134: DataFileContent},
) as reader:
Expand Down Expand Up @@ -996,7 +1017,7 @@ def read_manifest_list(input_file: InputFile) -> Iterator[ManifestFile]:
"""
with AvroFile[ManifestFile](
input_file,
MANIFEST_LIST_FILE_SCHEMAS[DEFAULT_READ_VERSION],
MANIFEST_LIST_FILE_SCHEMAS[_LATEST_MANIFEST_READ_VERSION],
read_types={-1: ManifestFile, 508: PartitionFieldSummary},
read_enums={517: ManifestContent},
) as reader:
Expand Down
8 changes: 7 additions & 1 deletion tests/avro/test_file.py
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,8 @@ def test_write_manifest_entry_with_iceberg_read_with_fastavro_v1() -> None:
del v2_entry["file_sequence_number"]
del v2_entry["data_file"]["content"]
del v2_entry["data_file"]["equality_ids"]
for field in ("first_row_id", "referenced_data_file", "content_offset", "content_size_in_bytes"):
del v2_entry["data_file"][field]

# Required in V1
v2_entry["data_file"]["block_size_in_bytes"] = DEFAULT_BLOCK_SIZE
Expand Down Expand Up @@ -222,7 +224,11 @@ def test_write_manifest_entry_with_iceberg_read_with_fastavro_v2() -> None:

fa_entry = next(it)

assert todict(entry) == fa_entry
v2_entry = todict(entry)
for field in ("first_row_id", "referenced_data_file", "content_offset", "content_size_in_bytes"):
del v2_entry["data_file"][field]

assert v2_entry == fa_entry


@pytest.mark.parametrize("format_version", [1, 2])
Expand Down
2 changes: 2 additions & 0 deletions tests/integration/test_rest_manifest.py
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,8 @@ def test_write_sample_manifest(table_test_all_types: Table, compression: AvroCom
wrapped_entry_v2 = copy(entry)
wrapped_entry_v2.data_file = wrapped_data_file_v2_debug
wrapped_entry_v2_dict = todict(wrapped_entry_v2, [field.name for field in test_spec.fields])
for field in ("first_row_id", "referenced_data_file", "content_offset", "content_size_in_bytes"):
del wrapped_entry_v2_dict["data_file"][field]

with TemporaryDirectory() as tmpdir:
tmp_avro_file = tmpdir + "/test_write_manifest.avro"
Expand Down
108 changes: 108 additions & 0 deletions tests/utils/test_manifest.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,12 @@

import pyiceberg.manifest as manifest_module
from pyiceberg.avro.codecs import AvroCompressionCodec
from pyiceberg.avro.file import AvroOutputFile
from pyiceberg.io import load_file_io
from pyiceberg.io.pyarrow import PyArrowFileIO
from pyiceberg.manifest import (
MANIFEST_ENTRY_SCHEMAS,
MANIFEST_LIST_FILE_SCHEMAS,
DataFile,
DataFileContent,
FileFormat,
Expand Down Expand Up @@ -92,6 +95,10 @@ def test_read_manifest_entry(generated_manifest_entry_file: str) -> None:
assert repr(data_file.partition) == "Record[1, 1925]"
assert data_file.record_count == 19513
assert data_file.file_size_in_bytes == 388872
assert data_file.first_row_id is None
assert data_file.referenced_data_file is None
assert data_file.content_offset is None
assert data_file.content_size_in_bytes is None
assert data_file.column_sizes == {
1: 53,
2: 98153,
Expand Down Expand Up @@ -192,6 +199,73 @@ def test_read_manifest_entry(generated_manifest_entry_file: str) -> None:
assert data_file.sort_order_id == 0


def test_read_manifest_entry_v3_fields(tmp_path: Path) -> None:
io = PyArrowFileIO()

def write_and_read(file_name: str, data_file: DataFile) -> DataFile:
manifest_path = str(tmp_path / file_name)
entry = ManifestEntry.from_args(
_table_format_version=3,
status=ManifestEntryStatus.ADDED,
snapshot_id=25,
sequence_number=1,
file_sequence_number=1,
data_file=data_file,
)
with AvroOutputFile[ManifestEntry](
output_file=io.new_output(manifest_path),
file_schema=MANIFEST_ENTRY_SCHEMAS[3],
record_schema=MANIFEST_ENTRY_SCHEMAS[3],
schema_name="manifest_entry",
metadata={"format-version": "3"},
) as writer:
writer.write_block([entry])

manifest = ManifestFile.from_args(
manifest_path=manifest_path,
manifest_length=0,
partition_spec_id=0,
added_snapshot_id=25,
sequence_number=1,
min_sequence_number=1,
)
return manifest.fetch_manifest_entry(io)[0].data_file

data_file = write_and_read(
"data-manifest.avro",
DataFile.from_args(
_table_format_version=3,
content=DataFileContent.DATA,
file_path="s3://bucket/data.parquet",
file_format=FileFormat.PARQUET,
partition=Record(),
record_count=10,
file_size_in_bytes=1024,
first_row_id=34,
),
)
assert data_file.first_row_id == 34

delete_file = write_and_read(
"delete-manifest.avro",
DataFile.from_args(
_table_format_version=3,
content=DataFileContent.POSITION_DELETES,
file_path="s3://bucket/deletes.puffin",
file_format=FileFormat.PUFFIN,
partition=Record(),
record_count=3,
file_size_in_bytes=47,
referenced_data_file="s3://bucket/data.parquet",
content_offset=1,
content_size_in_bytes=46,
),
)
assert delete_file.referenced_data_file == "s3://bucket/data.parquet"
assert delete_file.content_offset == 1
assert delete_file.content_size_in_bytes == 46


def test_read_manifest_list(generated_manifest_file_file_v1: str) -> None:
input_file = PyArrowFileIO().new_input(generated_manifest_file_file_v1)
manifest_list = list(read_manifest_list(input_file))[0]
Expand All @@ -216,6 +290,40 @@ def test_read_manifest_list(generated_manifest_file_file_v1: str) -> None:
assert manifest_list.added_rows_count == 237993
assert manifest_list.existing_rows_count == 0
assert manifest_list.deleted_rows_count == 0
assert manifest_list.first_row_id is None


def test_read_manifest_list_v3_fields(tmp_path: Path) -> None:
io = PyArrowFileIO()
path = str(tmp_path / "manifest-list.avro")
manifest = ManifestFile.from_args(
_table_format_version=3,
manifest_path="s3://bucket/manifest.avro",
manifest_length=1024,
partition_spec_id=0,
content=ManifestContent.DATA,
sequence_number=1,
min_sequence_number=1,
added_snapshot_id=25,
added_files_count=1,
existing_files_count=0,
deleted_files_count=0,
added_rows_count=10,
existing_rows_count=0,
deleted_rows_count=0,
first_row_id=34,
)
with AvroOutputFile[ManifestFile](
output_file=io.new_output(path),
file_schema=MANIFEST_LIST_FILE_SCHEMAS[3],
record_schema=MANIFEST_LIST_FILE_SCHEMAS[3],
schema_name="manifest_file",
metadata={"format-version": "3"},
) as writer:
writer.write_block([manifest])

read_manifest = list(read_manifest_list(io.new_input(path)))[0]
assert read_manifest.first_row_id == 34


def test_read_manifest_v1(generated_manifest_file_file_v1: str) -> None:
Expand Down