You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.06 14:02:38