使用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" } }
额外配置检查
- Docker网络地址问题:Sink的
connection.url中使用localhost错误,容器内部localhost指向自身,需替换为目标SQL Server在docker-compose中的服务名称(比如sqlserver-target)。 - Transform顺序:当前
dropPrefix(修改topic名)和unwrap(提取新记录状态)的顺序没问题,前者不影响记录内容,后者处理数据结构。
内容的提问来源于stack exchange,提问作者Cesar Ushiro
相关产品推荐
相关产品推荐

