Flink Avro序列化使用GenericRecord时遇非序列化错误求助
解决Flink与Confluent Schema Registry Avro序列化的不可序列化问题
你碰到的这个报错核心原因其实很明确:Avro的Schema类本身并没有实现Serializable接口,而你的GenericSerializer直接持有了schemaKey和schemaValue这两个Schema对象。Flink在提交作业时会序列化所有算子相关的对象(包括自定义序列化器),当它尝试序列化Schema对象时,就触发了NotSerializableException——这就是报错里提到的org.apache.avro.Schema$Field不可序列化的根源。
下面是针对你的Flink 1.9.1版本的具体解决方案:
方案1:改用RichKafkaSerializationSchema并标记Schema为transient
Flink的RichKafkaSerializationSchema继承了普通的KafkaSerializationSchema,同时提供了生命周期方法(比如open),我们可以利用这个特性延迟加载Schema,彻底避免序列化Schema对象的问题:
修改后的GenericSerializer代码
package com.reeeliance.flink; import org.apache.avro.Schema; import org.apache.avro.generic.GenericRecord; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.connectors.kafka.RichKafkaSerializationSchema; import org.apache.kafka.clients.producer.ProducerRecord; import flinkfix.ConfluentRegistryAvroSerializationSchema; public class GenericSerializer extends RichKafkaSerializationSchema<Tuple2<GenericRecord,GenericRecord>>{ private String topic; // 标记为transient,让Flink序列化时跳过这些字段 private transient Schema schemaKey; private transient Schema schemaValue; private String registryUrl; // 提前初始化序列化器实例,避免每次serialize都创建新对象 private transient ConfluentRegistryAvroSerializationSchema<GenericRecord> keySerializer; private transient ConfluentRegistryAvroSerializationSchema<GenericRecord> valueSerializer; public GenericSerializer(String topic, Schema schemaK, Schema schemaV, String url) { super(); this.topic = topic; this.schemaKey = schemaK; this.schemaValue = schemaV; this.registryUrl = url; } public GenericSerializer() { super(); } @Override public void open(Configuration parameters) throws Exception { super.open(parameters); // 在算子启动的open阶段初始化序列化器,此时不需要序列化Schema keySerializer = ConfluentRegistryAvroSerializationSchema.forGeneric(topic + "-key", schemaKey, registryUrl); valueSerializer = ConfluentRegistryAvroSerializationSchema.forGeneric(topic + "-value", schemaValue, registryUrl); } @Override public ProducerRecord<byte[], byte[]> serialize(Tuple2<GenericRecord,GenericRecord> element, Long timestamp) { byte[] key = keySerializer.serialize(element.f0); byte[] value = valueSerializer.serialize(element.f1); return new ProducerRecord<byte[], byte[]>(topic, key, value); } }
关键优化说明
- 把
schemaKey、schemaValue标记为transient,这样Flink序列化GenericSerializer实例时会跳过这两个不可序列化的对象; - 改用
RichKafkaSerializationSchema,在open方法中初始化实际的序列化器,既避免了序列化问题,又提升了性能(不用每次序列化都新建实例); - 如果你担心传入的Schema在序列化后丢失,还可以改成在
open方法中直接从Confluent Schema Registry拉取Schema(通过topic对应的subject名称),这样连Schema对象都不需要在构造方法中传递,彻底规避序列化风险。
方案2:升级Flink版本(长期推荐)
你的Flink 1.9.1版本确实比较老旧,后续的Flink版本(比如1.11及以上)官方提供了对Confluent Schema Registry的原生支持,不需要自己引入第三方的ConfluentRegistryAvroSerializationSchema类。比如可以直接使用官方的KafkaAvroSerializer配合Flink Kafka连接器,代码会简化很多,也能从根源上避免这类序列化问题。
额外注意点
- 确保你的
ConfluentRegistryAvroSerializationSchema类本身是可序列化的,或者同样在open方法中初始化它的实例,不要让它成为序列化器的非transient字段; - 如果需要传递Schema,优先传递Schema的字符串形式,然后在
open方法中用Schema.Parser解析成Schema对象,这也是规避Schema序列化问题的常用技巧。
内容的提问来源于stack exchange,提问作者kopaka
相关产品推荐
相关产品推荐

