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

MSK Kafka Connect Sink连接器因Schema Registry报错求助

问题描述

Kafka主题原始数据

{"tran_slip":"00002060","tran_amount":"111.22"}
{"tran_slip":"00000005","tran_amount":"123"}
{"tran_slip":"00000006","tran_amount":"123"}
{"tran_slip":"00000007","tran_amount":"123"}

自定义Avro Schema(已注册到AWS Glue Schema Registry)

{
  "type": "record",
  "namespace": "int_trans",
  "name": "transaction",
  "fields": [
    {
      "name": "tran_slip",
      "type": "string"
    },
    {
      "name": "tran_amount",
      "type": "string"
    }
  ]
}

MSK Kafka Connect Sink连接器配置

connector.class=io.confluent.connect.jdbc.JdbcSinkConnector
value.converter.schemaAutoRegistrationEnabled=true
connection.password=******
transforms.extractKeyFromStruct.type=org.apache.kafka.connect.transforms.ExtractField$Key
tasks.max=1
key.converter.region=*******
transforms=RenameField
key.converter.schemaName=KeySchema
value.converter.avroRecordType=GENERIC_RECORD
internal.key.converter.schemas.enable=false
value.converter.schemaName=ValueSchema
auto.evolve=false
transforms.RenameField.type=org.apache.kafka.connect.transforms.ReplaceField$Value
key.converter.avroRecordType=GENERIC_RECORD
value.converter=com.amazonaws.services.schemaregistry.kafkaconnect.AWSKafkaAvroConverter
insert.mode=upsert
key.converter=org.apache.kafka.connect.storage.StringConverter
transforms.RenameField.renames=tran_slip:TRAN_SLIP, tran_amount:TRAN_AMOUNT
table.name.format=abc.transactions_sink
topics=aws-db.abc.transactions
batch.size=1
value.converter.registry.name=registry_transactions
value.converter.region=*****
key.converter.registry.name=registry_transactions
key.converter.schemas.enable=false
internal.key.converter=com.amazonaws.services.schemaregistry.kafkaconnect.AWSKafkaAvroConverter
delete.enabled=false
key.converter.schemaAutoRegistrationEnabled=true
connection.user=*******
internal.value.converter.schemas.enable=false
value.converter.schemas.enable=true
internal.value.converter=com.amazonaws.services.schemaregistry.kafkaconnect.AWSKafkaAvroConverter
auto.create=false
connection.url=*********
pk.mode=record_value
pk.fields=tran_slip

错误日志

org.apache.kafka.connect.errors.DataException: Converting byte[] to Kafka Connect data failed due to serialization error:
    at com.amazonaws.services.schemaregistry.kafkaconnect.AWSKafkaAvroConverter.toConnectData(AWSKafkaAvroConverter.java:118)
    at org.apache.kafka.connect.storage.Converter.toConnectData(Converter.java:87)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.convertValue(WorkerSinkTask.java:545)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.lambda$convertAndTransformRecord$1(WorkerSinkTask.java:501)
    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:156)
    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:190)
    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:132)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.convertAndTransformRecord(WorkerSinkTask.java:501)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:478)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:328)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:232)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:201)
    at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:189)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:238)
    at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
    at java.base/java.lang.Thread.run(Thread.java:829)
Caused by: com.amazonaws.services.schemaregistry.exception.AWSSchemaRegistryException: Didn't find secondary deserializer.
    at com.amazonaws.services.schemaregistry.deserializers.SecondaryDeserializer.deserialize(SecondaryDeserializer.java:65)
    at com.amazonaws.services.schemaregistry.deserializers.avro.AWSKafkaAvroDeserializer.deserializeByHeaderVersionByte(AWSKafkaAvroDeserializer.java:150)
    at com.amazonaws.services.schemaregistry.deserializers.avro.AWSKafkaAvroDeserializer.deserialize(AWSKafkaAvroDeserializer.java:114)
    at com.amazonaws.services.schemaregistry.kafkaconnect.AWSKafkaAvroConverter.toConnectData(AWSKafkaAvroConverter.java:116)
问题分析与修复方案

错误核心是Didn't find secondary deserializer,原因是你的Kafka主题数据是无Schema的JSON格式,但配置的AWSKafkaAvroConverter默认只处理带有AWS Schema Registry头部标识的Avro序列化数据。当它遇到原始JSON数据时无法识别,需要指定secondary deserializer处理非Avro格式的原始数据,并关联已注册的Avro Schema完成转换。

具体修复步骤

  1. 添加secondary deserializer配置
    给value.converter补充以下配置,指定JSON转换器作为secondary deserializer读取原始JSON数据,并绑定已注册的Avro Schema版本ID(可在AWS Glue控制台的Schema Registry页面获取):

    value.converter.secondaryDeserializer=org.apache.kafka.connect.json.JsonConverter
    value.converter.secondaryDeserializer.schemas.enable=false
    value.converter.schemaVersionId=<你的Avro Schema版本ID>
    
  2. 清理冗余配置
    移除无效或不必要的配置项:

    • transforms.extractKeyFromStruct:未在transforms列表中启用,属于无效配置
    • internal.key.converter、internal.value.converter:无特殊内部转换需求时无需指定,Kafka Connect会使用默认值
    • 所有key.converter关联的Schema Registry配置:key.converter使用的是StringConverter,不需要Avro相关的schemaName、registry.name等参数
  3. 简化key转换器配置
    由于pk.mode=record_value,主键从value中提取,key转换器可简化为:

    key.converter=org.apache.kafka.connect.storage.StringConverter
    key.converter.schemas.enable=false
    

修复后的核心配置片段

connector.class=io.confluent.connect.jdbc.JdbcSinkConnector
value.converter=com.amazonaws.services.schemaregistry.kafkaconnect.AWSKafkaAvroConverter
value.converter.schemas.enable=true
value.converter.registry.name=registry_transactions
value.converter.region=*****
value.converter.schemaAutoRegistrationEnabled=true
value.converter.schemaName=ValueSchema
value.converter.avroRecordType=GENERIC_RECORD
# 新增secondary deserializer配置
value.converter.secondaryDeserializer=org.apache.kafka.connect.json.JsonConverter
value.converter.secondaryDeserializer.schemas.enable=false
value.converter.schemaVersionId=<你的Schema版本ID>

key.converter=org.apache.kafka.connect.storage.StringConverter
key.converter.schemas.enable=false

transforms=RenameField
transforms.RenameField.type=org.apache.kafka.connect.transforms.ReplaceField$Value
transforms.RenameField.renames=tran_slip:TRAN_SLIP, tran_amount:TRAN_AMOUNT

insert.mode=upsert
pk.mode=record_value
pk.fields=tran_slip
table.name.format=abc.transactions_sink
topics=aws-db.abc.transactions
tasks.max=1
batch.size=1
auto.create=false
auto.evolve=false
delete.enabled=false
connection.url=*********
connection.user=*******
connection.password=******

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 22:30:51