【问题标题】:Kafka stream append to JSON as event enrichmentKafka 流作为事件丰富附加到 JSON
【发布时间】:2020-11-20 23:21:17
【问题描述】:

我有一个生产者将 json 文件写入主题以供 kafka 消费者流读取。它的简单键值对。 我想通过添加/连接更多 JSON 键值行并发布到另一个主题来流式传输主题并丰富事件。 顺便说一句,没有任何值或键有任何共同之处。 我可能想多了,但是我该如何绕过这个逻辑?

【问题讨论】:

    标签: json apache-kafka kafka-consumer-api apache-kafka-streams


    【解决方案1】:

    我想你想在消费者端解码 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,因为您在同一主题中有多个架构。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-06-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多