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

使用含嵌套类的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 13:54:21