Flink Kafka集成Schema Registry及读取Confluent平台Avro数据咨询
关于Flink对接Confluent Schema Registry及读取Avro数据的解答
Flink是否支持与Confluent Schema Registry集成
完全支持,属于生产环境非常成熟的用法,不需要自行开发适配逻辑。
Flink官方已经提供了Confluent生态的专属对接组件,不管是DataStream API还是Flink SQL/Table API,都可以直接对接Schema Registry,支持Schema自动拉取、本地缓存、版本自动适配,完全兼容Confluent的Schema演进规则。
如何从Confluent平台读取AVRO格式数据
整体流程和普通Kafka消费逻辑差异不大,核心是替换反序列化实现、配置Schema Registry地址即可,分两种常用场景说明:
前置依赖准备
需要在项目中引入和自身Flink版本、Confluent版本匹配的两个核心依赖:
flink-connector-kafka:Flink对接Kafka的官方连接器flink-avro-confluent-registry:Flink对接Confluent Avro和Schema Registry的适配包
注意不要混用版本,比如Confluent用6.2版本,适配包就要选对应适配6.2的版本,否则会出现反序列化异常、Schema接口调用失败的问题。
DataStream API 实现方式
如果是Flink 1.14之前的版本,可以用旧版FlinkKafkaConsumer,1.14及之后版本推荐用新版KafkaSource,二者配置逻辑一致,核心配置示例:
import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.connector.kafka.source.KafkaSource; import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer; import org.apache.flink.formats.avro.registry.confluent.ConfluentRegistryAvroDeserializationSchema; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import java.util.Properties; public class ReadConfluentAvro { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); Properties kafkaProps = new Properties(); kafkaProps.put("bootstrap.servers", "kafka-broker1:9092,kafka-broker2:9092"); kafkaProps.put("schema.registry.url", "http://schema-registry:8081"); // 如果Schema Registry开了认证,在这里加对应鉴权配置即可,和原生Confluent消费者配置一致 KafkaSource<UserAvro> source = KafkaSource.<UserAvro>builder() .setTopics("target_topic") .setGroupId("flink-consumer-group") .setStartingOffsets(OffsetsInitializer.earliest()) .setProperties(kafkaProps) // 指定反序列化Schema,UserAvro是提前根据Avsc生成的SpecificRecord实体类 .setDeserializer(ConfluentRegistryAvroDeserializationSchema.forSpecific(UserAvro.class, "http://schema-registry:8081")) .build(); env.fromSource(source, WatermarkStrategy.noWatermarks(), "confluent-avro-source") .print(); env.execute(); } }
如果不需要用生成的SpecificRecord类,也可以替换成forGeneric方法,直接返回通用的GenericRecord类型,自行解析字段即可。
Flink SQL 实现方式
不需要写Java/Scala代码,建表时直接指定格式为avro-confluent即可,配置示例:
CREATE TABLE confluent_avro_source ( id BIGINT, username STRING, create_time TIMESTAMP(3) ) WITH ( 'connector' = 'kafka', 'topic' = 'target_topic', 'properties.bootstrap.servers' = 'kafka-broker1:9092,kafka-broker2:9092', 'properties.group.id' = 'flink-sql-group', 'format' = 'avro-confluent', 'avro-confluent.schema-registry.url' = 'http://schema-registry:8081', 'scan.startup.mode' = 'earliest-offset' ); -- 直接查询即可获取Avro格式的数据,字段自动映射 SELECT * FROM confluent_avro_source;
常见踩坑点
- 如果消费时报错Magic byte不匹配,先检查是不是Topic里存的不是Confluent序列化的Avro数据(Confluent Avro序列化会在消息头加5字节的Schema ID标识,普通Avro序列化没有这个头,会解析失败)
- Schema演进时只要符合Confluent配置的兼容策略(比如向后兼容),Flink作业不需要修改代码重启,会自动适配新Schema
- 如果Schema Registry开启了HTTPS访问,记得在配置里加上对应的SSL证书路径参数即可
内容的提问来源于stack exchange,提问作者user1112259
相关产品推荐
相关产品推荐

