如何在PySpark中展平含嵌套数组的JSON并拆分存储至不同文件?
问题描述
从API获取到如下JSON数据:
{ "field1": "value1", "field2": "value2", "message_records": [ { "field3": "value3", "field4": "value4" }, { "field5": "value5", "field6": "value6" } ], "messages": [ { "field7": "value3", "field8": "value4" }, { "field9": "value5", "field10": "value6" }, { "field11": "value5", "field12": "value6" } ] }
需用Python结合PySpark完成以下操作:
- 将嵌套数组
message_records和messages展开为独立记录 - 保留
field1、field2作为每条记录的公共字段 - 将
message_records对应的数据写入单独文件,messages对应的数据写入另一单独文件
解决方案
1. 初始化PySpark会话
先创建SparkSession实例:
from pyspark.sql import SparkSession from pyspark.sql.functions import explode # 初始化SparkSession spark = SparkSession.builder \ .appName("FlattenNestedJSON") \ .getOrCreate()
2. 加载JSON数据
将给定JSON数据转为PySpark DataFrame:
# 定义JSON数据 json_data = { "field1": "value1", "field2": "value2", "message_records": [ {"field3": "value3", "field4": "value4"}, {"field5": "value5", "field6": "value6"} ], "messages": [ {"field7": "value3", "field8": "value4"}, {"field9": "value5", "field10": "value6"}, {"field11": "value5", "field12": "value6"} ] } # 转为DataFrame df = spark.createDataFrame([json_data])
3. 处理并导出message_records数据
用explode展开嵌套数组,保留公共字段后写入文件:
# 展开message_records数组 message_records_df = df.select( "field1", "field2", explode("message_records").alias("record") ).select( "field1", "field2", "record.*" ) # 写入文件(示例用Parquet格式,可替换为csv/json等) message_records_df.write.mode("overwrite").parquet("./message_records_output")
4. 处理并导出messages数据
同样用explode处理messages数组:
# 展开messages数组 messages_df = df.select( "field1", "field2", explode("messages").alias("msg") ).select( "field1", "field2", "msg.*" ) # 写入文件 messages_df.write.mode("overwrite").parquet("./messages_output")
补充说明
explode函数负责将数组中每个元素转为独立行select("record.*")可以把嵌套结构体的字段展开为顶级列- 写入文件时可根据需求调整格式,输出路径也可自定义
mode("overwrite")表示覆盖已有文件,可替换为append等模式
内容的提问来源于stack exchange,提问作者Venkatesh
相关产品推荐
相关产品推荐

