Skip to content

Commit cd846f3

Browse files
committed
feat(manifest): support writing delete manifests in v2
ManifestWriterV2 hardcoded content() to ManifestContent.DATA and wrote "content": "data" into the Avro metadata, so PyIceberg could not produce a delete manifest even though it reads them and ManifestListWriterV1 already guards against storing one in a v1 table. Thread an optional content through write_manifest and ManifestWriterV2, defaulting to DATA so existing callers are unaffected, and reject a non-data content for format version 1.
1 parent 0d58407 commit cd846f3

2 files changed

Lines changed: 112 additions & 3 deletions

File tree

‎pyiceberg/manifest.py‎

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1297,18 +1297,22 @@ def prepare_entry(self, entry: ManifestEntry) -> ManifestEntry:
12971297

12981298

12991299
class ManifestWriterV2(ManifestWriter):
1300+
_content: ManifestContent
1301+
13001302
def __init__(
13011303
self,
13021304
spec: PartitionSpec,
13031305
schema: Schema,
13041306
output_file: OutputFile,
13051307
snapshot_id: int,
13061308
avro_compression: AvroCompressionCodec,
1309+
content: ManifestContent = ManifestContent.DATA,
13071310
):
13081311
super().__init__(spec, schema, output_file, snapshot_id, avro_compression)
1312+
self._content = content
13091313

13101314
def content(self) -> ManifestContent:
1311-
return ManifestContent.DATA
1315+
return self._content
13121316

13131317
@property
13141318
def version(self) -> TableVersion:
@@ -1318,7 +1322,7 @@ def version(self) -> TableVersion:
13181322
def _meta(self) -> dict[str, str]:
13191323
return {
13201324
**super()._meta,
1321-
"content": "data",
1325+
"content": "data" if self._content == ManifestContent.DATA else "deletes",
13221326
}
13231327

13241328
def prepare_entry(self, entry: ManifestEntry) -> ManifestEntry:
@@ -1337,11 +1341,14 @@ def write_manifest(
13371341
output_file: OutputFile,
13381342
snapshot_id: int,
13391343
avro_compression: AvroCompressionCodec,
1344+
content: ManifestContent = ManifestContent.DATA,
13401345
) -> ManifestWriter:
13411346
if format_version == 1:
1347+
if content != ManifestContent.DATA:
1348+
raise ValidationError("Cannot write delete manifests in a v1 table")
13421349
return ManifestWriterV1(spec, schema, output_file, snapshot_id, avro_compression)
13431350
elif format_version == 2:
1344-
return ManifestWriterV2(spec, schema, output_file, snapshot_id, avro_compression)
1351+
return ManifestWriterV2(spec, schema, output_file, snapshot_id, avro_compression, content)
13451352
else:
13461353
raise ValueError(f"Cannot write manifest for table version: {format_version}")
13471354

‎tests/utils/test_manifest.py‎

