Joey Pereira created FLINK-40370:
------------------------------------

             Summary: Bytes-backed decimal logical types corrupt PyFlink 
GenericRecordAvroTypeInfo JVM/Python serialization
                 Key: FLINK-40370
                 URL: https://issues.apache.org/jira/browse/FLINK-40370
             Project: Flink
          Issue Type: Bug
          Components: API / Python, Formats (JSON, Avro, Parquet, ORC, 
SequenceFile)
    Affects Versions: 2.1.3, 2.2.1, 2.3.0, 2.0.2
            Reporter: Joey Pereira


After FLINK-37192 replaced `avro-python3` with `avro>=1.12.0`, PyFlink began 
encoding bytes-backed Avro decimal logical types incompatibly with Flink’s 
internal Java serializer when generic records cross the JVM/Python boundary.

This does **not** affect standard Avro interoperability. Standard Python Avro, 
Java Avro, `fastavro`, and ordinary Kafka Avro payloads remain compatible. The 
regression affects Flink’s separate internal serializer for 
`GenericRecordAvroTypeInfo`.

For a record containing `amount = Decimal("12.34")` followed by `tail = 7`:

| Serialization boundary | Encoded bytes |
|---|---|
| Standard Python Avro and `fastavro` | `0404d20e` |
| Flink internal Java serializer | `0000000204d200000007` |
| Flink internal Python serializer with modern Avro | 
`000000000000000204d200000007` |

The actual Java `GenericDatumReader` successfully decodes the Python-generated 
bytes as a record with an empty decimal payload and `tail = 2`, leaving unread 
bytes. This is **silent data corruption**, not merely a decoding failure.

### Reproduction

```python
from decimal import Decimal
from io import BytesIO

import avro.io
import avro.schema
import fastavro
from pyflink.fn_execution.formats.avro import (
    FlinkAvroDatumReader,
    FlinkAvroDatumWriter,
    FlinkAvroDecoder,
    FlinkAvroEncoder,
)

schema_json = """
{
  "type": "record",
  "name": "DecimalRecord",
  "fields": [
    {
      "name": "amount",
      "type": {
        "type": "bytes",
        "logicalType": "decimal",
        "precision": 8,
        "scale": 2
      }
    },
    {"name": "tail", "type": "int"}
  ]
}
"""

schema = avro.schema.parse(schema_json)
record = {"amount": Decimal("12.34"), "tail": 7}

standard = BytesIO()
avro.io.DatumWriter(schema).write(
    record,
    avro.io.BinaryEncoder(standard),
)

fast = BytesIO()
fastavro.schemaless_writer(fast, schema.to_json(), record)

python_flink = BytesIO()
FlinkAvroDatumWriter(schema).write(
    record,
    FlinkAvroEncoder(python_flink),
)

java_flink = BytesIO()
java_compatible_encoder = FlinkAvroEncoder(java_flink)
java_compatible_encoder.write_bytes(bytes.fromhex("04d2"))
java_compatible_encoder.write_int(7)

print("Standard Python Avro:", standard.getvalue().hex())
print("Standard fastavro:", fast.getvalue().hex())
print("Expected Flink JVM:", java_flink.getvalue().hex())
print("Actual PyFlink:", python_flink.getvalue().hex())

try:
    result = FlinkAvroDatumReader(schema, schema).read(
        FlinkAvroDecoder(BytesIO(java_flink.getvalue()))
    )
    print("Python reads JVM:", result)
except Exception as exc:
    print("Python reads JVM:", type(exc).__name__, str(exc))

stream = BytesIO(python_flink.getvalue())
decoder = FlinkAvroDecoder(stream)

amount = decoder.read_bytes()
tail = decoder.read_int()

print(
    "JVM-compatible read:",
    {
        "amount_length": len(amount),
        "tail": tail,
        "unread": stream.read().hex(),
    },
)
```

Observed using Python `avro==1.12.1` and `fastavro==1.10.0`:

