使用含嵌套类的Avro无法注册Schema问题求助
Avro嵌套对象导致Schema无法保存,抛出"Can't redefine: io.confluent.connect.avro.ConnectDefault"异常
当Avro类包含嵌套对象时,Schema无法正常保存,抛出如下异常:
org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler connect_1 | at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:223) connect_1 | at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:149) connect_1 | at org.apache.kafka.connect.runtime.WorkerSourceTask.convertTransformedRecord(WorkerSourceTask.java:330) connect_1 | at org.apache.kafka.connect.runtime.WorkerSourceTask.sendRecords(WorkerSourceTask.java:356) connect_1 | at org.apache.kafka.connect.runtime.WorkerSourceTask.execute(WorkerSourceTask.java:258) connect_1 | at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:188) connect_1 | at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:243) connect_1 | at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) connect_1 | at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) connect_1 | at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) connect_1 | at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) connect_1 | at java.base/java.lang.Thread.run(Thread.java:829) connect_1 | Caused by: org.apache.avro.SchemaParseException: Can't redefine: io.confluent.connect.avro.ConnectDefault connect_1 | at org.apache.avro.Schema$Names.put(Schema.java:1550) connect_1 | at org.apache.avro.Schema$NamedSchema.writeNameRef(Schema.java:813) connect_1 | at org.apache.avro.Schema$RecordSchema.toJson(Schema.java:975) connect_1 | at org.apache.avro.Schema$UnionSchema.toJson(Schema.java:1242) connect_1 | at org.apache.avro.Schema$RecordSchema.fieldsToJson(Schema.java:1003) connect_1 | at org.apache.avro.Schema$RecordSchema.toJson(Schema.java:987) connect_1 | at org.apache.avro.Schema.toString(Schema.java:426) connect_1 | at org.apache.avro.Schema.toString(Schema.java:398) connect_1 | at org.apache.avro.Schema.toString(Schema.java:389) connect_1 | at io.apicurio.registry.serde.avro.AvroKafkaSerializer.getSchemaFromData(AvroKafkaSerializer.java:108) connect_1 | at io.apicurio.registry.serde.AbstractKafkaSerializer.lambda$serialize$0(AbstractKafkaSerializer.java:90) connect_1 | at io.apicurio.registry.serde.LazyLoadedParsedSchema.getRawSchema(LazyLoadedParsedSchema.java:55) connect_1 | at io.apicurio.registry.serde.DefaultSchemaResolver.resolveSchema(DefaultSchemaResolver.java:81) connect_1 | at io.apicurio.registry.serde.AbstractKafkaSerializer.serialize(AbstractKafkaSerializer.java:92) connect_1 | at io.apicurio.registry.serde.AbstractKafkaSerializer.serialize(AbstractKafkaSerializer.java:79) connect_1 | at io.apicurio.registry.utils.converter.SerdeBasedConverter.fromConnectData(SerdeBasedConverter.java:111) connect_1 | at org.apache.kafka.connect.storage.Converter.fromConnectData(Converter.java:64) connect_1 | at org.apache.kafka.connect.runtime.WorkerSourceTask.lambda$convertTransformedRecord$3(WorkerSourceTask.java:330) connect_1 | at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:173) connect_1 | at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:207) connect_1 | ... 11 more
使用工具:
- Debezium
- Apicurio Schema Registry
- Avro格式
问题原因
该异常源于Avro Schema中重复定义了io.confluent.connect.avro.ConnectDefault类型,通常是Debezium生成的嵌套结构Connect数据转换为Avro时,默认类型处理逻辑引发Schema冲突,或是Apicurio Registry的Schema注册策略导致重复定义。
解决方案
1. 调整Debezium转换器配置
使用Apicurio的Avro转换器,指定RecordIdStrategy作为全局ID策略,避免同一Topic下不同嵌套结构的Schema冲突:
key.converter=io.apicurio.registry.utils.converter.AvroConverter key.converter.apicurio.registry.url=http://<你的Apicurio地址> key.converter.apicurio.registry.global-id-strategy=io.apicurio.registry.utils.converter.strategy.RecordIdStrategy value.converter=io.apicurio.registry.utils.converter.AvroConverter value.converter.apicurio.registry.url=http://<你的Apicurio地址> value.converter.apicurio.registry.global-id-strategy=io.apicurio.registry.utils.converter.strategy.RecordIdStrategy
2. 清理嵌套结构冗余层级
使用Debezium的ExtractNewRecordState转换,去掉默认生成的before/after外层嵌套结构,降低Schema复杂度:
transforms=unwrap transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState transforms.unwrap.drop.tombstones=false transforms.unwrap.delete.handling.mode=rewrite
3. 清理冲突Schema并重新注册
登录Apicurio Registry管理界面,删除已存在的io.confluent.connect.avro.ConnectDefault相关Schema,重启Debezium连接器让其重新生成并注册正确的Schema。
4. 禁用默认元数据生成(可选)
若不需要默认的Connect类型元数据,可在转换器配置中禁用:
value.converter.connect.meta.data=false
内容的提问来源于stack exchange,提问作者G.Khandal
相关产品推荐
相关产品推荐