Lines changed: 102 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@
2626
import pyiceberg.manifest as manifest_module
2727
from pyiceberg.avro.codecs import AvroCompressionCodec
2828
from pyiceberg.avro.file import AvroFile, AvroOutputFile
29+
from pyiceberg.exceptions import ValidationError
2930
from pyiceberg.io import load_file_io
3031
from pyiceberg.io.pyarrow import PyArrowFileIO
3132
from pyiceberg.manifest import (
@@ -1376,3 +1377,104 @@ def test_negative_manifest_cache_size_raises_value_error(monkeypatch: pytest.Mon
13761377
finally:
13771378
monkeypatch.delenv("PYICEBERG_MANIFEST_CACHE_SIZE", raising=False)
13781379
importlib.reload(manifest_module)
1380+
1381+
1382+
@pytest.mark.parametrize("content", [ManifestContent.DATA, ManifestContent.DELETES])
1383+
def test_write_manifest_content(
1384+
generated_manifest_file_file_v2: str,
1385+
test_schema: Schema,
1386+
test_partition_spec: PartitionSpec,
1387+
content: ManifestContent,
1388+
) -> None:
1389+
"""A v2 manifest must record the content it was asked to write.
1390+
1391+
Without this the writer always claims `data`, so a delete manifest written by
1392+
PyIceberg would be read back as a data manifest and its entries applied as
1393+
live data files.
1394+
"""
1395+
io = load_file_io()
1396+
snapshot = Snapshot(
1397+
snapshot_id=25,
1398+
parent_snapshot_id=19,
1399+
timestamp_ms=1602638573590,
1400+
manifest_list=generated_manifest_file_file_v2,
1401+
summary=Summary(Operation.APPEND),
1402+
schema_id=3,
1403+
)
1404+
manifest_entries = snapshot.manifests(io)[0].fetch_manifest_entry(io)
1405+
1406+
with TemporaryDirectory() as tmpdir:
1407+
tmp_avro_file = tmpdir + "/test_write_manifest_content.avro"
1408+
with write_manifest(
1409+
format_version=2,
1410+
spec=test_partition_spec,
1411+
schema=test_schema,
1412+
output_file=io.new_output(tmp_avro_file),
1413+
snapshot_id=8744736658442914487,
1414+
avro_compression="deflate",
1415+
content=content,
1416+
) as writer:
1417+
for entry in manifest_entries:
1418+
writer.add_entry(entry)
1419+
new_manifest = writer.to_manifest_file()
1420+
1421+
assert new_manifest.content == content
1422+
_verify_metadata_with_fastavro(
1423+
tmp_avro_file,
1424+
{"content": "data" if content == ManifestContent.DATA else "deletes"},
1425+
)
1426+
1427+
1428+
def test_write_manifest_defaults_to_data_content(
1429+
generated_manifest_file_file_v2: str,
1430+
test_schema: Schema,
1431+
test_partition_spec: PartitionSpec,
1432+
) -> None:
1433+
"""Callers that do not ask for a content type still get a data manifest."""
1434+
io = load_file_io()
1435+
snapshot = Snapshot(
1436+
snapshot_id=25,
1437+
parent_snapshot_id=19,
1438+
timestamp_ms=1602638573590,
1439+
manifest_list=generated_manifest_file_file_v2,
1440+
summary=Summary(Operation.APPEND),
1441+
schema_id=3,
1442+
)
1443+
manifest_entries = snapshot.manifests(io)[0].fetch_manifest_entry(io)
1444+
1445+
with TemporaryDirectory() as tmpdir:
1446+
tmp_avro_file = tmpdir + "/test_write_manifest_default_content.avro"
1447+
with write_manifest(
1448+
format_version=2,
1449+
spec=test_partition_spec,
1450+
schema=test_schema,
1451+
output_file=io.new_output(tmp_avro_file),
1452+
snapshot_id=8744736658442914487,
1453+
avro_compression="deflate",
1454+
) as writer:
1455+
for entry in manifest_entries:
1456+
writer.add_entry(entry)
1457+
assert writer.to_manifest_file().content == ManifestContent.DATA
1458+
1459+
1460+
def test_write_manifest_v1_rejects_delete_content(
1461+
test_schema: Schema,
1462+
test_partition_spec: PartitionSpec,
1463+
) -> None:
1464+
"""v1 has no delete files, so asking for a delete manifest must fail loudly.
1465+
1466+
`ManifestListWriterV1` already refuses to store one; rejecting it at the
1467+
writer keeps a v1 table from producing a manifest it could never reference.
1468+
"""
1469+
io = load_file_io()
1470+
with TemporaryDirectory() as tmpdir:
1471+
with pytest.raises(ValidationError, match="Cannot write delete manifests in a v1 table"):
1472+
write_manifest(
1473+
format_version=1,
1474+
spec=test_partition_spec,
1475+
schema=test_schema,
1476+
output_file=io.new_output(tmpdir + "/test_write_manifest_v1_deletes.avro"),
1477+
snapshot_id=8744736658442914487,
1478+
avro_compression="deflate",
1479+
content=ManifestContent.DELETES,
1480+
)

0 commit comments

Comments
 (0)