```text
Standard Python Avro: 0404d20e
Standard fastavro: 0404d20e
Expected Flink JVM: 0000000204d200000007
Actual PyFlink: 000000000000000204d200000007

Python reads JVM: InvalidAvroBinaryEncoding Read 2 bytes, expected 8670806016 
bytes

JVM-compatible read: {
    "amount_length": 0,
    "tail": 2,
    "unread": "04d200000007"
}
```

Additionally, compiling and running Flink’s actual Java `DataOutputEncoder`, 
`DataInputDecoder`, and Avro `GenericDatumReader` produced:

```text
actual Java DataOutputEncoder:
0000000204d200000007

actual Java DataInputDecoder:
amount_length=0, tail=2, unread=04d200000007

actual Java GenericDatumReader:
amount_length=0, tail=2, unread=04d200000007
```

### Expected behavior

Flink’s Python and Java internal generic-record serializers agree on the 
four-byte byte-array length prefix. Decimal-containing records round-trip 
without corrupting subsequent fields.

### Root cause

Modern Python Avro’s decimal-specialized writer calls `write_long`, and its 
decimal-specialized reader calls `read_long`.

PyFlink overrides those methods using **eight-byte** fixed-width operations. 
However, Flink’s Java `DataOutputEncoder.writeBytes` and 
`DataInputDecoder.readBytes` use **four-byte** lengths.

Consequently:

- Python to Java can silently corrupt the decimal and subsequent fields.
- Java to Python fails decoding.
- Python to Python succeeds, so symmetric Python-only tests miss the 
incompatibility.

### Affected workflows

- `GenericRecordAvroTypeInfo` with bytes-backed decimal fields.
- `AvroBulkWriters.for_generic_record(...)`.
- `AvroParquetWriters.for_generic_record(...)`.
- Generic Avro readers crossing from Java into Python.
- Explicitly Avro-typed Python state.

Ordinary Kafka Avro payloads, `fastavro`, standard Java Avro, plain bytes 
fields, and unrelated SQL decimal coders are unaffected.

### Suggested fix

Ensure PyFlink’s decimal-specialized encoder and decoder use Flink’s four-byte 
internal byte-array framing.

Because affected Flink versions may already have checkpointed the eight-byte 
representation in explicitly Avro-typed Python state, the decoder should 
consider accepting both historical eight-byte and corrected four-byte 
representations, while new writes always use four-byte framing.

Add coverage for bidirectional JVM/Python compatibility, adjacent-field 
corruption, negative and high-precision decimals, nested and nullable records, 
Avro/Parquet sinks, both Python/Cython coder implementations, and previously 
persisted Python state.

### References

- [FLINK-37192: Dependency change introducing the 
regression](https://issues.apache.org/jira/browse/FLINK-37192)
- [PyFlink custom decimal-related 
encoder](https://github.com/apache/flink/blob/release-2.3.0/flink-python/pyflink/fn_execution/formats/avro.py#L174-L199)
- [PyFlink custom 
decoder](https://github.com/apache/flink/blob/release-2.3.0/flink-python/pyflink/fn_execution/formats/avro.py#L50-L88)
- [Java internal 
encoder](https://github.com/apache/flink/blob/release-2.3.0/flink-formats/flink-avro/src/main/java/org/apache/flink/formats/avro/utils/DataOutputEncoder.java#L86-L101)
- [Java internal 
decoder](https://github.com/apache/flink/blob/release-2.3.0/flink-formats/flink-avro/src/main/java/org/apache/flink/formats/avro/utils/DataInputDecoder.java#L85-L97)
- [Apache Avro decimal 
writer](https://github.com/apache/avro/blob/release-1.12.1/lang/py/avro/io.py#L468-L493)
- [Avro generic-record file 
sink](https://github.com/apache/flink/blob/release-2.3.0/flink-python/pyflink/datastream/formats/avro.py#L132-L166)
- [Parquet generic-record file 
sink](https://github.com/apache/flink/blob/release-2.3.0/flink-python/pyflink/datastream/formats/parquet.py#L76-L110)



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to