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

Debezium 3.0.8 JDBC Sink连接MSSQL varbinary字段失败求助

解决Debezium JdbcSinkConnector写入MSSQL VARBINARY字段的类型转换问题

方法1:自定义SMT转换HeapByteBuffer为byte[]

编写Kafka Connect的Single Message Transform(SMT),将消息中java.nio.HeapByteBuffer类型的字段转换为byte[],适配MSSQL JDBC驱动的类型要求。

示例SMT代码

import org.apache.kafka.common.config.ConfigDef;
import org.apache.kafka.connect.connector.ConnectRecord;
import org.apache.kafka.connect.data.Schema;
import org.apache.kafka.connect.data.SchemaBuilder;
import org.apache.kafka.connect.transforms.Transformation;
import org.apache.kafka.connect.transforms.util.SimpleConfig;

import java.nio.ByteBuffer;
import java.util.Map;

public class ByteBufferToByteArray implements Transformation<ConnectRecord<?>> {
    @Override
    public ConfigDef config() {
        return new ConfigDef();
    }

    @Override
    public void configure(Map<String, ?> configs) {}

    @Override
    public ConnectRecord<?> apply(ConnectRecord<?> record) {
        if (!(record.value() instanceof Map)) {
            return record;
        }

        Map<String, Object> valueMap = (Map<String, Object>) record.value();
        for (Map.Entry<String, Object> entry : valueMap.entrySet()) {
            if (entry.getValue() instanceof ByteBuffer buffer) {
                byte[] byteArray = new byte[buffer.remaining()];
                buffer.get(byteArray);
                entry.setValue(byteArray);
            }
        }

        return record.newRecord(
                record.topic(),
                record.kafkaPartition(),
                record.keySchema(),
                record.key(),
                SchemaBuilder.map(Schema.STRING_SCHEMA, Schema.BYTES_SCHEMA).build(),
                valueMap,
                record.timestamp()
        );
    }

    @Override
    public void close() {}
}

部署与配置

  1. 将编译后的JAR包放入Kafka Connect的plugin.path指定目录
  2. 在Debezium Source Connector配置中添加以下SMT配置:
transforms=byteBufferConvert
transforms.byteBufferConvert.type=com.your.package.ByteBufferToByteArray

方法2:配置JdbcSinkConnector的类型转换

通过调整JdbcSinkConnector的配置,结合内置转换处理类型适配:

  1. 启用自动建表/表结构演化,确保目标表字段类型为VARBINARY
  2. 添加类型转换配置,强制将字段转为bytes类型:
insert.mode=upsert
pk.mode=record_key
auto.create=true
auto.evolve=true
# 启用字段类型转换
transforms=castField0
transforms.castField0.type=org.apache.kafka.connect.transforms.Cast$Value
transforms.castField0.spec=Field0:bytes

如果内置Cast转换不生效,需结合方法1的自定义SMT完成ByteBuffer到byte[]的转换。

方法3:调整Debezium Source Connector输出格式

修改Source Connector配置,强制将binary类型字段输出为byte[]而非HeapByteBuffer:

# 配置Avro转换器确保Schema映射正确
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=http://your-schema-registry:8081
value.converter.connect.meta.data=true
# 禁用不必要的类型包装
decimal.handling.mode=string
database.history.kafka.value.converter=org.apache.kafka.connect.json.JsonConverter
database.history.kafka.value.converter.schemas.enable=true

验证步骤

  1. 重启Kafka Connect服务加载新配置/插件
  2. 检查Kafka主题中的消息,确认Field0字段类型为byte[]
  3. 查看JdbcSinkConnector日志,确认无类型转换错误
  4. 验证目标MSSQL表中Field0字段数据写入正常

内容的提问来源于stack exchange,提问作者Yuriy Vikulov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 20:22:44