【发布时间】:2020-11-20 23:21:17
【问题描述】:
我有一个生产者将 json 文件写入主题以供 kafka 消费者流读取。它的简单键值对。 我想通过添加/连接更多 JSON 键值行并发布到另一个主题来流式传输主题并丰富事件。 顺便说一句,没有任何值或键有任何共同之处。 我可能想多了,但是我该如何绕过这个逻辑?
【问题讨论】:
标签: json apache-kafka kafka-consumer-api apache-kafka-streams
我有一个生产者将 json 文件写入主题以供 kafka 消费者流读取。它的简单键值对。 我想通过添加/连接更多 JSON 键值行并发布到另一个主题来流式传输主题并丰富事件。 顺便说一句,没有任何值或键有任何共同之处。 我可能想多了,但是我该如何绕过这个逻辑?
【问题讨论】:
标签: json apache-kafka kafka-consumer-api apache-kafka-streams
我想你想在消费者端解码 JSON 消息。
如果您不关心架构,而只想将 JSON 作为 Map 处理,您可以使用 Jackson 库将 JSON 字符串读取为 Map<String,Object>。为此,您可以添加所需的字段,将其转换回 JSON 字符串并将其推送到新主题。
如果您想要一个架构,您需要存储有关它映射到哪个类或 JSON 架构或映射到此的某个 id 的信息,那么以下方法可以工作。
将架构信息存储在标头中
例如,您可以将 JSON 模式或 Java 类名存储在消息的标头中,同时生成和编写反序列化器以从标头中提取该信息并对其进行解码。
Deserializer#deserialize() 具有 Headers 参数。
default T deserialize(java.lang.String topic,
Headers headers,
byte[] data)
你可以做类似的事情..
objectMapper.readValue(data,
new Class.forName(
new String(headers.lastHeader("classname").value()
))
使用架构注册表
除此之外,还有一个来自 Confluent 的模式注册表,可以维护不同版本的模式。不过,您需要为此运行另一个进程。如果您打算使用它,您可能需要查看主题 naming strategy 并将其设置为 RecordNameStrategy,因为您在同一主题中有多个架构。
【讨论】: