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

Kafka Connect JDBC源连接器无法通过SMT写入NULL墓碑记录

解决Kafka Connect AvroConverter返回墓碑记录(Null值)的Schema校验错误

在使用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,让连接器直接使用:

  1. 通过Schema Registry的REST API导出JDBC源自动生成的原始Schema
  2. 修改Schema,将顶级类型改为["null", 原始Schema]的Union结构
  3. 手动将修改后的Schema注册到Schema Registry
  4. 在连接器配置中指定使用该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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 07:54:56