【发布时间】: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