在Apache Kafka中使用AVRO序列化如何处理嵌套源数据?
错误根因
你遇到的报错不是Kafka Connect Avro转换器不支持嵌套结构,而是以下两个常见原因导致:
- Confluent Avro转换器处理可选字段的默认值时,会自动生成名为
io.confluent.connect.avro.ConnectDefault的全局公共类型,如果你手动预注册的Schema中重复定义了该类型,或者旧版本转换器在处理嵌套结构时重复生成该类型,就会触发重定义报错 - 你手动预注册的Schema结构不符合Connect数据结构转Avro的映射规则,导致转换器生成的运行时Schema和预注册Schema出现类型冲突
解决方案
Kafka Connect完全支持保留嵌套结构序列化为Avro格式,无需扁平化处理,按照以下步骤配置即可:
方案1:优先使用转换器自动注册Schema(推荐)
不需要提前手动在Schema Registry创建subject,直接在连接器配置中添加以下参数,让转换器根据你生成的Struct结构自动注册匹配的Avro Schema:
"value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "http://你的Schema Registry服务地址:8081", "value.converter.auto.register.schemas": true, "value.converter.use.latest.version": false, // 可选配置:关闭Connect元字段注入,避免Schema中额外添加__connect_meta等字段 "value.converter.connect.meta.data": false
配置完成后重启连接器,即可正常写入带嵌套结构的Avro数据到Kafka主题。
方案2:必须手动预注册Schema的处理方法
如果业务要求必须提前预注册Schema,需要严格遵循Connect到Avro的映射规则调整你创建的Schema:
- 所有可选字段必须设置
default: null,不要自定义非null的默认值 - 不要在Schema中自定义任何名为
ConnectDefault的类型定义 - 嵌套Struct对应的Avro Record名称不要和全局定义的其他类型重名
调整完成后在连接器配置中关闭自动注册即可:
"value.converter.auto.register.schemas": false, "value.converter.use.latest.version": true
版本兼容修复
如果你使用的Confluent Platform版本低于6.2.0,该报错是已知的版本bug,升级到6.2.0及以上版本即可直接解决嵌套结构下的ConnectDefault重复定义问题。
内容的提问来源于stack exchange,提问作者Schedie
相关产品推荐
相关产品推荐

