PySpark解析嵌套JSON:如何实现目标结构化输出?
PySpark实现嵌套JSON数据扁平化
实现思路
核心是逐层拆解嵌套结构:先展开IDArray获取有效ID,再根据ID动态提取IDStruct中的对应数据,最后展开内层的LegCorr数组,最终筛选出目标字段。
完整代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import explode, col, expr # 初始化SparkSession spark = SparkSession.builder.appName("FlattenNestedJSON").getOrCreate() # 示例数据(也可通过spark.read.json读取外部文件) data = [ { "MainTag": { "GroupId": "10C81", "IDArray": ["ABC-XYZ-123"], "IDStruct": { "DSA-ASA-211": None, "BSA-ASA-211": None, "ABC-XYZ-123": [ { "BagId": "42425fsdfs", "TravelerId": "1234567", "LegCorr": [ {"DelID": "SQH", "SegID": "PQR-UVW"}, {"DelID": "GFS", "SegID": "GHS-UVW"} ] } ] } } } ] # 创建初始DataFrame df = spark.createDataFrame(data) # 步骤1:展开IDArray,得到每个有效ID,同时保留GroupId和IDStruct df_step1 = df.select( col("MainTag.GroupId"), explode(col("MainTag.IDArray")).alias("ID"), col("MainTag.IDStruct") ) # 步骤2:根据ID动态提取IDStruct中对应的数据,过滤空值 df_step2 = df_step1.withColumn( "id_data", expr(f"IDStruct.`{col('ID')}`") # 反引号处理带特殊字符的字段名 ).filter(col("id_data").isNotNull()) # 步骤3:展开id_data数组,获取单条Bag数据 df_step3 = df_step2.select( col("GroupId"), col("ID"), explode(col("id_data")).alias("id_struct") ) # 步骤4:展开LegCorr数组,提取所有目标字段 final_df = df_step3.select( col("GroupId"), col("ID"), col("id_struct.BagId"), col("id_struct.TravelerId"), explode(col("id_struct.LegCorr")).alias("leg_corr") ).select( col("GroupId"), col("ID"), col("BagId"), col("TravelerId"), col("leg_corr.DelID"), col("leg_corr.SegID") ) # 查看结果 final_df.show(truncate=False)
代码说明
- 步骤1:用
explode拆分IDArray,将每个有效ID转为单独行,同时保留关联的GroupId和完整IDStruct。 - 步骤2:通过
expr动态引用IDStruct中与当前行ID匹配的字段,过滤掉空值(对应IDArray外的无效ID)。 - 步骤3:再次用
explode拆分ID对应的结构体数组,得到每个Bag的详细数据。 - 步骤4:最后拆分
LegCorr数组,提取所有目标字段,完成扁平化。
输出结果
+-------+-----------+-----------+----------+------+---------+ |GroupId|ID |BagId |TravelerId|DelID |SegID | +-------+-----------+-----------+----------+------+---------+ |10C81 |ABC-XYZ-123|42425fsdfs |1234567 |SQH |PQR-UVW | |10C81 |ABC-XYZ-123|42425fsdfs |1234567 |GFS |GHS-UVW | +-------+-----------+-----------+----------+------+---------+
内容的提问来源于stack exchange,提问作者Vaibhav
相关产品推荐
相关产品推荐

