【问题标题】:Can I write custom kafka connect transform for converting JSON to AVRO?我可以编写自定义 kafka 连接转换以将 JSON 转换为 AVRO 吗?
【发布时间】:2018-02-13 09:49:59
【问题描述】:

我想使用 kafka-connect-hdfs 将无模式 json 记录从 kafka 写入 hdfs 文件。 如果我使用 JsonConvertor 作为键/值转换器,那么它不起作用。但是如果我使用的是 StringConvertor,那么它会将 json 写为转义字符串。

例如:

实际的 json -

{"name":"test"}

数据写入 hdfs 文件 -

"{\"name\":\"test\"}"

预期输出到 hdfs 文件 -

{"name":"test"}

有没有什么方法或替代方法可以实现这一点,或者我只能将它与模式一起使用?

以下是我尝试使用 JSONConvertor 时遇到的异常:

[2017-09-06 14:40:19,344] ERROR Task hdfs-sink-0 threw an uncaught and unrecoverable exception (org.apache.kafka.connect.runtime.WorkerTask:148)
org.apache.kafka.connect.errors.DataException: JsonConverter with schemas.enable requires "schema" and "payload" fields and may not contain additional fields. If you are trying to deserialize plain JSON data, set schemas.enable=false in your converter configuration.
    at org.apache.kafka.connect.json.JsonConverter.toConnectData(JsonConverter.java:308)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:406)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:250)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:180)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:148)
    at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:146)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:190)
    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)

quickstart-hdfs.properties的配置:

name=hdfs-sink
connector.class=io.confluent.connect.hdfs.HdfsSinkConnector
tasks.max=1
topics=test_hdfs_avro
hdfs.url=hdfs://localhost:9000
flush.size=1
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter

connect-avro-standalone.properties的配置:

bootstrap.servers=localhost:9092
schemas.enable=false
key.converter.schemas.enable=false
value.converter.schemas.enable=false

【问题讨论】:

  • "如果我使用 JsonConvertor 作为键/值转换器,那么它不起作用。" -> 你能描述一下你看到的问题吗?你得到的错误?
  • @RobinMoffatt 我已经添加了我为此连接器配置的异常和属性。

标签: apache-kafka apache-kafka-connect


【解决方案1】:

当您在连接器的配置属性中指定转换器时,您需要包含与此转换器相关的所有属性,无论这些属性是否也包含在工作器的配置中。

在上面的示例中,您需要同时指定两者:

value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false

快速启动-hdfs.properties中。

仅供参考,JSON 导出即将在 HDFS 连接器中推出。在此处跟踪相关的拉取请求:https://github.com/confluentinc/kafka-connect-hdfs/pull/196

更新:JsonFormat 已合并到 master 分支。

【讨论】:

  • 感谢您的信息。当请求合并到 master 并可供使用时,您能否更新此线程?
  • JsonFormat 已合并。
  • 非常感谢。所以要使用这段代码。我必须下载master并手动构建它还是有其他方法?
  • 没错。它将包含在即将发布的 Confluent 平台中。在此之前,早期采用者需要从源代码构建。
  • 好的,谢谢,但是否有手册可以找到我需要构建的所有项目以及手动构建源代码的步骤?我尝试构建 common 和 hdfs,但 hdfs 项目显示其他依赖项错误。它无法找到对融合 Maven 存储库的依赖项。有什么想法吗?
猜你喜欢
  • 2021-04-07
  • 2019-12-27
  • 2019-10-03
  • 2010-09-20
  • 2022-08-17
  • 2012-02-16
  • 1970-01-01
  • 1970-01-01
  • 2011-02-10
相关资源
最近更新 更多