如何将嵌套JSON转换为指定格式的PySpark DataFrame?
嵌套JSON转PySpark DataFrame格式优化
问题场景
需要将以下嵌套JSON转换为PySpark DataFrame,当前代码输出为数组列形式,无法得到预期的行展开结果。
原始JSON
{ "key1": 0.75, "values":[ { "id": 2313, "val1": 350, "val2": 6000 }, { "id": 2477, "val1": 340, "val2": 6500 } ] }
当前代码
import json from pyspark.sql import SparkSession spark = SparkSession.builder.appName("CreateDataFrame").getOrCreate() json_string = json.dumps({ "key1": 0.75, "values":[ { "id": 2313, "val1": 350, "val2": 6000 }, { "id": 2477, "val1": 340, "val2": 6500 } ] }) df = spark.read.json(spark.sparkContext.parallelize([json_string])) df = df.select("key1", "values.id", "values.val1", "values.val2") df.show()
当前输出
+----+-------------+-------------+-------------+ |key1| id| val1| val2| +----+-------------+-------------+-------------+ |0.75| [2313, 2477]| [350, 340]| [6000, 6500]| +----+-------------+-------------+-------------+
预期输出
+----+----+----+----+ |key1| id|val1|val2| +----+----+----+----+ |0.75|2313| 350|6000| |0.75|2477| 340|6500| +----+----+----+----+
解决方案
问题核心是values为数组类型,直接选取字段会生成数组列。需要用explode函数拆分数组为多行,再提取结构体字段。
优化后代码
import json from pyspark.sql import SparkSession from pyspark.sql.functions import explode, col spark = SparkSession.builder.appName("CreateDataFrame").getOrCreate() json_string = json.dumps({ "key1": 0.75, "values":[ { "id": 2313, "val1": 350, "val2": 6000 }, { "id": 2477, "val1": 340, "val2": 6500 } ] }) df = spark.read.json(spark.sparkContext.parallelize([json_string])) # 拆分values数组为多行,保留key1 df_exploded = df.select("key1", explode(col("values")).alias("value")) # 从结构体中提取各字段 df_result = df_exploded.select("key1", col("value.id"), col("value.val1"), col("value.val2")) df_result.show()
关键说明
explode(col("values")):将数组中的每个元素拆分为独立行,key1的值会自动复制到每一行。alias("value"):给拆分后的结构体列命名,便于后续字段提取。col("value.id"):从结构体列中提取指定字段,最终得到扁平化的DataFrame。
内容的提问来源于stack exchange,提问作者Leonard
相关产品推荐
相关产品推荐

