【问题标题】:Processing json data from kafka using structured streaming使用结构化流处理来自 kafka 的 json 数据
【发布时间】:2020-10-17 04:55:05
【问题描述】:

我想将来自 Kafka 的传入 JSON 数据转换为数据帧。

我正在使用Scala 2.12 的结构化流媒体

大部分人都添加了硬编码的schema,但是如果json可以有额外的字段,就需要每次都修改代码库,比较繁琐。

一种方法是将其写入文件并推断它,但我宁愿避免这样做。

还有其他方法可以解决这个问题吗?

编辑:找到了一种将json字符串转换为数据帧的方法,但无法从流源中提取它,可以提取它吗?

【问题讨论】:

  • 编辑:我通过模式注册表api成功创建了一个模式对象,它返回了一个模式对象,可以将它转换为一个StructType对象吗??

标签: json scala apache-kafka schema spark-structured-streaming


【解决方案1】:

将数据读取为字符串,然后将其转换为 map[string,String],这样您就可以在不知道其架构的情况下处理任何 json

【讨论】:

  • 解决方案看起来不错,但是如果数据包含其他类型,我会收到编码错误java.lang.ClassCastException: java.lang.Double cannot be cast to java.lang.String 是否有解决方法/您能提供一个基本的代码示例吗?
  • 当您尝试将 json 字符串转换为 Map[String,String] 时,会出现任何错误。当您尝试访问特定值(例如 Double )时,您将面临此类问题。解决方法将使用 'String.valueOf()'
【解决方案2】:
  1. 一种方法是将架构本身存储在消息头中(而不是键或值)。

    虽然这会增加消息大小,但解析 JSON 值会很容易,而无需任何外部资源,例如文件或模式注册表。

    新消息可以具有新架构,同时旧消息仍可以使用其旧架构本身进行处理,因为架构位于消息本身内。

  2. 或者,您可以 version 模式并在消息头中为每个模式包含一个 id(或)在键或值中包含一个 magic byte 并推断那里的架构。

    这种方法后面跟着Confluent Schema registry。它允许您基本上浏览同一架构的不同版本,并查看您的架构如何随着时间的推移而演变。

【讨论】:

  • 我有一个模式注册表,用于存储数据的 Avro 模式。 (JsonString , AvroSchema )=> Dataframe 不会有序列化问题吗?
  • @coding_potato 您可能想要编写一个自定义的反序列化器来检测数据的类型(avro 或 json)并相应地反序列化它。 Avro 模式,如果你在融合的 avro 序列化程序中使用,通常会放置一个 magic byte 来标记 avro 模式,如果你找到那个魔术字节,那么你可以将它反序列化为 avro,任何机会如果在反序列化时遇到异常,请尝试将其反序列化为 json。
【解决方案3】:

基于 JavaTechnical answer ,最好的方法是使用模式注册表和 avro 数据而不是 json,没有硬编码模式(目前)。

包含您的架构名称和 ID 作为标头,并使用它们从架构注册表中读取架构。

使用from_avro 函数将该数据转换为df!

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-09-30
    • 1970-01-01
    • 2019-11-25
    • 2019-08-24
    • 2019-10-03
    • 1970-01-01
    • 2019-01-14
    • 2019-11-28
    相关资源
    最近更新 更多