[
https://issues.apache.org/jira/browse/FLINK-18800?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Mohammad Hossein Gerami updated FLINK-18800:
--------------------------------------------
Description:
{color:#ff8b00}AvroSerializationSchema{color} and
{color:#ff8b00}ConfluentRegistryAvroSerializationSchema{color} doesn't support
Kafka key/value serialization. I implemented a custom Avro serialization schema
for solving this problem.
please consensus to implement new class to support kafka key/value
serialization.
for example in the Flink must implement a class like this:
{code:java}
public class KafkaAvroRegistrySchemaSerializationSchema extends
RegistryAvroSerializationSchema<GenericRecord> implements
KafkaSerializationSchema<GenericRecord>{code}
was:
{color:#ff8b00}AvroSerializationSchema{color} and
{color:#ff8b00}ConfluentRegistryAvroSerializationSchema{color} doesn't support
Kafka key/value serialization. I implemented a custom Avro serialization schema
for solving this problem.
for example in the Flink must implement a class like this.
{code:java}
public class KafkaAvroRegistrySchemaSerializationSchema extends
RegistryAvroSerializationSchema<GenericRecord> implements
KafkaSerializationSchema<GenericRecord>{code}
> Avro serialization schema doesn't support Kafka key/value serialization
> ------------------------------------------------------------------------
>
> Key: FLINK-18800
> URL: https://issues.apache.org/jira/browse/FLINK-18800
> Project: Flink
> Issue Type: Improvement
> Components: Connectors / Kafka, Formats (JSON, Avro, Parquet, ORC,
> SequenceFile)
> Affects Versions: 1.11.0, 1.11.1
> Reporter: Mohammad Hossein Gerami
> Priority: Major
>
> {color:#ff8b00}AvroSerializationSchema{color} and
> {color:#ff8b00}ConfluentRegistryAvroSerializationSchema{color} doesn't
> support Kafka key/value serialization. I implemented a custom Avro
> serialization schema for solving this problem.
> please consensus to implement new class to support kafka key/value
> serialization.
> for example in the Flink must implement a class like this:
> {code:java}
> public class KafkaAvroRegistrySchemaSerializationSchema extends
> RegistryAvroSerializationSchema<GenericRecord> implements
> KafkaSerializationSchema<GenericRecord>{code}
--
This message was sent by Atlassian Jira
(v8.3.4#803005)