【发布时间】: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、警告电子邮件……)。
我的问题是:实现这一目标的最佳方法是什么。
我的一些想法:
-
我能想到的第一件事是有一个UDF,它检查
NULLs,如果有任何问题,它会引发异常。我相信不可能通过 PySpark 将日志发送到 log4j,因为“spark”上下文无法在 UDF 中(在工作人员上)启动,并且想要使用默认值:log4jLogger = sc._jvm.org.apache.log4j logger = log4jLogger.LogManager.getLogger('PySpark Logger')
我能想到的第二件事是使用“foreach/foreachBatch”函数并将这个检查逻辑放在那里。
但我觉得这两种方法都像......像太多的自定义 - 我希望 Spark 内置了一些用于这些目的的东西。
【问题讨论】:
标签: json pyspark pyspark-sql spark-structured-streaming azure-eventhub