Kafka Connect JDBC源连接器无法通过SMT写入NULL墓碑记录
在使用Kafka Connect JDBC源连接器结合io.confluent.connect.avro.AvroConverter和Schema Registry时,自定义SMT返回Null值(墓碑记录)以清理旧数据,触发如下错误:
org.apache.kafka.connect.errors.DataException: Found null value for non-optional schema
该逻辑在JsonSchemaConverter下正常,但AvroConverter因Schema不可空校验失败,即使尝试将整个Schema或内部元素设为可空也无效。
核心原因
AvroConverter的AvroData类会执行严格校验:当value为null时,对应的Schema必须是可空的(即包含null的Union类型)。JDBC源连接器默认生成的Schema基于数据库表结构,非空字段对应不可空的Schema,且自动注册到Schema Registry的也是该不可空Schema;SMT中直接使用原始的不可空Schema返回Null值,就会触发校验错误。
解决方案
方案1:修改自定义SMT,动态转换Schema为可空Union类型
在SMT返回Null值时,将原始Schema包装成包含null的Union类型,这样schema.isOptional()会返回true,通过AvroData的校验。
修改后的SMT核心代码:
public R apply(R record) { Struct origStruct = Requirements.requireStruct(record.value(), "flat"); Struct targetStruct; Schema valueSchema = record.valueSchema(); if (String.valueOf(origStruct.get(this.checkField)).equals(this.checkContent)) { targetStruct = null; // 将原始Schema包装为包含null的Union类型 valueSchema = SchemaBuilder.unionOf() .nullType() .and(valueSchema) .build(); } else { targetStruct = origStruct; } return record.newRecord( record.topic(), record.kafkaPartition(), record.keySchema(), record.key(), valueSchema, targetStruct, record.timestamp()); }
方案2:提前在Schema Registry注册可空的Union Schema
如果不想修改SMT代码,可以手动注册包含null的Union Schema,让连接器直接使用:
- 通过Schema Registry的REST API导出JDBC源自动生成的原始Schema
- 修改Schema,将顶级类型改为
["null", 原始Schema]的Union结构 - 手动将修改后的Schema注册到Schema Registry
- 在连接器配置中指定使用该Schema:
"value.converter.schema.id": "注册的Schema ID", "value.converter.auto.register.schemas": "false"
关键配置调整
确保Schema Registry的兼容性设置允许从非可空Schema升级为可空Union类型,修正连接器配置中的兼容性参数(原配置的compatibility.strict为错误项,应使用compatibility.level):
"value.converter.compatibility.level": "BACKWARD"
该设置允许新增的Schema兼容旧版本,旧消费者可正常处理包含Null的记录。
注意事项
- 墓碑记录的Key必须与目标清理记录的Key完全一致,否则无法触发日志压缩(日志压缩依赖相同Key的Null记录)
- 若数据库表字段本身允许Null,JDBC源会自动生成包含Null的Union Schema,但业务逻辑的墓碑记录通常对应数据库非空字段,仍需通过SMT或手动注册Schema处理
内容的提问来源于stack exchange,提问作者mauam

