Confluent Kafka Connect下沉流至RDBMS时Avro反序列化失败求助
从你提供的错误日志和配置来看,核心问题是Kafka Connect尝试用Avro Converter去反序列化非Confluent Avro格式的数据,具体表现为Unknown magic byte!和Error deserializing Avro message for id -1——这两个错误说明数据开头没有Confluent Avro序列化格式要求的magic byte(0x0)和schema ID前缀,通常是因为你的Sink配置中错误地为key指定了Avro Converter,但topic的key实际是用其他序列化方式(比如原始Long类型)存储的。
问题分析
你用kafka-avro-console-consumer查看的是topic的value数据,看起来是正常的Avro格式,但错误发生在反序列化环节,大概率是topic的key不是Avro序列化的。因为你在Sink配置中同时设置了key.converter=io.confluent.connect.avro.AvroConverter,但如果这个topic的key是KSQL rekey操作后生成的原始Long类型(而非Avro格式),Connect就会无法解析key,抛出上述错误。
排查与解决步骤
1. 确认Topic的Key格式
先运行以下命令查看topic的key内容和序列化格式:
./bin/kafka-console-consumer --bootstrap-server localhost:9092 \ --from-beginning --max-messages 1 --topic STR_VAHG_REKEY_02 \ --property print.key=true \ --property key.deserializer=org.apache.kafka.common.serialization.LongDeserializer \ --property value.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer \ --property schema.registry.url=http://localhost:8081
如果能正常打印出key的数值(比如你看到的3941641584970777000),说明key是Long类型,而非Avro格式。
2. 修改Sink Connector配置
针对key是Long类型的情况,修改你的Sink配置文件,把key.converter换成对应的Long转换器,同时补全数据库连接的完整信息:
connector.class=io.confluent.connect.jdbc.JdbcSinkConnector tasks.max=1 topics=STR_VAHG_REKEY_02 # 补全完整的数据库连接地址 connection.url=jdbc:mysql://<你的数据库地址>:<端口>/<目标库名> connection.user=<数据库用户名> connection.password=<数据库密码> auto.create=true timestamp.column.name=create_at validate.non.null=false # 针对Long类型的key使用专用转换器 key.converter=org.apache.kafka.connect.converters.LongConverter # value仍然使用AvroConverter value.converter=io.confluent.connect.avro.AvroConverter value.converter.schema.registry.url=http://localhost:8081
3. 额外注意事项
- 如果你的KSQL rekey操作是通过
PARTITION BY语句完成的,KSQL默认会用原始类型序列化key(比如Long),而非Avro。如果需要将key也序列化为Avro,可以在KSQL中先执行SET KEY_SERIALIZATION_FORMAT='AVRO';,再执行rekey语句。 - 修改配置后,记得重启对应的Sink Connector,确保配置生效。
按照上述步骤调整后,应该就能解决这个反序列化错误,顺利将数据下沉到RDBMS了。
内容的提问来源于stack exchange,提问作者AnmP

