【问题标题】:org.apache.kafka.connect.errors.DataException: Invalid JSON for record default value: nullorg.apache.kafka.connect.errors.DataException:记录默认值的 JSON 无效:null
【发布时间】:2018-12-06 17:43:15
【问题描述】:

我有一个使用 KafkaAvroSerializer 生成的 Kafka Avro 主题。
我的独立属性如下。
我正在使用 Confluent 4.0.0 运行 Kafka 连接。

key.converter=io.confluent.connect.avro.AvroConverter
value.converter=io.confluent.connect.avro.AvroConverter
key.converter.schema.registry.url=<schema_registry_hostname>:8081
value.converter.schema.registry.url=<schema_registry_hostname>:8081
key.converter.schemas.enable=true
value.converter.schemas.enable=true
internal.key.converter=org.apache.kafka.connect.json.JsonConverter
internal.value.converter=org.apache.kafka.connect.json.JsonConverter
internal.key.converter.schemas.enable=false
internal.value.converter.schemas.enable=false

当我在独立模式下为 hdfs sink 运行 Kafka 连接器时,我收到以下错误消息:

[2018-06-27 17:47:41,746] ERROR WorkerSinkTask{id=camus-email-service-0} Task threw an uncaught and unrecoverable exception (org.apache.kafka.connect.runtime.WorkerTask)
org.apache.kafka.connect.errors.DataException: Invalid JSON for record default value: null
    at io.confluent.connect.avro.AvroData.defaultValueFromAvro(AvroData.java:1640)
    at io.confluent.connect.avro.AvroData.toConnectSchema(AvroData.java:1527)
    at io.confluent.connect.avro.AvroData.toConnectSchema(AvroData.java:1410)
    at io.confluent.connect.avro.AvroData.toConnectSchema(AvroData.java:1290)
    at io.confluent.connect.avro.AvroData.toConnectData(AvroData.java:1014)
    at io.confluent.connect.avro.AvroConverter.toConnectData(AvroConverter.java:88)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:454)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:287)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:198)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:166)
    at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:170)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:214)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
    at java.lang.Thread.run(Thread.java:745)
[2018-06-27 17:47:41,748] ERROR WorkerSinkTask{id=camus-email-service-0} Task is being killed and will not recover until manually restarted (    org.apache.kafka.connect.runtime.WorkerTask)
[2018-06-27 17:52:19,554] INFO Kafka Connect stopping (org.apache.kafka.connect.runtime.Connect).

当我使用 kafka-avro-console-consumer 传递模式注册表时,我会反序列化 Kafka 消息。

即:

/usr/bin/kafka-avro-console-consumer --bootstrap-server <kafka-host>:9092 --topic <KafkaTopicName> --property schema.registry.url=<schema_registry_hostname>:8081

【问题讨论】:

标签: apache-kafka avro apache-kafka-connect confluent-platform confluent-schema-registry


【解决方案1】:

我认为您的 Kafka 密钥为空,这不是 Avro。

或者它是某种其他类型但格式错误,并且未转换为RECORD 数据类型。见 AvroData 源代码

case RECORD: {
    if (!jsonValue.isObject()) {
      throw new DataException("Invalid JSON for record default value: " + jsonValue.toString());
    }

更新根据您的评论,您可以看到这是真的

$ curl -X GET localhost:8081/subjects/<kafka-topic>-key/versions/latest    
{"subject":"<kafka-topic>-key","version":2,"id":625,"schema":"\"bytes\""}

无论如何,HDFS Connect 本身并不存储密钥,因此尽量不要反序列化密钥,而不要使用 Avro。

key.converter=org.apache.kafka.connect.converters.ByteArrayConverter

此外,您的控制台使用者没有打印密钥,因此您的测试不充分。您需要添加--property print.key=true

【讨论】:

  • 谢谢@cricket_007。当我使用 --print.key=true 运行 kafka-avro-console-consumer 时,我看到了密钥。 /usr/bin/kafka-avro-console-consumer --bootstrap-server :9092 --topic --property schema.registry.url=:8081 --from-beginning - -max-messages 1 --property print.key=true "¶2âwÌS@¼T%\u0005ãé" {"id":{"string":"b632e277-cc53-4093-9cbc-86542505e3e9 ........ }
  • 所以我看到了一个键,我们在 Producer 使用 KafkaAvroSerializer 来获取键和值。在模式注册表中,虽然我将数据类型视为“字节”。 Producer 在模式注册表中生成这个键 $ curl -X GET localhost:8081/subjects/<kafka-topic>-key/versions/latest {"subject":"-key","version":2,"id":625,"schema":"\" bytes\""} 你认为这是 Key 的问题吗?
  • 另外你如何忽略独立属性中的键。当我没有指定密钥时,我收到以下错误消息:线程“主”org.apache.kafka.common.config.ConfigException 中的异常:缺少所需的配置“key.converter”,它没有默认值。
  • @RupeshMore 用key.converter 更新了答案,你需要bytes 数据
  • 我们将类型更改为 UNION 数据类型,并且 AvroConverters 能够反序列化该列。 {\"name\":\"subscription\",\"type\":[\"null\",{\"type\":\"record\",\"name\":\"Subscription\" ,\"doc\":\"模板订阅信息\",\"fields\":[{\"name\":\"subscriptionId\",\"type\":[\"null\",\" int\"],\"default\":null},{\"name\":\"channelId\",\"type\":[\"null\",\"int\"],\"default \":null}]}],\"default\":null}
【解决方案2】:

将“订阅”列的数据类型更改为联合数据类型解决了该问题。 Avroconverters 能够反序列化消息。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-08-17
    • 2018-01-05
    • 1970-01-01
    • 1970-01-01
    • 2020-06-27
    • 2020-12-01
    • 2014-12-07
    相关资源
    最近更新 更多