【问题标题】:How to capture incorrect (corrupt) JSON records in (Py)Spark Structured Streaming?如何在(Py)Spark Structured Streaming 中捕获不正确(损坏)的 JSON 记录?
【发布时间】:2019-11-10 12:25:44
【问题描述】:

我有一个 Azure Eventhub,它是流数据(JSON 格式)。 我将其读取为 Spark 数据帧,使用 from_json(col("body"), schema) 解析传入的“正文”,其中 schema 是预定义的。在代码中,它看起来像:

from pyspark.sql.functions import col, from_json
from pyspark.sql.types import *

schema = StructType().add(...) # define the incoming JSON schema 

df_stream_input = (spark
.readStream
.format("eventhubs")
.options(**ehConfInput)
.load()
.select(from_json(col("body").cast("string"), schema)
)

现在 = 如果传入 JSON 的架构与定义的架构之间存在一些不一致(例如,源 eventthub 开始以新格式发送数据,恕不另行通知),from_json() 函数不会抛出错误 = 相反,它会将 NULL 放入字段中,这些字段存在于我的 schema 定义中,但不存在于 JSONs eventthub 发送中。

我想捕获这些信息并将其记录在某处(Spark 的 log4j、Azure Monitor、警告电子邮件……)。

我的问题是:实现这一目标的最佳方法是什么。

我的一些想法:

  1. 我能想到的第一件事是有一个UDF,它检查NULLs,如果有任何问题,它会引发异常。我相信不可能通过 PySpark 将日志发送到 log4j,因为“spark”上下文无法在 UDF 中(在工作人员上)启动,并且想要使用默认值:

    log4jLogger = sc._jvm.org.apache.log4j logger = log4jLogger.LogManager.getLogger('PySpark Logger')

  2. 我能想到的第二件事是使用“foreach/foreachBatch”函数并将这个检查逻辑放在那里。

但我觉得这两种方法都像......像太多的自定义 - 我希望 Spark 内置了一些用于这些目的的东西。

【问题讨论】:

    标签: json pyspark pyspark-sql spark-structured-streaming azure-eventhub


    【解决方案1】:

    tl;dr您必须使用foreach 或foreachBatch 运算符自己执行此检查逻辑。


    原来我错误地认为columnNameOfCorruptRecord 选项可能是一个答案。它不会起作用。

    首先,由于this,它不起作用:

    case _: BadRecordException => null
    

    其次,由于this 只是禁用了任何其他解析模式(包括似乎与columnNameOfCorruptRecord 选项一起使用的PERMISSIVE):

    new JSONOptions(options + ("mode" -> FailFastMode.name), timeZoneId.get))
    

    换句话说,您唯一的选择是使用列表中的第二项,即foreach 或foreachBatch,并自己处理损坏的记录。

    解决方案可以使用from_json,同时保留初始body 列。任何带有不正确 JSON 的记录都会以null 和foreach* 的结果列结束,例如

    def handleCorruptRecords:
      // if json == null the body was corrupt
      // handle it
    
    df_stream_input = (spark
      .readStream
      .format("eventhubs")
      .options(**ehConfInput)
      .load()
      .select("body", from_json(col("body").cast("string"), schema).as("json"))
    ).foreach(handleCorruptRecords).start()
    

    【讨论】:

    • 好的 :) 感谢您的回答 - 也许将来会有一些变化
    • @mlC 我想我可能找到了另一种解决方案。介意检查一下吗? :) 您只需使用foreachBatch 并使用spark.read.json(batchDataFrame.rdd)。那可以工作。感谢您接受我的回答。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-10-26
    • 2020-09-12
    • 1970-01-01
    • 2020-03-08
    • 2020-03-19
    • 1970-01-01
    • 2021-03-06
    相关资源
    最近更新 更多