PySpark解析JSON字符串列报错:无需硬编码Schema展开为多列
PySpark 动态解析JSON列(无需硬编码Schema)
问题场景
DataFrame的value列存储如下格式的JSON字符串:
{ "result":{ "version":"1.2", "timeStamp":"2023-08-14 14:00:12", "description":"", "data":{ "DateTime_Received":"2023-08-14T14:01:10.4516457+01:00", "DateTime_Actual":"2023-08-14T14:00:12", "OtherInfo":null, "main":[ { "Status":0, "ID":111, "details":null } ] }, "tn":"aaa" } }
需要将该JSON列展开为多列,尝试使用schema_of_json自动生成Schema时触发错误:
df_decoded = df_decoded.withColumn("json_column", F.when(F.col("value").isNotNull(), F.col("value")).otherwise("{}")) json_schema = df_decoded.select(F.schema_of_json(F.col("json_column"))).collect()[0][0]
错误信息:
AnalysisException: cannot resolve 'schema_of_json(json_column)' due to data type mismatch: The input json should be a foldable string expression and not null; however, got json_column.;
错误原因
schema_of_json函数要求输入是可折叠的常量字符串(即编译阶段就能确定的固定值),不能直接传入DataFrame列引用——因为列值是运行时才会生成的,函数无法在查询计划阶段推断出Schema,因此报错。
解决方案
方法1:通过样本JSON生成Schema(推荐,适合复杂结构)
先从DataFrame中提取一个非空的JSON字符串作为样本,用这个样本生成Schema后,再解析整个列:
from pyspark.sql import functions as F # 提取一个非空的JSON样本(确保DataFrame中有非空值) sample_json = df_decoded.filter(F.col("value").isNotNull()).select("value").first()[0] # 用样本生成完整Schema json_schema = F.schema_of_json(sample_json) # 解析JSON列 df_decoded = df_decoded.withColumn("parsed_json", F.from_json(F.col("value"), json_schema)) # 逐层展开嵌套结构:先展开result层 df_flattened = df_decoded.select("*", "parsed_json.result.*") # 再展开data子层 df_flattened = df_flattened.select("*", "data.*") # 展开数组类型的main字段(如果有多个元素,用explode拆分) df_flattened = df_flattened.withColumn("main", F.explode(F.col("main"))) df_flattened = df_flattened.select("*", "main.*") # 清理中间冗余列 df_final = df_flattened.drop("value", "parsed_json", "data", "main")
方法2:使用get_json_object逐层提取(适合简单结构)
如果不想依赖样本生成Schema,可以用get_json_object按JSON路径逐层提取字段:
from pyspark.sql import functions as F # 提取顶层result对象 df_decoded = df_decoded.withColumn("result", F.get_json_object(F.col("value"), "$.result")) # 提取result下的基础字段 df_decoded = df_decoded.withColumn("version", F.get_json_object(F.col("result"), "$.version")) df_decoded = df_decoded.withColumn("timeStamp", F.get_json_object(F.col("result"), "$.timeStamp")) df_decoded = df_decoded.withColumn("tn", F.get_json_object(F.col("result"), "$.tn")) df_decoded = df_decoded.withColumn("description", F.get_json_object(F.col("result"), "$.description")) # 提取data子对象 df_decoded = df_decoded.withColumn("data", F.get_json_object(F.col("result"), "$.data")) # 提取data下的时间字段 df_decoded = df_decoded.withColumn("DateTime_Received", F.get_json_object(F.col("data"), "$.DateTime_Received")) df_decoded = df_decoded.withColumn("DateTime_Actual", F.get_json_object(F.col("data"), "$.DateTime_Actual")) df_decoded = df_decoded.withColumn("OtherInfo", F.get_json_object(F.col("data"), "$.OtherInfo")) # 解析并展开main数组 df_decoded = df_decoded.withColumn("main", F.from_json(F.get_json_object(F.col("data"), "$.main"), "array<struct<Status:int,ID:int,details:string>>")) df_decoded = df_decoded.withColumn("main", F.explode(F.col("main"))) df_decoded = df_decoded.withColumn("Status", F.col("main.Status")) df_decoded = df_decoded.withColumn("ID", F.col("main.ID")) df_decoded = df_decoded.withColumn("details", F.col("main.details")) # 清理中间列 df_final = df_decoded.drop("value", "result", "data", "main")
内容的提问来源于stack exchange,提问作者mpr
相关产品推荐
相关产品推荐

