【问题标题】:org.apache.kafka.connect.errors.DataException: Invalid JSON for array default value: "null"org.apache.kafka.connect.errors.DataException:数组默认值的 JSON 无效:“null”
【发布时间】:2018-10-29 11:21:00
【问题描述】:

我正在尝试使用 confluent Kafka s3 连接器,使用 confluent-4.1.1

s3-sink

"value.converter.schema.registry.url": "http://localhost:8081",
"value.converter": "io.confluent.connect.avro.AvroConverter",
"key.converter": "org.apache.kafka.connect.storage.StringConverter"

当我为 s3 接收器运行 Kafka 连接器时,我收到以下错误消息:

ERROR WorkerSinkTask{id=singular-s3-sink-0} Task threw an uncaught and unrecoverable exception (org.apache.kafka.connect.runtime.WorkerTask:172)
org.apache.kafka.connect.errors.DataException: Invalid JSON for array default value: "null"
        at io.confluent.connect.avro.AvroData.defaultValueFromAvro(AvroData.java:1649)
        at io.confluent.connect.avro.AvroData.toConnectSchema(AvroData.java:1562)
        at io.confluent.connect.avro.AvroData.toConnectSchema(AvroData.java:1443)
        at io.confluent.connect.avro.AvroData.toConnectSchema(AvroData.java:1443)
        at io.confluent.connect.avro.AvroData.toConnectSchema(AvroData.java:1323)
        at io.confluent.connect.avro.AvroData.toConnectData(AvroData.java:1047)
        at io.confluent.connect.avro.AvroConverter.toConnectData(AvroConverter.java:87)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:468)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:301)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:205)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:173)
        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:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)

我的 Schema 只包含 1 个数组类型字段,其架构是这样的

{"name":"item_id","type":{"type":"array","items":["null","string"]},"default":[]}

我可以使用 kafka-avro-console-consumer 命令查看反序列化消息。我见过similar question,但在他的情况下,他也在使用 Avro 序列化程序作为密钥。

./confluent-4.1.1/bin/kafka-avro-console-consumer --topic singular_custom_postback --bootstrap-server localhost:9092  -max-messages 2

"item_id":[{"string":"15552"},{"string":"37810"},{"string":"38061"}]
"item_id":[]

我不能把我从控制台消费者那里得到的全部输出,因为它包含敏感的用户信息,所以我在我的架构中添加了唯一的数组类型字段。

提前致谢。

【问题讨论】:

  • 您使用 5.0 S3 连接是否有同样的错误?
  • 你能把你的问题编辑成你从控制台消费者看到的输出吗?
  • @cricket_007 我已经添加了控制台消费者的输出。不,我没有在 5.0 s3 连接器上测试它,因为我们在生产中使用 4.1.1。这是特定于版本的问题吗?

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


【解决方案1】:

调用io.confluent.connect.avro.AvroData.defaultValueFromAvro(AvroData.java:1649) 将您读取的消息的 avro 模式转换为连接接收器的内部模式。我相信这与您的消息数据无关。这就是为什么AbstractKafkaAvroDeserializer 可以成功反序列化您的消息(例如通过kafka-avro-console-consumer),因为您的消息是有效的 avro 消息。如果您的默认值为null,而null 不是您的字段的有效值,则可能会发生上述异常。例如

{
   "name":"item_id",
   "type":{
      "type":"array",
      "items":[
         "string"
      ]
   },
   "default": null
}

我建议你远程调试连接,看看到底是什么失败了。

【讨论】:

    【解决方案2】:

    与您链接到的问题相同。

    In the source code,你可以看到这个条件。

      case ARRAY: {
        if (!jsonValue.isArray()) {
          throw new DataException("Invalid JSON for array default value: " + jsonValue.toString());
        }
    

    当架构类型在您的情况下定义为type:"array" 时,可能会引发异常,但有效负载本身具有null 值(或任何其他值类型)而不是实际上的数组,尽管您拥有定义为您的架构默认值。仅当 items 元素根本不存在时才应用默认值,而不是在 "items":null 时应用


    除此之外,我建议使用这样的模式,即记录对象,而不仅仅是命名数组,默认为空数组,而不是 null

    {
      "type" : "record",
      "name" : "Items",
      "namespace" : "com.example.avro",
      "fields" : [ {
        "name" : "item_id",
        "type" : {
          "type" : "array",
          "items" : [ "null", "string" ]
        },
        "default": []
      } ]
    }
    

    【讨论】:

    • 但是,如果 items 的值为 null,我们确保发送一个空数组,所以 "items": null 永远不会发生
    • @cricket_007 我遇到了同样的问题。下周我也会调查。但是有人可能知道为什么 kafka-avro-console-consumer 能够反序列化 s3 sink 无法反序列化的东西吗? Avro 是否有效。我错过了什么吗?谢谢!
    • @Vassilis Connect API 报错,控制台消费者不使用
    猜你喜欢
    • 2018-12-06
    • 1970-01-01
    • 2021-08-17
    • 2018-01-05
    • 1970-01-01
    • 1970-01-01
    • 2020-06-27
    • 2020-12-01
    • 2014-12-07
    相关资源
    最近更新 更多