配置Confluent S3 Sink Connector连接MSK时序列化错误求助
场景1:AWS Glue Schema Registry 报「Didn't find secondary deserializer」
错误原因
Glue Schema Registry转换器要求指定次要反序列化器,用于处理未注册Schema、或与Registry中Schema不匹配的消息。仅修改key.converter为AWSKafkaAvroConverter,未配置次要反序列化器的话,仍会触发该错误。
解决方案
在Connector配置中,为value.converter(若key也需处理则添加key.converter)补充次要反序列化器相关配置:
# 使用Glue的JSON Schema转换器为例 value.converter=software.amazon.kafka.connect.schemaregistry.AWSKafkaJsonSchemaConverter value.converter.schema.registry.name=testregistry1 value.converter.region=<你的AWS区域> # 启用次要反序列化器 value.converter.use.secondary.deserializer=true # 指定次要反序列化器(适配你的消息格式,此处为原始JSON) value.converter.secondary.deserializer=org.apache.kafka.connect.json.JsonConverter # 若消息无内嵌Schema,关闭该配置 value.converter.secondary.deserializer.schemas.enable=false
如果使用Avro转换器,只需将value.converter改为software.amazon.kafka.connect.schemaregistry.AWSKafkaAvroConverter,次要反序列化器根据实际消息格式调整即可。
场景2:Confluent Schema Registry 报「Unknown magic byte」
错误原因
kafka-console-producer.sh发送的是原始JSON消息,但io.confluent.connect.json.JsonSchemaConverter仅能识别带Confluent Schema格式的序列化消息(消息开头包含magic byte和Schema ID,需通过Confluent的JsonSchemaSerializer生成)。原始JSON不符合该格式,因此触发错误。
解决方案
有两种可行方案:
方案1:改用普通JSON转换器处理原始消息
将value.converter替换为支持原始JSON的转换器,并手动指定Schema:value.converter=org.apache.kafka.connect.json.JsonConverter # 启用Schema支持 value.converter.schemas.enable=true # 手动指定消息对应的JSON Schema value.converter.schema={"type":"object","properties":{"field1":{"type":"string"},"field2":{"type":"integer"}}}方案2:发送符合Confluent Schema格式的消息
使用kafka-console-producer.sh时,指定Confluent的JsonSchemaSerializer,确保消息带Schema ID:- 创建
producer.properties配置文件:
key.serializer=org.apache.kafka.common.serialization.StringSerializer value.serializer=io.confluent.kafka.serializers.json.JsonSchemaSerializer schema.registry.url=http://<你的EC2上Confluent Registry地址>:8081 value.subject.name.strategy=io.confluent.kafka.serializers.subject.TopicNameStrategy- 用该配置发送消息:
kafka-console-producer.sh --broker-list <MSK集群Broker地址> --topic <你的Topic名> --producer.config producer.properties此时发送的
{"field1":"test4578_01","field2":1}会被正确序列化,Sink Connector可正常解析。- 创建
内容的提问来源于stack exchange,提问作者Homer

