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完成转换。
具体修复步骤
添加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>清理冗余配置
移除无效或不必要的配置项:transforms.extractKeyFromStruct:未在transforms列表中启用,属于无效配置internal.key.converter、internal.value.converter:无特殊内部转换需求时无需指定,Kafka Connect会使用默认值- 所有
key.converter关联的Schema Registry配置:key.converter使用的是StringConverter,不需要Avro相关的schemaName、registry.name等参数
简化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
相关产品推荐
相关产品推荐

