[
https://issues.apache.org/jira/browse/FLINK-16048?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=17036670#comment-17036670
]
Leonard Xu commented on FLINK-16048:
------------------------------------
Compare to AvroRowSerializationSchema[1], KafkaAvroSerializer[2] writes
external magic bytes and id field in result byte array. So,the example code[3]
can not read correct record and throws following Exception:
{code:java}
Caused by: java.io.IOException: Failed to deserialize Avro record.Caused by:
java.io.IOException: Failed to deserialize Avro record. at
org.apache.flink.formats.avro.AvroRowDeserializationSchema.deserialize(AvroRowDeserializationSchema.java:170)
at
org.apache.flink.formats.avro.AvroRowDeserializationSchema.deserialize(AvroRowDeserializationSchema.java:78)
at
org.apache.flink.streaming.connectors.kafka.internals.KafkaDeserializationSchemaWrapper.deserialize(KafkaDeserializationSchemaWrapper.java:45)
at
org.apache.flink.streaming.connectors.kafka.internal.Kafka010Fetcher.runFetchLoop(Kafka010Fetcher.java:146)
at
org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerBase.run(FlinkKafkaConsumerBase.java:715)
at
org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:100)
at
org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:63)
at
org.apache.flink.streaming.runtime.tasks.SourceStreamTask$LegacySourceFunctionThread.run(SourceStreamTask.java:208)Caused
by: org.apache.avro.AvroRuntimeException: Malformed data. Length is negative:
-1 at org.apache.avro.io.BinaryDecoder.doReadBytes(BinaryDecoder.java:336) at
org.apache.avro.io.BinaryDecoder.readString(BinaryDecoder.java:263) at
org.apache.avro.io.ResolvingDecoder.readString(ResolvingDecoder.java:201) at
org.apache.avro.generic.GenericDatumReader.readString(GenericDatumReader.java:422)
at
org.apache.avro.generic.GenericDatumReader.readString(GenericDatumReader.java:414)
at
org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:181)
at
org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:153) at
org.apache.avro.generic.GenericDatumReader.readField(GenericDatumReader.java:232)
at
org.apache.avro.specific.SpecificDatumReader.readField(SpecificDatumReader.java:122)
at
org.apache.avro.generic.GenericDatumReader.readRecord(GenericDatumReader.java:222)
at
org.apache.avro.generic.GenericDatumReader.readWithoutConversion(GenericDatumReader.java:175)
at
org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:153) at
org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:145) at
org.apache.flink.formats.avro.AvroRowDeserializationSchema.deserialize(AvroRowDeserializationSchema.java:167)
... 7 more
{code}
[1] confluent avro format serialize():
[https://github.com/confluentinc/schema-registry/blob/c19da74f7c4326438dd96d022a7e3d38dd3c68af/avro-ser[1]ializer/src/main/java/io/confluent/kafka/serializers/AbstractKafkaAvroSerializer.java#L58|https://github.com/confluentinc/schema-registry/blob/c19da74f7c4326438dd96d022a7e3d38dd3c68af/avro-serializer/src/main/java/io/confluent/kafka/serializers/AbstractKafkaAvroSerializer.java#L58]
[2] Flink Table avro format serialize() :
[https://github.com/apache/flink/blob/master/flink-formats/flink-avro/src/main/java/org/apache/flink/formats/avro/AvroRowSerializationSchema.java#L138]
[3][https://github.com/leonardBang/flink-sql-etl/blob/master/etl-job/src/main/java/kafka2kafka_1/ConsumeConfluentAvroTest.java]
> Support read/write confluent schema registry avro data from Kafka
> ------------------------------------------------------------------
>
> Key: FLINK-16048
> URL: https://issues.apache.org/jira/browse/FLINK-16048
> Project: Flink
> Issue Type: Improvement
> Components: Formats (JSON, Avro, Parquet, ORC, SequenceFile)
> Affects Versions: 1.11.0
> Reporter: Leonard Xu
> Priority: Major
> Fix For: 1.11.0
>
>
> KafkaAvroSerializer and AvroRowSerializationSchema
> I found SQL Kafka connector can not consume avro data that was serialized by
> `KafkaAvroSerializer` and only can consume Row data with avro schema because
> we use `AvroRowDeserializationSchema/AvroRowSerializationSchema` to se/de
> data in `AvroRowFormatFactory`.
> I think we should support this because `KafkaAvroSerializer` is very common
> in Kafka.
> and someone met same question in stackoverflow[1].
> [[1]https://stackoverflow.com/questions/56452571/caused-by-org-apache-avro-avroruntimeexception-malformed-data-length-is-negat/56478259|https://stackoverflow.com/questions/56452571/caused-by-org-apache-avro-avroruntimeexception-malformed-data-length-is-negat/56478259]
--
This message was sent by Atlassian Jira
(v8.3.4#803005)