【问题标题】:Kafka Connect JSON format卡夫卡连接 JSON 格式
【发布时间】:2018-08-29 21:44:57
【问题描述】:

我在 Kafka 中有一个标题为 newtest 的主题,其中包含三条消息:

Hello 
Is anybody out there
Can you hear me

...我有以下连接作业的配置:

{
    "name":"connect-test-9",
    "config":
    {
        "connector.class":"FileStreamSink",
        "file":"connector-test",
        "topics":"newtest",
        "name":"connect-test-9",
        "value.converter":"org.apache.kafka.connect.storage.StringConverter",
        "value.converter.schemas.enable":"false",
        "key.converter":"org.apache.kafka.connect.storage.StringConverter",
        "key.converter.schemas.enable":"false",
        "transforms":"Hoist, AddTimestamp",
        "transforms.Hoist.type":"org.apache.kafka.connect.transforms.HoistField$Value",
        "transforms.Hoist.field":"line",
        "transforms.AddTimestamp.type":"org.apache.kafka.connect.transforms.InsertField$Value",
        "transforms.AddTimestamp.timestamp.field":"Timestamp"
    }
}

我在文件 connector-test 中得到以下输出:

Struct{line=Hello,Timestamp=Mon Mar 12 14:50:34 PDT 2018}
Struct{line=Is anybody out there,Timestamp=Mon Mar 12 14:50:44 PDT 2018}
Struct{line=Can you hear me,Timestamp=Mon Mar 12 14:50:52 PDT 2018}

我想得到这个:

{"line":"Hello","Timestamp":"Mon Mar 12 14:50:34 PDT 2018"}
{"line":"Is anybody out there","Timestamp":"Mon Mar 12 14:50:44 PDT 2018"}
{"line":"Can you hear me","Timestamp":"Mon Mar 12 14:50:52 PDT 2018"}

我试过改变 value.converter,不好(解析异常)。我还有另一个主题,其中消息已经是 Json,并且在那里解析成功,并且我可以在没有 Hoist 的情况下添加时间戳。但是我的输出是相同的非Json格式{key1=value1,key2=value2}

有什么方法可以得到正确的 JSON 输出?


这是我看到的解析异常:

com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'Can': was expecting ('true', 'false', or 'null')

【问题讨论】:

  • 如果你想要 JSON,为什么要使用字符串转换器?你能在问题中添加解析异常吗?
  • 已编辑问题解析异常。源主题是一个纯字符串。转换器不必与源相匹配吗?
  • 您的源数据不是 JSON。这是一个字符串。如果要使用 JSON,则需要发送 JSON。 Struct{line=message, Timestamp=...} 输出对于您提供的数据是正确的。
  • 我的来源不是文件。这是一个话题。该主题具有三个字符串(非 JSON)消息。我想使用 Connect 将它们转换为 JSON 并添加时间戳。这可能吗?
  • 顺便更新了我的答案

标签: apache-kafka apache-kafka-connect


【解决方案1】:

要获得 JSON 输出,您需要使用 JsonConverter 而不是 StringConverter the converter happens before the sink, and after the consumer deserialization

您的数据已经在 Kafka 中,并且您使用了 Source Connector,可能带有 StringConveter 来摄取,然后您转换为内部 Struct,可以使用 Sink Connector 和不同的 Converter 类型进行设置

【讨论】:

    猜你喜欢
    • 2018-06-23
    • 2017-11-14
    • 1970-01-01
    • 2023-02-15
    • 2020-08-27
    • 2016-09-03
    • 2018-01-18
    • 2017-03-02
    相关资源
    最近更新 更多