Skip to content

Commit b52f97b

Browse files
committed
Tolerate legacy long equality_ids in resolver and write int in manifest (#3840)
1 parent eabbb72 commit b52f97b

5 files changed

Lines changed: 169 additions & 116 deletions

File tree

pyiceberg/avro/resolver.py

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -14,8 +14,6 @@
1414
# KIND, either express or implied. See the License for the
1515
# specific language governing permissions and limitations
1616
# under the License.
17-
# pylint: disable=arguments-renamed,unused-argument
18-
import warnings
1917
from collections.abc import Callable
2018
from enum import Enum
2119

@@ -379,6 +377,9 @@ def visit_geography(self, geography_type: "GeographyType", partner: IcebergType
379377
return BinaryWriter()
380378

381379

380+
_MANIFEST_DATA_FILE_EQUALITY_IDS_ELEMENT_ID = 136
381+
382+
382383
class ReadSchemaResolver(PrimitiveWithPartnerVisitor[IcebergType, Reader]):
383384
__slots__ = ("read_types", "read_enums", "context")
384385
read_types: dict[int, Callable[..., StructProtocol]]
@@ -462,13 +463,14 @@ def primitive(self, primitive: PrimitiveType, expected_primitive: IcebergType |
462463

463464
# ensure that the type can be projected to the expected
464465
if primitive != expected_primitive:
465-
if isinstance(primitive, LongType) and isinstance(expected_primitive, IntegerType):
466-
warnings.warn(
467-
"Encountered non-compliant manifest with long equality_ids (spec requires int). "
468-
"Support for legacy long equality_ids is deprecated and will be removed in a future release.",
469-
DeprecationWarning,
470-
stacklevel=2,
471-
)
466+
is_manifest_data_file = getattr(self.read_types.get(2), "__name__", "") == "DataFile"
467+
if (
468+
is_manifest_data_file
469+
and _MANIFEST_DATA_FILE_EQUALITY_IDS_ELEMENT_ID in self.context
470+
and isinstance(primitive, LongType)
471+
and isinstance(expected_primitive, IntegerType)
472+
):
473+
pass
472474
else:
473475
promote(primitive, expected_primitive)
474476

pyiceberg/manifest.py

Lines changed: 18 additions & 39 deletions
Original file line numberDiff line numberDiff line change
@@ -433,9 +433,7 @@ def __repr__(self) -> str:
433433
}
434434

435435

436-
def data_file_with_partition(
437-
partition_type: StructType, format_version: TableVersion, legacy_equality_ids: bool = False
438-
) -> StructType:
436+
def data_file_with_partition(partition_type: StructType, format_version: TableVersion) -> StructType:
439437
data_file_partition_type = StructType(
440438
*[
441439
NestedField(
@@ -448,32 +446,20 @@ def data_file_with_partition(
448446
]
449447
)
450448

451-
fields = []
452-
for field in DATA_FILE_TYPE[format_version].fields:
453-
if field.field_id == 102:
454-
fields.append(
455-
NestedField(
456-
field_id=102,
457-
name="partition",
458-
field_type=data_file_partition_type,
459-
required=True,
460-
doc="Partition data tuple, schema based on the partition spec",
461-
)
462-
)
463-
elif field.field_id == 135 and legacy_equality_ids:
464-
fields.append(
465-
NestedField(
466-
field_id=135,
467-
name="equality_ids",
468-
field_type=ListType(element_id=136, element_type=LongType(), element_required=True),
469-
required=False,
470-
doc="Field ids used to determine row equality in equality delete files.",
471-
)
449+
return StructType(
450+
*[
451+
NestedField(
452+
field_id=102,
453+
name="partition",
454+
field_type=data_file_partition_type,
455+
required=True,
456+
doc="Partition data tuple, schema based on the partition spec",
472457
)
473-
else:
474-
fields.append(field)
475-
476-
return StructType(*fields)
458+
if field.field_id == 102
459+
else field
460+
for field in DATA_FILE_TYPE[format_version].fields
461+
]
462+
)
477463

478464

479465
class DataFile(Record):
@@ -1072,7 +1058,6 @@ class ManifestWriter(ABC):
10721058
_min_sequence_number: int | None
10731059
_partitions: list[Record]
10741060
_compression: AvroCompressionCodec
1075-
_legacy_equality_ids: bool
10761061

10771062
def __init__(
10781063
self,
@@ -1081,7 +1066,6 @@ def __init__(
10811066
output_file: OutputFile,
10821067
snapshot_id: int,
10831068
avro_compression: AvroCompressionCodec,
1084-
legacy_equality_ids: bool = False,
10851069
) -> None:
10861070
self.closed = False
10871071
self._spec = spec
@@ -1098,7 +1082,6 @@ def __init__(
10981082
self._min_sequence_number = None
10991083
self._partitions = []
11001084
self._compression = avro_compression
1101-
self._legacy_equality_ids = legacy_equality_ids
11021085

11031086
def __enter__(self) -> ManifestWriter:
11041087
"""Open the writer."""
@@ -1144,7 +1127,6 @@ def _with_partition(self, format_version: TableVersion) -> Schema:
11441127
data_file_type = data_file_with_partition(
11451128
format_version=format_version,
11461129
partition_type=self._spec.partition_type(self._schema),
1147-
legacy_equality_ids=self._legacy_equality_ids,
11481130
)
11491131
return manifest_entry_schema_with_data_file(format_version=format_version, data_file=data_file_type)
11501132

@@ -1257,9 +1239,8 @@ def __init__(
12571239
output_file: OutputFile,
12581240
snapshot_id: int,
12591241
avro_compression: AvroCompressionCodec,
1260-
legacy_equality_ids: bool = False,
12611242
):
1262-
super().__init__(spec, schema, output_file, snapshot_id, avro_compression, legacy_equality_ids=legacy_equality_ids)
1243+
super().__init__(spec, schema, output_file, snapshot_id, avro_compression)
12631244

12641245
def content(self) -> ManifestContent:
12651246
return ManifestContent.DATA
@@ -1280,9 +1261,8 @@ def __init__(
12801261
output_file: OutputFile,
12811262
snapshot_id: int,
12821263
avro_compression: AvroCompressionCodec,
1283-
legacy_equality_ids: bool = False,
12841264
):
1285-
super().__init__(spec, schema, output_file, snapshot_id, avro_compression, legacy_equality_ids=legacy_equality_ids)
1265+
super().__init__(spec, schema, output_file, snapshot_id, avro_compression)
12861266

12871267
def content(self) -> ManifestContent:
12881268
return ManifestContent.DATA
@@ -1314,12 +1294,11 @@ def write_manifest(
13141294
output_file: OutputFile,
13151295
snapshot_id: int,
13161296
avro_compression: AvroCompressionCodec,
1317-
legacy_equality_ids: bool = False,
13181297
) -> ManifestWriter:
13191298
if format_version == 1:
1320-
return ManifestWriterV1(spec, schema, output_file, snapshot_id, avro_compression, legacy_equality_ids=legacy_equality_ids)
1299+
return ManifestWriterV1(spec, schema, output_file, snapshot_id, avro_compression)
13211300
elif format_version == 2:
1322-
return ManifestWriterV2(spec, schema, output_file, snapshot_id, avro_compression, legacy_equality_ids=legacy_equality_ids)
1301+
return ManifestWriterV2(spec, schema, output_file, snapshot_id, avro_compression)
13231302
else:
13241303
raise ValueError(f"Cannot write manifest for table version: {format_version}")
13251304

pyiceberg/table/__init__.py

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -167,9 +167,6 @@ class TableProperties:
167167
WRITE_AVRO_COMPRESSION = "write.avro.compression-codec"
168168
WRITE_AVRO_COMPRESSION_DEFAULT = "gzip"
169169

170-
WRITE_MANIFEST_LEGACY_LONG_EQUALITY_IDS = "write.manifest.legacy-long-equality-ids"
171-
WRITE_MANIFEST_LEGACY_LONG_EQUALITY_IDS_DEFAULT = False
172-
173170
DEFAULT_WRITE_METRICS_MODE = "write.metadata.metrics.default"
174171
DEFAULT_WRITE_METRICS_MODE_DEFAULT = "truncate(16)"
175172

tests/avro/test_resolver.py

Lines changed: 85 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -409,3 +409,88 @@ def test_writer_missing_optional_in_read_schema() -> None:
409409
expected = StructWriter(field_writers=((None, OptionWriter(option=StringWriter())),))
410410

411411
assert actual == expected
412+
413+
414+
def test_resolver_long_to_int_isolated_to_manifest_equality_ids() -> None:
415+
from pyiceberg.manifest import DataFile
416+
417+
# 1. User table schema: Long -> Int must fail even if field ID is 136
418+
user_file_schema = Schema(NestedField(field_id=136, name="col", field_type=LongType(), required=True))
419+
user_read_schema = Schema(NestedField(field_id=136, name="col", field_type=IntegerType(), required=True))
420+
421+
with pytest.raises(ResolveError, match="Cannot promote long to int"):
422+
resolve_reader(file_schema=user_file_schema, read_schema=user_read_schema)
423+
424+
# 2. Manifest DataFile schema with equality_ids (field 135, element 136): Long -> Int must succeed
425+
manifest_file_schema = Schema(
426+
NestedField(
427+
field_id=2,
428+
name="data_file",
429+
field_type=StructType(
430+
NestedField(
431+
field_id=135,
432+
name="equality_ids",
433+
field_type=ListType(element_id=136, element_type=LongType(), element_required=True),
434+
required=False,
435+
)
436+
),
437+
required=True,
438+
)
439+
)
440+
manifest_read_schema = Schema(
441+
NestedField(
442+
field_id=2,
443+
name="data_file",
444+
field_type=StructType(
445+
NestedField(
446+
field_id=135,
447+
name="equality_ids",
448+
field_type=ListType(element_id=136, element_type=IntegerType(), element_required=True),
449+
required=False,
450+
)
451+
),
452+
required=True,
453+
)
454+
)
455+
456+
# Without read_types={2: DataFile}, it fails because it's not recognized as a manifest data_file
457+
with pytest.raises(ResolveError, match="Cannot promote long to int"):
458+
resolve_reader(file_schema=manifest_file_schema, read_schema=manifest_read_schema)
459+
460+
# With read_types={2: DataFile}, equality_ids Long -> Int succeeds!
461+
reader = resolve_reader(
462+
file_schema=manifest_file_schema,
463+
read_schema=manifest_read_schema,
464+
read_types={2: DataFile},
465+
)
466+
assert reader is not None
467+
468+
# 3. Manifest DataFile schema with other field (e.g. record_count field 103):
469+
# Long -> Int must FAIL even with read_types={2: DataFile}
470+
manifest_record_count_file_schema = Schema(
471+
NestedField(
472+
field_id=2,
473+
name="data_file",
474+
field_type=StructType(
475+
NestedField(field_id=103, name="record_count", field_type=LongType(), required=True),
476+
),
477+
required=True,
478+
)
479+
)
480+
manifest_record_count_read_schema = Schema(
481+
NestedField(
482+
field_id=2,
483+
name="data_file",
484+
field_type=StructType(
485+
NestedField(field_id=103, name="record_count", field_type=IntegerType(), required=True),
486+
),
487+
required=True,
488+
)
489+
)
490+
491+
with pytest.raises(ResolveError, match="Cannot promote long to int"):
492+
resolve_reader(
493+
file_schema=manifest_record_count_file_schema,
494+
read_schema=manifest_record_count_read_schema,
495+
read_types={2: DataFile},
496+
)

0 commit comments

Comments
 (0)