如何在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
相关产品推荐
相关产品推荐

