dalelane commented on code in PR #28932:
URL: https://github.com/apache/flink/pull/28932#discussion_r4207195503
##########
flink-formats/flink-avro/src/main/java/org/apache/flink/formats/avro/AvroDeserializationSchema.java:
##########
@@ -181,7 +181,33 @@ public T deserialize(@Nullable byte[] message) throws
IOException {
((JsonDecoder) this.decoder).configure(inputStream);
}
- return datumReader.read(null, decoder);
+ try {
+ return datumReader.read(null, decoder);
+ } catch (IOException | RuntimeException e) {
+ // FLINK-34474: a failed read can leave the pooled decoder in an
+ // inconsistent internal state, so that subsequent reads return
+ // corrupted data even after the input buffer is reset. Discard
+ // the poisoned decoder so the next message starts from a clean
state.
+ resetDecoder();
+ throw e;
+ }
+ }
+
+ private void resetDecoder() {
+ try {
+ if (encoding == AvroEncoding.JSON) {
+ this.decoder =
DecoderFactory.get().jsonDecoder(getReaderSchema(), inputStream);
Review Comment:
Is this needed? `configure()` already resets the JsonDecoder on each call,
and without this change I would expect JSON to recover after corrupt input.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]