【问题标题】:Kafka s3 json connector卡夫卡 s3 json 连接器
【发布时间】:2018-06-23 07:17:40
【问题描述】:

我尝试使用最新的 kafka (confluent-platform-2.11) 连接将 Json 放入 s3。我在 quickstart-s3.properties 文件中设置了 format.class=io.confluent.connect.s3.format.json.JsonFormat

和加载连接器:

~$ confluent load s3-sink  {   "name": "s3-sink",   "config": {
    "connector.class": "io.confluent.connect.s3.S3SinkConnector",
    "tasks.max": "1",
    "topics": "s3_hose",
    "s3.region": "us-east-1",
    "s3.bucket.name": "some-bucket-name",
    "s3.part.size": "5242880",
    "flush.size": "1",
    "storage.class": "io.confluent.connect.s3.storage.S3Storage",
    "format.class": "io.confluent.connect.s3.format.json.JsonFormat",
    "schema.generator.class": "io.confluent.connect.storage.hive.schema.DefaultSchemaGenerator",
    "partitioner.class": "io.confluent.connect.storage.partitioner.DefaultPartitioner",
    "schema.compatibility": "NONE",
    "name": "s3-sink"   },   "tasks": [
    {
      "connector": "s3-sink",
      "task": 0
    }   ],   "type": null }

然后我给卡夫卡发了一行:

~$ kafka-console-producer --broker-list localhost:9092 --topic s3_hose

{"q":1}

我在连接器日志中看到 Avro 转换异常:

[2018-01-14 14:41:30,832] ERROR WorkerSinkTask{id=s3-sink-0} Task threw an uncaught and unrecoverable exception (org.apache.kafka.connect.runti me.WorkerTask:172) org.apache.kafka.connect.errors.DataException: s3_hose
        at io.confluent.connect.avro.AvroConverter.toConnectData(AvroConverter.java:96)
        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:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748) Caused by: org.apache.kafka.common.errors.SerializationException: Error deserializing Avro message for id -1 Caused by: org.apache.kafka.common.errors.SerializationException: Unknown magic byte!

如果我设置 format.class=io.confluent.connect.s3.format.json.JsonFormat,为什么会尝试使用某些 Avro 转换器?

【问题讨论】:

    标签: apache-kafka-connect


    【解决方案1】:

    此消息与转换器有关。这与最终的输出格式不同。它用于将 Kafka 中的数据转换为连接数据 API 格式,因此连接器有一些标准可以使用。要设置转换器,您可以

    1) 将 key.converter 和 value.converter 设置为工作器属性文件中的内置 JsonConverter,使其成为工作器中运行的所有连接器的默认值

    2) 在连接器级别设置 key.converter 和 value.converter 属性以覆盖在工作人员级别设置的内容

    注意,由于这是一个接收器连接器,因此您非常希望将转换器与主题中的数据类型相匹配,以便正确转换。

    【讨论】:

    • connect-distributed.propertiesconnect-standalone.properties 具有 key.converter 值。转换器设置为org.apache.kafka.connect.json.JsonConverter,但它不起作用。添加 key.converter=org.apache.kafka.connect.json.JsonConverter, value.converter=org.apache.kafka.connect.json.JsonConverter, key.converter.schemas.enable=false, value.converter。 schemas.enable=false 到 quickstart-s3.properties 完成了这项工作!谢谢!
    猜你喜欢
    • 2018-08-29
    • 2020-08-27
    • 2016-09-03
    • 2017-11-14
    • 1970-01-01
    • 2023-02-15
    • 2020-06-20
    • 2018-01-18
    • 2017-03-02
    相关资源
    最近更新 更多