【问题标题】:Kafka s3 connect "Value is not Struct type" errorKafka s3 连接“值不是结构类型”错误
【发布时间】:2018-10-26 15:01:56
【问题描述】:

我使用以下参数加载 s3 连接器:

confluent load s3-sink
{
  "name": "s3-sink",
  "config": {
    "connector.class": "io.confluent.connect.s3.S3SinkConnector",
    "tasks.max": "1",
    "topics": "s3_topic",
    "s3.region": "us-east-1",
    "s3.bucket.name": "some_bucket",
    "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.FieldPartitioner",
    "schema.compatibility": "NONE",
    "partition.field.name": "f1",
    "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",
    "name": "s3-sink"
  },
  "tasks": [
    {
      "connector": "s3-sink",
      "task": 0
    }
  ],
  "type": null
}

接下来我用 kafka-console-producer JSON 发送它:

{"f1":"partition","data":"some data"}

我在连接日志中收到以下错误:

[2018-05-16 16:32:05,150] ERROR Value is not Struct type. (io.confluent.connect.storage.partitioner.FieldPartitioner:67)
[2018-05-16 16:32:05,150] ERROR WorkerSinkTask{id=s3-sink-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not re
cover until manually restarted. (org.apache.kafka.connect.runtime.WorkerSinkTask:515)
io.confluent.connect.storage.errors.PartitionException: Error encoding partition.

我记得它在前一段时间有效。
现在我使用 Confluent Open Source v. 4.1

【问题讨论】:

    标签: apache-kafka apache-kafka-connect confluent-platform


    【解决方案1】:

    从 Confluent 4.1 版本开始,FieldPartitioner does not support JSON 字段提取。

    You could instead use kafka-avro-console-producer 使用 Avro Schema 发送相同的 JSON blob,那么它应该可以工作

    这是您要使用的属性

    --property value.schema='{"type":"record","name":"myrecord","fields":[{"name":"f1","type":"string"},{"name":"data","type":"string"}]}'

    然后就可以发送了

    {"f1":"partition","data":"some data"}
    

    您需要在 Connect 中使用这些属性

    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schemas.enable": "true",
    

    【讨论】:

    • 谢谢!如何使用 Java Kafka 驱动程序做到这一点?
    • 发送 Avro?在此处查看示例。 github.com/confluentinc/examples/tree/4.0.x/kafka-clients
    • 嘿@cricket_007,你有partition.field.name 的示例(.properties 格式),使用多个分区,基于字段名称(例如:field1=value/field2=value/字段3=值)?谢谢
    • @Jul partition.field.name=field1,field2,field3。不过,请再次使用 Connect 分布式模式
    猜你喜欢
    • 2017-08-30
    • 2018-11-30
    • 2021-09-10
    • 2020-06-02
    • 2018-01-31
    • 2015-01-08
    • 1970-01-01
    • 2022-11-30
    • 2021-02-25
    相关资源
    最近更新 更多