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

使用Kafka Connect实现SQL Server间ETL遇序列化错误求排查

问题分析与解决方案

核心错误原因

Unknown magic byte! 错误说明Sink连接器用Avro转换器解析的内容并非Avro格式。Debezium SQL Server连接器默认使用JSON序列化输出数据,而你的Sink连接器配置了io.confluent.connect.avro.AvroConverter,两者格式不匹配导致解析失败。

修复方案(二选一)

方案1:让源连接器输出Avro格式

修改Debezium源连接器配置,添加Avro转换器及Schema Registry地址:

{
   "name":"db-connector",
   "config":{
      "connector.class":"io.debezium.connector.sqlserver.SqlServerConnector",
      "tasks.max":"1",    
      "database.server.name":"sqlserver",
      "database.hostname":"sqlserver",
      "database.port":"1433",
      "database.user":"sa",
      "database.password":"P@ssw0rd",
      "database.dbname":"db",
      "database.history.kafka.bootstrap.servers":"kafka:9092",
      "database.history.kafka.topic":"schema-changes.db",
      // 新增Avro转换器配置
      "key.converter":"io.confluent.connect.avro.AvroConverter",
      "key.converter.schema.registry.url":"http://schema-registry:8081",
      "value.converter":"io.confluent.connect.avro.AvroConverter",
      "value.converter.schema.registry.url":"http://schema-registry:8081"
   }
}

注意:Docker环境下要使用Schema Registry的服务名(对应docker-compose里的服务名称),而非localhost,否则容器内部无法访问。

方案2:让Sink连接器使用JSON转换器

如果不需要Avro格式,将Sink的转换器改为JSON类型:

{
   "name":"jdbc-sqlserver-sink",
   "config":{
      "connector.class":"io.confluent.connect.jdbc.JdbcSinkConnector",
      "tasks.max":"1",
      // 修改为JSON转换器
      "key.converter":"org.apache.kafka.connect.json.JsonConverter",
      "key.converter.schemas.enable":"true",
      "value.converter":"org.apache.kafka.connect.json.JsonConverter",
      "value.converter.schemas.enable":"true",
      "header.converter":"org.apache.kafka.connect.storage.SimpleHeaderConverter",
      "topics":"sqlserver.dbo.t1",      
      // Docker环境下替换目标SQL Server地址为服务名
      "connection.url":"jdbc:sqlserver://sqlserver-target;instance=instanceName;databaseName=db",
      "connection.user":"kafka",
      "connection.password":"kafka",
      "insert.mode":"upsert",
      "pk.mode":"record_key",
      "auto.create":"true",
      "auto.evolve":"false",
      "max.retries":"1",
      "delete.enabled":"true",
      "transforms":"dropPrefix,unwrap",
      "transforms.dropPrefix.type":"org.apache.kafka.connect.transforms.RegexRouter",
      "transforms.dropPrefix.regex":"sqlserver.dbo.(.*)",
      "transforms.dropPrefix.replacement":"$1",
      "transforms.unwrap.type":"io.debezium.transforms.ExtractNewRecordState",
      "transforms.unwrap.drop.tombstones":"false"
   }
}

额外配置检查

  1. Docker网络地址问题:Sink的connection.url中使用localhost错误,容器内部localhost指向自身,需替换为目标SQL Server在docker-compose中的服务名称(比如sqlserver-target)。
  2. Transform顺序:当前dropPrefix(修改topic名)和unwrap(提取新记录状态)的顺序没问题,前者不影响记录内容,后者处理数据结构。

内容的提问来源于stack exchange,提问作者Cesar Ushiro

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 23:47:42