Skip to content

Commit 71d6eff

Browse files
committed
fix(avro): advance the decoder when skipping an enum
EnumReader.skip() was a no-op, so the enum's bytes stayed in the stream and the next field was decoded from them. Reading a manifest with the status field projected out returns wrong values for every field after it. Delegate to the wrapped reader, which is what read() already does.
1 parent 0d58407 commit 71d6eff

2 files changed

Lines changed: 33 additions & 3 deletions

File tree

‎pyiceberg/avro/resolver.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -272,7 +272,7 @@ def read(self, decoder: BinaryDecoder) -> Enum:
272272
return self.enum(self.reader.read(decoder))
273273

274274
def skip(self, decoder: BinaryDecoder) -> None:
275-
pass
275+
self.reader.skip(decoder)
276276

277277

278278
class WriteSchemaResolver(PrimitiveWithPartnerVisitor[IcebergType, Writer]):

‎tests/avro/test_resolver.py‎

Lines changed: 32 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,11 +15,14 @@
1515
# specific language governing permissions and limitations
1616
# under the License.
1717

18+
from collections.abc import Callable
1819
from tempfile import TemporaryDirectory
1920

2021
import pytest
2122
from pydantic import Field
2223

24+
from pyiceberg.avro.decoder import BinaryDecoder, StreamingBinaryDecoder
25+
from pyiceberg.avro.decoder_fast import CythonBinaryDecoder
2326
from pyiceberg.avro.file import AvroFile
2427
from pyiceberg.avro.reader import (
2528
DecimalReader,
@@ -32,7 +35,7 @@
3235
StringReader,
3336
StructReader,
3437
)
35-
from pyiceberg.avro.resolver import resolve_reader, resolve_writer
38+
from pyiceberg.avro.resolver import EnumReader, resolve_reader, resolve_writer
3639
from pyiceberg.avro.writer import (
3740
BinaryWriter,
3841
DefaultWriter,
@@ -46,7 +49,7 @@
4649
)
4750
from pyiceberg.exceptions import ResolveError
4851
from pyiceberg.io.pyarrow import PyArrowFileIO
49-
from pyiceberg.manifest import MANIFEST_ENTRY_SCHEMAS
52+
from pyiceberg.manifest import MANIFEST_ENTRY_SCHEMAS, ManifestEntryStatus
5053
from pyiceberg.schema import Schema
5154
from pyiceberg.typedef import Record
5255
from pyiceberg.types import (
@@ -418,3 +421,30 @@ def test_writer_missing_optional_in_read_schema() -> None:
418421
expected = StructWriter(field_writers=((None, OptionWriter(option=StringWriter())),))
419422

420423
assert actual == expected
424+
425+
426+
@pytest.mark.parametrize("decoder_class", [StreamingBinaryDecoder, CythonBinaryDecoder])
427+
def test_enum_reader_skip_advances_the_decoder(decoder_class: Callable[[bytes], BinaryDecoder]) -> None:
428+
"""Skipping an enum field must consume its bytes, or the next field is read from them."""
429+
# Two ints, zigzag encoded: the enum's ordinal 1, then the next field's value 12.
430+
decoder = decoder_class(b"\x02\x18")
431+
reader = EnumReader(ManifestEntryStatus, IntegerReader())
432+
433+
reader.skip(decoder)
434+
435+
assert IntegerReader().read(decoder) == 12
436+
437+
438+
@pytest.mark.parametrize("decoder_class", [StreamingBinaryDecoder, CythonBinaryDecoder])
439+
def test_enum_reader_skip_matches_read(decoder_class: Callable[[bytes], BinaryDecoder]) -> None:
440+
"""Reading and skipping must leave the decoder at the same position."""
441+
encoded = b"\x02\x18"
442+
reader = EnumReader(ManifestEntryStatus, IntegerReader())
443+
444+
read_decoder = decoder_class(encoded)
445+
assert reader.read(read_decoder) == ManifestEntryStatus.ADDED
446+
447+
skip_decoder = decoder_class(encoded)
448+
reader.skip(skip_decoder)
449+
450+
assert IntegerReader().read(read_decoder) == IntegerReader().read(skip_decoder)

0 commit comments

Comments
 (0)