Spark DataFrame中嵌套JSON无法使用explode函数问题求助
解决Spark中explode字符串类型JSON数组的问题
你的问题根源很明确:explode函数仅支持数组(ARRAY)或映射(MAP)类型,但你的ApprovalJSON是字符串类型,直接调用自然会报错。要实现目标的扁平结构,需要分三步处理:
将字符串类型的JSON数组转换为Spark数组类型
使用from_json函数,先定义目标JSON的Schema,把字符串解析成Spark能识别的数组结构。假设你的ApprovalJSON格式是类似[{"ApprovalType":"经理审批","Status":"通过"},{"ApprovalType":"财务审批","Status":"待处理"}]的嵌套JSON数组,对应的Schema定义如下:from pyspark.sql.types import StructType, StructField, StringType, ArrayType approval_schema = ArrayType( StructType([ StructField("ApprovalType", StringType(), nullable=True), StructField("Status", StringType(), nullable=True) ]) )接着用
from_json转换列:df_with_array = df.withColumn("ApprovalArray", from_json(df["ApprovalJSON"], approval_schema))用explode展开数组
此时ApprovalArray已经是数组类型,可以正常使用explode:df_exploded = df_with_array.withColumn("ApprovalItem", explode("ApprovalArray"))提取嵌套字段并整理最终结构
从展开的ApprovalItem中提取ApprovalType和Status,保留原ID列,最后删除中间过渡列:final_df = df_exploded.select( "ID", df_exploded["ApprovalItem.ApprovalType"].alias("ApprovalType"), df_exploded["ApprovalItem.Status"].alias("Status") ).drop("ApprovalJSON", "ApprovalArray", "ApprovalItem")
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, explode from pyspark.sql.types import StructType, StructField, StringType, ArrayType # 初始化SparkSession(未初始化时执行) spark = SparkSession.builder.appName("FlattenJSON").getOrCreate() # 模拟原始DataFrame数据 data = [ (1, '[{"ApprovalType":"经理审批","Status":"通过"},{"ApprovalType":"财务审批","Status":"待处理"}]'), (2, '[{"ApprovalType":"人事审批","Status":"驳回"}]') ] df = spark.createDataFrame(data, ["ID", "ApprovalJSON"]) # 定义JSON Schema approval_schema = ArrayType( StructType([ StructField("ApprovalType", StringType(), nullable=True), StructField("Status", StringType(), nullable=True) ]) ) # 转换字符串为数组类型 df_with_array = df.withColumn("ApprovalArray", from_json(df["ApprovalJSON"], approval_schema)) # 展开数组 df_exploded = df_with_array.withColumn("ApprovalItem", explode("ApprovalArray")) # 提取字段并生成最终表 final_df = df_exploded.select( "ID", "ApprovalItem.ApprovalType", "ApprovalItem.Status" ).drop("ApprovalJSON", "ApprovalArray", "ApprovalItem") # 查看结果 final_df.show()
运行后就能得到ID、ApprovalType、Status的扁平结构,完全符合你的需求。
内容的提问来源于stack exchange,提问作者Mark
相关产品推荐
相关产品推荐

