[ 
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)

Reply via email to