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

如何在Apache Flink中消费Kafka Topic的Avro格式数据?

嘿,刚好之前做过Flink消费Kafka Avro消息的场景,来给你分享下实操方案~你现在用自定义反序列化器的思路是可行的,不过也可以试试Flink自带的Avro工具类,能省不少代码。下面分两种方式给你示例:

方式一:自定义Avro反序列化器(贴合你当前的思路)

首先得确保你已经通过Avro Schema生成了对应的UserInfo Java类(用avro-maven-plugin或者命令行工具都可以)。然后编写自定义反序列化器:

import org.apache.flink.api.common.serialization.DeserializationSchema;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.avro.io.DatumReader;
import org.apache.avro.io.Decoder;
import org.apache.avro.io.DecoderFactory;
import org.apache.avro.specific.SpecificDatumReader;
import com.example.flink.avro.UserInfo;

import java.io.IOException;

public class UserInfoAvroDeserializer implements DeserializationSchema<UserInfo> {

    // 懒加载DatumReader,避免重复初始化
    private transient DatumReader<UserInfo> datumReader;

    @Override
    public UserInfo deserialize(byte[] message) throws IOException {
        if (datumReader == null) {
            datumReader = new SpecificDatumReader<>(UserInfo.getClassSchema());
        }
        Decoder decoder = DecoderFactory.get().binaryDecoder(message, null);
        return datumReader.read(null, decoder);
    }

    @Override
    public boolean isEndOfStream(UserInfo nextElement) {
        return false;
    }

    @Override
    public TypeInformation<UserInfo> getProducedType() {
        return TypeInformation.of(UserInfo.class);
    }
}

然后在Flink消费作业里配置这个反序列化器:

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;
import java.util.Properties;

public class KafkaAvroConsumer {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        Properties kafkaProps = new Properties();
        kafkaProps.setProperty("bootstrap.servers", "你的Kafka Broker地址:9092");
        kafkaProps.setProperty("group.id", "flink-avro-consumer-group");

        FlinkKafkaConsumer<UserInfo> consumer = new FlinkKafkaConsumer<>(
                "test-topic",
                new UserInfoAvroDeserializer(),
                kafkaProps
        );

        env.addSource(consumer)
           .print()
           .setParallelism(1);

        env.execute("Flink 消费Kafka Avro消息作业");
    }
}
方式二:用Flink内置的Avro序列化工具(更简洁)

Flink从1.11版本开始内置了Avro的支持,直接用AvroDeserializationSchema就能搞定,不用自己写反序列化器:

import org.apache.flink.formats.avro.AvroDeserializationSchema;
// 其他导入和上面一致

public class KafkaAvroConsumer {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        Properties kafkaProps = new Properties();
        kafkaProps.setProperty("bootstrap.servers", "你的Kafka Broker地址:9092");
        kafkaProps.setProperty("group.id", "flink-avro-consumer-group");

        FlinkKafkaConsumer<UserInfo> consumer = new FlinkKafkaConsumer<>(
                "test-topic",
                AvroDeserializationSchema.forSpecific(UserInfo.class),
                kafkaProps
        );

        env.addSource(consumer)
           .print()
           .setParallelism(1);

        env.execute("Flink 消费Kafka Avro消息作业");
    }
}

如果是用GenericRecord格式发送的消息(没生成具体Java类),可以换成AvroDeserializationSchema.forGeneric(UserInfo.getClassSchema())来解析。

几个关键注意点
  • 项目依赖要加全,Maven的话需要这两个依赖:
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-avro</artifactId>
    <version>${flink.version}</version>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-kafka</artifactId>
    <version>${flink.version}</version>
</dependency>
  • 消费端的Schema要和生产端兼容,如果有Schema演进需求,建议用Schema Registry来统一管理,这样能避免Schema不匹配的问题。

内容的提问来源于stack exchange,提问作者Sumit Nekar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 03:32:37