You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 07:15:52