如何在PySpark中展平嵌套JSON数据?
问题描述
我有如下格式的JSON数据:
[ { "student_id": 1234, "room_id": "abc", "enrolled": false }, { "student_id": 4321, "room_id": "def", "enrolled": true, "enrollment": { "type": "home", "date": "01-01-2020" } }, { "student_id": 678, "room_id": "htf", "sports": { "team": "hockey", "position": "forward" } } ]
我通过以下代码实现了部分展平:
df = sc.parallelize(data).map(lambda x: json.dumps(x))
得到的结果如下:
| student_id | room_id | enrolled | enrollment | sports |
|---|---|---|---|---|
| 1234 | abc | false | NULL | NULL |
| 4321 | def | true | {home, 01-01-2020} | NULL |
| 678 | htf | NULL | NULL | {hockey, forward} |
需要进一步展平,得到如下格式的结果:
| student_id | room_id | enrolled | type | date | team | position |
|---|---|---|---|---|---|---|
| 1234 | abc | false | NULL | NULL | NULL | NULL |
| 4321 | def | true | home | 01-01-2020 | NULL | NULL |
| 678 | htf | NULL | NULL | NULL | hockey | forward |
解决方案
你可以利用PySpark的结构化数据操作直接展开嵌套字段,以下是两种可行的实现方式:
方式一:手动定义Schema解析JSON
如果需要严格控制数据类型,可以先定义Schema,再解析JSON字符串:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, IntegerType, StringType, BooleanType from pyspark.sql.functions import from_json, col # 初始化SparkSession(若未初始化) spark = SparkSession.builder.appName("FlattenNestedJSON").getOrCreate() # 定义JSON结构的Schema schema = StructType([ StructField("student_id", IntegerType(), nullable=True), StructField("room_id", StringType(), nullable=True), StructField("enrolled", BooleanType(), nullable=True), StructField("enrollment", StructType([ StructField("type", StringType(), nullable=True), StructField("date", StringType(), nullable=True) ]), nullable=True), StructField("sports", StructType([ StructField("team", StringType(), nullable=True), StructField("position", StringType(), nullable=True) ]), nullable=True) ]) # 将RDD转为DataFrame并解析JSON字符串 json_rdd = sc.parallelize(data).map(lambda x: json.dumps(x)) df = json_rdd.toDF("json_str") df = df.select(from_json(col("json_str"), schema).alias("data")).select("data.*") # 展平嵌套字段 flattened_df = df.select( "student_id", "room_id", "enrolled", col("enrollment.type").alias("type"), col("enrollment.date").alias("date"), col("sports.team").alias("team"), col("sports.position").alias("position") ) # 查看结果 flattened_df.show()
方式二:自动推断Schema(简化写法)
如果不需要严格指定数据类型,可以让Spark自动推断JSON结构,省去手动定义Schema的步骤:
from pyspark.sql import SparkSession # 初始化SparkSession spark = SparkSession.builder.appName("FlattenNestedJSON").getOrCreate() # 直接将JSON数据转为结构化DataFrame df = spark.read.json(sc.parallelize(data)) # 直接提取嵌套字段完成展平 flattened_df = df.select( "student_id", "room_id", "enrolled", "enrollment.type", "enrollment.date", "sports.team", "sports.position" ) # 查看结果 flattened_df.show()
两种方式都会生成你需要的展平结果,缺失的字段会自动填充为NULL。
内容的提问来源于stack exchange,提问作者unlocknew
相关产品推荐
相关产品推荐

