【问题标题】:Validating Schema of Column with StructType in Pyspark 2.4在 Pyspark 2.4 中使用 StructType 验证列的模式
【发布时间】:2023-04-11 01:04:01
【问题描述】:

我有一个数据框,其中有一列是 JSON 字符串

from pyspark.sql import SparkSession
from pyspark.sql.types import *
import pyspark.sql.functions as F

sc = SparkSession.builder.getOrCreate()

l = [
    (1, """{"key1": true, "nested_key": {"mylist": ["foo", "bar"], "mybool": true}})"""),
    (2, """{"key1": true, "nested_key": {"mylist": "", "mybool": true}})"""),
    ]
df = sc.createDataFrame(l, ["id", "json_str"])

并希望使用架构解析json_str 列和from_json


schema = StructType([
    StructField("key1", BooleanType(), False),
    StructField("nested_key", StructType([
        StructField("mylist", ArrayType(StringType()), False),
        StructField("mybool", BooleanType(), False)
     ]))
])

df = df.withColumn("data", F.from_json(F.col("json_str"), schema))
df.show(truncate=False)
+---+--------------------------+
|id |data                      |
+---+--------------------------+
|1  |[true, [[foo, bar], true]]|
|2  |[true, [, true]]          |
+---+--------------------------+

正如所见,第二行不符合schema 中的架构,因此即使我在StructField 中将False 传递给nullable,它也是空的。对我的管道来说,重要的是,如果存在不符合架构定义的数据,则会以某种方式引发警报,但我不确定在 Pyspark 中执行此操作的最佳方法。真实数据有很多很多键,其中一些是嵌套的,因此使用某种形式的isNan 检查每个键是不可行的,因为我们已经定义了模式,感觉应该可以利用它。

如果重要的话,我不一定需要检查整个数据框的架构,我真的是在检查 StructType 列的架构之后

【问题讨论】:

  • 这似乎仍然是一个悬而未决的问题。我能找到验证自定义 JSON 的唯一方法是通过 here 所描述的自定义读取器/写入器,请查看脚注:)

标签: apache-spark pyspark apache-spark-sql


【解决方案1】:

查看options 参数: https://spark.apache.org/docs/2.3.1/api/python/pyspark.sql.html?highlight=from_json#pyspark.sql.functions.from_json

这有点模糊,但它允许您将dict 传递给这里的底层方法: https://spark.apache.org/docs/2.3.1/api/python/pyspark.sql.html?highlight=from_json#pyspark.sql.DataFrameReader.json

您可能会成功通过 options={'mode' : 'FAILFAST'} 之类的内容。

【讨论】:

  • 这是一个很好的建议,也是一个可靠的发现,但我在帖子中的示例中尝试过,但没有成功
猜你喜欢
  • 2021-03-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-10-21
  • 2011-11-20
  • 2015-09-03
  • 1970-01-01
  • 2021-01-15
相关资源
最近更新 更多