【发布时间】: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