【问题标题】:Remove (corrupt) rows from Spark Streaming DataFrame that don't fit schema (incoming JSON data from Kafka)从 Spark Streaming DataFrame 中删除(损坏)不适合模式的行(来自 Kafka 的传入 JSON 数据)
【发布时间】:2023-03-05 20:27:01
【问题描述】:

我有一个从 Kafka 读取的 spark 结构化蒸汽应用程序。 这是我的代码的基本结构。

我创建了 Spark 会话。

val spark = SparkSession
  .builder
  .appName("app_name")
  .getOrCreate()

然后我从流中读取

val data_stream = spark
  .readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "server_list")
  .option("subscribe", "topic")
  .load()

在 Kafka 记录中,我将“值”转换为字符串。它从二进制转换为字符串。此时数据框中有1列

val df = data_stream
    .select($"value".cast("string") as "json")

基于预定义的架构,我尝试将 JSON 结构解析为列。然而,这里的问题是如果数据是“坏的”,或者不同的格式,那么它与定义的模式不匹配。因此,下一个数据帧 (df2) 将空值放入列中。

val df2 = df.select(from_json($"json", schema) as "data")
  .select("data.*")

我希望能够从 df2 中过滤掉在某个列(我用作数据库中的主键)中具有“null”的行,即忽略与架构不匹配的错误数据?

编辑:我在某种程度上能够做到这一点,但不是我想要的方式。 在我的流程中,我使用了一个使用.foreach(writer) 流程的查询。它的作用是打开与数据库的连接,处理每一行,然后关闭连接。 structured streaming 的文档提到了此过程所需的必需品。在 process 方法中,我从每一行获取值并检查我的主键是否为空,如果为空,我不将其插入数据库。

【问题讨论】:

    标签: apache-spark apache-kafka spark-structured-streaming


    【解决方案1】:

    只需过滤掉你不想要的任何空值:

    df2
      .filter(row => row("colName") != null)
    

    【讨论】:

      【解决方案2】:

      Kafka 将数据存储为原始字节数组格式。数据生产者和消费者需要就处理数据的结构达成一致。

      如果产生的消息格式发生变化,消费者需要调整以读取相同的格式。当您的数据结构不断发展时,问题就出现了,您可能需要在消费者端兼容。

      通过 Protobuff 定义消息格式解决了这个问题。

      【讨论】:

        猜你喜欢
        • 2015-09-13
        • 2017-02-10
        • 2019-05-25
        • 2016-11-03
        • 2021-10-12
        • 1970-01-01
        • 1970-01-01
        • 2016-11-04
        • 2018-04-04
        相关资源
        最近更新 更多