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

配置Confluent S3 Sink Connector连接MSK时序列化错误求助

Confluent S3 Sink Connector 序列化错误解决方案

场景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:

    1. 创建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
    
    1. 用该配置发送消息:
    kafka-console-producer.sh --broker-list <MSK集群Broker地址> --topic <你的Topic名> --producer.config producer.properties
    

    此时发送的{"field1":"test4578_01","field2":1}会被正确序列化,Sink Connector可正常解析。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 14:02:31